From 6d076c281a3c3890159e277245dfc4540be26603 Mon Sep 17 00:00:00 2001 From: KyuubiYoru Date: Thu, 16 Jul 2026 07:37:02 +0200 Subject: [PATCH] feat: implement authenticated NAT mediator (#11) Closes #11 --- .../0005-session-lease-lifecycle.md | 14 +- ...0009-authenticated-litenet-nat-mediator.md | 68 +++ docs/architecture/README.md | 1 + docs/contracts/udp-v1.md | 59 ++- src/FinalFactory.Rendezvous.Client/README.md | 15 + .../ContractLimits.cs | 2 + .../Udp/NatPunchRequestTokenCodec.cs | 131 ++++++ .../JoinAttempts/JoinAttemptService.cs | 3 + src/FinalFactory.Rendezvous.Server/Program.cs | 11 +- .../State/InMemoryEphemeralRendezvousStore.cs | 59 ++- .../Transport/LiteNetNatRequestCodec.cs | 73 ++++ .../Transport/NatMediationProcessor.cs | 323 ++++++++++++++ .../Transport/UdpMediatorOptions.cs | 8 + .../Transport/UdpMediatorService.cs | 225 +++++----- .../appsettings.json | 4 +- .../NatPunchRequestTokenCodecTests.cs | 61 +++ .../Server/NatMediationProcessorTests.cs | 404 ++++++++++++++++++ .../Server/UdpMediatorServiceTests.cs | 357 +++++++++++++--- .../Contracts/v1/contracts-public-api.txt | 16 + 19 files changed, 1660 insertions(+), 174 deletions(-) create mode 100644 docs/architecture/0009-authenticated-litenet-nat-mediator.md create mode 100644 src/FinalFactory.Rendezvous.Contracts/Udp/NatPunchRequestTokenCodec.cs create mode 100644 src/FinalFactory.Rendezvous.Server/Transport/LiteNetNatRequestCodec.cs create mode 100644 src/FinalFactory.Rendezvous.Server/Transport/NatMediationProcessor.cs create mode 100644 tests/FinalFactory.Rendezvous.Tests/Contracts/NatPunchRequestTokenCodecTests.cs create mode 100644 tests/FinalFactory.Rendezvous.Tests/Server/NatMediationProcessorTests.cs diff --git a/docs/architecture/0005-session-lease-lifecycle.md b/docs/architecture/0005-session-lease-lifecycle.md index 5d98813..ca47bd7 100644 --- a/docs/architecture/0005-session-lease-lifecycle.md +++ b/docs/architecture/0005-session-lease-lifecycle.md @@ -60,12 +60,14 @@ the supplied ID. ### UDP presence -Only a structurally valid `HostPresence` datagram with the issued capability can -refresh presence. The public endpoint is the UDP packet's observed source on the -host's gameplay socket; the HTTP API never accepts one. The bounded local candidate -comes from the authenticated datagram. Invalid, unknown, or client-presence packets -receive no response. Presence expiry demotes public visibility but keeps the lease, -so the same handle can restore visibility without changing session identity. +Only a structurally valid frozen `HostPresence` envelope or native LiteNetLib +host-presence request with the issued capability can refresh presence. The public +endpoint is the UDP packet's observed source on the host's gameplay socket; the +HTTP API never accepts one. The bounded local candidate comes from the authenticated +packet. Invalid or unknown inputs receive no response. ADR 0009 defines the later +attempt-role use of frozen `ClientPresence` and native host/client requests. +Presence expiry demotes public visibility but keeps the lease, so the same handle +can restore visibility without changing session identity. Public listing responses contain bounded listing data only. They never contain public/local endpoints, lease tokens, presence capabilities, fingerprints, store diff --git a/docs/architecture/0009-authenticated-litenet-nat-mediator.md b/docs/architecture/0009-authenticated-litenet-nat-mediator.md new file mode 100644 index 0000000..1716513 --- /dev/null +++ b/docs/architecture/0009-authenticated-litenet-nat-mediator.md @@ -0,0 +1,68 @@ +# ADR 0009: authenticated bounded LiteNetLib NAT mediator + +- Status: Accepted +- Date: 2026-07-16 +- Tracking: #11 + +## Decision + +The server owns one LiteNetLib `NetManager` and its `NatPunchModule` on the +configured UDP endpoint. It runs in manual mode with a configured maximum number +of datagrams per poll and a short caller-owned poll interval. LiteNetLib events +are unsynchronized so authenticated requests are processed immediately on that +single polling path rather than accumulated in an unbounded event queue. The +mediator never accepts a LiteNetLib gameplay connection or handles application +payloads. + +The packet layer also consumes the frozen v1 presence envelope on the same +socket. Native NAT requests use a canonical fixed-size 192-character token that +binds a role (`HostPresence`, attempt `Host`, or attempt `Client`), mediation +handle, and the already-issued capability. Both transports enter one processor +and the same atomic store operations. No transport-supplied public address is +trusted; the socket source is authoritative. + +LiteNetLib's native NAT packet family also contains introduction-response and +punch frames that are appropriate for peers but unsafe on a public mediator: a +forged response can name arbitrary destinations. The packet layer therefore +decodes only the pinned `NatIntroduceRequest` wire shape and consumes every +inbound packet before `NatPunchModule` sees it. The module is outbound-only and +may send introductions solely from a completed authorized plan. + +Listing presence refreshes authorize no response. Attempt contributions bind the +first observed endpoint for exactly one capability role. Exact duplicates are +idempotent; a different endpoint, the opposite role, an expired/cancelled +attempt, or a stale listing presence cannot replace it. The introduction is +consumed atomically only after both roles bind and their observed address +families match, preventing concurrent attempts for one listing from cross-wiring. + +A reported local candidate is eligible only when it is RFC 1918 IPv4 or IPv6 +unique-local unicast, matches the observed family, and both peers have the same +observed public address. Otherwise `NatIntroduce` receives the observed public +endpoint in the local slot. Loopback, link-local, multicast, unspecified, +documentation IPv6, global-address claims, and cross-family claims are never +disclosed as local targets. IPv4 is required; observed global IPv6 can be used +when both peers contribute IPv6, without claiming guaranteed IPv6 NAT traversal. + +The introduction carries only the distinct connection ticket and is emitted at +most once to each verified observed endpoint. The fixed authenticated native +request and bounded frozen envelope keep the combined response bytes within the +2.0 verified amplification budget; unauthenticated inputs receive zero bytes. +Malformed, truncated, oversized, spoofed, or unrelated LiteNetLib packets do not +grow Rendezvous state. Raw endpoints and credentials are never logged or exposed +through diagnostic string representations. + +Frozen IPv6 listing-presence refresh remains valid because it emits no response. +IPv6 attempt roles require the fixed-size native LiteNetLib request; accepting the +short frozen envelope would exceed the 2.0 byte budget for two IPv6 introduction +frames. The required IPv4 listen address and optional IPv6 listen address are +configured separately so enabling one family never widens the other family to a +wildcard bind. + +## Consequences + +- Hosts refresh listing presence and answer invitations from their actual + gameplay socket; a separate mediator socket would observe the wrong mapping. +- Caller-owned SDK coordination in #12 must poll the host invitation endpoint, + send the corresponding native role token, and consume the returned ticket. +- UDP loss can prevent traversal, but it cannot cause an arbitrary destination, + replay, role substitution, or cross-attempt introduction. diff --git a/docs/architecture/README.md b/docs/architecture/README.md index 7f708ac..c1c84ad 100644 --- a/docs/architecture/README.md +++ b/docs/architecture/README.md @@ -11,6 +11,7 @@ decision requires a superseding ADR and corresponding contract/test updates. - [ADR 0006: bounded compatible session browser](0006-compatible-session-browser.md) - [ADR 0007: caller-owned .NET publisher and browser SDK](0007-caller-owned-dotnet-client-sdk.md) - [ADR 0008: scoped join attempts and one-time connection tickets](0008-scoped-join-attempts-and-tickets.md) +- [ADR 0009: authenticated bounded LiteNetLib NAT mediator](0009-authenticated-litenet-nat-mediator.md) - [Threat model](../security/threat-model.md) - [Security promise and test matrix](../security/control-matrix.md) - [Versioned HTTP and UDP contracts](../contracts/README.md) diff --git a/docs/contracts/udp-v1.md b/docs/contracts/udp-v1.md index 65bede1..9ce1b7a 100644 --- a/docs/contracts/udp-v1.md +++ b/docs/contracts/udp-v1.md @@ -1,11 +1,11 @@ -# UDP presence contract v1 +# UDP presence and NAT-punch contract v1 -Tracking: #4 +Tracking: #4, #11 -The UDP mediator accepts a single bounded presence envelope from a host or -client. It associates the authenticated mediation handle with the packet's -observed public source endpoint and the sender's reported local endpoint. It -does not carry gameplay packets. +The UDP mediator accepts the frozen bounded presence envelope below and native +LiteNetLib NAT-introduction requests. Both forms associate an authenticated +mediation handle with the packet's observed public source endpoint and the +sender's reported local endpoint. Neither form carries gameplay packets. All multi-byte integers use network byte order. UUID bytes use the canonical RFC 4122 textual order (the byte pairs from the 32 hexadecimal digits), not the @@ -49,5 +49,48 @@ Capabilities are short-lived, single-purpose, scoped to one mediation handle, and compared without exposing them in logs. A valid-looking packet does not prove authorization until the capability is checked. Invalid packets receive no UDP response, preventing the mediator from becoming an amplification oracle. -Replay, expiry, pairing, and rate-limit policy are defined by later mediator -issues; the v1 envelope deliberately leaves no unbounded or reflected payload. +For the frozen envelope, `HostPresence` is resolved against either the listing's +host-presence capability or an attempt's host-role capability. `ClientPresence` +is resolved only against the attempt's client-role capability. Handles are +globally distinct in the active store, so this does not permit role confusion. + +## Native LiteNetLib request token + +A game using LiteNetLib sends `NatPunchModule.SendNatIntroduceRequest` from its +gameplay `NetManager`. The `additionalInfo` value is produced by +`NatPunchRequestTokenCodec` and is exactly 192 ASCII characters: + +```text +rv1::<32 lowercase handle hex>:<43-character capability> +``` + +`role` is `p` for listing host-presence refresh, `h` for the host side of a join +attempt, or `c` for its client side. Padding is canonical and leaves the token +below LiteNetLib's 256-character ceiling. Its fixed size also ensures that the +two authenticated introduction responses remain within the 2.0 response-byte +budget. Tokens with a wrong length, role, handle, capability, or padding receive +no response. + +The mediator runs LiteNetLib in bounded manual-poll mode. Its packet layer admits +only the pinned native `NatIntroduceRequest` frame, consumes every inbound frame +before LiteNetLib can act on it, and uses `NatPunchModule` only to emit authorized +introductions. Native and frozen v1 inputs reach the same atomic role/capability +checks. Only the packet source is +used as the public endpoint. A claimed private candidate is retained only when +it is private unicast, matches the observed address family, and both authorized +peers were observed behind the same public address; otherwise the observed +public endpoint is substituted. IPv4 punching is required. IPv6 sources must be +observed global unicast and both roles must use IPv6; IPv6 NAT traversal remains +best-effort rather than a v1 release requirement. + +The second valid contribution atomically consumes the introduction and starts +the connection-ticket lifetime. `NatIntroduce` is called once with the distinct +43-character connection ticket. Reordered and exact duplicate requests are +idempotent. Endpoint substitution, cross-role use, stale host presence, expired +or cancelled attempts, malformed packets, and gameplay payloads produce no +introduction and create no mediator queue or endpoint state. + +Frozen envelopes may refresh listing presence over IPv6 because that operation +has no response. IPv6 attempt contributions must use the fixed-size native token; +the shorter frozen IPv6 envelope cannot fund two IPv6 introduction frames within +the 2.0 response-byte ceiling and is therefore dropped without response. diff --git a/src/FinalFactory.Rendezvous.Client/README.md b/src/FinalFactory.Rendezvous.Client/README.md index ac3ff80..19db5f7 100644 --- a/src/FinalFactory.Rendezvous.Client/README.md +++ b/src/FinalFactory.Rendezvous.Client/README.md @@ -42,6 +42,21 @@ if (!registered.IsSuccess || registered.Value is null) Load `publisherCredential` from the game's deployment secret boundary; never embed it in a client build or source control. A successful registration returns a `PublishedSession` containing the lease and host-presence capabilities. +Send a periodic presence request from the host's gameplay `NetManager` using the +server-controlled refresh interval and the fixed-size native token: + +```csharp +string presenceToken = NatPunchRequestTokenCodec.Encode( + NatPunchPeerRole.HostPresence, + session.HostPresenceHandle, + session.HostPresenceCapability); +gameplayNetManager.NatPunchModule.SendNatIntroduceRequest(mediator, presenceToken); +``` + +The same codec creates `Host` tokens for host-polled invitations and `Client` +tokens for a created join attempt. Always send them from the same LiteNetLib +socket that will carry the direct game connection; the mediator ignores any +caller-supplied public endpoint. Lease renewal is explicit and caller-controlled: diff --git a/src/FinalFactory.Rendezvous.Contracts/ContractLimits.cs b/src/FinalFactory.Rendezvous.Contracts/ContractLimits.cs index ead6b77..f29d8d4 100644 --- a/src/FinalFactory.Rendezvous.Contracts/ContractLimits.cs +++ b/src/FinalFactory.Rendezvous.Contracts/ContractLimits.cs @@ -23,6 +23,8 @@ public static class ContractLimits public const int OpaqueHttpCredentialMaxCharacters = 1_024; public const int UdpCapabilityMaxCharacters = 192; public const int ConnectionTicketMaxCharacters = 192; + public const int DerivedCredentialCharacters = 43; + public const int NatPunchRequestTokenCharacters = 192; public const int LiteNetLibNatTokenMaxCharacters = 256; public const int SessionCapacityMaxPlayers = 10_000; } diff --git a/src/FinalFactory.Rendezvous.Contracts/Udp/NatPunchRequestTokenCodec.cs b/src/FinalFactory.Rendezvous.Contracts/Udp/NatPunchRequestTokenCodec.cs new file mode 100644 index 0000000..d9ac15b --- /dev/null +++ b/src/FinalFactory.Rendezvous.Contracts/Udp/NatPunchRequestTokenCodec.cs @@ -0,0 +1,131 @@ +namespace FinalFactory.Rendezvous.Contracts; + +public enum NatPunchPeerRole +{ + HostPresence = 1, + Host = 2, + Client = 3, +} + +public sealed class NatPunchRequestToken +{ + public NatPunchPeerRole Role { get; set; } + public MediationHandle MediationHandle { get; set; } + public string Capability { get; set; } = string.Empty; + + public override string ToString() => "[NatPunchRequestToken: capability redacted]"; +} + +public static class NatPunchRequestTokenCodec +{ + public const int EncodedLength = ContractLimits.NatPunchRequestTokenCharacters; + + private const string VersionPrefix = "rv1:"; + private const int HandleLength = 32; + private const int CapabilityLength = ContractLimits.DerivedCredentialCharacters; + private const char Separator = ':'; + private const char Padding = '.'; + + public static string Encode( + NatPunchPeerRole role, + MediationHandle mediationHandle, + string capability) + { + if (!TryGetRoleCode(role, out char roleCode) + || mediationHandle.Value == Guid.Empty + || capability is null + || capability.Length != CapabilityLength + || !ContractValidation.IsCapabilityValid(capability)) + { + throw new ArgumentException("The NAT punch request token fields are invalid."); + } + + string payload = string.Concat( + VersionPrefix, + roleCode, + Separator, + mediationHandle.Value.ToString("N"), + Separator, + capability); + return payload.PadRight(EncodedLength, Padding); + } + + public static bool TryDecode(string? encoded, out NatPunchRequestToken? token) + { + token = null; + if (encoded is null + || encoded.Length != EncodedLength + || !encoded.StartsWith(VersionPrefix, StringComparison.Ordinal) + || !TryParseRole(encoded[VersionPrefix.Length], out NatPunchPeerRole role)) + { + return false; + } + + int roleSeparator = VersionPrefix.Length + 1; + int handleOffset = roleSeparator + 1; + int capabilitySeparator = handleOffset + HandleLength; + int capabilityOffset = capabilitySeparator + 1; + int paddingOffset = capabilityOffset + CapabilityLength; + string handleText = encoded.Substring(handleOffset, HandleLength); + if (encoded[roleSeparator] != Separator + || encoded[capabilitySeparator] != Separator + || !Guid.TryParseExact(handleText, "N", out Guid handle) + || handle == Guid.Empty + || !string.Equals(handleText, handle.ToString("N"), StringComparison.Ordinal) + || !ContainsOnlyPadding(encoded, paddingOffset)) + { + return false; + } + + string capability = encoded.Substring(capabilityOffset, CapabilityLength); + if (!ContractValidation.IsCapabilityValid(capability)) + { + return false; + } + + token = new NatPunchRequestToken + { + Role = role, + MediationHandle = new MediationHandle(handle), + Capability = capability, + }; + return true; + } + + private static bool TryGetRoleCode(NatPunchPeerRole role, out char code) + { + code = role switch + { + NatPunchPeerRole.HostPresence => 'p', + NatPunchPeerRole.Host => 'h', + NatPunchPeerRole.Client => 'c', + _ => default, + }; + return code != default; + } + + private static bool TryParseRole(char code, out NatPunchPeerRole role) + { + role = code switch + { + 'p' => NatPunchPeerRole.HostPresence, + 'h' => NatPunchPeerRole.Host, + 'c' => NatPunchPeerRole.Client, + _ => default, + }; + return role != default; + } + + private static bool ContainsOnlyPadding(string value, int offset) + { + for (int index = offset; index < value.Length; index++) + { + if (value[index] != Padding) + { + return false; + } + } + + return true; + } +} diff --git a/src/FinalFactory.Rendezvous.Server/JoinAttempts/JoinAttemptService.cs b/src/FinalFactory.Rendezvous.Server/JoinAttempts/JoinAttemptService.cs index 1cb042f..90227ad 100644 --- a/src/FinalFactory.Rendezvous.Server/JoinAttempts/JoinAttemptService.cs +++ b/src/FinalFactory.Rendezvous.Server/JoinAttempts/JoinAttemptService.cs @@ -316,6 +316,9 @@ internal sealed class JoinAttemptService( ContractValidation.IsCapabilityValid(hostCapability) && ContractValidation.IsCapabilityValid(clientCapability) && ContractValidation.IsConnectionTicketValid(ticket) + && hostCapability.Length == ContractLimits.DerivedCredentialCharacters + && clientCapability.Length == ContractLimits.DerivedCredentialCharacters + && ticket.Length == ContractLimits.DerivedCredentialCharacters && hostCapability.Length <= ContractLimits.LiteNetLibNatTokenMaxCharacters && clientCapability.Length <= ContractLimits.LiteNetLibNatTokenMaxCharacters; diff --git a/src/FinalFactory.Rendezvous.Server/Program.cs b/src/FinalFactory.Rendezvous.Server/Program.cs index e596aa2..b88ff2f 100644 --- a/src/FinalFactory.Rendezvous.Server/Program.cs +++ b/src/FinalFactory.Rendezvous.Server/Program.cs @@ -133,12 +133,19 @@ builder.Services .BindConfiguration(UdpMediatorOptions.SectionName) .ValidateDataAnnotations() .Validate( - options => IPAddress.TryParse(options.ListenAddress, out _), - $"{UdpMediatorOptions.SectionName}:ListenAddress must be an IP address.") + options => IPAddress.TryParse(options.ListenAddress, out IPAddress? address) + && address.AddressFamily == System.Net.Sockets.AddressFamily.InterNetwork, + $"{UdpMediatorOptions.SectionName}:ListenAddress must be an IPv4 address.") + .Validate( + options => string.IsNullOrWhiteSpace(options.Ipv6ListenAddress) + || (IPAddress.TryParse(options.Ipv6ListenAddress, out IPAddress? address) + && address.AddressFamily == System.Net.Sockets.AddressFamily.InterNetworkV6), + $"{UdpMediatorOptions.SectionName}:Ipv6ListenAddress must be an IPv6 address when configured.") .ValidateOnStart(); builder.Services.AddSingleton(); if (!isOpenApiGeneration) { + builder.Services.AddSingleton(); builder.Services.AddHostedService(static services => services.GetRequiredService()); } diff --git a/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs b/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs index 0e2c337..9237714 100644 --- a/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs +++ b/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs @@ -4,6 +4,7 @@ namespace FinalFactory.Rendezvous.Server.State; internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousStore { + private static readonly TimeSpan UdpMaintenanceInterval = TimeSpan.FromSeconds(1); private readonly object _gate = new(); private readonly EphemeralStoreOptions _options; private readonly IMonotonicClock _monotonicClock; @@ -19,6 +20,8 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto private readonly Dictionary _replay = new(StringComparer.Ordinal); private readonly Dictionary _revocations = new(StringComparer.Ordinal); private TimeSpan? _drainDeadline; + private TimeSpan _nextUdpMaintenance; + private long _maintenanceSweepCount; private bool _available = true; public InMemoryEphemeralRendezvousStore( @@ -38,6 +41,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto } public Guid InstanceId { get; } + internal long MaintenanceSweepCount => Interlocked.Read(ref _maintenanceSweepCount); public bool IsAvailable { @@ -270,6 +274,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto if (!_presenceHandles.TryGetValue(command.Handle, out SessionListingId listingId) || !_listings.TryGetValue(listingId, out ListingEntry? entry) + || entry.LeaseDeadline <= now || entry.Definition.HostPresenceFingerprint != command.CapabilityFingerprint) { return new(StoreResultCode.NotFound); @@ -285,7 +290,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto command.LocalEndpoint, now + _options.PresenceLifetime); return new(StoreResultCode.Success, Snapshot(entry)); - }, cancellationToken); + }, cancellationToken, eagerCleanup: false); public StoreResult> BrowseVisibleListings( VisibleListingQuery query, @@ -446,13 +451,18 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto return new(StoreResultCode.NotFound); } + if (attempt.IntroductionConsumed) + { + return new(StoreResultCode.Conflict); + } + RemoveAttempt(command.AttemptId); return new(StoreResultCode.Success, true); }, cancellationToken); public StoreResult BindAttemptEndpoint( BindAttemptEndpointCommand command, - CancellationToken cancellationToken = default) => Atomic(_ => + CancellationToken cancellationToken = default) => Atomic(now => { ArgumentNullException.ThrowIfNull(command); if (command.Handle.Value == Guid.Empty @@ -469,7 +479,13 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto } if (!_attemptHandles.TryGetValue(command.Handle, out JoinAttemptId attemptId) - || !_attempts.TryGetValue(attemptId, out AttemptEntry? attempt)) + || !_attempts.TryGetValue(attemptId, out AttemptEntry? attempt) + || attempt.Deadline <= now) + { + return new(StoreResultCode.NotFound); + } + + if (!HasFreshHostPresence(attempt, now)) { return new(StoreResultCode.NotFound); } @@ -503,7 +519,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto } return new(StoreResultCode.Success, Snapshot(attempt)); - }, cancellationToken); + }, cancellationToken, eagerCleanup: false); public StoreResult ConsumeIntroduction( MediationHandle handle, @@ -515,7 +531,13 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto } if (!_attemptHandles.TryGetValue(handle, out JoinAttemptId attemptId) - || !_attempts.TryGetValue(attemptId, out AttemptEntry? attempt)) + || !_attempts.TryGetValue(attemptId, out AttemptEntry? attempt) + || attempt.Deadline <= now) + { + return new(StoreResultCode.NotFound); + } + + if (!HasFreshHostPresence(attempt, now)) { return new(StoreResultCode.NotFound); } @@ -540,7 +562,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto Snapshot(attempt), attempt.HostEndpoint, attempt.ClientEndpoint)); - }, cancellationToken); + }, cancellationToken, eagerCleanup: false); public StoreResult ConsumeConnectionTicket( ConsumeConnectionTicketCommand command, @@ -687,18 +709,38 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto } } - private StoreResult Atomic(Func> operation, CancellationToken cancellationToken) + private StoreResult Atomic( + Func> operation, + CancellationToken cancellationToken, + bool eagerCleanup = true) { cancellationToken.ThrowIfCancellationRequested(); lock (_gate) { cancellationToken.ThrowIfCancellationRequested(); TimeSpan now = _monotonicClock.Elapsed; - Cleanup(now); + // Authenticated UDP duplicates need O(1) store work. Their operations + // check exact resource deadlines and amortize physical expiry removal. + bool drainExpired = _drainDeadline is TimeSpan drainDeadline + && now >= drainDeadline; + if (drainExpired || eagerCleanup || now >= _nextUdpMaintenance) + { + Cleanup(now); + _nextUdpMaintenance = now + UdpMaintenanceInterval; + } + return operation(now); } } + private bool HasFreshHostPresence(AttemptEntry attempt, TimeSpan now) => + _listings.TryGetValue(attempt.Command.ListingId, out ListingEntry? listing) + && listing.LeaseDeadline > now + && _presence.TryGetValue( + listing.Definition.HostPresenceHandle, + out PresenceEntry? presence) + && presence.Deadline > now; + private StoreResult? CheckNewWorkAdmission(string subject) { if (!_available) @@ -718,6 +760,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto private void Cleanup(TimeSpan now) { + _maintenanceSweepCount++; if (_drainDeadline is TimeSpan drainDeadline && now >= drainDeadline) { ClearActiveState(); diff --git a/src/FinalFactory.Rendezvous.Server/Transport/LiteNetNatRequestCodec.cs b/src/FinalFactory.Rendezvous.Server/Transport/LiteNetNatRequestCodec.cs new file mode 100644 index 0000000..0b40026 --- /dev/null +++ b/src/FinalFactory.Rendezvous.Server/Transport/LiteNetNatRequestCodec.cs @@ -0,0 +1,73 @@ +using System.Buffers.Binary; +using System.Net; +using System.Text; +using FinalFactory.Rendezvous.Contracts; + +namespace FinalFactory.Rendezvous.Server.Transport; + +internal static class LiteNetNatRequestCodec +{ + private const byte NatMessageProperty = 17; + private const int TypeIdentifierLength = 8; + private const int TokenLengthPrefix = NatPunchRequestTokenCodec.EncodedLength + 1; + // LiteNetLib 2.1.4's private NatIntroduceRequest type ID. The native socket + // integration test deliberately fails if a package upgrade changes this wire value. + private static ReadOnlySpan RequestTypeIdentifier => + [0x88, 0xbe, 0x10, 0x26, 0xbf, 0xb1, 0x66, 0x9c]; + + public static bool TryDecode( + ReadOnlySpan datagram, + out IPEndPoint? claimedLocalEndpoint, + out string? token) + { + claimedLocalEndpoint = null; + token = null; + if (datagram.Length < 1 + TypeIdentifierLength + 1 + 4 + 2 + 2 + + NatPunchRequestTokenCodec.EncodedLength + || datagram[0] != NatMessageProperty + || !datagram.Slice(1, TypeIdentifierLength).SequenceEqual(RequestTypeIdentifier)) + { + return false; + } + + int offset = 1 + TypeIdentifierLength; + int addressLength = datagram[offset++] switch + { + 0 => 4, + 1 => 16, + _ => 0, + }; + int expectedLength = offset + addressLength + 2 + 2 + + NatPunchRequestTokenCodec.EncodedLength; + if (addressLength == 0 || datagram.Length != expectedLength) + { + return false; + } + + IPAddress localAddress = new(datagram.Slice(offset, addressLength)); + offset += addressLength; + int localPort = BinaryPrimitives.ReadUInt16LittleEndian(datagram.Slice(offset, 2)); + offset += 2; + int encodedTokenLength = BinaryPrimitives.ReadUInt16LittleEndian(datagram.Slice(offset, 2)); + offset += 2; + if (localPort == 0 || encodedTokenLength != TokenLengthPrefix) + { + return false; + } + + ReadOnlySpan tokenBytes = datagram.Slice( + offset, + NatPunchRequestTokenCodec.EncodedLength); + for (int index = 0; index < tokenBytes.Length; index++) + { + if (tokenBytes[index] > 0x7f) + { + return false; + } + } + + claimedLocalEndpoint = new(localAddress, localPort); + token = Encoding.ASCII.GetString(tokenBytes); + return true; + } +} diff --git a/src/FinalFactory.Rendezvous.Server/Transport/NatMediationProcessor.cs b/src/FinalFactory.Rendezvous.Server/Transport/NatMediationProcessor.cs new file mode 100644 index 0000000..1e935f6 --- /dev/null +++ b/src/FinalFactory.Rendezvous.Server/Transport/NatMediationProcessor.cs @@ -0,0 +1,323 @@ +using System.Net; +using System.Net.Sockets; +using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.JoinAttempts; +using FinalFactory.Rendezvous.Server.Sessions; +using FinalFactory.Rendezvous.Server.State; + +namespace FinalFactory.Rendezvous.Server.Transport; + +internal interface INatIntroductionSink +{ + void Introduce(NatIntroductionPlan plan); +} + +internal sealed record NatIntroductionPlan( + IPEndPoint HostLocal, + IPEndPoint HostPublic, + IPEndPoint ClientLocal, + IPEndPoint ClientPublic, + string ConnectionTicket) +{ + public override string ToString() => "[NatIntroductionPlan: endpoints and ticket redacted]"; +} + +internal enum NatMediationResult +{ + Dropped = 0, + HostPresenceAccepted = 1, + HostPresenceRejected = 2, + WaitingForPeer = 3, + Introduced = 4, + Duplicate = 5, + Rejected = 6, +} + +internal sealed class NatMediationProcessor( + IEphemeralRendezvousStore store, + ISessionCapabilityService capabilities, + JoinAttemptService joinAttempts) +{ + public NatMediationResult ProcessDatagram( + ReadOnlySpan encoded, + IPEndPoint observedPublicEndpoint, + INatIntroductionSink introductionSink, + CancellationToken cancellationToken = default) + { + if (!RendezvousUdpCodec.TryDecode(encoded, out PresenceDatagram? datagram, out _) + || datagram is null + || datagram.Capability.Length != ContractLimits.DerivedCredentialCharacters + || !IPAddress.TryParse(datagram.LocalAddress, out IPAddress? localAddress)) + { + return NatMediationResult.Dropped; + } + + IPEndPoint claimedLocalEndpoint = new(localAddress, datagram.LocalPort); + NatPunchPeerRole role = datagram.MessageType == UdpPresenceMessageType.ClientPresence + ? NatPunchPeerRole.Client + : NatPunchPeerRole.HostPresence; + bool observedIpv6 = observedPublicEndpoint.AddressFamily == AddressFamily.InterNetworkV6 + && !observedPublicEndpoint.Address.IsIPv4MappedToIPv6; + if (role == NatPunchPeerRole.Client && observedIpv6) + { + return NatMediationResult.Dropped; + } + + NatMediationResult result = ProcessRequest( + claimedLocalEndpoint, + observedPublicEndpoint, + NatPunchRequestTokenCodec.Encode(role, datagram.MediationHandle, datagram.Capability), + introductionSink, + cancellationToken); + if (role != NatPunchPeerRole.HostPresence + || result != NatMediationResult.HostPresenceRejected + || observedIpv6) + { + return result; + } + + return ProcessRequest( + claimedLocalEndpoint, + observedPublicEndpoint, + NatPunchRequestTokenCodec.Encode( + NatPunchPeerRole.Host, + datagram.MediationHandle, + datagram.Capability), + introductionSink, + cancellationToken); + } + + public NatMediationResult ProcessRequest( + IPEndPoint claimedLocalEndpoint, + IPEndPoint observedPublicEndpoint, + string token, + INatIntroductionSink introductionSink, + CancellationToken cancellationToken = default) + { + ArgumentNullException.ThrowIfNull(claimedLocalEndpoint); + ArgumentNullException.ThrowIfNull(observedPublicEndpoint); + ArgumentNullException.ThrowIfNull(introductionSink); + + if (!NatPunchRequestTokenCodec.TryDecode(token, out NatPunchRequestToken? request) + || request is null + || !TryCreateObservedEndpoint(observedPublicEndpoint, out ObservedEndpoint publicEndpoint) + || !capabilities.TryFingerprint(request.Capability, out SecretFingerprint fingerprint)) + { + return NatMediationResult.Dropped; + } + + ObservedEndpoint? localEndpoint = TryCreatePrivateCandidate( + claimedLocalEndpoint, + publicEndpoint.AddressFamily, + out ObservedEndpoint candidate) + ? candidate + : null; + + if (request.Role == NatPunchPeerRole.HostPresence) + { + StoreResult presence = store.BindHostPresence(new( + request.MediationHandle, + fingerprint, + publicEndpoint, + localEndpoint), cancellationToken); + return presence.Succeeded + ? NatMediationResult.HostPresenceAccepted + : NatMediationResult.HostPresenceRejected; + } + + AttemptPeerRole role = request.Role switch + { + NatPunchPeerRole.Host => AttemptPeerRole.Host, + NatPunchPeerRole.Client => AttemptPeerRole.Client, + _ => default, + }; + if (role == default) + { + return NatMediationResult.Dropped; + } + + StoreResult bound = store.BindAttemptEndpoint(new( + request.MediationHandle, + role, + fingerprint, + publicEndpoint, + localEndpoint), cancellationToken); + if (!bound.Succeeded || bound.Value is null) + { + return bound.Code == StoreResultCode.ReplayRejected + ? NatMediationResult.Rejected + : NatMediationResult.Dropped; + } + + StoredJoinAttempt attempt = bound.Value; + if (attempt.IntroductionConsumed) + { + return NatMediationResult.Duplicate; + } + + if (attempt.HostEndpoint is null || attempt.ClientEndpoint is null) + { + return NatMediationResult.WaitingForPeer; + } + + if (attempt.HostEndpoint.PublicEndpoint.AddressFamily + != attempt.ClientEndpoint.PublicEndpoint.AddressFamily) + { + return NatMediationResult.Rejected; + } + + StoreResult consumed = store.ConsumeIntroduction( + request.MediationHandle, + cancellationToken); + if (!consumed.Succeeded || consumed.Value is null) + { + return consumed.Code == StoreResultCode.ReplayRejected + ? NatMediationResult.Duplicate + : NatMediationResult.Rejected; + } + + JoinAttemptServiceResult ticket = joinAttempts.IssueConnectionTicket( + consumed.Value.Attempt); + if (!ticket.Succeeded || ticket.Value is null) + { + return NatMediationResult.Rejected; + } + + try + { + introductionSink.Introduce(CreatePlan(consumed.Value, ticket.Value.Ticket)); + return NatMediationResult.Introduced; + } + catch (Exception exception) when (exception is SocketException + or InvalidOperationException + or ArgumentException) + { + return NatMediationResult.Rejected; + } + } + + private static NatIntroductionPlan CreatePlan( + IntroductionEndpoints endpoints, + string connectionTicket) + { + IPEndPoint hostPublic = ToIpEndpoint(endpoints.Host.PublicEndpoint); + IPEndPoint clientPublic = ToIpEndpoint(endpoints.Client.PublicEndpoint); + bool sameNat = hostPublic.Address.Equals(clientPublic.Address); + IPEndPoint hostLocal = sameNat && endpoints.Host.LocalEndpoint is { } hostCandidate + ? ToIpEndpoint(hostCandidate) + : hostPublic; + IPEndPoint clientLocal = sameNat && endpoints.Client.LocalEndpoint is { } clientCandidate + ? ToIpEndpoint(clientCandidate) + : clientPublic; + return new(hostLocal, hostPublic, clientLocal, clientPublic, connectionTicket); + } + + private static bool TryCreateObservedEndpoint( + IPEndPoint source, + out ObservedEndpoint endpoint) + { + endpoint = default; + if (source.Port is < 1 or > 65_535) + { + return false; + } + + IPAddress address = source.Address.IsIPv4MappedToIPv6 + ? source.Address.MapToIPv4() + : source.Address; + if (address.Equals(IPAddress.Any) + || address.Equals(IPAddress.IPv6Any) + || address.IsIPv6Multicast + || IsIpv4MulticastOrBroadcast(address) + || (address.AddressFamily == AddressFamily.InterNetworkV6 + && !IsGlobalIpv6(address))) + { + return false; + } + + AddressFamilyKind family = address.AddressFamily switch + { + AddressFamily.InterNetwork => AddressFamilyKind.Ipv4, + AddressFamily.InterNetworkV6 => AddressFamilyKind.Ipv6, + _ => default, + }; + if (family == default) + { + return false; + } + + endpoint = new(family, address.ToString(), source.Port); + return true; + } + + private static bool TryCreatePrivateCandidate( + IPEndPoint source, + AddressFamilyKind publicFamily, + out ObservedEndpoint endpoint) + { + endpoint = default; + if (source.Port is < 1 or > 65_535) + { + return false; + } + + IPAddress address = source.Address.IsIPv4MappedToIPv6 + ? source.Address.MapToIPv4() + : source.Address; + AddressFamilyKind family = address.AddressFamily switch + { + AddressFamily.InterNetwork => AddressFamilyKind.Ipv4, + AddressFamily.InterNetworkV6 => AddressFamilyKind.Ipv6, + _ => default, + }; + if (family != publicFamily || !IsPrivateUnicast(address)) + { + return false; + } + + endpoint = new(family, address.ToString(), source.Port); + return true; + } + + private static bool IsPrivateUnicast(IPAddress address) + { + byte[] bytes = address.GetAddressBytes(); + return address.AddressFamily switch + { + AddressFamily.InterNetwork => bytes[0] == 10 + || (bytes[0] == 172 && bytes[1] is >= 16 and <= 31) + || (bytes[0] == 192 && bytes[1] == 168), + AddressFamily.InterNetworkV6 => (bytes[0] & 0xfe) == 0xfc, + _ => false, + }; + } + + private static bool IsGlobalIpv6(IPAddress address) => + !address.Equals(IPAddress.IPv6Loopback) + && !address.Equals(IPAddress.IPv6Any) + && !address.IsIPv6LinkLocal + && !address.IsIPv6Multicast + && !address.IsIPv6SiteLocal + && !IsPrivateUnicast(address) + && !IsDocumentationIpv6(address); + + private static bool IsIpv4MulticastOrBroadcast(IPAddress address) + { + if (address.AddressFamily != AddressFamily.InterNetwork) + { + return false; + } + + byte[] bytes = address.GetAddressBytes(); + return bytes[0] >= 224 || bytes.All(static value => value == byte.MaxValue); + } + + private static bool IsDocumentationIpv6(IPAddress address) + { + byte[] bytes = address.GetAddressBytes(); + return bytes[0] == 0x20 && bytes[1] == 0x01 && bytes[2] == 0x0d && bytes[3] == 0xb8; + } + + private static IPEndPoint ToIpEndpoint(ObservedEndpoint endpoint) => + new(IPAddress.Parse(endpoint.Address), endpoint.Port); +} diff --git a/src/FinalFactory.Rendezvous.Server/Transport/UdpMediatorOptions.cs b/src/FinalFactory.Rendezvous.Server/Transport/UdpMediatorOptions.cs index 2e8b389..a91b199 100644 --- a/src/FinalFactory.Rendezvous.Server/Transport/UdpMediatorOptions.cs +++ b/src/FinalFactory.Rendezvous.Server/Transport/UdpMediatorOptions.cs @@ -18,9 +18,17 @@ public sealed class UdpMediatorOptions [Required] public string ListenAddress { get; set; } = "0.0.0.0"; + public string? Ipv6ListenAddress { get; set; } + /// /// Gets or sets the UDP port. Zero requests an ephemeral port for tests. /// [Range(0, 65_535)] public int Port { get; set; } = 9050; + + [Range(1, 4_096)] + public int MaxDatagramsPerPoll { get; set; } = 256; + + [Range(1, 100)] + public int PollIntervalMilliseconds { get; set; } = 2; } diff --git a/src/FinalFactory.Rendezvous.Server/Transport/UdpMediatorService.cs b/src/FinalFactory.Rendezvous.Server/Transport/UdpMediatorService.cs index 5a5422f..722d2e7 100644 --- a/src/FinalFactory.Rendezvous.Server/Transport/UdpMediatorService.cs +++ b/src/FinalFactory.Rendezvous.Server/Transport/UdpMediatorService.cs @@ -1,161 +1,136 @@ +using System.Diagnostics; using System.Net; using System.Net.Sockets; using FinalFactory.Rendezvous.Contracts; -using FinalFactory.Rendezvous.Server.Sessions; -using FinalFactory.Rendezvous.Server.State; +using LiteNetLib; +using LiteNetLib.Layers; using Microsoft.Extensions.Options; namespace FinalFactory.Rendezvous.Server.Transport; -/// -/// Owns the cancellable UDP socket used by the future NAT mediator. -/// internal sealed partial class UdpMediatorService : BackgroundService { private readonly ILogger _logger; private readonly UdpMediatorOptions _options; - private readonly IEphemeralRendezvousStore _store; - private readonly ISessionCapabilityService _capabilities; - private UdpClient? _udpClient; + private readonly NatMediationProcessor _processor; + private LiteNetManager? _manager; + private LiteNetIntroductionSink? _introductionSink; - /// - /// Initializes a new UDP mediator service. - /// public UdpMediatorService( IOptions options, ILogger logger, - IEphemeralRendezvousStore store, - ISessionCapabilityService capabilities) + NatMediationProcessor processor) { _options = options.Value; _logger = logger; - _store = store; - _capabilities = capabilities; + _processor = processor; } - /// - /// Gets the bound endpoint after startup completes. - /// public IPEndPoint? LocalEndpoint { get; private set; } + public IPEndPoint? LocalIpv6Endpoint { get; private set; } - /// public override Task StartAsync(CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); - - if (_udpClient is not null) + if (_manager is not null) { throw new InvalidOperationException("The UDP mediator is already running."); } IPAddress listenAddress = IPAddress.Parse(_options.ListenAddress); - UdpClient udpClient = new(new IPEndPoint(listenAddress, _options.Port)); - _udpClient = udpClient; - IPEndPoint localEndpoint = - (IPEndPoint?)udpClient.Client.LocalEndPoint - ?? throw new InvalidOperationException("The UDP socket did not expose its bound endpoint."); - LocalEndpoint = localEndpoint; + if (listenAddress.AddressFamily != AddressFamily.InterNetwork) + { + throw new InvalidOperationException("The required UDP listen address must be IPv4."); + } - LogMediatorListening(_logger, localEndpoint.Address, localEndpoint.Port); + IPAddress? ipv6ListenAddress = string.IsNullOrWhiteSpace(_options.Ipv6ListenAddress) + ? null + : IPAddress.Parse(_options.Ipv6ListenAddress); + if (ipv6ListenAddress is not null + && ipv6ListenAddress.AddressFamily != AddressFamily.InterNetworkV6) + { + throw new InvalidOperationException("The optional UDP IPv6 listen address must be IPv6."); + } + EventBasedLiteNetListener listener = new(); + RendezvousPacketLayer packetLayer = new(_processor); + LiteNetManager manager = new(listener, packetLayer) + { + NatPunchEnabled = true, + IPv6Enabled = ipv6ListenAddress is not null, + UnsyncedEvents = true, + MaxPacketPerManualReceive = _options.MaxDatagramsPerPoll, + }; + manager.NatPunchModule.UnsyncedEvents = true; + _introductionSink = new(manager.NatPunchModule); + packetLayer.Attach(_introductionSink); + + if (!manager.StartInManualMode( + listenAddress, + ipv6ListenAddress ?? IPAddress.IPv6Any, + _options.Port)) + { + _introductionSink = null; + manager.Stop(); + throw new InvalidOperationException("The UDP mediator could not bind its LiteNetLib socket."); + } + + _manager = manager; + LocalEndpoint = new(listenAddress, manager.LocalPort); + LocalIpv6Endpoint = ipv6ListenAddress is null + ? null + : new(ipv6ListenAddress, manager.LocalPort); + LogMediatorListening(_logger, listenAddress, manager.LocalPort); return base.StartAsync(cancellationToken); } - /// public override async Task StopAsync(CancellationToken cancellationToken) { await base.StopAsync(cancellationToken).ConfigureAwait(false); - _udpClient?.Dispose(); - _udpClient = null; - LocalEndpoint = null; + StopManager(); LogMediatorStopped(_logger); } - /// public override void Dispose() { - _udpClient?.Dispose(); - _udpClient = null; - LocalEndpoint = null; + StopManager(); base.Dispose(); } - /// protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - UdpClient udpClient = _udpClient + LiteNetManager manager = _manager ?? throw new InvalidOperationException("The UDP mediator socket was not initialized."); - + long previous = Stopwatch.GetTimestamp(); try { while (!stoppingToken.IsCancellationRequested) { - UdpReceiveResult received = await udpClient - .ReceiveAsync(stoppingToken) - .ConfigureAwait(false); - ProcessDatagram(received.Buffer, received.RemoteEndPoint, stoppingToken); - // Bootstrap deliberately emits no UDP response. Protocol handling lands in #11. + manager.PollEvents(); + manager.NatPunchModule.PollEvents(); + long current = Stopwatch.GetTimestamp(); + manager.ManualUpdate((float)Stopwatch.GetElapsedTime(previous, current).TotalMilliseconds); + previous = current; + await Task.Delay(_options.PollIntervalMilliseconds, stoppingToken).ConfigureAwait(false); } } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { - // Expected during normal shutdown. - } - catch (ObjectDisposedException) when (stoppingToken.IsCancellationRequested) - { - // Disposing the socket is the fallback that releases a blocked receive. } finally { LocalEndpoint = null; + LocalIpv6Endpoint = null; } } - internal UdpPresenceProcessingResult ProcessDatagram( - ReadOnlySpan encoded, - IPEndPoint observedSource, - CancellationToken cancellationToken = default) + private void StopManager() { - ArgumentNullException.ThrowIfNull(observedSource); - if (!RendezvousUdpCodec.TryDecode(encoded, out PresenceDatagram? datagram, out _) - || datagram is null - || !_capabilities.TryFingerprint(datagram.Capability, out SecretFingerprint fingerprint)) - { - return UdpPresenceProcessingResult.Dropped; - } - - if (datagram.MessageType != UdpPresenceMessageType.HostPresence) - { - return UdpPresenceProcessingResult.ClientPresenceDeferred; - } - - AddressFamilyKind publicFamily = observedSource.AddressFamily switch - { - AddressFamily.InterNetwork => AddressFamilyKind.Ipv4, - AddressFamily.InterNetworkV6 => AddressFamilyKind.Ipv6, - _ => 0, - }; - if (publicFamily == 0) - { - return UdpPresenceProcessingResult.Dropped; - } - - ObservedEndpoint publicEndpoint = new( - publicFamily, - observedSource.Address.ToString(), - observedSource.Port); - ObservedEndpoint localEndpoint = new( - datagram.AddressFamily, - datagram.LocalAddress, - datagram.LocalPort); - StoreResult bound = _store.BindHostPresence(new( - datagram.MediationHandle, - fingerprint, - publicEndpoint, - localEndpoint), cancellationToken); - return bound.Succeeded - ? UdpPresenceProcessingResult.HostPresenceAccepted - : UdpPresenceProcessingResult.HostPresenceRejected; + LiteNetManager? manager = Interlocked.Exchange(ref _manager, null); + _introductionSink = null; + LocalEndpoint = null; + LocalIpv6Endpoint = null; + manager?.Stop(); } [LoggerMessage( @@ -172,12 +147,62 @@ internal sealed partial class UdpMediatorService : BackgroundService Level = LogLevel.Information, Message = "UDP mediator stopped")] private static partial void LogMediatorStopped(ILogger logger); -} -internal enum UdpPresenceProcessingResult -{ - Dropped = 0, - HostPresenceAccepted = 1, - HostPresenceRejected = 2, - ClientPresenceDeferred = 3, + private sealed class LiteNetIntroductionSink(NatPunchModule module) : INatIntroductionSink + { + public void Introduce(NatIntroductionPlan plan) => module.NatIntroduce( + plan.HostLocal, + plan.HostPublic, + plan.ClientLocal, + plan.ClientPublic, + plan.ConnectionTicket); + } + + private sealed class RendezvousPacketLayer(NatMediationProcessor processor) : PacketLayerBase(0) + { + private INatIntroductionSink? _sink; + + public void Attach(INatIntroductionSink sink) => _sink = sink; + + public override void ProcessInboundPacket( + ref IPEndPoint endPoint, + ref byte[] data, + ref int length) + { + bool isFrozenEnvelope = length >= 2 + && data[0] == RendezvousUdpCodec.MagicFirst + && data[1] == RendezvousUdpCodec.MagicSecond; + INatIntroductionSink? sink = _sink; + if (isFrozenEnvelope) + { + if (sink is not null) + { + _ = processor.ProcessDatagram(data.AsSpan(0, length), endPoint, sink); + } + } + else if (sink is not null + && LiteNetNatRequestCodec.TryDecode( + data.AsSpan(0, length), + out IPEndPoint? claimedLocalEndpoint, + out string? token) + && claimedLocalEndpoint is not null + && token is not null) + { + _ = processor.ProcessRequest(claimedLocalEndpoint, endPoint, token, sink); + } + + // Every inbound packet is consumed here. NatPunchModule is used only for outbound introductions. + Drop(ref length); + } + + public override void ProcessOutBoundPacket( + ref IPEndPoint endPoint, + ref byte[] data, + ref int offset, + ref int length) + { + } + + private static void Drop(ref int length) => length = 0; + } } diff --git a/src/FinalFactory.Rendezvous.Server/appsettings.json b/src/FinalFactory.Rendezvous.Server/appsettings.json index 966bfca..c286692 100644 --- a/src/FinalFactory.Rendezvous.Server/appsettings.json +++ b/src/FinalFactory.Rendezvous.Server/appsettings.json @@ -2,7 +2,9 @@ "Rendezvous": { "Udp": { "ListenAddress": "0.0.0.0", - "Port": 9050 + "Port": 9050, + "MaxDatagramsPerPoll": 256, + "PollIntervalMilliseconds": 2 } }, "Logging": { diff --git a/tests/FinalFactory.Rendezvous.Tests/Contracts/NatPunchRequestTokenCodecTests.cs b/tests/FinalFactory.Rendezvous.Tests/Contracts/NatPunchRequestTokenCodecTests.cs new file mode 100644 index 0000000..ef73cdb --- /dev/null +++ b/tests/FinalFactory.Rendezvous.Tests/Contracts/NatPunchRequestTokenCodecTests.cs @@ -0,0 +1,61 @@ +using FinalFactory.Rendezvous.Contracts; + +namespace FinalFactory.Rendezvous.Tests.Contracts; + +public sealed class NatPunchRequestTokenCodecTests +{ + [Theory] + [InlineData(NatPunchPeerRole.HostPresence)] + [InlineData(NatPunchPeerRole.Host)] + [InlineData(NatPunchPeerRole.Client)] + public void FixedSizeTokensRoundTripBelowLiteNetLibLimit(NatPunchPeerRole role) + { + MediationHandle handle = new(Guid.NewGuid()); + const string capability = "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"; + + string encoded = NatPunchRequestTokenCodec.Encode(role, handle, capability); + + Assert.Equal(NatPunchRequestTokenCodec.EncodedLength, encoded.Length); + Assert.True(encoded.Length <= ContractLimits.LiteNetLibNatTokenMaxCharacters); + Assert.True(NatPunchRequestTokenCodec.TryDecode(encoded, out NatPunchRequestToken? decoded)); + Assert.NotNull(decoded); + Assert.Equal(role, decoded.Role); + Assert.Equal(handle, decoded.MediationHandle); + Assert.Equal(capability, decoded.Capability); + Assert.DoesNotContain(capability, decoded.ToString(), StringComparison.Ordinal); + } + + [Fact] + public void MalformedAndNonCanonicalTokensAreRejected() + { + string valid = NatPunchRequestTokenCodec.Encode( + NatPunchPeerRole.Client, + new MediationHandle(Guid.Parse("00112233-4455-6677-8899-aabbccddeeff")), + "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"); + string uppercaseHandle = valid[..6] + + valid.Substring(6, 32).ToUpperInvariant() + + valid[38..]; + + Assert.False(NatPunchRequestTokenCodec.TryDecode(null, out _)); + Assert.False(NatPunchRequestTokenCodec.TryDecode(valid[..^1], out _)); + Assert.False(NatPunchRequestTokenCodec.TryDecode("x" + valid[1..], out _)); + Assert.False(NatPunchRequestTokenCodec.TryDecode(valid[..^1] + "x", out _)); + Assert.False(NatPunchRequestTokenCodec.TryDecode(uppercaseHandle, out _)); + Assert.Throws(() => NatPunchRequestTokenCodec.Encode( + (NatPunchPeerRole)99, + new MediationHandle(Guid.NewGuid()), + "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA")); + } + + [Fact] + public void FixedWireLengthHasADedicatedLiteNetSafeContractLimit() + { + Assert.Equal(192, ContractLimits.NatPunchRequestTokenCharacters); + Assert.Equal( + ContractLimits.NatPunchRequestTokenCharacters, + NatPunchRequestTokenCodec.EncodedLength); + Assert.True( + ContractLimits.NatPunchRequestTokenCharacters + <= ContractLimits.LiteNetLibNatTokenMaxCharacters); + } +} diff --git a/tests/FinalFactory.Rendezvous.Tests/Server/NatMediationProcessorTests.cs b/tests/FinalFactory.Rendezvous.Tests/Server/NatMediationProcessorTests.cs new file mode 100644 index 0000000..873da39 --- /dev/null +++ b/tests/FinalFactory.Rendezvous.Tests/Server/NatMediationProcessorTests.cs @@ -0,0 +1,404 @@ +using System.Collections.Concurrent; +using System.Net; +using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.JoinAttempts; +using FinalFactory.Rendezvous.Server.State; +using FinalFactory.Rendezvous.Server.Transport; +using FinalFactory.Rendezvous.Tests.JoinAttempts; + +namespace FinalFactory.Rendezvous.Tests.Server; + +public sealed class NatMediationProcessorTests +{ + [Fact] + public void AuthenticatedHostPresenceUsesTheObservedGameplaySocket() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(bindPresence: false); + NatMediationProcessor processor = CreateProcessor(fixture); + CaptureIntroductionSink sink = new(); + string token = NatPunchRequestTokenCodec.Encode( + NatPunchPeerRole.HostPresence, + registration.HostPresenceHandle, + registration.HostPresenceCapability); + + Assert.Equal( + NatMediationResult.HostPresenceAccepted, + processor.ProcessRequest( + Endpoint("192.168.1.50", 40_000), + Endpoint("203.0.113.77", 51_234), + token, + sink)); + Assert.Equal(registration.ListingId, Assert.Single(fixture.Sessions.Browse()).Definition.ListingId); + Assert.Empty(sink.Plans); + } + + [Fact] + public void MatchedPeersReceiveOneIntroductionAndSameNatPrivateCandidates() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + AttemptCredentials attempt = CreateAttempt(fixture, registration, "same-nat"); + NatMediationProcessor processor = CreateProcessor(fixture); + CaptureIntroductionSink sink = new(); + + Assert.Equal( + NatMediationResult.WaitingForPeer, + Process(processor, sink, attempt, NatPunchPeerRole.Host, + Endpoint("192.168.1.10", 41_000), Endpoint("203.0.113.20", 51_000))); + Assert.Equal( + NatMediationResult.Introduced, + Process(processor, sink, attempt, NatPunchPeerRole.Client, + Endpoint("192.168.1.11", 42_000), Endpoint("203.0.113.20", 52_000))); + + NatIntroductionPlan plan = Assert.Single(sink.Plans); + Assert.Equal(Endpoint("192.168.1.10", 41_000), plan.HostLocal); + Assert.Equal(Endpoint("192.168.1.11", 42_000), plan.ClientLocal); + Assert.Equal(Endpoint("203.0.113.20", 51_000), plan.HostPublic); + Assert.Equal(Endpoint("203.0.113.20", 52_000), plan.ClientPublic); + Assert.Equal(43, plan.ConnectionTicket.Length); + Assert.DoesNotContain(plan.ConnectionTicket, plan.ToString(), StringComparison.Ordinal); + + Assert.True(fixture.Sessions.Capabilities.TryFingerprint( + plan.ConnectionTicket, + out SecretFingerprint ticketFingerprint)); + Assert.True(fixture.Sessions.Store.ConsumeConnectionTicket(new( + attempt.AttemptId, + ticketFingerprint)).Succeeded); + + Assert.Equal( + NatMediationResult.Duplicate, + Process(processor, sink, attempt, NatPunchPeerRole.Client, + Endpoint("192.168.1.11", 42_000), Endpoint("203.0.113.20", 52_000))); + Assert.Single(sink.Plans); + } + + [Fact] + public void DifferentNatsAndInvalidLocalClaimsExposeOnlyObservedPublicEndpoints() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + AttemptCredentials attempt = CreateAttempt(fixture, registration, "different-nats"); + NatMediationProcessor processor = CreateProcessor(fixture); + CaptureIntroductionSink sink = new(); + + _ = Process(processor, sink, attempt, NatPunchPeerRole.Client, + Endpoint("8.8.8.8", 42_000), Endpoint("198.51.100.40", 52_000)); + Assert.Equal( + NatMediationResult.Introduced, + Process(processor, sink, attempt, NatPunchPeerRole.Host, + Endpoint("192.168.1.10", 41_000), Endpoint("203.0.113.20", 51_000))); + + NatIntroductionPlan plan = Assert.Single(sink.Plans); + Assert.Equal(plan.HostPublic, plan.HostLocal); + Assert.Equal(plan.ClientPublic, plan.ClientLocal); + Assert.NotEqual(IPAddress.Parse("8.8.8.8"), plan.ClientLocal.Address); + } + + [Fact] + public void RoleAndEndpointSubstitutionAreRejectedWithoutChangingTheFirstBinding() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + AttemptCredentials attempt = CreateAttempt(fixture, registration, "substitution"); + NatMediationProcessor processor = CreateProcessor(fixture); + CaptureIntroductionSink sink = new(); + + string crossRole = NatPunchRequestTokenCodec.Encode( + NatPunchPeerRole.Host, + attempt.Handle, + attempt.ClientCapability); + Assert.Equal( + NatMediationResult.Dropped, + processor.ProcessRequest( + Endpoint("192.168.1.10", 41_000), + Endpoint("203.0.113.20", 51_000), + crossRole, + sink)); + + Assert.Equal( + NatMediationResult.WaitingForPeer, + Process(processor, sink, attempt, NatPunchPeerRole.Host, + Endpoint("192.168.1.10", 41_000), Endpoint("203.0.113.20", 51_000))); + Assert.Equal( + NatMediationResult.Rejected, + Process(processor, sink, attempt, NatPunchPeerRole.Host, + Endpoint("192.168.1.99", 41_999), Endpoint("203.0.113.99", 51_999))); + Assert.Equal( + NatMediationResult.Introduced, + Process(processor, sink, attempt, NatPunchPeerRole.Client, + Endpoint("192.168.2.10", 42_000), Endpoint("198.51.100.40", 52_000))); + + Assert.Equal(Endpoint("203.0.113.20", 51_000), Assert.Single(sink.Plans).HostPublic); + } + + [Fact] + public void ConcurrentAttemptsForOneSessionNeverCrossWire() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + AttemptCredentials first = CreateAttempt(fixture, registration, "parallel-1"); + AttemptCredentials second = CreateAttempt(fixture, registration, "parallel-2"); + NatMediationProcessor processor = CreateProcessor(fixture); + CaptureIntroductionSink sink = new(); + + _ = Process(processor, sink, first, NatPunchPeerRole.Host, + Endpoint("10.0.0.10", 41_001), Endpoint("203.0.113.10", 51_001)); + _ = Process(processor, sink, second, NatPunchPeerRole.Host, + Endpoint("10.0.0.20", 41_002), Endpoint("203.0.113.20", 51_002)); + _ = Process(processor, sink, second, NatPunchPeerRole.Client, + Endpoint("10.0.0.21", 42_002), Endpoint("198.51.100.20", 52_002)); + _ = Process(processor, sink, first, NatPunchPeerRole.Client, + Endpoint("10.0.0.11", 42_001), Endpoint("198.51.100.10", 52_001)); + + Assert.Equal(2, sink.Plans.Count); + Assert.Contains(sink.Plans, plan => + plan.HostPublic.Equals(Endpoint("203.0.113.10", 51_001)) + && plan.ClientPublic.Equals(Endpoint("198.51.100.10", 52_001))); + Assert.Contains(sink.Plans, plan => + plan.HostPublic.Equals(Endpoint("203.0.113.20", 51_002)) + && plan.ClientPublic.Equals(Endpoint("198.51.100.20", 52_002))); + } + + [Fact] + public async Task ConcurrentDuplicateCompletionEmitsExactlyOneIntroduction() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + AttemptCredentials attempt = CreateAttempt(fixture, registration, "completion-race"); + NatMediationProcessor processor = CreateProcessor(fixture); + ConcurrentIntroductionSink sink = new(); + _ = Process(processor, sink, attempt, NatPunchPeerRole.Host, + Endpoint("192.168.1.10", 41_000), Endpoint("203.0.113.20", 51_000)); + using Barrier barrier = new(2); + + Task[] completions = Enumerable.Range(0, 2) + .Select(_ => Task.Run(() => + { + barrier.SignalAndWait(); + return Process(processor, sink, attempt, NatPunchPeerRole.Client, + Endpoint("192.168.1.11", 42_000), Endpoint("198.51.100.40", 52_000)); + })) + .ToArray(); + NatMediationResult[] results = await Task.WhenAll(completions); + + Assert.Single(results, result => result == NatMediationResult.Introduced); + Assert.Single(sink.Plans); + } + + [Fact] + public async Task CancellationCannotReportSuccessAfterIntroductionIsConsumed() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + AttemptCredentials attempt = CreateAttempt(fixture, registration, "cancel-race"); + NatMediationProcessor processor = CreateProcessor(fixture); + using BlockingIntroductionSink sink = new(); + _ = Process(processor, sink, attempt, NatPunchPeerRole.Host, + Endpoint("192.168.1.10", 41_000), Endpoint("203.0.113.20", 51_000)); + Task completion = Task.Run(() => Process( + processor, + sink, + attempt, + NatPunchPeerRole.Client, + Endpoint("192.168.1.11", 42_000), + Endpoint("198.51.100.40", 52_000))); + Assert.True(sink.WaitUntilEntered(TimeSpan.FromSeconds(2))); + + JoinAttemptServiceResult cancelled = fixture.Service.Cancel( + attempt.AttemptId, + attempt.ClientCapability); + + Assert.Equal(RendezvousErrorCode.Conflict, cancelled.Error); + sink.Release(); + Assert.Equal(NatMediationResult.Introduced, await completion); + } + + [Fact] + public void DuplicateFloodAmortizesGlobalExpiryMaintenance() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + AttemptCredentials attempt = CreateAttempt(fixture, registration, "maintenance-budget"); + NatMediationProcessor processor = CreateProcessor(fixture); + CaptureIntroductionSink sink = new(); + long before = fixture.Sessions.Store.MaintenanceSweepCount; + + for (int index = 0; index < 256; index++) + { + Assert.Equal( + NatMediationResult.WaitingForPeer, + Process(processor, sink, attempt, NatPunchPeerRole.Host, + Endpoint("192.168.1.10", 41_000), Endpoint("203.0.113.20", 51_000))); + } + + Assert.InRange(fixture.Sessions.Store.MaintenanceSweepCount - before, 0, 1); + Assert.Empty(sink.Plans); + } + + [Fact] + public void MissingStaleCancelledAndMalformedRequestsNeverIntroduce() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + AttemptCredentials stale = CreateAttempt(fixture, registration, "stale"); + AttemptCredentials cancelled = CreateAttempt(fixture, registration, "cancelled"); + NatMediationProcessor processor = CreateProcessor(fixture); + CaptureIntroductionSink sink = new(); + + Assert.Equal( + NatMediationResult.WaitingForPeer, + Process(processor, sink, stale, NatPunchPeerRole.Client, + Endpoint("192.168.1.11", 42_000), Endpoint("198.51.100.40", 52_000))); + Assert.True(fixture.Service.Cancel(cancelled.AttemptId, cancelled.ClientCapability).Succeeded); + Assert.Equal( + NatMediationResult.Dropped, + Process(processor, sink, cancelled, NatPunchPeerRole.Client, + Endpoint("192.168.1.12", 42_001), Endpoint("198.51.100.41", 52_001))); + + fixture.Sessions.Clock.Advance(TimeSpan.FromSeconds(21)); + Assert.Equal( + NatMediationResult.Dropped, + Process(processor, sink, stale, NatPunchPeerRole.Host, + Endpoint("192.168.1.10", 41_000), Endpoint("203.0.113.20", 51_000))); + Assert.Equal( + NatMediationResult.Dropped, + processor.ProcessRequest( + Endpoint("192.168.1.10", 41_000), + Endpoint("203.0.113.20", 51_000), + "malformed", + sink)); + Assert.Empty(sink.Plans); + } + + [Fact] + public void AddressFamiliesMustMatchAndOnlyGlobalIpv6SourcesAreAccepted() + { + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + AttemptCredentials mismatch = CreateAttempt(fixture, registration, "family-mismatch"); + AttemptCredentials ipv6 = CreateAttempt(fixture, registration, "ipv6"); + NatMediationProcessor processor = CreateProcessor(fixture); + CaptureIntroductionSink sink = new(); + + byte[] shortFrozenIpv6 = RendezvousUdpCodec.Encode(new PresenceDatagram + { + MessageType = UdpPresenceMessageType.ClientPresence, + MediationHandle = ipv6.Handle, + AddressFamily = AddressFamilyKind.Ipv6, + LocalAddress = "fd00::11", + LocalPort = 42_000, + Capability = ipv6.ClientCapability, + }); + Assert.Equal( + NatMediationResult.Dropped, + processor.ProcessDatagram( + shortFrozenIpv6, + Endpoint("2606:4700:4700::1001", 52_000), + sink)); + + _ = Process(processor, sink, mismatch, NatPunchPeerRole.Host, + Endpoint("192.168.1.10", 41_000), Endpoint("203.0.113.20", 51_000)); + Assert.Equal( + NatMediationResult.Rejected, + Process(processor, sink, mismatch, NatPunchPeerRole.Client, + Endpoint("fd00::11", 42_000), Endpoint("2606:4700:4700::1111", 52_000))); + + Assert.Equal( + NatMediationResult.Dropped, + Process(processor, sink, ipv6, NatPunchPeerRole.Host, + Endpoint("fd00::10", 41_000), Endpoint("2001:db8::10", 51_000))); + _ = Process(processor, sink, ipv6, NatPunchPeerRole.Host, + Endpoint("fd00::10", 41_000), Endpoint("2606:4700:4700::1000", 51_000)); + Assert.Equal( + NatMediationResult.Introduced, + Process(processor, sink, ipv6, NatPunchPeerRole.Client, + Endpoint("fd00::11", 42_000), Endpoint("2606:4700:4700::1001", 52_000))); + Assert.Single(sink.Plans); + } + + private static NatMediationProcessor CreateProcessor(JoinAttemptFixture fixture) => new( + fixture.Sessions.Store, + fixture.Sessions.Capabilities, + fixture.Service); + + private static AttemptCredentials CreateAttempt( + JoinAttemptFixture fixture, + RegisterSessionResponse registration, + string idempotencyKey) + { + CreateJoinAttemptResponse created = fixture.Create(registration.ListingId, idempotencyKey); + HostJoinAttempt host = fixture.Service.BrowseForHost( + registration.ListingId, + ContractLimits.ContractVersion, + registration.LeaseToken, + ContractLimits.BrowserPageMaxItems, + null).Value!.Items.Single(item => item.AttemptId == created.AttemptId); + return new( + created.AttemptId, + created.MediationHandle, + host.HostPunchCapability, + created.ClientPunchCapability); + } + + private static NatMediationResult Process( + NatMediationProcessor processor, + INatIntroductionSink sink, + AttemptCredentials attempt, + NatPunchPeerRole role, + IPEndPoint local, + IPEndPoint observed) => processor.ProcessRequest( + local, + observed, + NatPunchRequestTokenCodec.Encode( + role, + attempt.Handle, + role == NatPunchPeerRole.Client + ? attempt.ClientCapability + : attempt.HostCapability), + sink); + + private static IPEndPoint Endpoint(string address, int port) => + new(IPAddress.Parse(address), port); + + private sealed record AttemptCredentials( + JoinAttemptId AttemptId, + MediationHandle Handle, + string HostCapability, + string ClientCapability); + + private sealed class CaptureIntroductionSink : INatIntroductionSink + { + public List Plans { get; } = []; + + public void Introduce(NatIntroductionPlan plan) => Plans.Add(plan); + } + + private sealed class ConcurrentIntroductionSink : INatIntroductionSink + { + public ConcurrentBag Plans { get; } = []; + + public void Introduce(NatIntroductionPlan plan) => Plans.Add(plan); + } + + private sealed class BlockingIntroductionSink : INatIntroductionSink, IDisposable + { + private readonly ManualResetEventSlim _entered = new(); + private readonly ManualResetEventSlim _release = new(); + + public void Introduce(NatIntroductionPlan plan) + { + _entered.Set(); + _release.Wait(TimeSpan.FromSeconds(2)); + } + + public bool WaitUntilEntered(TimeSpan timeout) => _entered.Wait(timeout); + public void Release() => _release.Set(); + + public void Dispose() + { + _entered.Dispose(); + _release.Dispose(); + } + } +} diff --git a/tests/FinalFactory.Rendezvous.Tests/Server/UdpMediatorServiceTests.cs b/tests/FinalFactory.Rendezvous.Tests/Server/UdpMediatorServiceTests.cs index 28ad77c..8393850 100644 --- a/tests/FinalFactory.Rendezvous.Tests/Server/UdpMediatorServiceTests.cs +++ b/tests/FinalFactory.Rendezvous.Tests/Server/UdpMediatorServiceTests.cs @@ -1,10 +1,9 @@ using System.Net; +using System.Net.Sockets; using FinalFactory.Rendezvous.Contracts; -using FinalFactory.Rendezvous.Server.Sessions; -using FinalFactory.Rendezvous.Server.State; using FinalFactory.Rendezvous.Server.Transport; -using FinalFactory.Rendezvous.Tests.Sessions; -using FinalFactory.Rendezvous.Tests.State; +using FinalFactory.Rendezvous.Tests.JoinAttempts; +using LiteNetLib; using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; @@ -13,55 +12,24 @@ namespace FinalFactory.Rendezvous.Tests.Server; public sealed class UdpMediatorServiceTests { [Fact] - public void AuthenticatedHostDatagramGatesVisibilityUsingObservedGameplaySocket() - { - using SessionLeaseFixture fixture = new(); - RegisterSessionResponse registration = fixture.Register(); - using UdpMediatorService service = new( - Options.Create(new UdpMediatorOptions { ListenAddress = "127.0.0.1", Port = 0 }), - NullLogger.Instance, - fixture.Store, - fixture.Capabilities); - PresenceDatagram presence = new() - { - MessageType = UdpPresenceMessageType.HostPresence, - MediationHandle = registration.HostPresenceHandle, - AddressFamily = AddressFamilyKind.Ipv4, - LocalAddress = "192.168.1.50", - LocalPort = 40_000, - Capability = registration.HostPresenceCapability, - }; - IPEndPoint observedGameplaySocket = new(IPAddress.Parse("203.0.113.77"), 51_234); - presence.Capability = "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"; - Assert.Equal( - UdpPresenceProcessingResult.HostPresenceRejected, - service.ProcessDatagram(RendezvousUdpCodec.Encode(presence), observedGameplaySocket)); - Assert.Empty(fixture.Browse()); - - presence.Capability = registration.HostPresenceCapability; - Assert.Equal( - UdpPresenceProcessingResult.HostPresenceAccepted, - service.ProcessDatagram(RendezvousUdpCodec.Encode(presence), observedGameplaySocket)); - Assert.Equal(registration.ListingId, Assert.Single(fixture.Browse()).Definition.ListingId); - } - - [Fact] - public async Task ServiceBindsAnEphemeralUdpPortAndStopsCleanly() + public async Task ServiceBindsAnEphemeralLiteNetLibPortAndStopsCleanly() { using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(5)); - UdpMediatorOptions options = new() - { - ListenAddress = IPAddress.Loopback.ToString(), - Port = 0, - }; - ManualRendezvousClock clock = new(); - InMemoryEphemeralRendezvousStore store = new(new EphemeralStoreOptions(), clock, clock); - using EphemeralCapabilityIssuer capabilities = new(); + using JoinAttemptFixture fixture = new(); + NatMediationProcessor processor = new( + fixture.Sessions.Store, + fixture.Sessions.Capabilities, + fixture.Service); using UdpMediatorService service = new( - Options.Create(options), + Options.Create(new UdpMediatorOptions + { + ListenAddress = IPAddress.Loopback.ToString(), + Port = 0, + MaxDatagramsPerPoll = 8, + PollIntervalMilliseconds = 1, + }), NullLogger.Instance, - store, - capabilities); + processor); await service.StartAsync(timeout.Token); @@ -69,9 +37,300 @@ public sealed class UdpMediatorServiceTests Assert.NotNull(boundEndpoint); Assert.Equal(IPAddress.Loopback, boundEndpoint.Address); Assert.InRange(boundEndpoint.Port, 1, 65_535); + Assert.Null(service.LocalIpv6Endpoint); await service.StopAsync(timeout.Token); Assert.Null(service.LocalEndpoint); } + + [Fact] + public async Task OptionalIpv6BindingNeverWidensTheRequiredIpv4Binding() + { + if (!Socket.OSSupportsIPv6) + { + return; + } + + using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(5)); + using JoinAttemptFixture fixture = new(); + NatMediationProcessor processor = new( + fixture.Sessions.Store, + fixture.Sessions.Capabilities, + fixture.Service); + using UdpMediatorService service = new( + Options.Create(new UdpMediatorOptions + { + ListenAddress = IPAddress.Loopback.ToString(), + Ipv6ListenAddress = IPAddress.IPv6Loopback.ToString(), + Port = 0, + }), + NullLogger.Instance, + processor); + + await service.StartAsync(timeout.Token); + + Assert.Equal(IPAddress.Loopback, service.LocalEndpoint!.Address); + Assert.Equal(IPAddress.IPv6Loopback, service.LocalIpv6Endpoint!.Address); + Assert.Equal(service.LocalEndpoint.Port, service.LocalIpv6Endpoint.Port); + IPAddress? otherIpv4 = Dns.GetHostAddresses(Dns.GetHostName()) + .FirstOrDefault(address => + address.AddressFamily == AddressFamily.InterNetwork + && !IPAddress.IsLoopback(address)); + if (otherIpv4 is not null) + { + using UdpClient scopeProbe = new(new IPEndPoint(otherIpv4, service.LocalEndpoint.Port)); + Assert.Equal(otherIpv4, ((IPEndPoint)scopeProbe.Client.LocalEndPoint!).Address); + } + + await service.StopAsync(timeout.Token); + } + + [Fact] + public async Task NativeLiteNetLibRequestsIntroduceTheAuthorizedPair() + { + using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(5)); + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + CreateJoinAttemptResponse created = fixture.Create(registration.ListingId, "native-litenet"); + HostJoinAttempt hostAttempt = fixture.Service.BrowseForHost( + registration.ListingId, + ContractLimits.ContractVersion, + registration.LeaseToken, + ContractLimits.BrowserPageMaxItems, + null).Value!.Items.Single(item => item.AttemptId == created.AttemptId); + NatMediationProcessor processor = new( + fixture.Sessions.Store, + fixture.Sessions.Capabilities, + fixture.Service); + using UdpMediatorService service = new( + Options.Create(new UdpMediatorOptions + { + ListenAddress = IPAddress.Loopback.ToString(), + Port = 0, + MaxDatagramsPerPoll = 8, + PollIntervalMilliseconds = 1, + }), + NullLogger.Instance, + processor); + await service.StartAsync(timeout.Token); + + EventBasedNetListener hostListener = new(); + EventBasedNetListener clientListener = new(); + NetManager host = new(hostListener) { NatPunchEnabled = true }; + NetManager client = new(clientListener) { NatPunchEnabled = true }; + EventBasedNatPunchListener hostPunch = new(); + EventBasedNatPunchListener clientPunch = new(); + List hostTickets = []; + List clientTickets = []; + hostPunch.NatIntroductionSuccess += (_, _, ticket) => hostTickets.Add(ticket); + clientPunch.NatIntroductionSuccess += (_, _, ticket) => clientTickets.Add(ticket); + host.NatPunchModule.Init(hostPunch); + client.NatPunchModule.Init(clientPunch); + + try + { + Assert.True(host.Start(0)); + Assert.True(client.Start(0)); + IPEndPoint mediator = Assert.IsType(service.LocalEndpoint); + host.NatPunchModule.SendNatIntroduceRequest( + mediator, + NatPunchRequestTokenCodec.Encode( + NatPunchPeerRole.Host, + created.MediationHandle, + hostAttempt.HostPunchCapability)); + client.NatPunchModule.SendNatIntroduceRequest( + mediator, + NatPunchRequestTokenCodec.Encode( + NatPunchPeerRole.Client, + created.MediationHandle, + created.ClientPunchCapability)); + + while ((hostTickets.Count == 0 || clientTickets.Count == 0) + && !timeout.IsCancellationRequested) + { + host.PollEvents(); + host.NatPunchModule.PollEvents(); + client.PollEvents(); + client.NatPunchModule.PollEvents(); + await Task.Delay(5, timeout.Token); + } + + string hostTicket = Assert.Single(hostTickets.Distinct(StringComparer.Ordinal)); + string clientTicket = Assert.Single(clientTickets.Distinct(StringComparer.Ordinal)); + Assert.Equal(hostTicket, clientTicket); + Assert.Equal(43, hostTicket.Length); + } + finally + { + host.Stop(); + client.Stop(); + await service.StopAsync(CancellationToken.None); + } + } + + [Fact] + public async Task FrozenV1EnvelopeIsConsumedOnTheLiteNetSocketWithinAmplificationBudget() + { + using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(5)); + using JoinAttemptFixture fixture = new(); + (RegisterSessionResponse registration, _) = fixture.CreateHost(); + CreateJoinAttemptResponse created = fixture.Create(registration.ListingId, "v1-envelope"); + HostJoinAttempt hostAttempt = fixture.Service.BrowseForHost( + registration.ListingId, + ContractLimits.ContractVersion, + registration.LeaseToken, + ContractLimits.BrowserPageMaxItems, + null).Value!.Items.Single(item => item.AttemptId == created.AttemptId); + NatMediationProcessor processor = new( + fixture.Sessions.Store, + fixture.Sessions.Capabilities, + fixture.Service); + using UdpMediatorService service = new( + Options.Create(new UdpMediatorOptions + { + ListenAddress = IPAddress.Loopback.ToString(), + Port = 0, + MaxDatagramsPerPoll = 8, + PollIntervalMilliseconds = 1, + }), + NullLogger.Instance, + processor); + await service.StartAsync(timeout.Token); + using UdpClient host = new(new IPEndPoint(IPAddress.Loopback, 0)); + using UdpClient client = new(new IPEndPoint(IPAddress.Loopback, 0)); + IPEndPoint mediator = Assert.IsType(service.LocalEndpoint); + byte[] hostDatagram = RendezvousUdpCodec.Encode(new PresenceDatagram + { + MessageType = UdpPresenceMessageType.HostPresence, + MediationHandle = created.MediationHandle, + AddressFamily = AddressFamilyKind.Ipv4, + LocalAddress = "192.168.1.10", + LocalPort = 41_000, + Capability = hostAttempt.HostPunchCapability, + }); + byte[] clientDatagram = RendezvousUdpCodec.Encode(new PresenceDatagram + { + MessageType = UdpPresenceMessageType.ClientPresence, + MediationHandle = created.MediationHandle, + AddressFamily = AddressFamilyKind.Ipv4, + LocalAddress = "192.168.1.11", + LocalPort = 42_000, + Capability = created.ClientPunchCapability, + }); + + try + { + await host.SendAsync(hostDatagram, mediator, timeout.Token); + await client.SendAsync(clientDatagram, mediator, timeout.Token); + UdpReceiveResult hostIntroduction = await host.ReceiveAsync(timeout.Token); + UdpReceiveResult clientIntroduction = await client.ReceiveAsync(timeout.Token); + + Assert.True( + hostIntroduction.Buffer.Length + clientIntroduction.Buffer.Length + <= clientDatagram.Length * 2, + "The completing authenticated contribution exceeded the 2.0 response-byte budget."); + } + finally + { + await service.StopAsync(CancellationToken.None); + } + } + + [Fact] + public async Task OversizedMalformedAndGameplayDatagramsReceiveNoResponse() + { + using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(5)); + using JoinAttemptFixture fixture = new(); + NatMediationProcessor processor = new( + fixture.Sessions.Store, + fixture.Sessions.Capabilities, + fixture.Service); + using UdpMediatorService service = new( + Options.Create(new UdpMediatorOptions + { + ListenAddress = IPAddress.Loopback.ToString(), + Port = 0, + MaxDatagramsPerPoll = 8, + PollIntervalMilliseconds = 1, + }), + NullLogger.Instance, + processor); + await service.StartAsync(timeout.Token); + using UdpClient sender = new(new IPEndPoint(IPAddress.Loopback, 0)); + IPEndPoint mediator = Assert.IsType(service.LocalEndpoint); + byte[] oversized = new byte[ContractLimits.UdpDatagramMaxBytes + 1]; + oversized[0] = RendezvousUdpCodec.MagicFirst; + oversized[1] = RendezvousUdpCodec.MagicSecond; + byte[] gameplayPayload = [0x01, 0x02, 0x03, 0x04]; + byte[] malformedNative = [17, 0]; + + try + { + await sender.SendAsync(oversized, mediator, timeout.Token); + await sender.SendAsync(gameplayPayload, mediator, timeout.Token); + await sender.SendAsync(malformedNative, mediator, timeout.Token); + using CancellationTokenSource noResponse = new(TimeSpan.FromMilliseconds(150)); + await Assert.ThrowsAnyAsync(async () => + await sender.ReceiveAsync(noResponse.Token)); + } + finally + { + await service.StopAsync(CancellationToken.None); + } + } + + [Fact] + public async Task ForgedNativeIntroductionResponseCannotReflectToPayloadEndpoint() + { + using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(5)); + using JoinAttemptFixture fixture = new(); + NatMediationProcessor processor = new( + fixture.Sessions.Store, + fixture.Sessions.Capabilities, + fixture.Service); + using UdpMediatorService service = new( + Options.Create(new UdpMediatorOptions + { + ListenAddress = IPAddress.Loopback.ToString(), + Port = 0, + MaxDatagramsPerPoll = 8, + PollIntervalMilliseconds = 1, + }), + NullLogger.Instance, + processor); + await service.StartAsync(timeout.Token); + using UdpClient reflectedTarget = new(new IPEndPoint(IPAddress.Loopback, 0)); + using UdpClient responseCapture = new(new IPEndPoint(IPAddress.Loopback, 0)); + using UdpClient attacker = new(new IPEndPoint(IPAddress.Loopback, 0)); + LiteNetManager generator = new(new EventBasedLiteNetListener()) { NatPunchEnabled = true }; + + try + { + Assert.True(generator.Start(0)); + IPEndPoint target = (IPEndPoint)reflectedTarget.Client.LocalEndPoint!; + IPEndPoint capture = (IPEndPoint)responseCapture.Client.LocalEndPoint!; + generator.NatPunchModule.NatIntroduce( + target, + new IPEndPoint(IPAddress.Loopback, 9), + capture, + capture, + "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"); + byte[] forgedResponse = (await responseCapture.ReceiveAsync(timeout.Token)).Buffer; + + await attacker.SendAsync( + forgedResponse, + Assert.IsType(service.LocalEndpoint), + timeout.Token); + + using CancellationTokenSource noReflection = new(TimeSpan.FromMilliseconds(150)); + await Assert.ThrowsAnyAsync(async () => + await reflectedTarget.ReceiveAsync(noReflection.Token)); + } + finally + { + generator.Stop(); + await service.StopAsync(CancellationToken.None); + } + } } diff --git a/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/contracts-public-api.txt b/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/contracts-public-api.txt index 1642b76..eca1d73 100644 --- a/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/contracts-public-api.txt +++ b/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/contracts-public-api.txt @@ -49,6 +49,7 @@ TYPE FinalFactory.Rendezvous.Contracts.ContractLimits FIELD System.Int32 ConnectionTicketMaxCharacters=192 FIELD System.Int32 ContractVersion=1 FIELD System.Int32 CursorMaxCharacters=512 + FIELD System.Int32 DerivedCredentialCharacters=43 FIELD System.Int32 DiagnosticCodeMaxCharacters=64 FIELD System.Int32 DisplayNameMaxBytes=128 FIELD System.Int32 EnvironmentIdMaxCharacters=32 @@ -61,6 +62,7 @@ TYPE FinalFactory.Rendezvous.Contracts.ContractLimits FIELD System.Int32 MetadataMaxBytes=4096 FIELD System.Int32 MetadataMaxKeys=32 FIELD System.Int32 MetadataValueMaxBytes=256 + FIELD System.Int32 NatPunchRequestTokenCharacters=192 FIELD System.Int32 OpaqueHttpCredentialMaxCharacters=1024 FIELD System.Int32 RegionIdMaxCharacters=32 FIELD System.Int32 SessionCapacityMaxPlayers=10000 @@ -171,6 +173,20 @@ TYPE FinalFactory.Rendezvous.Contracts.MediationHandle METHOD System.Boolean TryParse(System.String value, FinalFactory.Rendezvous.Contracts.MediationHandle& id) METHOD System.Boolean op_Equality(FinalFactory.Rendezvous.Contracts.MediationHandle left, FinalFactory.Rendezvous.Contracts.MediationHandle right) METHOD System.Boolean op_Inequality(FinalFactory.Rendezvous.Contracts.MediationHandle left, FinalFactory.Rendezvous.Contracts.MediationHandle right) +TYPE FinalFactory.Rendezvous.Contracts.NatPunchPeerRole + ENUM HostPresence=1 + ENUM Host=2 + ENUM Client=3 +TYPE FinalFactory.Rendezvous.Contracts.NatPunchRequestToken + CTOR () + PROP System.String Capability {get;set;} + PROP FinalFactory.Rendezvous.Contracts.MediationHandle MediationHandle {get;set;} + PROP FinalFactory.Rendezvous.Contracts.NatPunchPeerRole Role {get;set;} + METHOD System.String ToString() +TYPE FinalFactory.Rendezvous.Contracts.NatPunchRequestTokenCodec + FIELD System.Int32 EncodedLength=192 + METHOD System.String Encode(FinalFactory.Rendezvous.Contracts.NatPunchPeerRole role, FinalFactory.Rendezvous.Contracts.MediationHandle mediationHandle, System.String capability) + METHOD System.Boolean TryDecode(System.String encoded, FinalFactory.Rendezvous.Contracts.NatPunchRequestToken& token) TYPE FinalFactory.Rendezvous.Contracts.NetworkEndpoint CTOR () PROP System.String Address {get;set;}