Compare commits

...

1 Commits

Author SHA1 Message Date
KyuubiYoru 6d076c281a feat: implement authenticated NAT mediator (#11)
quality-gate / quality (push) Successful in 59s
Closes #11
2026-07-16 07:37:02 +02:00
19 changed files with 1660 additions and 174 deletions
@@ -60,12 +60,14 @@ the supplied ID.
### UDP presence ### UDP presence
Only a structurally valid `HostPresence` datagram with the issued capability can Only a structurally valid frozen `HostPresence` envelope or native LiteNetLib
refresh presence. The public endpoint is the UDP packet's observed source on the host-presence request with the issued capability can refresh presence. The public
host's gameplay socket; the HTTP API never accepts one. The bounded local candidate endpoint is the UDP packet's observed source on the host's gameplay socket; the
comes from the authenticated datagram. Invalid, unknown, or client-presence packets HTTP API never accepts one. The bounded local candidate comes from the authenticated
receive no response. Presence expiry demotes public visibility but keeps the lease, packet. Invalid or unknown inputs receive no response. ADR 0009 defines the later
so the same handle can restore visibility without changing session identity. 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 listing responses contain bounded listing data only. They never contain
public/local endpoints, lease tokens, presence capabilities, fingerprints, store public/local endpoints, lease tokens, presence capabilities, fingerprints, store
@@ -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.
+1
View File
@@ -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 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 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 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) - [Threat model](../security/threat-model.md)
- [Security promise and test matrix](../security/control-matrix.md) - [Security promise and test matrix](../security/control-matrix.md)
- [Versioned HTTP and UDP contracts](../contracts/README.md) - [Versioned HTTP and UDP contracts](../contracts/README.md)
+51 -8
View File
@@ -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 The UDP mediator accepts the frozen bounded presence envelope below and native
client. It associates the authenticated mediation handle with the packet's LiteNetLib NAT-introduction requests. Both forms associate an authenticated
observed public source endpoint and the sender's reported local endpoint. It mediation handle with the packet's observed public source endpoint and the
does not carry gameplay packets. sender's reported local endpoint. Neither form carries gameplay packets.
All multi-byte integers use network byte order. UUID bytes use the canonical 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 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 and compared without exposing them in logs. A valid-looking packet does not
prove authorization until the capability is checked. Invalid packets receive prove authorization until the capability is checked. Invalid packets receive
no UDP response, preventing the mediator from becoming an amplification oracle. no UDP response, preventing the mediator from becoming an amplification oracle.
Replay, expiry, pairing, and rate-limit policy are defined by later mediator For the frozen envelope, `HostPresence` is resolved against either the listing's
issues; the v1 envelope deliberately leaves no unbounded or reflected payload. 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:<role>:<32 lowercase handle hex>:<43-character capability><dot padding>
```
`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.
@@ -42,6 +42,21 @@ if (!registered.IsSuccess || registered.Value is null)
Load `publisherCredential` from the game's deployment secret boundary; never Load `publisherCredential` from the game's deployment secret boundary; never
embed it in a client build or source control. A successful registration returns a embed it in a client build or source control. A successful registration returns a
`PublishedSession` containing the lease and host-presence capabilities. `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: Lease renewal is explicit and caller-controlled:
@@ -23,6 +23,8 @@ public static class ContractLimits
public const int OpaqueHttpCredentialMaxCharacters = 1_024; public const int OpaqueHttpCredentialMaxCharacters = 1_024;
public const int UdpCapabilityMaxCharacters = 192; public const int UdpCapabilityMaxCharacters = 192;
public const int ConnectionTicketMaxCharacters = 192; public const int ConnectionTicketMaxCharacters = 192;
public const int DerivedCredentialCharacters = 43;
public const int NatPunchRequestTokenCharacters = 192;
public const int LiteNetLibNatTokenMaxCharacters = 256; public const int LiteNetLibNatTokenMaxCharacters = 256;
public const int SessionCapacityMaxPlayers = 10_000; public const int SessionCapacityMaxPlayers = 10_000;
} }
@@ -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;
}
}
@@ -316,6 +316,9 @@ internal sealed class JoinAttemptService(
ContractValidation.IsCapabilityValid(hostCapability) ContractValidation.IsCapabilityValid(hostCapability)
&& ContractValidation.IsCapabilityValid(clientCapability) && ContractValidation.IsCapabilityValid(clientCapability)
&& ContractValidation.IsConnectionTicketValid(ticket) && ContractValidation.IsConnectionTicketValid(ticket)
&& hostCapability.Length == ContractLimits.DerivedCredentialCharacters
&& clientCapability.Length == ContractLimits.DerivedCredentialCharacters
&& ticket.Length == ContractLimits.DerivedCredentialCharacters
&& hostCapability.Length <= ContractLimits.LiteNetLibNatTokenMaxCharacters && hostCapability.Length <= ContractLimits.LiteNetLibNatTokenMaxCharacters
&& clientCapability.Length <= ContractLimits.LiteNetLibNatTokenMaxCharacters; && clientCapability.Length <= ContractLimits.LiteNetLibNatTokenMaxCharacters;
@@ -133,12 +133,19 @@ builder.Services
.BindConfiguration(UdpMediatorOptions.SectionName) .BindConfiguration(UdpMediatorOptions.SectionName)
.ValidateDataAnnotations() .ValidateDataAnnotations()
.Validate( .Validate(
options => IPAddress.TryParse(options.ListenAddress, out _), options => IPAddress.TryParse(options.ListenAddress, out IPAddress? address)
$"{UdpMediatorOptions.SectionName}:ListenAddress must be an IP 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(); .ValidateOnStart();
builder.Services.AddSingleton<UdpMediatorService>(); builder.Services.AddSingleton<UdpMediatorService>();
if (!isOpenApiGeneration) if (!isOpenApiGeneration)
{ {
builder.Services.AddSingleton<NatMediationProcessor>();
builder.Services.AddHostedService(static services => builder.Services.AddHostedService(static services =>
services.GetRequiredService<UdpMediatorService>()); services.GetRequiredService<UdpMediatorService>());
} }
@@ -4,6 +4,7 @@ namespace FinalFactory.Rendezvous.Server.State;
internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousStore internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousStore
{ {
private static readonly TimeSpan UdpMaintenanceInterval = TimeSpan.FromSeconds(1);
private readonly object _gate = new(); private readonly object _gate = new();
private readonly EphemeralStoreOptions _options; private readonly EphemeralStoreOptions _options;
private readonly IMonotonicClock _monotonicClock; private readonly IMonotonicClock _monotonicClock;
@@ -19,6 +20,8 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
private readonly Dictionary<string, TimeSpan> _replay = new(StringComparer.Ordinal); private readonly Dictionary<string, TimeSpan> _replay = new(StringComparer.Ordinal);
private readonly Dictionary<string, TimeSpan> _revocations = new(StringComparer.Ordinal); private readonly Dictionary<string, TimeSpan> _revocations = new(StringComparer.Ordinal);
private TimeSpan? _drainDeadline; private TimeSpan? _drainDeadline;
private TimeSpan _nextUdpMaintenance;
private long _maintenanceSweepCount;
private bool _available = true; private bool _available = true;
public InMemoryEphemeralRendezvousStore( public InMemoryEphemeralRendezvousStore(
@@ -38,6 +41,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
} }
public Guid InstanceId { get; } public Guid InstanceId { get; }
internal long MaintenanceSweepCount => Interlocked.Read(ref _maintenanceSweepCount);
public bool IsAvailable public bool IsAvailable
{ {
@@ -270,6 +274,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
if (!_presenceHandles.TryGetValue(command.Handle, out SessionListingId listingId) if (!_presenceHandles.TryGetValue(command.Handle, out SessionListingId listingId)
|| !_listings.TryGetValue(listingId, out ListingEntry? entry) || !_listings.TryGetValue(listingId, out ListingEntry? entry)
|| entry.LeaseDeadline <= now
|| entry.Definition.HostPresenceFingerprint != command.CapabilityFingerprint) || entry.Definition.HostPresenceFingerprint != command.CapabilityFingerprint)
{ {
return new(StoreResultCode.NotFound); return new(StoreResultCode.NotFound);
@@ -285,7 +290,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
command.LocalEndpoint, command.LocalEndpoint,
now + _options.PresenceLifetime); now + _options.PresenceLifetime);
return new(StoreResultCode.Success, Snapshot(entry)); return new(StoreResultCode.Success, Snapshot(entry));
}, cancellationToken); }, cancellationToken, eagerCleanup: false);
public StoreResult<IReadOnlyList<StoredListing>> BrowseVisibleListings( public StoreResult<IReadOnlyList<StoredListing>> BrowseVisibleListings(
VisibleListingQuery query, VisibleListingQuery query,
@@ -446,13 +451,18 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
return new(StoreResultCode.NotFound); return new(StoreResultCode.NotFound);
} }
if (attempt.IntroductionConsumed)
{
return new(StoreResultCode.Conflict);
}
RemoveAttempt(command.AttemptId); RemoveAttempt(command.AttemptId);
return new(StoreResultCode.Success, true); return new(StoreResultCode.Success, true);
}, cancellationToken); }, cancellationToken);
public StoreResult<StoredJoinAttempt> BindAttemptEndpoint( public StoreResult<StoredJoinAttempt> BindAttemptEndpoint(
BindAttemptEndpointCommand command, BindAttemptEndpointCommand command,
CancellationToken cancellationToken = default) => Atomic<StoredJoinAttempt>(_ => CancellationToken cancellationToken = default) => Atomic<StoredJoinAttempt>(now =>
{ {
ArgumentNullException.ThrowIfNull(command); ArgumentNullException.ThrowIfNull(command);
if (command.Handle.Value == Guid.Empty if (command.Handle.Value == Guid.Empty
@@ -469,7 +479,13 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
} }
if (!_attemptHandles.TryGetValue(command.Handle, out JoinAttemptId attemptId) 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); return new(StoreResultCode.NotFound);
} }
@@ -503,7 +519,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
} }
return new(StoreResultCode.Success, Snapshot(attempt)); return new(StoreResultCode.Success, Snapshot(attempt));
}, cancellationToken); }, cancellationToken, eagerCleanup: false);
public StoreResult<IntroductionEndpoints> ConsumeIntroduction( public StoreResult<IntroductionEndpoints> ConsumeIntroduction(
MediationHandle handle, MediationHandle handle,
@@ -515,7 +531,13 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
} }
if (!_attemptHandles.TryGetValue(handle, out JoinAttemptId attemptId) 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); return new(StoreResultCode.NotFound);
} }
@@ -540,7 +562,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
Snapshot(attempt), Snapshot(attempt),
attempt.HostEndpoint, attempt.HostEndpoint,
attempt.ClientEndpoint)); attempt.ClientEndpoint));
}, cancellationToken); }, cancellationToken, eagerCleanup: false);
public StoreResult<bool> ConsumeConnectionTicket( public StoreResult<bool> ConsumeConnectionTicket(
ConsumeConnectionTicketCommand command, ConsumeConnectionTicketCommand command,
@@ -687,18 +709,38 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
} }
} }
private StoreResult<T> Atomic<T>(Func<TimeSpan, StoreResult<T>> operation, CancellationToken cancellationToken) private StoreResult<T> Atomic<T>(
Func<TimeSpan, StoreResult<T>> operation,
CancellationToken cancellationToken,
bool eagerCleanup = true)
{ {
cancellationToken.ThrowIfCancellationRequested(); cancellationToken.ThrowIfCancellationRequested();
lock (_gate) lock (_gate)
{ {
cancellationToken.ThrowIfCancellationRequested(); cancellationToken.ThrowIfCancellationRequested();
TimeSpan now = _monotonicClock.Elapsed; TimeSpan now = _monotonicClock.Elapsed;
// 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); Cleanup(now);
_nextUdpMaintenance = now + UdpMaintenanceInterval;
}
return operation(now); 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<T>? CheckNewWorkAdmission<T>(string subject) private StoreResult<T>? CheckNewWorkAdmission<T>(string subject)
{ {
if (!_available) if (!_available)
@@ -718,6 +760,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
private void Cleanup(TimeSpan now) private void Cleanup(TimeSpan now)
{ {
_maintenanceSweepCount++;
if (_drainDeadline is TimeSpan drainDeadline && now >= drainDeadline) if (_drainDeadline is TimeSpan drainDeadline && now >= drainDeadline)
{ {
ClearActiveState(); ClearActiveState();
@@ -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<byte> RequestTypeIdentifier =>
[0x88, 0xbe, 0x10, 0x26, 0xbf, 0xb1, 0x66, 0x9c];
public static bool TryDecode(
ReadOnlySpan<byte> 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<byte> 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;
}
}
@@ -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<byte> 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<StoredListing> 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<StoredJoinAttempt> 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<IntroductionEndpoints> consumed = store.ConsumeIntroduction(
request.MediationHandle,
cancellationToken);
if (!consumed.Succeeded || consumed.Value is null)
{
return consumed.Code == StoreResultCode.ReplayRejected
? NatMediationResult.Duplicate
: NatMediationResult.Rejected;
}
JoinAttemptServiceResult<ConnectionTicketGrant> 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);
}
@@ -18,9 +18,17 @@ public sealed class UdpMediatorOptions
[Required] [Required]
public string ListenAddress { get; set; } = "0.0.0.0"; public string ListenAddress { get; set; } = "0.0.0.0";
public string? Ipv6ListenAddress { get; set; }
/// <summary> /// <summary>
/// Gets or sets the UDP port. Zero requests an ephemeral port for tests. /// Gets or sets the UDP port. Zero requests an ephemeral port for tests.
/// </summary> /// </summary>
[Range(0, 65_535)] [Range(0, 65_535)]
public int Port { get; set; } = 9050; 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;
} }
@@ -1,161 +1,136 @@
using System.Diagnostics;
using System.Net; using System.Net;
using System.Net.Sockets; using System.Net.Sockets;
using FinalFactory.Rendezvous.Contracts; using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Sessions; using LiteNetLib;
using FinalFactory.Rendezvous.Server.State; using LiteNetLib.Layers;
using Microsoft.Extensions.Options; using Microsoft.Extensions.Options;
namespace FinalFactory.Rendezvous.Server.Transport; namespace FinalFactory.Rendezvous.Server.Transport;
/// <summary>
/// Owns the cancellable UDP socket used by the future NAT mediator.
/// </summary>
internal sealed partial class UdpMediatorService : BackgroundService internal sealed partial class UdpMediatorService : BackgroundService
{ {
private readonly ILogger<UdpMediatorService> _logger; private readonly ILogger<UdpMediatorService> _logger;
private readonly UdpMediatorOptions _options; private readonly UdpMediatorOptions _options;
private readonly IEphemeralRendezvousStore _store; private readonly NatMediationProcessor _processor;
private readonly ISessionCapabilityService _capabilities; private LiteNetManager? _manager;
private UdpClient? _udpClient; private LiteNetIntroductionSink? _introductionSink;
/// <summary>
/// Initializes a new UDP mediator service.
/// </summary>
public UdpMediatorService( public UdpMediatorService(
IOptions<UdpMediatorOptions> options, IOptions<UdpMediatorOptions> options,
ILogger<UdpMediatorService> logger, ILogger<UdpMediatorService> logger,
IEphemeralRendezvousStore store, NatMediationProcessor processor)
ISessionCapabilityService capabilities)
{ {
_options = options.Value; _options = options.Value;
_logger = logger; _logger = logger;
_store = store; _processor = processor;
_capabilities = capabilities;
} }
/// <summary>
/// Gets the bound endpoint after startup completes.
/// </summary>
public IPEndPoint? LocalEndpoint { get; private set; } public IPEndPoint? LocalEndpoint { get; private set; }
public IPEndPoint? LocalIpv6Endpoint { get; private set; }
/// <inheritdoc />
public override Task StartAsync(CancellationToken cancellationToken) public override Task StartAsync(CancellationToken cancellationToken)
{ {
cancellationToken.ThrowIfCancellationRequested(); cancellationToken.ThrowIfCancellationRequested();
if (_manager is not null)
if (_udpClient is not null)
{ {
throw new InvalidOperationException("The UDP mediator is already running."); throw new InvalidOperationException("The UDP mediator is already running.");
} }
IPAddress listenAddress = IPAddress.Parse(_options.ListenAddress); IPAddress listenAddress = IPAddress.Parse(_options.ListenAddress);
UdpClient udpClient = new(new IPEndPoint(listenAddress, _options.Port)); if (listenAddress.AddressFamily != AddressFamily.InterNetwork)
_udpClient = udpClient; {
IPEndPoint localEndpoint = throw new InvalidOperationException("The required UDP listen address must be IPv4.");
(IPEndPoint?)udpClient.Client.LocalEndPoint }
?? throw new InvalidOperationException("The UDP socket did not expose its bound endpoint.");
LocalEndpoint = localEndpoint;
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); return base.StartAsync(cancellationToken);
} }
/// <inheritdoc />
public override async Task StopAsync(CancellationToken cancellationToken) public override async Task StopAsync(CancellationToken cancellationToken)
{ {
await base.StopAsync(cancellationToken).ConfigureAwait(false); await base.StopAsync(cancellationToken).ConfigureAwait(false);
_udpClient?.Dispose(); StopManager();
_udpClient = null;
LocalEndpoint = null;
LogMediatorStopped(_logger); LogMediatorStopped(_logger);
} }
/// <inheritdoc />
public override void Dispose() public override void Dispose()
{ {
_udpClient?.Dispose(); StopManager();
_udpClient = null;
LocalEndpoint = null;
base.Dispose(); base.Dispose();
} }
/// <inheritdoc />
protected override async Task ExecuteAsync(CancellationToken stoppingToken) protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{ {
UdpClient udpClient = _udpClient LiteNetManager manager = _manager
?? throw new InvalidOperationException("The UDP mediator socket was not initialized."); ?? throw new InvalidOperationException("The UDP mediator socket was not initialized.");
long previous = Stopwatch.GetTimestamp();
try try
{ {
while (!stoppingToken.IsCancellationRequested) while (!stoppingToken.IsCancellationRequested)
{ {
UdpReceiveResult received = await udpClient manager.PollEvents();
.ReceiveAsync(stoppingToken) manager.NatPunchModule.PollEvents();
.ConfigureAwait(false); long current = Stopwatch.GetTimestamp();
ProcessDatagram(received.Buffer, received.RemoteEndPoint, stoppingToken); manager.ManualUpdate((float)Stopwatch.GetElapsedTime(previous, current).TotalMilliseconds);
// Bootstrap deliberately emits no UDP response. Protocol handling lands in #11. previous = current;
await Task.Delay(_options.PollIntervalMilliseconds, stoppingToken).ConfigureAwait(false);
} }
} }
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) 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 finally
{ {
LocalEndpoint = null; LocalEndpoint = null;
LocalIpv6Endpoint = null;
} }
} }
internal UdpPresenceProcessingResult ProcessDatagram( private void StopManager()
ReadOnlySpan<byte> encoded,
IPEndPoint observedSource,
CancellationToken cancellationToken = default)
{ {
ArgumentNullException.ThrowIfNull(observedSource); LiteNetManager? manager = Interlocked.Exchange(ref _manager, null);
if (!RendezvousUdpCodec.TryDecode(encoded, out PresenceDatagram? datagram, out _) _introductionSink = null;
|| datagram is null LocalEndpoint = null;
|| !_capabilities.TryFingerprint(datagram.Capability, out SecretFingerprint fingerprint)) LocalIpv6Endpoint = null;
{ manager?.Stop();
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<StoredListing> bound = _store.BindHostPresence(new(
datagram.MediationHandle,
fingerprint,
publicEndpoint,
localEndpoint), cancellationToken);
return bound.Succeeded
? UdpPresenceProcessingResult.HostPresenceAccepted
: UdpPresenceProcessingResult.HostPresenceRejected;
} }
[LoggerMessage( [LoggerMessage(
@@ -172,12 +147,62 @@ internal sealed partial class UdpMediatorService : BackgroundService
Level = LogLevel.Information, Level = LogLevel.Information,
Message = "UDP mediator stopped")] Message = "UDP mediator stopped")]
private static partial void LogMediatorStopped(ILogger logger); private static partial void LogMediatorStopped(ILogger logger);
}
internal enum UdpPresenceProcessingResult private sealed class LiteNetIntroductionSink(NatPunchModule module) : INatIntroductionSink
{ {
Dropped = 0, public void Introduce(NatIntroductionPlan plan) => module.NatIntroduce(
HostPresenceAccepted = 1, plan.HostLocal,
HostPresenceRejected = 2, plan.HostPublic,
ClientPresenceDeferred = 3, 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;
}
} }
@@ -2,7 +2,9 @@
"Rendezvous": { "Rendezvous": {
"Udp": { "Udp": {
"ListenAddress": "0.0.0.0", "ListenAddress": "0.0.0.0",
"Port": 9050 "Port": 9050,
"MaxDatagramsPerPoll": 256,
"PollIntervalMilliseconds": 2
} }
}, },
"Logging": { "Logging": {
@@ -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<ArgumentException>(() => 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);
}
}
@@ -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<NatMediationResult>[] 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<NatMediationResult> 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<bool> 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<NatIntroductionPlan> Plans { get; } = [];
public void Introduce(NatIntroductionPlan plan) => Plans.Add(plan);
}
private sealed class ConcurrentIntroductionSink : INatIntroductionSink
{
public ConcurrentBag<NatIntroductionPlan> 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();
}
}
}
@@ -1,10 +1,9 @@
using System.Net; using System.Net;
using System.Net.Sockets;
using FinalFactory.Rendezvous.Contracts; using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Sessions;
using FinalFactory.Rendezvous.Server.State;
using FinalFactory.Rendezvous.Server.Transport; using FinalFactory.Rendezvous.Server.Transport;
using FinalFactory.Rendezvous.Tests.Sessions; using FinalFactory.Rendezvous.Tests.JoinAttempts;
using FinalFactory.Rendezvous.Tests.State; using LiteNetLib;
using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options; using Microsoft.Extensions.Options;
@@ -13,55 +12,24 @@ namespace FinalFactory.Rendezvous.Tests.Server;
public sealed class UdpMediatorServiceTests public sealed class UdpMediatorServiceTests
{ {
[Fact] [Fact]
public void AuthenticatedHostDatagramGatesVisibilityUsingObservedGameplaySocket() public async Task ServiceBindsAnEphemeralLiteNetLibPortAndStopsCleanly()
{
using SessionLeaseFixture fixture = new();
RegisterSessionResponse registration = fixture.Register();
using UdpMediatorService service = new(
Options.Create(new UdpMediatorOptions { ListenAddress = "127.0.0.1", Port = 0 }),
NullLogger<UdpMediatorService>.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()
{ {
using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(5)); using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(5));
UdpMediatorOptions options = new() 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(), ListenAddress = IPAddress.Loopback.ToString(),
Port = 0, Port = 0,
}; MaxDatagramsPerPoll = 8,
ManualRendezvousClock clock = new(); PollIntervalMilliseconds = 1,
InMemoryEphemeralRendezvousStore store = new(new EphemeralStoreOptions(), clock, clock); }),
using EphemeralCapabilityIssuer capabilities = new();
using UdpMediatorService service = new(
Options.Create(options),
NullLogger<UdpMediatorService>.Instance, NullLogger<UdpMediatorService>.Instance,
store, processor);
capabilities);
await service.StartAsync(timeout.Token); await service.StartAsync(timeout.Token);
@@ -69,9 +37,300 @@ public sealed class UdpMediatorServiceTests
Assert.NotNull(boundEndpoint); Assert.NotNull(boundEndpoint);
Assert.Equal(IPAddress.Loopback, boundEndpoint.Address); Assert.Equal(IPAddress.Loopback, boundEndpoint.Address);
Assert.InRange(boundEndpoint.Port, 1, 65_535); Assert.InRange(boundEndpoint.Port, 1, 65_535);
Assert.Null(service.LocalIpv6Endpoint);
await service.StopAsync(timeout.Token); await service.StopAsync(timeout.Token);
Assert.Null(service.LocalEndpoint); 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<UdpMediatorService>.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<UdpMediatorService>.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<string> hostTickets = [];
List<string> 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<IPEndPoint>(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<UdpMediatorService>.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<IPEndPoint>(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<UdpMediatorService>.Instance,
processor);
await service.StartAsync(timeout.Token);
using UdpClient sender = new(new IPEndPoint(IPAddress.Loopback, 0));
IPEndPoint mediator = Assert.IsType<IPEndPoint>(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<OperationCanceledException>(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<UdpMediatorService>.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<IPEndPoint>(service.LocalEndpoint),
timeout.Token);
using CancellationTokenSource noReflection = new(TimeSpan.FromMilliseconds(150));
await Assert.ThrowsAnyAsync<OperationCanceledException>(async () =>
await reflectedTarget.ReceiveAsync(noReflection.Token));
}
finally
{
generator.Stop();
await service.StopAsync(CancellationToken.None);
}
}
} }
@@ -49,6 +49,7 @@ TYPE FinalFactory.Rendezvous.Contracts.ContractLimits
FIELD System.Int32 ConnectionTicketMaxCharacters=192 FIELD System.Int32 ConnectionTicketMaxCharacters=192
FIELD System.Int32 ContractVersion=1 FIELD System.Int32 ContractVersion=1
FIELD System.Int32 CursorMaxCharacters=512 FIELD System.Int32 CursorMaxCharacters=512
FIELD System.Int32 DerivedCredentialCharacters=43
FIELD System.Int32 DiagnosticCodeMaxCharacters=64 FIELD System.Int32 DiagnosticCodeMaxCharacters=64
FIELD System.Int32 DisplayNameMaxBytes=128 FIELD System.Int32 DisplayNameMaxBytes=128
FIELD System.Int32 EnvironmentIdMaxCharacters=32 FIELD System.Int32 EnvironmentIdMaxCharacters=32
@@ -61,6 +62,7 @@ TYPE FinalFactory.Rendezvous.Contracts.ContractLimits
FIELD System.Int32 MetadataMaxBytes=4096 FIELD System.Int32 MetadataMaxBytes=4096
FIELD System.Int32 MetadataMaxKeys=32 FIELD System.Int32 MetadataMaxKeys=32
FIELD System.Int32 MetadataValueMaxBytes=256 FIELD System.Int32 MetadataValueMaxBytes=256
FIELD System.Int32 NatPunchRequestTokenCharacters=192
FIELD System.Int32 OpaqueHttpCredentialMaxCharacters=1024 FIELD System.Int32 OpaqueHttpCredentialMaxCharacters=1024
FIELD System.Int32 RegionIdMaxCharacters=32 FIELD System.Int32 RegionIdMaxCharacters=32
FIELD System.Int32 SessionCapacityMaxPlayers=10000 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 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_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) 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 TYPE FinalFactory.Rendezvous.Contracts.NetworkEndpoint
CTOR () CTOR ()
PROP System.String Address {get;set;} PROP System.String Address {get;set;}