feat(client): standardize connection outcomes (#13)
quality-gate / quality (push) Successful in 59s
quality-gate / quality (push) Successful in 59s
This commit is contained in:
@@ -0,0 +1,183 @@
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using LiteNetLib;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Client;
|
||||
|
||||
public enum RendezvousConnectionOutcomeSource
|
||||
{
|
||||
RendezvousService = 1,
|
||||
LocalTraversal = 2,
|
||||
RemoteHost = 3,
|
||||
Caller = 4,
|
||||
Lifecycle = 5,
|
||||
}
|
||||
|
||||
public enum RendezvousConnectionFailureCategory
|
||||
{
|
||||
None = 0,
|
||||
Directory = 1,
|
||||
Compatibility = 2,
|
||||
Authorization = 3,
|
||||
Capacity = 4,
|
||||
HostPresence = 5,
|
||||
Service = 6,
|
||||
Mediation = 7,
|
||||
NatTraversal = 8,
|
||||
DirectConnection = 9,
|
||||
Lifecycle = 10,
|
||||
}
|
||||
|
||||
public enum RendezvousConnectionPhase
|
||||
{
|
||||
Directory = 1,
|
||||
Authorization = 2,
|
||||
Mediation = 3,
|
||||
NatTraversal = 4,
|
||||
DirectConnection = 5,
|
||||
Complete = 6,
|
||||
}
|
||||
|
||||
public sealed class RendezvousConnectionOutcome
|
||||
{
|
||||
private readonly NetworkEndpoint? _dedicatedFallback;
|
||||
|
||||
private RendezvousConnectionOutcome(
|
||||
ConnectionOutcomeKind kind,
|
||||
RendezvousConnectionOutcomeSource source,
|
||||
RendezvousConnectionFailureCategory category,
|
||||
RendezvousConnectionPhase phase,
|
||||
TimeSpan elapsed,
|
||||
RendezvousErrorCode? serviceError,
|
||||
NetworkEndpoint? dedicatedFallback,
|
||||
NetPeer? peer)
|
||||
{
|
||||
if (elapsed < TimeSpan.Zero)
|
||||
{
|
||||
throw new ArgumentOutOfRangeException(nameof(elapsed));
|
||||
}
|
||||
|
||||
if (dedicatedFallback is not null
|
||||
&& !ContractValidation.IsNetworkEndpointValid(dedicatedFallback))
|
||||
{
|
||||
throw new ArgumentException("The dedicated fallback endpoint is invalid.", nameof(dedicatedFallback));
|
||||
}
|
||||
|
||||
Kind = kind;
|
||||
Source = source;
|
||||
Category = category;
|
||||
Phase = phase;
|
||||
Elapsed = elapsed;
|
||||
ServiceError = serviceError;
|
||||
_dedicatedFallback = RendezvousEndpoint.Copy(dedicatedFallback);
|
||||
Peer = peer;
|
||||
}
|
||||
|
||||
public ConnectionOutcomeKind Kind { get; }
|
||||
public RendezvousConnectionOutcomeSource Source { get; }
|
||||
public RendezvousConnectionFailureCategory Category { get; }
|
||||
public RendezvousConnectionPhase Phase { get; }
|
||||
public TimeSpan Elapsed { get; }
|
||||
public RendezvousErrorCode? ServiceError { get; }
|
||||
public NetworkEndpoint? DedicatedFallback => RendezvousEndpoint.Copy(_dedicatedFallback);
|
||||
public NetPeer? Peer { get; }
|
||||
public bool IsSuccess => Kind == ConnectionOutcomeKind.Connected;
|
||||
public bool HasDedicatedFallback => _dedicatedFallback is not null;
|
||||
|
||||
public static RendezvousConnectionOutcome FromServiceError(
|
||||
RendezvousErrorCode error,
|
||||
TimeSpan elapsed,
|
||||
NetworkEndpoint? dedicatedFallback = null)
|
||||
{
|
||||
if (error == RendezvousErrorCode.None)
|
||||
{
|
||||
throw new ArgumentException("A service failure outcome requires an error.", nameof(error));
|
||||
}
|
||||
|
||||
(ConnectionOutcomeKind kind, RendezvousConnectionFailureCategory category, RendezvousConnectionPhase phase) =
|
||||
error switch
|
||||
{
|
||||
RendezvousErrorCode.NotFound => (
|
||||
ConnectionOutcomeKind.DirectoryNotFound,
|
||||
RendezvousConnectionFailureCategory.Directory,
|
||||
RendezvousConnectionPhase.Directory),
|
||||
RendezvousErrorCode.Expired => (
|
||||
ConnectionOutcomeKind.AttemptExpired,
|
||||
RendezvousConnectionFailureCategory.Authorization,
|
||||
RendezvousConnectionPhase.Authorization),
|
||||
RendezvousErrorCode.IncompatibleProtocol => (
|
||||
ConnectionOutcomeKind.IncompatibleProtocol,
|
||||
RendezvousConnectionFailureCategory.Compatibility,
|
||||
RendezvousConnectionPhase.Directory),
|
||||
RendezvousErrorCode.AuthenticationRequired
|
||||
or RendezvousErrorCode.Forbidden
|
||||
or RendezvousErrorCode.ReplayRejected => (
|
||||
ConnectionOutcomeKind.Unauthorized,
|
||||
RendezvousConnectionFailureCategory.Authorization,
|
||||
RendezvousConnectionPhase.Authorization),
|
||||
RendezvousErrorCode.RateLimited
|
||||
or RendezvousErrorCode.CapacityExceeded => (
|
||||
ConnectionOutcomeKind.RateLimited,
|
||||
RendezvousConnectionFailureCategory.Capacity,
|
||||
RendezvousConnectionPhase.Authorization),
|
||||
RendezvousErrorCode.StaleHost => (
|
||||
ConnectionOutcomeKind.NoHostPresence,
|
||||
RendezvousConnectionFailureCategory.HostPresence,
|
||||
RendezvousConnectionPhase.Mediation),
|
||||
RendezvousErrorCode.ServiceUnavailable => (
|
||||
ConnectionOutcomeKind.ServiceUnavailable,
|
||||
RendezvousConnectionFailureCategory.Service,
|
||||
RendezvousConnectionPhase.Authorization),
|
||||
_ => (
|
||||
ConnectionOutcomeKind.ServiceRejected,
|
||||
RendezvousConnectionFailureCategory.Service,
|
||||
RendezvousConnectionPhase.Authorization),
|
||||
};
|
||||
return new(
|
||||
kind,
|
||||
RendezvousConnectionOutcomeSource.RendezvousService,
|
||||
category,
|
||||
phase,
|
||||
elapsed,
|
||||
error,
|
||||
dedicatedFallback,
|
||||
null);
|
||||
}
|
||||
|
||||
public static ConnectionElapsedBucket BucketElapsed(TimeSpan elapsed)
|
||||
{
|
||||
if (elapsed < TimeSpan.Zero)
|
||||
{
|
||||
throw new ArgumentOutOfRangeException(nameof(elapsed));
|
||||
}
|
||||
|
||||
return elapsed.TotalSeconds switch
|
||||
{
|
||||
< 1 => ConnectionElapsedBucket.UnderOneSecond,
|
||||
< 5 => ConnectionElapsedBucket.OneToFiveSeconds,
|
||||
< 15 => ConnectionElapsedBucket.FiveToFifteenSeconds,
|
||||
< 30 => ConnectionElapsedBucket.FifteenToThirtySeconds,
|
||||
_ => ConnectionElapsedBucket.ThirtySecondsOrMore,
|
||||
};
|
||||
}
|
||||
|
||||
public override string ToString() =>
|
||||
$"[RendezvousConnectionOutcome {Kind}; {Source}; credentials redacted]";
|
||||
|
||||
internal static RendezvousConnectionOutcome Create(
|
||||
ConnectionOutcomeKind kind,
|
||||
RendezvousConnectionOutcomeSource source,
|
||||
RendezvousConnectionFailureCategory category,
|
||||
RendezvousConnectionPhase phase,
|
||||
TimeSpan elapsed,
|
||||
NetworkEndpoint? dedicatedFallback = null,
|
||||
NetPeer? peer = null) => new(
|
||||
kind,
|
||||
source,
|
||||
category,
|
||||
phase,
|
||||
elapsed,
|
||||
null,
|
||||
dedicatedFallback,
|
||||
peer);
|
||||
|
||||
}
|
||||
+35
@@ -0,0 +1,35 @@
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Client;
|
||||
|
||||
public sealed class RendezvousConnectionStartResult
|
||||
{
|
||||
internal RendezvousConnectionStartResult(
|
||||
CreateJoinAttemptResponse? attempt,
|
||||
RendezvousConnectionOutcome? outcome)
|
||||
{
|
||||
if ((attempt is null) == (outcome is null))
|
||||
{
|
||||
throw new ArgumentException(
|
||||
"A connection start result requires exactly one attempt or terminal outcome.");
|
||||
}
|
||||
|
||||
Attempt = attempt;
|
||||
Outcome = outcome;
|
||||
}
|
||||
|
||||
public CreateJoinAttemptResponse? Attempt { get; }
|
||||
public RendezvousConnectionOutcome? Outcome { get; }
|
||||
public bool IsReadyForTraversal => Attempt is not null;
|
||||
public bool IsCompleted => Outcome is not null;
|
||||
|
||||
public static RendezvousConnectionStartResult ReadyForTraversal(
|
||||
CreateJoinAttemptResponse attempt) => new(
|
||||
attempt ?? throw new ArgumentNullException(nameof(attempt)),
|
||||
null);
|
||||
|
||||
public static RendezvousConnectionStartResult Completed(
|
||||
RendezvousConnectionOutcome outcome) => new(
|
||||
null,
|
||||
outcome ?? throw new ArgumentNullException(nameof(outcome)));
|
||||
}
|
||||
@@ -1,3 +1,4 @@
|
||||
using System.Diagnostics;
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Client;
|
||||
@@ -39,6 +40,45 @@ public sealed class RendezvousJoinClient : IRendezvousJoinClient
|
||||
cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<RendezvousConnectionStartResult> CreateConnectionAttemptAsync(
|
||||
CreateJoinAttemptRequest request,
|
||||
NetworkEndpoint? dedicatedFallback = null,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
if (dedicatedFallback is not null
|
||||
&& !ContractValidation.IsNetworkEndpointValid(dedicatedFallback))
|
||||
{
|
||||
throw new ArgumentException("The dedicated fallback endpoint is invalid.", nameof(dedicatedFallback));
|
||||
}
|
||||
|
||||
Stopwatch elapsed = Stopwatch.StartNew();
|
||||
try
|
||||
{
|
||||
RendezvousClientResult<CreateJoinAttemptResponse> result = await CreateAsync(
|
||||
request,
|
||||
cancellationToken).ConfigureAwait(false);
|
||||
elapsed.Stop();
|
||||
return result.IsSuccess && result.Value is not null
|
||||
? RendezvousConnectionStartResult.ReadyForTraversal(result.Value)
|
||||
: RendezvousConnectionStartResult.Completed(
|
||||
RendezvousConnectionOutcome.FromServiceError(
|
||||
result.Error,
|
||||
elapsed.Elapsed,
|
||||
dedicatedFallback));
|
||||
}
|
||||
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
|
||||
{
|
||||
elapsed.Stop();
|
||||
return RendezvousConnectionStartResult.Completed(
|
||||
RendezvousConnectionOutcome.Create(
|
||||
ConnectionOutcomeKind.Cancelled,
|
||||
RendezvousConnectionOutcomeSource.Caller,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
RendezvousConnectionPhase.Authorization,
|
||||
elapsed.Elapsed));
|
||||
}
|
||||
}
|
||||
|
||||
public Task<RendezvousClientResult<bool>> CancelAsync(
|
||||
CreateJoinAttemptResponse attempt,
|
||||
CancellationToken cancellationToken = default)
|
||||
@@ -130,6 +170,41 @@ public sealed class RendezvousJoinClient : IRendezvousJoinClient
|
||||
$"Host invitation polling exceeded the configured {maximumPages}-page limit.");
|
||||
}
|
||||
|
||||
public Task<RendezvousClientResult<ReportConnectionOutcomeResponse>> ReportOutcomeAsync(
|
||||
CreateJoinAttemptResponse attempt,
|
||||
RendezvousConnectionOutcome outcome,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
if (attempt is null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(attempt));
|
||||
}
|
||||
if (outcome is null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(outcome));
|
||||
}
|
||||
if (!ContractValidation.IsReportableConnectionOutcome(outcome.Kind))
|
||||
{
|
||||
throw new ArgumentException(
|
||||
"This outcome cannot be reported for an issued join attempt.",
|
||||
nameof(outcome));
|
||||
}
|
||||
|
||||
ReportConnectionOutcomeRequest body = new()
|
||||
{
|
||||
Outcome = outcome.Kind,
|
||||
ElapsedBucket = RendezvousConnectionOutcome.BucketElapsed(outcome.Elapsed),
|
||||
};
|
||||
return _transport.SendSafeAsync<ReportConnectionOutcomeResponse>(
|
||||
() => HeaderJsonRequest(
|
||||
HttpMethod.Post,
|
||||
$"v1/join-attempts/{attempt.AttemptId}/outcome",
|
||||
ClientPunchCapabilityHeader,
|
||||
RequireHeaderValue(attempt.ClientPunchCapability, nameof(attempt)),
|
||||
body),
|
||||
cancellationToken);
|
||||
}
|
||||
|
||||
private static HttpRequestMessage HeaderRequest(
|
||||
HttpMethod method,
|
||||
string uri,
|
||||
@@ -141,6 +216,18 @@ public sealed class RendezvousJoinClient : IRendezvousJoinClient
|
||||
return request;
|
||||
}
|
||||
|
||||
private static HttpRequestMessage HeaderJsonRequest<T>(
|
||||
HttpMethod method,
|
||||
string uri,
|
||||
string header,
|
||||
string value,
|
||||
T body)
|
||||
{
|
||||
HttpRequestMessage request = RendezvousHttpTransport.JsonRequest(method, uri, body);
|
||||
request.Headers.TryAddWithoutValidation(header, value);
|
||||
return request;
|
||||
}
|
||||
|
||||
private static string RequireHeaderValue(string value, string parameterName) =>
|
||||
!string.IsNullOrWhiteSpace(value)
|
||||
? value
|
||||
|
||||
@@ -29,6 +29,12 @@ RendezvousClientResult<PublishedSession> registered = await publisher.RegisterAs
|
||||
DisplayName = "My server",
|
||||
Visibility = ListingVisibility.Public,
|
||||
Capacity = new() { CurrentPlayers = 1, MaximumPlayers = 8 },
|
||||
DedicatedFallback = new()
|
||||
{
|
||||
AddressFamily = AddressFamilyKind.Ipv4,
|
||||
Address = "203.0.113.40",
|
||||
Port = 7777,
|
||||
},
|
||||
},
|
||||
publisherCredential,
|
||||
cancellationToken);
|
||||
@@ -97,14 +103,22 @@ the accepted peer as connected. Register ordinary gameplay callbacks on
|
||||
`networkEvents.GameplayEvents`; the routing listener reserves Rendezvous direct
|
||||
requests for ticket validation and forwards every other callback normally.
|
||||
|
||||
The joining game first creates the HTTP attempt, then uses its own already-started
|
||||
gameplay manager in the same frame loop:
|
||||
The joining game first requests an attempt through the typed start API. It returns
|
||||
exactly one issued attempt or one terminal service outcome, so service authority
|
||||
is not confused with a later locally observed traversal failure:
|
||||
|
||||
```csharp
|
||||
CreateJoinAttemptResponse attempt = (await joins.CreateAsync(
|
||||
RendezvousConnectionStartResult start = await joins.CreateConnectionAttemptAsync(
|
||||
createJoinRequest,
|
||||
cancellationToken)).Value
|
||||
?? throw new InvalidOperationException("Join issuance failed.");
|
||||
cancellationToken: cancellationToken);
|
||||
if (start.Outcome is { } serviceOutcome)
|
||||
{
|
||||
ShowConnectionFailure(serviceOutcome.Kind, serviceOutcome.Category);
|
||||
return;
|
||||
}
|
||||
|
||||
CreateJoinAttemptResponse attempt = start.Attempt
|
||||
?? throw new InvalidOperationException("The typed start result was invalid.");
|
||||
using RendezvousClientCoordinator client = new(
|
||||
gameplayNetManager,
|
||||
networkEvents,
|
||||
@@ -116,12 +130,35 @@ client.Poll();
|
||||
```
|
||||
|
||||
NAT introduction changes the client state to `Connecting`; it is not success.
|
||||
Only `Connected` supplies `ConnectedPeer`. Call `Cancel()` and then `Poll()` for
|
||||
local cancellation, or `CancelAsync(joins, cancellationToken)` to also revoke the
|
||||
service attempt. Terminal client paths release all event subscriptions. Disposing
|
||||
a coordinator never stops or disposes the caller-owned manager and does not touch
|
||||
an in-flight peer; call `Cancel()` followed by `Poll()` first when that peer must
|
||||
also be disconnected.
|
||||
Only a `Connected` outcome supplies `Peer`. Completion exposes a stable kind,
|
||||
source, category, phase, and elapsed duration. The default HTTP silence, punch,
|
||||
and direct-connect budgets are five, ten, and five seconds respectively; configure
|
||||
them through `RendezvousClientOptions` and `RendezvousCoordinatorOptions` when a
|
||||
game has measured reasons to do so. The signed attempt expiry is always the
|
||||
absolute upper bound.
|
||||
|
||||
Call `Cancel()` and then `Poll()` for local cancellation, or
|
||||
`CancelAsync(joins, cancellationToken)` to also revoke the service attempt.
|
||||
Terminal client paths complete exactly once and release all event subscriptions,
|
||||
so late packets and callbacks are inert. Disposing a coordinator never stops or
|
||||
disposes the caller-owned manager and does not touch an in-flight peer; call
|
||||
`Cancel()` followed by `Poll()` first when that peer must also be disconnected.
|
||||
|
||||
After terminal completion, reporting is explicit and safe to retry. It sends only
|
||||
the authenticated outcome enum and a coarse elapsed bucket—never the endpoint,
|
||||
exact duration, diagnostic text, metadata, or player identity:
|
||||
|
||||
```csharp
|
||||
RendezvousClientResult<ReportConnectionOutcomeResponse> report =
|
||||
await client.ReportOutcomeAsync(joins, cancellationToken);
|
||||
```
|
||||
|
||||
An optional `DedicatedFallback` is copied from the authoritative listing into the
|
||||
issued attempt and terminal outcome. A local deployment may replace it with
|
||||
`RendezvousCoordinatorOptions.DedicatedFallbackOverride`. The SDK only returns
|
||||
the endpoint; it never connects automatically. The game must explicitly decide
|
||||
whether to use it and then connect and authenticate through its own gameplay
|
||||
transport. If the outcome has no fallback, v1 offers no relay.
|
||||
|
||||
Lease renewal is explicit and caller-controlled:
|
||||
|
||||
@@ -155,5 +192,6 @@ apply its own player identity, capacity, ban, and gameplay admission rules. Revo
|
||||
the attempt on cancellation and dispose the validator during host shutdown so its
|
||||
keyed ticket digests are zeroed.
|
||||
|
||||
See the repository's ADR 0007 for HTTP ownership/retry semantics and ADR 0008 for
|
||||
join-capability and connection-ticket security semantics.
|
||||
See the repository's ADR 0007 for HTTP ownership/retry semantics, ADR 0008 for
|
||||
join-capability and connection-ticket security semantics, and ADR 0010 for typed
|
||||
outcomes, deadlines, reporting, and caller-owned fallback.
|
||||
|
||||
@@ -148,6 +148,11 @@ public interface IRendezvousSessionBrowserClient
|
||||
|
||||
public interface IRendezvousJoinClient
|
||||
{
|
||||
Task<RendezvousConnectionStartResult> CreateConnectionAttemptAsync(
|
||||
CreateJoinAttemptRequest request,
|
||||
NetworkEndpoint? dedicatedFallback = null,
|
||||
CancellationToken cancellationToken = default);
|
||||
|
||||
Task<RendezvousClientResult<CreateJoinAttemptResponse>> CreateAsync(
|
||||
CreateJoinAttemptRequest request,
|
||||
CancellationToken cancellationToken = default);
|
||||
@@ -166,6 +171,11 @@ public interface IRendezvousJoinClient
|
||||
PublishedSession session,
|
||||
int maximumPages = 100,
|
||||
CancellationToken cancellationToken = default);
|
||||
|
||||
Task<RendezvousClientResult<ReportConnectionOutcomeResponse>> ReportOutcomeAsync(
|
||||
CreateJoinAttemptResponse attempt,
|
||||
RendezvousConnectionOutcome outcome,
|
||||
CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
||||
public interface IRendezvousDelay
|
||||
@@ -176,6 +186,7 @@ public interface IRendezvousDelay
|
||||
public sealed class RendezvousClientOptions
|
||||
{
|
||||
public int MaximumSafeRetries { get; set; } = 2;
|
||||
public TimeSpan RequestTimeout { get; set; } = TimeSpan.FromSeconds(5);
|
||||
public TimeSpan InitialRetryDelay { get; set; } = TimeSpan.FromMilliseconds(200);
|
||||
public TimeSpan MaximumRetryDelay { get; set; } = TimeSpan.FromSeconds(2);
|
||||
public double JitterRatio { get; set; } = 0.2;
|
||||
@@ -183,6 +194,8 @@ public sealed class RendezvousClientOptions
|
||||
internal void Validate()
|
||||
{
|
||||
if (MaximumSafeRetries is < 0 or > 5
|
||||
|| RequestTimeout <= TimeSpan.Zero
|
||||
|| RequestTimeout > TimeSpan.FromSeconds(30)
|
||||
|| InitialRetryDelay < TimeSpan.Zero
|
||||
|| MaximumRetryDelay < InitialRetryDelay
|
||||
|| MaximumRetryDelay > TimeSpan.FromSeconds(30)
|
||||
@@ -198,3 +211,15 @@ internal sealed class SystemRendezvousDelay : IRendezvousDelay
|
||||
public Task DelayAsync(TimeSpan delay, CancellationToken cancellationToken) =>
|
||||
Task.Delay(delay, cancellationToken);
|
||||
}
|
||||
|
||||
internal static class RendezvousEndpoint
|
||||
{
|
||||
internal static NetworkEndpoint? Copy(NetworkEndpoint? endpoint) => endpoint is null
|
||||
? null
|
||||
: new NetworkEndpoint
|
||||
{
|
||||
AddressFamily = endpoint.AddressFamily,
|
||||
Address = endpoint.Address,
|
||||
Port = endpoint.Port,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -24,6 +24,7 @@ internal sealed class RendezvousHttpTransport
|
||||
_options = new RendezvousClientOptions
|
||||
{
|
||||
MaximumSafeRetries = suppliedOptions.MaximumSafeRetries,
|
||||
RequestTimeout = suppliedOptions.RequestTimeout,
|
||||
InitialRetryDelay = suppliedOptions.InitialRetryDelay,
|
||||
MaximumRetryDelay = suppliedOptions.MaximumRetryDelay,
|
||||
JitterRatio = suppliedOptions.JitterRatio,
|
||||
@@ -38,11 +39,15 @@ internal sealed class RendezvousHttpTransport
|
||||
for (int attempt = 0; ; attempt++)
|
||||
{
|
||||
cancellationToken.ThrowIfCancellationRequested();
|
||||
using CancellationTokenSource requestTimeout =
|
||||
CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
|
||||
requestTimeout.CancelAfter(_options.RequestTimeout);
|
||||
CancellationToken requestCancellation = requestTimeout.Token;
|
||||
try
|
||||
{
|
||||
using HttpRequestMessage request = requestFactory();
|
||||
using HttpResponseMessage response = await _httpClient
|
||||
.SendAsync(request, HttpCompletionOption.ResponseHeadersRead, cancellationToken)
|
||||
.SendAsync(request, HttpCompletionOption.ResponseHeadersRead, requestCancellation)
|
||||
.ConfigureAwait(false);
|
||||
if (response.IsSuccessStatusCode)
|
||||
{
|
||||
@@ -54,7 +59,7 @@ internal sealed class RendezvousHttpTransport
|
||||
byte[] payload;
|
||||
try
|
||||
{
|
||||
payload = await ReadBoundedAsync(response.Content, cancellationToken)
|
||||
payload = await ReadBoundedAsync(response.Content, requestCancellation)
|
||||
.ConfigureAwait(false);
|
||||
}
|
||||
catch (InvalidDataException)
|
||||
@@ -81,7 +86,7 @@ internal sealed class RendezvousHttpTransport
|
||||
: RendezvousClientResult.Success(value);
|
||||
}
|
||||
|
||||
ApiError error = await ReadErrorAsync(response, cancellationToken).ConfigureAwait(false);
|
||||
ApiError error = await ReadErrorAsync(response, requestCancellation).ConfigureAwait(false);
|
||||
int? retryAfter = error.RetryAfterSeconds ?? GetRetryAfterSeconds(response.Headers.RetryAfter);
|
||||
if (attempt < _options.MaximumSafeRetries && IsTransient(error.Code))
|
||||
{
|
||||
|
||||
@@ -88,6 +88,7 @@ public sealed class RendezvousPublisherClient : IRendezvousPublisherClient
|
||||
DisplayName = request.DisplayName,
|
||||
Capacity = CopyCapacity(request.Capacity),
|
||||
Metadata = CopyMetadata(request.Metadata),
|
||||
DedicatedFallback = RendezvousEndpoint.Copy(request.DedicatedFallback),
|
||||
};
|
||||
return _transport.SendSafeAsync<bool>(
|
||||
() => RendezvousHttpTransport.JsonRequest(
|
||||
@@ -138,6 +139,7 @@ public sealed class RendezvousPublisherClient : IRendezvousPublisherClient
|
||||
Visibility = request.Visibility,
|
||||
Capacity = CopyCapacity(request.Capacity),
|
||||
Metadata = CopyMetadata(request.Metadata),
|
||||
DedicatedFallback = RendezvousEndpoint.Copy(request.DedicatedFallback),
|
||||
};
|
||||
|
||||
private static SessionCapacity CopyCapacity(SessionCapacity capacity) => new()
|
||||
@@ -148,4 +150,5 @@ public sealed class RendezvousPublisherClient : IRendezvousPublisherClient
|
||||
|
||||
private static Dictionary<string, string> CopyMetadata(Dictionary<string, string> metadata) =>
|
||||
new(metadata, StringComparer.Ordinal);
|
||||
|
||||
}
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
using System.Net;
|
||||
using System.Net.Sockets;
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using LiteNetLib;
|
||||
|
||||
@@ -12,12 +13,21 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
private readonly IPEndPoint _mediator;
|
||||
private readonly CreateJoinAttemptResponse _attempt;
|
||||
private readonly IRendezvousCoordinatorClock _clock;
|
||||
private readonly RendezvousCoordinatorOptions _options;
|
||||
private readonly RendezvousPunchRetrySchedule _retry;
|
||||
private readonly object _completionGate = new();
|
||||
private readonly TimeSpan _startedAt;
|
||||
private readonly TimeSpan _attemptDeadline;
|
||||
private readonly TimeSpan _punchDeadline;
|
||||
private readonly NetworkEndpoint? _dedicatedFallback;
|
||||
private NetPeer? _connectingPeer;
|
||||
private IPEndPoint? _directEndpoint;
|
||||
private TimeSpan? _directDeadline;
|
||||
private RendezvousConnectionOutcome? _outcome;
|
||||
private bool _cancelRequested;
|
||||
private int _polling;
|
||||
private bool _subscriptionsReleased;
|
||||
private bool _disposed;
|
||||
private int _disposed;
|
||||
|
||||
public RendezvousClientCoordinator(
|
||||
NetManager manager,
|
||||
@@ -49,23 +59,31 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
_mediator = mediator ?? throw new ArgumentNullException(nameof(mediator));
|
||||
_attempt = attempt ?? throw new ArgumentNullException(nameof(attempt));
|
||||
_clock = clock ?? throw new ArgumentNullException(nameof(clock));
|
||||
RendezvousCoordinatorOptions validated = (options ?? new RendezvousCoordinatorOptions())
|
||||
_options = (options ?? new RendezvousCoordinatorOptions())
|
||||
.CopyAndValidate();
|
||||
_retry = new(validated, _clock);
|
||||
_retry = new(_options, _clock);
|
||||
|
||||
RendezvousManagerGuard.Validate(_manager, _networkEvents);
|
||||
DateTimeOffset startedUtc = _clock.UtcNow;
|
||||
if (_mediator.Port is < 1 or > 65_535
|
||||
|| _attempt.AttemptId.Value == Guid.Empty
|
||||
|| _attempt.MediationHandle.Value == Guid.Empty
|
||||
|| !ContractValidation.IsCapabilityValid(_attempt.ClientPunchCapability)
|
||||
|| !ContractValidation.IsConnectionTicketValid(_attempt.ConnectionTicketDigest)
|
||||
|| _attempt.ExpiresAt <= _clock.UtcNow)
|
||||
|| _attempt.ExpiresAt <= startedUtc)
|
||||
{
|
||||
throw new ArgumentException("The client traversal inputs are invalid.");
|
||||
}
|
||||
|
||||
_startedAt = _clock.Elapsed;
|
||||
_attemptDeadline = _startedAt + (_attempt.ExpiresAt - startedUtc);
|
||||
_punchDeadline = Min(_attemptDeadline, _startedAt + _options.PunchTimeout);
|
||||
_dedicatedFallback = RendezvousEndpoint.Copy(
|
||||
_options.DedicatedFallbackOverride ?? _attempt.DedicatedFallback);
|
||||
|
||||
_networkEvents.RendezvousPeerConnected += OnPeerConnected;
|
||||
_networkEvents.RendezvousPeerDisconnected += OnPeerDisconnected;
|
||||
_networkEvents.RendezvousNetworkError += OnNetworkError;
|
||||
_punchEvents.NatIntroductionSuccess += OnNatIntroductionSuccess;
|
||||
}
|
||||
|
||||
@@ -73,7 +91,8 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
|
||||
public RendezvousConnectionState State { get; private set; } = RendezvousConnectionState.Punching;
|
||||
public NetPeer? ConnectedPeer { get; private set; }
|
||||
public bool IsCompleted => IsTerminal(State);
|
||||
public RendezvousConnectionOutcome? Outcome => Volatile.Read(ref _outcome);
|
||||
public bool IsCompleted => Outcome is not null;
|
||||
|
||||
public void Cancel() => Volatile.Write(ref _cancelRequested, true);
|
||||
|
||||
@@ -91,6 +110,23 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
return await joinClient.CancelAsync(_attempt, cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
|
||||
public Task<RendezvousClientResult<ReportConnectionOutcomeResponse>> ReportOutcomeAsync(
|
||||
IRendezvousJoinClient joinClient,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
if (joinClient is null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(joinClient));
|
||||
}
|
||||
ThrowIfDisposed();
|
||||
if (Outcome is null)
|
||||
{
|
||||
throw new InvalidOperationException("The connection attempt has not completed.");
|
||||
}
|
||||
|
||||
return joinClient.ReportOutcomeAsync(_attempt, Outcome, cancellationToken);
|
||||
}
|
||||
|
||||
public void Poll()
|
||||
{
|
||||
ThrowIfDisposed();
|
||||
@@ -109,13 +145,18 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
if (Volatile.Read(ref _cancelRequested))
|
||||
{
|
||||
DisconnectPendingPeer();
|
||||
Complete(RendezvousConnectionState.Cancelled);
|
||||
Complete(
|
||||
RendezvousConnectionState.Cancelled,
|
||||
ConnectionOutcomeKind.Cancelled,
|
||||
RendezvousConnectionOutcomeSource.Caller,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
CurrentPhase());
|
||||
return;
|
||||
}
|
||||
|
||||
if (!_manager.IsRunning)
|
||||
{
|
||||
Complete(RendezvousConnectionState.ManagerStopped);
|
||||
CompleteManagerStopped();
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -127,35 +168,67 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
}
|
||||
|
||||
DateTimeOffset now = _clock.UtcNow;
|
||||
TimeSpan elapsed = _clock.Elapsed;
|
||||
if (Volatile.Read(ref _cancelRequested))
|
||||
{
|
||||
DisconnectPendingPeer();
|
||||
Complete(RendezvousConnectionState.Cancelled);
|
||||
Complete(
|
||||
RendezvousConnectionState.Cancelled,
|
||||
ConnectionOutcomeKind.Cancelled,
|
||||
RendezvousConnectionOutcomeSource.Caller,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
CurrentPhase());
|
||||
}
|
||||
else if (!_manager.IsRunning)
|
||||
{
|
||||
Complete(RendezvousConnectionState.ManagerStopped);
|
||||
CompleteManagerStopped();
|
||||
}
|
||||
else if (now >= _attempt.ExpiresAt)
|
||||
else if (now >= _attempt.ExpiresAt || elapsed >= _attemptDeadline)
|
||||
{
|
||||
DisconnectPendingPeer();
|
||||
Complete(RendezvousConnectionState.TimedOut);
|
||||
Complete(
|
||||
RendezvousConnectionState.TimedOut,
|
||||
ConnectionOutcomeKind.AttemptExpired,
|
||||
RendezvousConnectionOutcomeSource.RendezvousService,
|
||||
RendezvousConnectionFailureCategory.Authorization,
|
||||
RendezvousConnectionPhase.Authorization);
|
||||
}
|
||||
else if (State == RendezvousConnectionState.Punching && _retry.IsDue(now))
|
||||
else if (State == RendezvousConnectionState.Punching)
|
||||
{
|
||||
if (_retry.IsExhausted)
|
||||
if (elapsed >= _punchDeadline
|
||||
|| _retry.IsExhausted && _retry.IsDue(elapsed))
|
||||
{
|
||||
Complete(RendezvousConnectionState.TimedOut);
|
||||
Complete(
|
||||
RendezvousConnectionState.TimedOut,
|
||||
ConnectionOutcomeKind.PunchTimedOut,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.NatTraversal,
|
||||
RendezvousConnectionPhase.NatTraversal);
|
||||
return;
|
||||
}
|
||||
|
||||
_manager.NatPunchModule.SendNatIntroduceRequest(
|
||||
_mediator,
|
||||
NatPunchRequestTokenCodec.Encode(
|
||||
NatPunchPeerRole.Client,
|
||||
_attempt.MediationHandle,
|
||||
_attempt.ClientPunchCapability));
|
||||
_retry.RecordRequest();
|
||||
if (_retry.IsDue(elapsed))
|
||||
{
|
||||
_manager.NatPunchModule.SendNatIntroduceRequest(
|
||||
_mediator,
|
||||
NatPunchRequestTokenCodec.Encode(
|
||||
NatPunchPeerRole.Client,
|
||||
_attempt.MediationHandle,
|
||||
_attempt.ClientPunchCapability));
|
||||
_retry.RecordRequest();
|
||||
}
|
||||
}
|
||||
else if (State == RendezvousConnectionState.Connecting
|
||||
&& _directDeadline is TimeSpan directDeadline
|
||||
&& directDeadline <= elapsed)
|
||||
{
|
||||
DisconnectPendingPeer();
|
||||
Complete(
|
||||
RendezvousConnectionState.TimedOut,
|
||||
ConnectionOutcomeKind.DirectConnectTimedOut,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.DirectConnection,
|
||||
RendezvousConnectionPhase.DirectConnection);
|
||||
}
|
||||
}
|
||||
finally
|
||||
@@ -166,18 +239,22 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
if (_disposed)
|
||||
if (Interlocked.Exchange(ref _disposed, 1) != 0)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
if (!IsCompleted)
|
||||
{
|
||||
Complete(RendezvousConnectionState.Disposed);
|
||||
Complete(
|
||||
RendezvousConnectionState.Disposed,
|
||||
ConnectionOutcomeKind.Disposed,
|
||||
RendezvousConnectionOutcomeSource.Lifecycle,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
CurrentPhase());
|
||||
}
|
||||
|
||||
ReleaseSubscriptions();
|
||||
_disposed = true;
|
||||
}
|
||||
|
||||
public override string ToString() =>
|
||||
@@ -205,16 +282,25 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
byte[] connectionData = DirectConnectionRequestCodec.Encode(
|
||||
introduction.AttemptId,
|
||||
introduction.ConnectionTicket);
|
||||
_directEndpoint = target;
|
||||
_connectingPeer = _manager.Connect(target, connectionData);
|
||||
if (_connectingPeer is null
|
||||
|| _connectingPeer.ConnectionState != ConnectionState.Outgoing)
|
||||
{
|
||||
_connectingPeer = null;
|
||||
Complete(RendezvousConnectionState.Rejected);
|
||||
Complete(
|
||||
RendezvousConnectionState.Rejected,
|
||||
ConnectionOutcomeKind.TransportError,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.DirectConnection,
|
||||
RendezvousConnectionPhase.DirectConnection);
|
||||
return;
|
||||
}
|
||||
|
||||
State = RendezvousConnectionState.Connecting;
|
||||
_directDeadline = Min(
|
||||
_attemptDeadline,
|
||||
_clock.Elapsed + _options.DirectConnectTimeout);
|
||||
}
|
||||
|
||||
private void OnPeerConnected(NetPeer peer)
|
||||
@@ -225,17 +311,60 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
return;
|
||||
}
|
||||
|
||||
ConnectedPeer = peer;
|
||||
Complete(RendezvousConnectionState.Connected, peer);
|
||||
Complete(
|
||||
RendezvousConnectionState.Connected,
|
||||
ConnectionOutcomeKind.Connected,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.None,
|
||||
RendezvousConnectionPhase.Complete,
|
||||
peer);
|
||||
}
|
||||
|
||||
private void OnPeerDisconnected(NetPeer peer, DisconnectInfo disconnectInfo)
|
||||
{
|
||||
_ = disconnectInfo;
|
||||
if (State == RendezvousConnectionState.Connecting
|
||||
&& ReferenceEquals(peer, _connectingPeer))
|
||||
{
|
||||
Complete(RendezvousConnectionState.Rejected);
|
||||
ConnectionOutcomeKind kind = disconnectInfo.Reason == DisconnectReason.Timeout
|
||||
? ConnectionOutcomeKind.DirectConnectTimedOut
|
||||
: disconnectInfo.Reason == DisconnectReason.ConnectionFailed
|
||||
? ConnectionOutcomeKind.TransportError
|
||||
: ConnectionOutcomeKind.HostRejected;
|
||||
Complete(
|
||||
kind == ConnectionOutcomeKind.DirectConnectTimedOut
|
||||
? RendezvousConnectionState.TimedOut
|
||||
: RendezvousConnectionState.Rejected,
|
||||
kind,
|
||||
kind == ConnectionOutcomeKind.HostRejected
|
||||
? RendezvousConnectionOutcomeSource.RemoteHost
|
||||
: RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.DirectConnection,
|
||||
RendezvousConnectionPhase.DirectConnection);
|
||||
}
|
||||
}
|
||||
|
||||
private void OnNetworkError(IPEndPoint endpoint, SocketError socketError)
|
||||
{
|
||||
_ = socketError;
|
||||
if (State == RendezvousConnectionState.Punching && endpoint.Equals(_mediator))
|
||||
{
|
||||
Complete(
|
||||
RendezvousConnectionState.Rejected,
|
||||
ConnectionOutcomeKind.MediatorUnavailable,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.Mediation,
|
||||
RendezvousConnectionPhase.Mediation);
|
||||
}
|
||||
else if (State == RendezvousConnectionState.Connecting
|
||||
&& endpoint.Equals(_directEndpoint))
|
||||
{
|
||||
DisconnectPendingPeer();
|
||||
Complete(
|
||||
RendezvousConnectionState.Rejected,
|
||||
ConnectionOutcomeKind.TransportError,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.DirectConnection,
|
||||
RendezvousConnectionPhase.DirectConnection);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -247,42 +376,85 @@ public sealed class RendezvousClientCoordinator : IDisposable
|
||||
}
|
||||
}
|
||||
|
||||
private void Complete(RendezvousConnectionState terminalState, NetPeer? peer = null)
|
||||
private void Complete(
|
||||
RendezvousConnectionState terminalState,
|
||||
ConnectionOutcomeKind kind,
|
||||
RendezvousConnectionOutcomeSource source,
|
||||
RendezvousConnectionFailureCategory category,
|
||||
RendezvousConnectionPhase phase,
|
||||
NetPeer? peer = null)
|
||||
{
|
||||
if (IsCompleted)
|
||||
RendezvousConnectionCompletedEventArgs completion;
|
||||
lock (_completionGate)
|
||||
{
|
||||
return;
|
||||
if (_outcome is not null)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
RendezvousConnectionOutcome outcome = RendezvousConnectionOutcome.Create(
|
||||
kind,
|
||||
source,
|
||||
category,
|
||||
phase,
|
||||
_clock.Elapsed - _startedAt,
|
||||
ShouldOfferFallback(kind) ? _dedicatedFallback : null,
|
||||
peer);
|
||||
State = terminalState;
|
||||
if (kind == ConnectionOutcomeKind.Connected)
|
||||
{
|
||||
ConnectedPeer = peer;
|
||||
}
|
||||
Volatile.Write(ref _outcome, outcome);
|
||||
ReleaseSubscriptions();
|
||||
completion = new(terminalState, outcome);
|
||||
}
|
||||
|
||||
State = terminalState;
|
||||
ReleaseSubscriptions();
|
||||
Completed?.Invoke(this, new(terminalState, peer));
|
||||
Completed?.Invoke(this, completion);
|
||||
}
|
||||
|
||||
private void ReleaseSubscriptions()
|
||||
{
|
||||
if (_subscriptionsReleased)
|
||||
lock (_completionGate)
|
||||
{
|
||||
return;
|
||||
}
|
||||
if (_subscriptionsReleased)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
_networkEvents.RendezvousPeerConnected -= OnPeerConnected;
|
||||
_networkEvents.RendezvousPeerDisconnected -= OnPeerDisconnected;
|
||||
_punchEvents.NatIntroductionSuccess -= OnNatIntroductionSuccess;
|
||||
_subscriptionsReleased = true;
|
||||
_networkEvents.RendezvousPeerConnected -= OnPeerConnected;
|
||||
_networkEvents.RendezvousPeerDisconnected -= OnPeerDisconnected;
|
||||
_networkEvents.RendezvousNetworkError -= OnNetworkError;
|
||||
_punchEvents.NatIntroductionSuccess -= OnNatIntroductionSuccess;
|
||||
_subscriptionsReleased = true;
|
||||
}
|
||||
}
|
||||
|
||||
private static bool IsTerminal(RendezvousConnectionState state) => state is
|
||||
RendezvousConnectionState.Connected
|
||||
or RendezvousConnectionState.Cancelled
|
||||
or RendezvousConnectionState.TimedOut
|
||||
or RendezvousConnectionState.Rejected
|
||||
or RendezvousConnectionState.ManagerStopped
|
||||
or RendezvousConnectionState.Disposed;
|
||||
private void CompleteManagerStopped() => Complete(
|
||||
RendezvousConnectionState.ManagerStopped,
|
||||
ConnectionOutcomeKind.ManagerStopped,
|
||||
RendezvousConnectionOutcomeSource.Lifecycle,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
CurrentPhase());
|
||||
|
||||
private RendezvousConnectionPhase CurrentPhase() => State switch
|
||||
{
|
||||
RendezvousConnectionState.Punching => RendezvousConnectionPhase.NatTraversal,
|
||||
RendezvousConnectionState.Connecting => RendezvousConnectionPhase.DirectConnection,
|
||||
_ => RendezvousConnectionPhase.Complete,
|
||||
};
|
||||
|
||||
private static bool ShouldOfferFallback(ConnectionOutcomeKind kind) => kind is not (
|
||||
ConnectionOutcomeKind.Connected
|
||||
or ConnectionOutcomeKind.Cancelled
|
||||
or ConnectionOutcomeKind.Disposed);
|
||||
|
||||
private static TimeSpan Min(TimeSpan left, TimeSpan right) =>
|
||||
left <= right ? left : right;
|
||||
|
||||
private void ThrowIfDisposed()
|
||||
{
|
||||
if (_disposed)
|
||||
if (Volatile.Read(ref _disposed) != 0)
|
||||
{
|
||||
throw new ObjectDisposedException(nameof(RendezvousClientCoordinator));
|
||||
}
|
||||
|
||||
@@ -1,4 +1,6 @@
|
||||
using System.Diagnostics;
|
||||
using System.Security.Cryptography;
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using LiteNetLib;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Client;
|
||||
@@ -15,12 +17,97 @@ public enum RendezvousConnectionState
|
||||
Disposed = 8,
|
||||
}
|
||||
|
||||
public sealed class RendezvousConnectionCompletedEventArgs(
|
||||
RendezvousConnectionState state,
|
||||
NetPeer? peer = null) : EventArgs
|
||||
public sealed class RendezvousConnectionCompletedEventArgs : EventArgs
|
||||
{
|
||||
public RendezvousConnectionState State { get; } = state;
|
||||
public NetPeer? Peer { get; } = peer;
|
||||
[Obsolete("Completion events now expose a typed Outcome. Construct these arguments only for legacy test doubles.")]
|
||||
public RendezvousConnectionCompletedEventArgs(
|
||||
RendezvousConnectionState state,
|
||||
NetPeer? peer)
|
||||
: this(state, RendezvousCompletionInvariant.FromLegacy(state, peer))
|
||||
{
|
||||
}
|
||||
|
||||
internal RendezvousConnectionCompletedEventArgs(
|
||||
RendezvousConnectionState state,
|
||||
RendezvousConnectionOutcome outcome)
|
||||
{
|
||||
RendezvousCompletionInvariant.Validate(state, outcome);
|
||||
State = state;
|
||||
Outcome = outcome;
|
||||
}
|
||||
|
||||
public RendezvousConnectionState State { get; }
|
||||
public RendezvousConnectionOutcome Outcome { get; }
|
||||
public NetPeer? Peer => Outcome.Peer;
|
||||
}
|
||||
|
||||
internal static class RendezvousCompletionInvariant
|
||||
{
|
||||
internal static RendezvousConnectionOutcome FromLegacy(
|
||||
RendezvousConnectionState state,
|
||||
NetPeer? peer) => state switch
|
||||
{
|
||||
RendezvousConnectionState.Connected when peer is not null => RendezvousConnectionOutcome.Create(
|
||||
ConnectionOutcomeKind.Connected,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.None,
|
||||
RendezvousConnectionPhase.Complete,
|
||||
TimeSpan.Zero,
|
||||
peer: peer),
|
||||
RendezvousConnectionState.Cancelled => RendezvousConnectionOutcome.Create(
|
||||
ConnectionOutcomeKind.Cancelled,
|
||||
RendezvousConnectionOutcomeSource.Caller,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
RendezvousConnectionPhase.Complete,
|
||||
TimeSpan.Zero),
|
||||
RendezvousConnectionState.TimedOut => RendezvousConnectionOutcome.Create(
|
||||
ConnectionOutcomeKind.DirectConnectTimedOut,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.DirectConnection,
|
||||
RendezvousConnectionPhase.DirectConnection,
|
||||
TimeSpan.Zero),
|
||||
RendezvousConnectionState.Rejected => RendezvousConnectionOutcome.Create(
|
||||
ConnectionOutcomeKind.HostRejected,
|
||||
RendezvousConnectionOutcomeSource.RemoteHost,
|
||||
RendezvousConnectionFailureCategory.Authorization,
|
||||
RendezvousConnectionPhase.Authorization,
|
||||
TimeSpan.Zero),
|
||||
RendezvousConnectionState.ManagerStopped => RendezvousConnectionOutcome.Create(
|
||||
ConnectionOutcomeKind.ManagerStopped,
|
||||
RendezvousConnectionOutcomeSource.Lifecycle,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
RendezvousConnectionPhase.Complete,
|
||||
TimeSpan.Zero),
|
||||
RendezvousConnectionState.Disposed => RendezvousConnectionOutcome.Create(
|
||||
ConnectionOutcomeKind.Disposed,
|
||||
RendezvousConnectionOutcomeSource.Lifecycle,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
RendezvousConnectionPhase.Complete,
|
||||
TimeSpan.Zero),
|
||||
RendezvousConnectionState.Connected => throw new ArgumentNullException(
|
||||
nameof(peer),
|
||||
"A connected completion requires a peer."),
|
||||
_ => throw new ArgumentOutOfRangeException(
|
||||
nameof(state),
|
||||
state,
|
||||
"A completion event requires a terminal connection state."),
|
||||
};
|
||||
|
||||
internal static void Validate(
|
||||
RendezvousConnectionState state,
|
||||
RendezvousConnectionOutcome outcome)
|
||||
{
|
||||
if (outcome is null)
|
||||
{
|
||||
throw new ArgumentNullException(nameof(outcome));
|
||||
}
|
||||
if ((state == RendezvousConnectionState.Connected) != outcome.IsSuccess)
|
||||
{
|
||||
throw new ArgumentException(
|
||||
"The connection state and typed outcome contradict each other.",
|
||||
nameof(outcome));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public sealed class RendezvousCoordinatorOptions
|
||||
@@ -29,8 +116,11 @@ public sealed class RendezvousCoordinatorOptions
|
||||
public int MaximumAttemptChecksPerPoll { get; set; } = 128;
|
||||
public TimeSpan InitialPunchRetryDelay { get; set; } = TimeSpan.FromMilliseconds(200);
|
||||
public TimeSpan MaximumPunchRetryDelay { get; set; } = TimeSpan.FromSeconds(2);
|
||||
public TimeSpan PunchTimeout { get; set; } = TimeSpan.FromSeconds(10);
|
||||
public TimeSpan DirectConnectTimeout { get; set; } = TimeSpan.FromSeconds(5);
|
||||
public TimeSpan ConnectionTicketLifetime { get; set; } = TimeSpan.FromSeconds(20);
|
||||
public double JitterRatio { get; set; } = 0.2;
|
||||
public NetworkEndpoint? DedicatedFallbackOverride { get; set; }
|
||||
|
||||
internal RendezvousCoordinatorOptions CopyAndValidate()
|
||||
{
|
||||
@@ -39,9 +129,15 @@ public sealed class RendezvousCoordinatorOptions
|
||||
|| InitialPunchRetryDelay < TimeSpan.FromMilliseconds(10)
|
||||
|| MaximumPunchRetryDelay < InitialPunchRetryDelay
|
||||
|| MaximumPunchRetryDelay > TimeSpan.FromSeconds(10)
|
||||
|| PunchTimeout <= TimeSpan.Zero
|
||||
|| PunchTimeout > TimeSpan.FromSeconds(30)
|
||||
|| DirectConnectTimeout <= TimeSpan.Zero
|
||||
|| DirectConnectTimeout > TimeSpan.FromSeconds(30)
|
||||
|| ConnectionTicketLifetime <= TimeSpan.Zero
|
||||
|| ConnectionTicketLifetime > TimeSpan.FromSeconds(20)
|
||||
|| JitterRatio is < 0 or > 1)
|
||||
|| JitterRatio is < 0 or > 1
|
||||
|| DedicatedFallbackOverride is not null
|
||||
&& !ContractValidation.IsNetworkEndpointValid(DedicatedFallbackOverride))
|
||||
{
|
||||
throw new ArgumentOutOfRangeException(nameof(RendezvousCoordinatorOptions));
|
||||
}
|
||||
@@ -52,8 +148,11 @@ public sealed class RendezvousCoordinatorOptions
|
||||
MaximumAttemptChecksPerPoll = MaximumAttemptChecksPerPoll,
|
||||
InitialPunchRetryDelay = InitialPunchRetryDelay,
|
||||
MaximumPunchRetryDelay = MaximumPunchRetryDelay,
|
||||
PunchTimeout = PunchTimeout,
|
||||
DirectConnectTimeout = DirectConnectTimeout,
|
||||
ConnectionTicketLifetime = ConnectionTicketLifetime,
|
||||
JitterRatio = JitterRatio,
|
||||
DedicatedFallbackOverride = RendezvousEndpoint.Copy(DedicatedFallbackOverride),
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -61,11 +160,16 @@ public sealed class RendezvousCoordinatorOptions
|
||||
internal interface IRendezvousCoordinatorClock
|
||||
{
|
||||
DateTimeOffset UtcNow { get; }
|
||||
TimeSpan Elapsed { get; }
|
||||
}
|
||||
|
||||
internal sealed class SystemRendezvousCoordinatorClock : IRendezvousCoordinatorClock
|
||||
{
|
||||
private readonly long _origin = Stopwatch.GetTimestamp();
|
||||
|
||||
public DateTimeOffset UtcNow => DateTimeOffset.UtcNow;
|
||||
public TimeSpan Elapsed => TimeSpan.FromSeconds(
|
||||
(Stopwatch.GetTimestamp() - _origin) / (double)Stopwatch.Frequency);
|
||||
}
|
||||
|
||||
internal static class RendezvousManagerGuard
|
||||
@@ -95,11 +199,11 @@ internal sealed class RendezvousPunchRetrySchedule(
|
||||
IRendezvousCoordinatorClock clock)
|
||||
{
|
||||
public int RequestsSent { get; private set; }
|
||||
public DateTimeOffset NextRequestAt { get; private set; } = DateTimeOffset.MinValue;
|
||||
public TimeSpan NextRequestAt { get; private set; } = TimeSpan.Zero;
|
||||
|
||||
public bool IsExhausted => RequestsSent >= options.MaximumPunchRequests;
|
||||
|
||||
public bool IsDue(DateTimeOffset now) => now >= NextRequestAt;
|
||||
public bool IsDue(TimeSpan elapsed) => elapsed >= NextRequestAt;
|
||||
|
||||
public void RecordRequest()
|
||||
{
|
||||
@@ -119,6 +223,6 @@ internal sealed class RendezvousPunchRetrySchedule(
|
||||
options.MaximumPunchRetryDelay.TotalMilliseconds);
|
||||
}
|
||||
|
||||
NextRequestAt = clock.UtcNow + TimeSpan.FromMilliseconds(milliseconds);
|
||||
NextRequestAt = clock.Elapsed + TimeSpan.FromMilliseconds(milliseconds);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
using System.Net;
|
||||
using System.Net.Sockets;
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using LiteNetLib;
|
||||
|
||||
@@ -11,14 +12,37 @@ public enum RendezvousHostState
|
||||
Disposed = 3,
|
||||
}
|
||||
|
||||
public sealed class RendezvousHostAttemptCompletedEventArgs(
|
||||
JoinAttemptId attemptId,
|
||||
RendezvousConnectionState state,
|
||||
NetPeer? peer = null) : EventArgs
|
||||
public sealed class RendezvousHostAttemptCompletedEventArgs : EventArgs
|
||||
{
|
||||
public JoinAttemptId AttemptId { get; } = attemptId;
|
||||
public RendezvousConnectionState State { get; } = state;
|
||||
public NetPeer? Peer { get; } = peer;
|
||||
[Obsolete("Completion events now expose a typed Outcome. Construct these arguments only for legacy test doubles.")]
|
||||
public RendezvousHostAttemptCompletedEventArgs(
|
||||
JoinAttemptId attemptId,
|
||||
RendezvousConnectionState state,
|
||||
NetPeer? peer)
|
||||
: this(attemptId, state, RendezvousCompletionInvariant.FromLegacy(state, peer))
|
||||
{
|
||||
}
|
||||
|
||||
internal RendezvousHostAttemptCompletedEventArgs(
|
||||
JoinAttemptId attemptId,
|
||||
RendezvousConnectionState state,
|
||||
RendezvousConnectionOutcome outcome)
|
||||
{
|
||||
if (attemptId.Value == Guid.Empty)
|
||||
{
|
||||
throw new ArgumentException("The completed attempt ID is invalid.", nameof(attemptId));
|
||||
}
|
||||
|
||||
RendezvousCompletionInvariant.Validate(state, outcome);
|
||||
AttemptId = attemptId;
|
||||
State = state;
|
||||
Outcome = outcome;
|
||||
}
|
||||
|
||||
public JoinAttemptId AttemptId { get; }
|
||||
public RendezvousConnectionState State { get; }
|
||||
public RendezvousConnectionOutcome Outcome { get; }
|
||||
public NetPeer? Peer => Outcome.Peer;
|
||||
}
|
||||
|
||||
public sealed class RendezvousHostCoordinator : IDisposable
|
||||
@@ -37,6 +61,7 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
private readonly Dictionary<JoinAttemptId, DeferredConnectionRequest> _deferredRequests = [];
|
||||
private readonly Dictionary<JoinAttemptId, DateTimeOffset> _terminalAttempts = [];
|
||||
private readonly Queue<JoinAttemptId> _attemptSchedule = [];
|
||||
private readonly SortedDictionary<long, Queue<HostAttemptDeadline>> _deadlines = [];
|
||||
private readonly List<JoinAttemptId> _cleanupScratch = [];
|
||||
private HostJoinAttempt[]? _latestSnapshot;
|
||||
private DateTimeOffset _nextPresenceAt = DateTimeOffset.MinValue;
|
||||
@@ -90,6 +115,7 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
_networkEvents.RendezvousConnectionRequest += OnConnectionRequest;
|
||||
_networkEvents.RendezvousPeerConnected += OnPeerConnected;
|
||||
_networkEvents.RendezvousPeerDisconnected += OnPeerDisconnected;
|
||||
_networkEvents.RendezvousNetworkError += OnNetworkError;
|
||||
_punchEvents.NatIntroductionSuccess += OnNatIntroductionSuccess;
|
||||
}
|
||||
|
||||
@@ -161,7 +187,10 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
ApplySnapshots();
|
||||
if (!_manager.IsRunning)
|
||||
{
|
||||
Stop(RendezvousHostState.ManagerStopped, RendezvousConnectionState.ManagerStopped);
|
||||
Stop(
|
||||
RendezvousHostState.ManagerStopped,
|
||||
RendezvousConnectionState.ManagerStopped,
|
||||
ConnectionOutcomeKind.ManagerStopped);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -174,13 +203,23 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
}
|
||||
|
||||
DateTimeOffset now = _clock.UtcNow;
|
||||
TimeSpan elapsed = _clock.Elapsed;
|
||||
if (!_manager.IsRunning)
|
||||
{
|
||||
Stop(RendezvousHostState.ManagerStopped, RendezvousConnectionState.ManagerStopped);
|
||||
Stop(
|
||||
RendezvousHostState.ManagerStopped,
|
||||
RendezvousConnectionState.ManagerStopped,
|
||||
ConnectionOutcomeKind.ManagerStopped);
|
||||
return;
|
||||
}
|
||||
|
||||
RefreshPresence(now);
|
||||
ProcessDueDeadlines(elapsed);
|
||||
if (State != RendezvousHostState.Active)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
int checks = Math.Min(
|
||||
_attemptSchedule.Count,
|
||||
_options.MaximumAttemptChecksPerPoll);
|
||||
@@ -192,18 +231,22 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
continue;
|
||||
}
|
||||
|
||||
if (now >= attempt.Invitation.ExpiresAt)
|
||||
if (attempt.State != RendezvousConnectionState.Punching)
|
||||
{
|
||||
CompleteAttempt(attemptId, RendezvousConnectionState.TimedOut);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (attempt.State == RendezvousConnectionState.Punching
|
||||
&& attempt.Retry.IsDue(now))
|
||||
if (attempt.Retry.IsDue(elapsed))
|
||||
{
|
||||
if (attempt.Retry.IsExhausted)
|
||||
{
|
||||
CompleteAttempt(attemptId, RendezvousConnectionState.TimedOut);
|
||||
CompleteAttempt(
|
||||
attemptId,
|
||||
RendezvousConnectionState.TimedOut,
|
||||
ConnectionOutcomeKind.PunchTimedOut,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.NatTraversal,
|
||||
RendezvousConnectionPhase.NatTraversal);
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -251,9 +294,13 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
return;
|
||||
}
|
||||
|
||||
Stop(RendezvousHostState.Disposed, RendezvousConnectionState.Disposed);
|
||||
Stop(
|
||||
RendezvousHostState.Disposed,
|
||||
RendezvousConnectionState.Disposed,
|
||||
ConnectionOutcomeKind.Disposed);
|
||||
Interlocked.Exchange(ref _latestSnapshot, null);
|
||||
_attemptSchedule.Clear();
|
||||
_deadlines.Clear();
|
||||
_terminalAttempts.Clear();
|
||||
_cleanupScratch.Clear();
|
||||
_tickets.Dispose();
|
||||
@@ -272,6 +319,7 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
}
|
||||
|
||||
DateTimeOffset now = _clock.UtcNow;
|
||||
TimeSpan elapsed = _clock.Elapsed;
|
||||
foreach (HostJoinAttempt invitation in latest)
|
||||
{
|
||||
if (invitation.AttemptId.Value == Guid.Empty
|
||||
@@ -289,7 +337,11 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
{
|
||||
CompleteAttempt(
|
||||
invitation.AttemptId,
|
||||
RendezvousConnectionState.Cancelled);
|
||||
RendezvousConnectionState.Cancelled,
|
||||
ConnectionOutcomeKind.Cancelled,
|
||||
RendezvousConnectionOutcomeSource.RendezvousService,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
RendezvousConnectionPhase.Authorization);
|
||||
}
|
||||
|
||||
_terminalAttempts[invitation.AttemptId] = invitation.ExpiresAt;
|
||||
@@ -303,11 +355,23 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
continue;
|
||||
}
|
||||
|
||||
TimeSpan attemptDeadline = elapsed + (invitation.ExpiresAt - now);
|
||||
TimeSpan punchDeadline = Min(
|
||||
attemptDeadline,
|
||||
elapsed + _options.PunchTimeout);
|
||||
_attempts.Add(
|
||||
invitation.AttemptId,
|
||||
new PendingHostAttempt(
|
||||
CopyAttempt(invitation),
|
||||
new RendezvousPunchRetrySchedule(_options, _clock)));
|
||||
new RendezvousPunchRetrySchedule(_options, _clock),
|
||||
elapsed,
|
||||
attemptDeadline,
|
||||
punchDeadline));
|
||||
EnqueueDeadline(
|
||||
new HostAttemptDeadline(
|
||||
invitation.AttemptId,
|
||||
RendezvousConnectionState.Punching,
|
||||
punchDeadline));
|
||||
_attemptSchedule.Enqueue(invitation.AttemptId);
|
||||
}
|
||||
}
|
||||
@@ -354,6 +418,13 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
}
|
||||
|
||||
attempt.State = RendezvousConnectionState.Connecting;
|
||||
attempt.DirectDeadline = Min(
|
||||
attempt.AttemptDeadline,
|
||||
_clock.Elapsed + _options.DirectConnectTimeout);
|
||||
EnqueueDeadline(new HostAttemptDeadline(
|
||||
introduction.AttemptId,
|
||||
RendezvousConnectionState.Connecting,
|
||||
attempt.DirectDeadline.Value));
|
||||
if (_deferredRequests.Remove(
|
||||
introduction.AttemptId,
|
||||
out DeferredConnectionRequest? deferred))
|
||||
@@ -410,7 +481,14 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
{
|
||||
if (_acceptedPeers.TryGetValue(peer, out JoinAttemptId attemptId))
|
||||
{
|
||||
CompleteAttempt(attemptId, RendezvousConnectionState.Connected, peer);
|
||||
CompleteAttempt(
|
||||
attemptId,
|
||||
RendezvousConnectionState.Connected,
|
||||
ConnectionOutcomeKind.Connected,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.None,
|
||||
RendezvousConnectionPhase.Complete,
|
||||
peer);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -419,32 +497,113 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
_ = disconnectInfo;
|
||||
if (_acceptedPeers.TryGetValue(peer, out JoinAttemptId attemptId))
|
||||
{
|
||||
CompleteAttempt(attemptId, RendezvousConnectionState.Rejected);
|
||||
ConnectionOutcomeKind kind = disconnectInfo.Reason == DisconnectReason.Timeout
|
||||
? ConnectionOutcomeKind.DirectConnectTimedOut
|
||||
: ConnectionOutcomeKind.TransportError;
|
||||
CompleteAttempt(
|
||||
attemptId,
|
||||
kind == ConnectionOutcomeKind.DirectConnectTimedOut
|
||||
? RendezvousConnectionState.TimedOut
|
||||
: RendezvousConnectionState.Rejected,
|
||||
kind,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.DirectConnection,
|
||||
RendezvousConnectionPhase.DirectConnection);
|
||||
}
|
||||
}
|
||||
|
||||
private void OnNetworkError(IPEndPoint endpoint, SocketError socketError)
|
||||
{
|
||||
_ = socketError;
|
||||
if (!endpoint.Equals(_mediator))
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
foreach (JoinAttemptId attemptId in _attempts
|
||||
.Where(static item => item.Value.State == RendezvousConnectionState.Punching)
|
||||
.Select(static item => item.Key)
|
||||
.ToArray())
|
||||
{
|
||||
CompleteAttempt(
|
||||
attemptId,
|
||||
RendezvousConnectionState.Rejected,
|
||||
ConnectionOutcomeKind.MediatorUnavailable,
|
||||
RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
RendezvousConnectionFailureCategory.Mediation,
|
||||
RendezvousConnectionPhase.Mediation);
|
||||
}
|
||||
}
|
||||
|
||||
private void CompleteAttempt(
|
||||
JoinAttemptId attemptId,
|
||||
RendezvousConnectionState state,
|
||||
ConnectionOutcomeKind kind,
|
||||
RendezvousConnectionOutcomeSource source,
|
||||
RendezvousConnectionFailureCategory category,
|
||||
RendezvousConnectionPhase phase,
|
||||
NetPeer? peer = null)
|
||||
{
|
||||
if (TryCompleteAttempt(
|
||||
attemptId,
|
||||
state,
|
||||
kind,
|
||||
source,
|
||||
category,
|
||||
phase,
|
||||
peer,
|
||||
out RendezvousHostAttemptCompletedEventArgs? completion))
|
||||
{
|
||||
AttemptCompleted?.Invoke(this, completion!);
|
||||
}
|
||||
}
|
||||
|
||||
private bool TryCompleteAttempt(
|
||||
JoinAttemptId attemptId,
|
||||
RendezvousConnectionState state,
|
||||
ConnectionOutcomeKind kind,
|
||||
RendezvousConnectionOutcomeSource source,
|
||||
RendezvousConnectionFailureCategory category,
|
||||
RendezvousConnectionPhase phase,
|
||||
NetPeer? peer,
|
||||
out RendezvousHostAttemptCompletedEventArgs? completion)
|
||||
{
|
||||
completion = null;
|
||||
if (!_attempts.Remove(attemptId, out PendingHostAttempt? attempt))
|
||||
{
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
|
||||
if (attempt.AcceptedPeer is not null)
|
||||
{
|
||||
_acceptedPeers.Remove(attempt.AcceptedPeer);
|
||||
if (kind != ConnectionOutcomeKind.Connected)
|
||||
{
|
||||
attempt.AcceptedPeer.Disconnect();
|
||||
}
|
||||
}
|
||||
|
||||
_deferredRequests.Remove(attemptId);
|
||||
if (_deferredRequests.Remove(attemptId, out DeferredConnectionRequest? deferred))
|
||||
{
|
||||
deferred.Request.RejectForce([]);
|
||||
}
|
||||
_tickets.Revoke(attemptId);
|
||||
_terminalAttempts[attemptId] = attempt.Invitation.ExpiresAt;
|
||||
AttemptCompleted?.Invoke(this, new(attemptId, state, peer));
|
||||
RendezvousConnectionOutcome outcome = RendezvousConnectionOutcome.Create(
|
||||
kind,
|
||||
source,
|
||||
category,
|
||||
phase,
|
||||
_clock.Elapsed - attempt.StartedAt,
|
||||
peer: peer);
|
||||
completion = new(attemptId, state, outcome);
|
||||
return true;
|
||||
}
|
||||
|
||||
private void Stop(RendezvousHostState hostState, RendezvousConnectionState attemptState)
|
||||
private void Stop(
|
||||
RendezvousHostState hostState,
|
||||
RendezvousConnectionState attemptState,
|
||||
ConnectionOutcomeKind outcomeKind)
|
||||
{
|
||||
if (State != RendezvousHostState.Active)
|
||||
{
|
||||
@@ -452,12 +611,32 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
}
|
||||
|
||||
State = hostState;
|
||||
List<RendezvousHostAttemptCompletedEventArgs> completions = [];
|
||||
foreach (JoinAttemptId attemptId in _attempts.Keys.ToArray())
|
||||
{
|
||||
CompleteAttempt(attemptId, attemptState);
|
||||
RendezvousConnectionPhase phase = _attempts[attemptId].State
|
||||
== RendezvousConnectionState.Connecting
|
||||
? RendezvousConnectionPhase.DirectConnection
|
||||
: RendezvousConnectionPhase.NatTraversal;
|
||||
if (TryCompleteAttempt(
|
||||
attemptId,
|
||||
attemptState,
|
||||
outcomeKind,
|
||||
RendezvousConnectionOutcomeSource.Lifecycle,
|
||||
RendezvousConnectionFailureCategory.Lifecycle,
|
||||
phase,
|
||||
null,
|
||||
out RendezvousHostAttemptCompletedEventArgs? completion))
|
||||
{
|
||||
completions.Add(completion!);
|
||||
}
|
||||
}
|
||||
|
||||
ReleaseSubscriptions();
|
||||
foreach (RendezvousHostAttemptCompletedEventArgs completion in completions)
|
||||
{
|
||||
AttemptCompleted?.Invoke(this, completion);
|
||||
}
|
||||
}
|
||||
|
||||
private void ReleaseSubscriptions()
|
||||
@@ -470,6 +649,7 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
_networkEvents.RendezvousConnectionRequest -= OnConnectionRequest;
|
||||
_networkEvents.RendezvousPeerConnected -= OnPeerConnected;
|
||||
_networkEvents.RendezvousPeerDisconnected -= OnPeerDisconnected;
|
||||
_networkEvents.RendezvousNetworkError -= OnNetworkError;
|
||||
_punchEvents.NatIntroductionSuccess -= OnNatIntroductionSuccess;
|
||||
_subscriptionsReleased = true;
|
||||
}
|
||||
@@ -496,9 +676,78 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
ExpiresAt = attempt.ExpiresAt,
|
||||
};
|
||||
|
||||
private static TimeSpan Min(TimeSpan left, TimeSpan right) =>
|
||||
left <= right ? left : right;
|
||||
|
||||
private static DateTimeOffset Min(DateTimeOffset left, DateTimeOffset right) =>
|
||||
left <= right ? left : right;
|
||||
|
||||
private void EnqueueDeadline(HostAttemptDeadline deadline)
|
||||
{
|
||||
if (!_deadlines.TryGetValue(deadline.Deadline.Ticks, out Queue<HostAttemptDeadline>? bucket))
|
||||
{
|
||||
bucket = new Queue<HostAttemptDeadline>();
|
||||
_deadlines.Add(deadline.Deadline.Ticks, bucket);
|
||||
}
|
||||
|
||||
bucket.Enqueue(deadline);
|
||||
}
|
||||
|
||||
private void ProcessDueDeadlines(TimeSpan elapsed)
|
||||
{
|
||||
while (_deadlines.Count > 0)
|
||||
{
|
||||
KeyValuePair<long, Queue<HostAttemptDeadline>> first = _deadlines.First();
|
||||
if (first.Key > elapsed.Ticks)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
HostAttemptDeadline deadline = first.Value.Dequeue();
|
||||
if (first.Value.Count == 0)
|
||||
{
|
||||
_deadlines.Remove(first.Key);
|
||||
}
|
||||
|
||||
if (!_attempts.TryGetValue(deadline.AttemptId, out PendingHostAttempt? attempt)
|
||||
|| attempt.State != deadline.ExpectedState
|
||||
|| (deadline.ExpectedState == RendezvousConnectionState.Punching
|
||||
? attempt.PunchDeadline
|
||||
: attempt.DirectDeadline) != deadline.Deadline)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
bool expired = elapsed >= attempt.AttemptDeadline;
|
||||
CompleteAttempt(
|
||||
deadline.AttemptId,
|
||||
RendezvousConnectionState.TimedOut,
|
||||
expired
|
||||
? ConnectionOutcomeKind.AttemptExpired
|
||||
: deadline.ExpectedState == RendezvousConnectionState.Punching
|
||||
? ConnectionOutcomeKind.PunchTimedOut
|
||||
: ConnectionOutcomeKind.DirectConnectTimedOut,
|
||||
expired
|
||||
? RendezvousConnectionOutcomeSource.RendezvousService
|
||||
: RendezvousConnectionOutcomeSource.LocalTraversal,
|
||||
expired
|
||||
? RendezvousConnectionFailureCategory.Authorization
|
||||
: deadline.ExpectedState == RendezvousConnectionState.Punching
|
||||
? RendezvousConnectionFailureCategory.NatTraversal
|
||||
: RendezvousConnectionFailureCategory.DirectConnection,
|
||||
expired
|
||||
? RendezvousConnectionPhase.Authorization
|
||||
: deadline.ExpectedState == RendezvousConnectionState.Punching
|
||||
? RendezvousConnectionPhase.NatTraversal
|
||||
: RendezvousConnectionPhase.DirectConnection);
|
||||
|
||||
if (State != RendezvousHostState.Active)
|
||||
{
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void AcceptAuthorizedRequest(
|
||||
JoinAttemptId attemptId,
|
||||
PendingHostAttempt attempt,
|
||||
@@ -529,14 +778,31 @@ public sealed class RendezvousHostCoordinator : IDisposable
|
||||
|
||||
private sealed class PendingHostAttempt(
|
||||
HostJoinAttempt invitation,
|
||||
RendezvousPunchRetrySchedule retry)
|
||||
RendezvousPunchRetrySchedule retry,
|
||||
TimeSpan startedAt,
|
||||
TimeSpan attemptDeadline,
|
||||
TimeSpan punchDeadline)
|
||||
{
|
||||
internal HostJoinAttempt Invitation { get; } = invitation;
|
||||
internal RendezvousPunchRetrySchedule Retry { get; } = retry;
|
||||
internal TimeSpan StartedAt { get; } = startedAt;
|
||||
internal TimeSpan AttemptDeadline { get; } = attemptDeadline;
|
||||
internal TimeSpan PunchDeadline { get; } = punchDeadline;
|
||||
internal TimeSpan? DirectDeadline { get; set; }
|
||||
internal RendezvousConnectionState State { get; set; } = RendezvousConnectionState.Punching;
|
||||
internal NetPeer? AcceptedPeer { get; set; }
|
||||
}
|
||||
|
||||
private sealed class HostAttemptDeadline(
|
||||
JoinAttemptId attemptId,
|
||||
RendezvousConnectionState expectedState,
|
||||
TimeSpan deadline)
|
||||
{
|
||||
internal JoinAttemptId AttemptId { get; } = attemptId;
|
||||
internal RendezvousConnectionState ExpectedState { get; } = expectedState;
|
||||
internal TimeSpan Deadline { get; } = deadline;
|
||||
}
|
||||
|
||||
private sealed class DeferredConnectionRequest(
|
||||
ConnectionRequest request,
|
||||
string connectionTicket)
|
||||
|
||||
@@ -29,6 +29,7 @@ public sealed class RendezvousNetListener : INetEventListener
|
||||
internal event Action<NetPeer>? RendezvousPeerConnected;
|
||||
internal event Action<NetPeer, DisconnectInfo>? RendezvousPeerDisconnected;
|
||||
internal event Action<ConnectionRequest>? RendezvousConnectionRequest;
|
||||
internal event Action<IPEndPoint, SocketError>? RendezvousNetworkError;
|
||||
|
||||
internal void ValidateManager(NetManager manager)
|
||||
{
|
||||
@@ -51,8 +52,11 @@ public sealed class RendezvousNetListener : INetEventListener
|
||||
((INetEventListener)GameplayEvents).OnPeerDisconnected(peer, disconnectInfo);
|
||||
}
|
||||
|
||||
public void OnNetworkError(IPEndPoint endPoint, SocketError socketError) =>
|
||||
public void OnNetworkError(IPEndPoint endPoint, SocketError socketError)
|
||||
{
|
||||
RendezvousNetworkError?.Invoke(endPoint, socketError);
|
||||
((INetEventListener)GameplayEvents).OnNetworkError(endPoint, socketError);
|
||||
}
|
||||
|
||||
public void OnNetworkReceive(
|
||||
NetPeer peer,
|
||||
|
||||
@@ -49,6 +49,27 @@ public enum ConnectionOutcomeKind
|
||||
HostRejected = 7,
|
||||
TransportFailed = 8,
|
||||
FallbackOffered = 9,
|
||||
DirectoryNotFound = 10,
|
||||
AttemptExpired = 11,
|
||||
Unauthorized = 12,
|
||||
RateLimited = 13,
|
||||
NoHostPresence = 14,
|
||||
ServiceUnavailable = 15,
|
||||
MediatorUnavailable = 16,
|
||||
PunchTimedOut = 17,
|
||||
DirectConnectTimedOut = 18,
|
||||
TransportError = 19,
|
||||
ManagerStopped = 20,
|
||||
Disposed = 21,
|
||||
}
|
||||
|
||||
public enum ConnectionElapsedBucket
|
||||
{
|
||||
UnderOneSecond = 1,
|
||||
OneToFiveSeconds = 2,
|
||||
FiveToFifteenSeconds = 3,
|
||||
FifteenToThirtySeconds = 4,
|
||||
ThirtySecondsOrMore = 5,
|
||||
}
|
||||
|
||||
public enum UdpPresenceMessageType : byte
|
||||
|
||||
@@ -45,6 +45,23 @@ public static class ContractValidation
|
||||
public static bool IsDiagnosticCodeValid(string? value) =>
|
||||
value is null || IsVisibleAsciiWithin(value, ContractLimits.DiagnosticCodeMaxCharacters);
|
||||
|
||||
public static bool IsReportableConnectionOutcome(ConnectionOutcomeKind outcome) => outcome is
|
||||
ConnectionOutcomeKind.Connected
|
||||
or ConnectionOutcomeKind.Cancelled
|
||||
or ConnectionOutcomeKind.TimedOut
|
||||
or ConnectionOutcomeKind.StaleHost
|
||||
or ConnectionOutcomeKind.TransportFailed
|
||||
or ConnectionOutcomeKind.FallbackOffered
|
||||
or ConnectionOutcomeKind.AttemptExpired
|
||||
or ConnectionOutcomeKind.NoHostPresence
|
||||
or ConnectionOutcomeKind.MediatorUnavailable
|
||||
or ConnectionOutcomeKind.PunchTimedOut
|
||||
or ConnectionOutcomeKind.DirectConnectTimedOut
|
||||
or ConnectionOutcomeKind.HostRejected
|
||||
or ConnectionOutcomeKind.TransportError
|
||||
or ConnectionOutcomeKind.ManagerStopped
|
||||
or ConnectionOutcomeKind.Disposed;
|
||||
|
||||
public static bool IsBuildVersionValid(string? value) =>
|
||||
!string.IsNullOrWhiteSpace(value)
|
||||
&& IsUtf8LengthWithin(value, ContractLimits.BuildVersionMaxBytes);
|
||||
|
||||
@@ -0,0 +1,33 @@
|
||||
using System.Text.Json.Serialization;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Contracts;
|
||||
|
||||
public sealed class ReportConnectionOutcomeRequest
|
||||
{
|
||||
[JsonRequired]
|
||||
public int ContractVersion { get; set; } = ContractLimits.ContractVersion;
|
||||
|
||||
[JsonRequired]
|
||||
public ConnectionOutcomeKind Outcome { get; set; }
|
||||
|
||||
public ConnectionElapsedBucket ElapsedBucket { get; set; }
|
||||
|
||||
[Obsolete("Use ElapsedBucket. Exact elapsed time is accepted only for v1 compatibility and is not retained.")]
|
||||
[JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingDefault)]
|
||||
public int ElapsedMilliseconds { get; set; }
|
||||
|
||||
[Obsolete("Diagnostic codes are accepted only for v1 compatibility and are not retained.")]
|
||||
public string? DiagnosticCode { get; set; }
|
||||
}
|
||||
|
||||
public sealed class ReportConnectionOutcomeResponse
|
||||
{
|
||||
[JsonRequired]
|
||||
public int ContractVersion { get; set; } = ContractLimits.ContractVersion;
|
||||
|
||||
[JsonRequired]
|
||||
public bool Accepted { get; set; }
|
||||
|
||||
[JsonRequired]
|
||||
public bool IsDuplicate { get; set; }
|
||||
}
|
||||
@@ -76,26 +76,3 @@ public sealed class BrowseHostJoinAttemptsResponse
|
||||
|
||||
public string? NextCursor { get; set; }
|
||||
}
|
||||
|
||||
public sealed class ReportConnectionOutcomeRequest
|
||||
{
|
||||
[JsonRequired]
|
||||
public int ContractVersion { get; set; } = ContractLimits.ContractVersion;
|
||||
|
||||
[JsonRequired]
|
||||
public ConnectionOutcomeKind Outcome { get; set; }
|
||||
|
||||
[JsonRequired]
|
||||
public int ElapsedMilliseconds { get; set; }
|
||||
|
||||
public string? DiagnosticCode { get; set; }
|
||||
}
|
||||
|
||||
public sealed class ReportConnectionOutcomeResponse
|
||||
{
|
||||
[JsonRequired]
|
||||
public int ContractVersion { get; set; } = ContractLimits.ContractVersion;
|
||||
|
||||
[JsonRequired]
|
||||
public bool Accepted { get; set; }
|
||||
}
|
||||
|
||||
@@ -39,6 +39,8 @@ public sealed class SessionListing
|
||||
|
||||
[JsonRequired]
|
||||
public Dictionary<string, string> Metadata { get; set; } = new(StringComparer.Ordinal);
|
||||
|
||||
public NetworkEndpoint? DedicatedFallback { get; set; }
|
||||
}
|
||||
|
||||
public sealed class RegisterSessionRequest
|
||||
@@ -75,6 +77,8 @@ public sealed class RegisterSessionRequest
|
||||
|
||||
[JsonRequired]
|
||||
public Dictionary<string, string> Metadata { get; set; } = new(StringComparer.Ordinal);
|
||||
|
||||
public NetworkEndpoint? DedicatedFallback { get; set; }
|
||||
}
|
||||
|
||||
public sealed class RegisterSessionResponse
|
||||
@@ -147,6 +151,8 @@ public sealed class UpdateSessionRequest
|
||||
|
||||
[JsonRequired]
|
||||
public Dictionary<string, string> Metadata { get; set; } = new(StringComparer.Ordinal);
|
||||
|
||||
public NetworkEndpoint? DedicatedFallback { get; set; }
|
||||
}
|
||||
|
||||
public sealed class DeleteSessionRequest
|
||||
|
||||
@@ -25,7 +25,9 @@ public static class ContractJson
|
||||
|
||||
options.AllowTrailingCommas = false;
|
||||
options.DefaultIgnoreCondition = JsonIgnoreCondition.WhenWritingNull;
|
||||
options.MaxDepth = 8;
|
||||
// Nine is the minimum that lets ASP.NET generate the nullable fallback
|
||||
// OpenAPI schema; the 16 KiB HTTP body limit still bounds parser work.
|
||||
options.MaxDepth = 9;
|
||||
options.NumberHandling = JsonNumberHandling.Strict;
|
||||
options.PropertyNameCaseInsensitive = false;
|
||||
options.PropertyNamingPolicy = JsonNamingPolicy.CamelCase;
|
||||
|
||||
@@ -153,5 +153,6 @@ internal sealed class SessionBrowserService(
|
||||
static item => item.Key,
|
||||
static item => item.Value,
|
||||
StringComparer.Ordinal),
|
||||
DedicatedFallback = StoredListing.CopyEndpoint(stored.Definition.DedicatedFallback),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using FinalFactory.Rendezvous.Server.Sessions;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.ConnectionOutcomes;
|
||||
|
||||
internal sealed record ConnectionOutcomeServiceResult(
|
||||
RendezvousErrorCode Error,
|
||||
ReportConnectionOutcomeResponse? Value = null)
|
||||
{
|
||||
public bool Succeeded => Error == RendezvousErrorCode.None;
|
||||
}
|
||||
|
||||
internal sealed class ConnectionOutcomeMetrics
|
||||
{
|
||||
private readonly object _gate = new();
|
||||
private readonly Dictionary<(ConnectionOutcomeKind, ConnectionElapsedBucket), long> _counts = [];
|
||||
|
||||
internal void Record(ConnectionOutcomeKind outcome, ConnectionElapsedBucket elapsedBucket)
|
||||
{
|
||||
lock (_gate)
|
||||
{
|
||||
(ConnectionOutcomeKind, ConnectionElapsedBucket) key = (outcome, elapsedBucket);
|
||||
_counts.TryGetValue(key, out long count);
|
||||
_counts[key] = count + 1;
|
||||
}
|
||||
}
|
||||
|
||||
internal long GetCount(ConnectionOutcomeKind outcome, ConnectionElapsedBucket elapsedBucket)
|
||||
{
|
||||
lock (_gate)
|
||||
{
|
||||
return _counts.GetValueOrDefault((outcome, elapsedBucket));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
internal sealed class ConnectionOutcomeService(
|
||||
IEphemeralRendezvousStore store,
|
||||
ISessionCapabilityService capabilities,
|
||||
ConnectionOutcomeMetrics metrics)
|
||||
{
|
||||
internal ConnectionOutcomeServiceResult Report(
|
||||
JoinAttemptId attemptId,
|
||||
string? clientPunchCapability,
|
||||
ReportConnectionOutcomeRequest request,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
ArgumentNullException.ThrowIfNull(request);
|
||||
RendezvousErrorCode version = ContractValidation.ValidateContractVersion(
|
||||
request.ContractVersion);
|
||||
if (version != RendezvousErrorCode.None)
|
||||
{
|
||||
return new(version);
|
||||
}
|
||||
|
||||
if (attemptId.Value == Guid.Empty
|
||||
|| !ContractValidation.IsCapabilityValid(clientPunchCapability)
|
||||
|| !TryNormalizeReport(request, out ConnectionOutcomeKind outcome, out ConnectionElapsedBucket elapsedBucket)
|
||||
|| !capabilities.TryFingerprint(
|
||||
clientPunchCapability,
|
||||
out SecretFingerprint capabilityFingerprint))
|
||||
{
|
||||
return new(RendezvousErrorCode.InvalidRequest);
|
||||
}
|
||||
|
||||
StoreResult<StoredConnectionOutcome> reported = store.ReportConnectionOutcome(new(
|
||||
attemptId,
|
||||
capabilityFingerprint,
|
||||
outcome,
|
||||
elapsedBucket), cancellationToken);
|
||||
if (!reported.Succeeded)
|
||||
{
|
||||
return new(reported.Code.ToContractError());
|
||||
}
|
||||
|
||||
if (!reported.IsIdempotentReplay)
|
||||
{
|
||||
metrics.Record(outcome, elapsedBucket);
|
||||
}
|
||||
|
||||
return new(RendezvousErrorCode.None, new ReportConnectionOutcomeResponse
|
||||
{
|
||||
Accepted = true,
|
||||
IsDuplicate = reported.IsIdempotentReplay,
|
||||
});
|
||||
}
|
||||
|
||||
private static bool TryNormalizeReport(
|
||||
ReportConnectionOutcomeRequest request,
|
||||
out ConnectionOutcomeKind outcome,
|
||||
out ConnectionElapsedBucket elapsedBucket)
|
||||
{
|
||||
outcome = request.Outcome switch
|
||||
{
|
||||
ConnectionOutcomeKind.TimedOut => ConnectionOutcomeKind.PunchTimedOut,
|
||||
ConnectionOutcomeKind.StaleHost => ConnectionOutcomeKind.NoHostPresence,
|
||||
ConnectionOutcomeKind.TransportFailed => ConnectionOutcomeKind.TransportError,
|
||||
_ => request.Outcome,
|
||||
};
|
||||
if (!ContractValidation.IsReportableConnectionOutcome(request.Outcome))
|
||||
{
|
||||
elapsedBucket = default;
|
||||
return false;
|
||||
}
|
||||
|
||||
if (Enum.IsDefined(request.ElapsedBucket))
|
||||
{
|
||||
elapsedBucket = request.ElapsedBucket;
|
||||
return true;
|
||||
}
|
||||
|
||||
#pragma warning disable CS0618 // Frozen v1 compatibility input; never retained at exact precision.
|
||||
if (request.ElapsedBucket == default && request.ElapsedMilliseconds >= 0)
|
||||
{
|
||||
elapsedBucket = BucketElapsedMilliseconds(request.ElapsedMilliseconds);
|
||||
return true;
|
||||
}
|
||||
#pragma warning restore CS0618
|
||||
|
||||
elapsedBucket = default;
|
||||
return false;
|
||||
}
|
||||
|
||||
private static ConnectionElapsedBucket BucketElapsedMilliseconds(int elapsedMilliseconds) =>
|
||||
elapsedMilliseconds switch
|
||||
{
|
||||
< 1_000 => ConnectionElapsedBucket.UnderOneSecond,
|
||||
< 5_000 => ConnectionElapsedBucket.OneToFiveSeconds,
|
||||
< 15_000 => ConnectionElapsedBucket.FiveToFifteenSeconds,
|
||||
< 30_000 => ConnectionElapsedBucket.FifteenToThirtySeconds,
|
||||
_ => ConnectionElapsedBucket.ThirtySecondsOrMore,
|
||||
};
|
||||
}
|
||||
@@ -1,6 +1,7 @@
|
||||
using System.Net;
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using FinalFactory.Rendezvous.Server.Browser;
|
||||
using FinalFactory.Rendezvous.Server.ConnectionOutcomes;
|
||||
using FinalFactory.Rendezvous.Server.JoinAttempts;
|
||||
using FinalFactory.Rendezvous.Server.Provisioning;
|
||||
using FinalFactory.Rendezvous.Server.Sessions;
|
||||
@@ -11,8 +12,6 @@ namespace FinalFactory.Rendezvous.Server.Http;
|
||||
|
||||
internal static class ContractEndpoints
|
||||
{
|
||||
private const int NotImplementedStatus = StatusCodes.Status501NotImplemented;
|
||||
|
||||
public static IEndpointRouteBuilder MapRendezvousContractEndpoints(
|
||||
this IEndpointRouteBuilder endpoints)
|
||||
{
|
||||
@@ -83,6 +82,7 @@ internal static class ContractEndpoints
|
||||
.Produces<ApiError>(StatusCodes.Status400BadRequest)
|
||||
.Produces<ApiError>(StatusCodes.Status404NotFound)
|
||||
.Produces<ApiError>(StatusCodes.Status409Conflict)
|
||||
.Produces<ApiError>(StatusCodes.Status410Gone)
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.Produces<ApiError>(StatusCodes.Status503ServiceUnavailable)
|
||||
.WithName("CreateJoinAttempt");
|
||||
@@ -95,7 +95,10 @@ internal static class ContractEndpoints
|
||||
attempts.MapPost("/{attemptId}/outcome", ReportConnectionOutcome)
|
||||
.Accepts<ReportConnectionOutcomeRequest>("application/json")
|
||||
.Produces<ReportConnectionOutcomeResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status501NotImplemented)
|
||||
.Produces<ApiError>(StatusCodes.Status400BadRequest)
|
||||
.Produces<ApiError>(StatusCodes.Status404NotFound)
|
||||
.Produces<ApiError>(StatusCodes.Status409Conflict)
|
||||
.Produces<ApiError>(StatusCodes.Status503ServiceUnavailable)
|
||||
.WithName("ReportConnectionOutcome");
|
||||
|
||||
return endpoints;
|
||||
@@ -334,16 +337,20 @@ internal static class ContractEndpoints
|
||||
|
||||
private static IResult ReportConnectionOutcome(
|
||||
JoinAttemptId attemptId,
|
||||
[FromBody] ReportConnectionOutcomeRequest request) => NotImplemented();
|
||||
|
||||
private static IResult NotImplemented() => Results.Json(
|
||||
new ApiError
|
||||
{
|
||||
Code = RendezvousErrorCode.ServiceUnavailable,
|
||||
Message = "The v1 contract is reserved; implementation is tracked by subsequent issues.",
|
||||
},
|
||||
ContractJson.Options,
|
||||
statusCode: NotImplementedStatus);
|
||||
[FromHeader(Name = "X-Rendezvous-Client-Punch-Capability")] string clientPunchCapability,
|
||||
[FromBody] ReportConnectionOutcomeRequest request,
|
||||
[FromServices] ConnectionOutcomeService outcomes,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
ConnectionOutcomeServiceResult result = outcomes.Report(
|
||||
attemptId,
|
||||
clientPunchCapability,
|
||||
request,
|
||||
cancellationToken);
|
||||
return result.Succeeded && result.Value is not null
|
||||
? Results.Ok(result.Value)
|
||||
: Error(result.Error);
|
||||
}
|
||||
|
||||
private static bool TryAuthenticatePublisher(
|
||||
string? authorizationHeader,
|
||||
@@ -389,9 +396,11 @@ internal static class ContractEndpoints
|
||||
{
|
||||
RendezvousErrorCode.AuthenticationRequired => StatusCodes.Status401Unauthorized,
|
||||
RendezvousErrorCode.Forbidden => StatusCodes.Status403Forbidden,
|
||||
RendezvousErrorCode.NotFound or RendezvousErrorCode.StaleHost => StatusCodes.Status404NotFound,
|
||||
RendezvousErrorCode.Conflict or RendezvousErrorCode.ReplayRejected => StatusCodes.Status409Conflict,
|
||||
RendezvousErrorCode.Expired => StatusCodes.Status410Gone,
|
||||
RendezvousErrorCode.NotFound => StatusCodes.Status404NotFound,
|
||||
RendezvousErrorCode.Conflict
|
||||
or RendezvousErrorCode.IncompatibleProtocol
|
||||
or RendezvousErrorCode.ReplayRejected => StatusCodes.Status409Conflict,
|
||||
RendezvousErrorCode.Expired or RendezvousErrorCode.StaleHost => StatusCodes.Status410Gone,
|
||||
RendezvousErrorCode.RateLimited or RendezvousErrorCode.CapacityExceeded =>
|
||||
StatusCodes.Status429TooManyRequests,
|
||||
RendezvousErrorCode.ServiceUnavailable => StatusCodes.Status503ServiceUnavailable,
|
||||
@@ -406,6 +415,7 @@ internal static class ContractEndpoints
|
||||
RendezvousErrorCode.NotFound => "The session was not found or is not owned by this publisher.",
|
||||
RendezvousErrorCode.Conflict => "The session changed concurrently; retry with current state.",
|
||||
RendezvousErrorCode.Expired => "The session lease has expired.",
|
||||
RendezvousErrorCode.StaleHost => "The session has no fresh host presence.",
|
||||
RendezvousErrorCode.IncompatibleProtocol => "The gameplay protocol is not enabled for this game.",
|
||||
RendezvousErrorCode.CapacityExceeded => "The configured session capacity is currently exhausted.",
|
||||
RendezvousErrorCode.ServiceUnavailable => "Session state is temporarily unavailable.",
|
||||
|
||||
@@ -131,6 +131,7 @@ internal sealed class JoinAttemptService(
|
||||
ConnectionTicketDigest = NatIntroductionTokenCodec.ComputeDigest(
|
||||
CreateConnectionTicket(persisted)),
|
||||
ExpiresAt = persisted.ExpiresAt,
|
||||
DedicatedFallback = StoredListing.CopyEndpoint(persisted.DedicatedFallback),
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
using System.Net;
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using FinalFactory.Rendezvous.Server.Browser;
|
||||
using FinalFactory.Rendezvous.Server.ConnectionOutcomes;
|
||||
using FinalFactory.Rendezvous.Server.Http;
|
||||
using FinalFactory.Rendezvous.Server.JoinAttempts;
|
||||
using FinalFactory.Rendezvous.Server.Provisioning;
|
||||
@@ -50,6 +51,14 @@ builder.Services.AddOpenApi("v1", static options =>
|
||||
BearerFormat = "rv1 publisher credential",
|
||||
Description = "Tenant-scoped publisher credential issued during game provisioning.",
|
||||
};
|
||||
const string attemptSchemeName = "JoinAttemptCapability";
|
||||
document.Components.SecuritySchemes[attemptSchemeName] = new OpenApiSecurityScheme
|
||||
{
|
||||
Type = SecuritySchemeType.ApiKey,
|
||||
Name = "X-Rendezvous-Client-Punch-Capability",
|
||||
In = ParameterLocation.Header,
|
||||
Description = "Attempt-scoped client capability returned only to the joining caller.",
|
||||
};
|
||||
|
||||
HashSet<string> securedOperations = new(StringComparer.Ordinal)
|
||||
{
|
||||
@@ -59,6 +68,7 @@ builder.Services.AddOpenApi("v1", static options =>
|
||||
"DeleteSession",
|
||||
};
|
||||
OpenApiSecuritySchemeReference reference = new(schemeName, document, null);
|
||||
OpenApiSecuritySchemeReference attemptReference = new(attemptSchemeName, document, null);
|
||||
foreach (OpenApiPathItem path in document.Paths.Values)
|
||||
{
|
||||
if (path.Operations is null)
|
||||
@@ -75,6 +85,17 @@ builder.Services.AddOpenApi("v1", static options =>
|
||||
[reference] = [],
|
||||
});
|
||||
}
|
||||
|
||||
foreach (OpenApiOperation operation in path.Operations.Values.Where(
|
||||
operation => operation.OperationId is
|
||||
"CancelJoinAttempt" or "ReportConnectionOutcome"))
|
||||
{
|
||||
operation.Security ??= [];
|
||||
operation.Security.Add(new OpenApiSecurityRequirement
|
||||
{
|
||||
[attemptReference] = [],
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
return Task.CompletedTask;
|
||||
@@ -125,6 +146,8 @@ else
|
||||
builder.Services.AddSingleton<SessionBrowserService>();
|
||||
builder.Services.AddSingleton<JoinAttemptCursorCodec>();
|
||||
builder.Services.AddSingleton<JoinAttemptService>();
|
||||
builder.Services.AddSingleton<ConnectionOutcomeMetrics>();
|
||||
builder.Services.AddSingleton<ConnectionOutcomeService>();
|
||||
builder.Services.AddSingleton(new ProvisioningReadiness(true));
|
||||
}
|
||||
|
||||
|
||||
@@ -55,6 +55,11 @@ internal sealed class SessionLeaseService(
|
||||
}
|
||||
|
||||
AuthorizedPublisherContext context = authorized.Context;
|
||||
if (!IsFallbackAllowed(context.Policy, request.DedicatedFallback))
|
||||
{
|
||||
return new(RendezvousErrorCode.Forbidden);
|
||||
}
|
||||
|
||||
string requestFingerprint = ComputeRegistrationFingerprint(request);
|
||||
string derivationSalt = capabilities.CreateDerivationSalt();
|
||||
string leaseToken = capabilities.DeriveCapability(
|
||||
@@ -119,6 +124,7 @@ internal sealed class SessionLeaseService(
|
||||
CurrentPlayers = request.Capacity.CurrentPlayers,
|
||||
MaximumPlayers = request.Capacity.MaximumPlayers,
|
||||
Metadata = request.Metadata,
|
||||
DedicatedFallback = request.DedicatedFallback,
|
||||
LeaseFingerprint = leaseFingerprint,
|
||||
HostPresenceHandle = presenceHandle,
|
||||
HostPresenceFingerprint = presenceFingerprint,
|
||||
@@ -235,10 +241,14 @@ internal sealed class SessionLeaseService(
|
||||
|
||||
StoredListing ownedListing = listing!;
|
||||
PublisherAuthorizationResult authorized = AuthorizeExisting(principal, ownedListing, request.Metadata);
|
||||
if (!authorized.IsAllowed)
|
||||
if (!authorized.IsAllowed || authorized.Context is null)
|
||||
{
|
||||
return new(MapAuthorization(authorized.Error));
|
||||
}
|
||||
if (!IsFallbackAllowed(authorized.Context.Policy, request.DedicatedFallback))
|
||||
{
|
||||
return new(RendezvousErrorCode.Forbidden);
|
||||
}
|
||||
|
||||
capabilities.TryFingerprint(request.LeaseToken, out SecretFingerprint fingerprint);
|
||||
StoreResult<StoredListing> updated = store.UpdateListing(new(
|
||||
@@ -250,7 +260,8 @@ internal sealed class SessionLeaseService(
|
||||
request.DisplayName,
|
||||
request.Capacity.CurrentPlayers,
|
||||
request.Capacity.MaximumPlayers,
|
||||
request.Metadata), cancellationToken);
|
||||
request.Metadata,
|
||||
request.DedicatedFallback), cancellationToken);
|
||||
return updated.Succeeded
|
||||
? new(RendezvousErrorCode.None, true)
|
||||
: new(updated.Code.ToContractError());
|
||||
@@ -341,6 +352,9 @@ internal sealed class SessionLeaseService(
|
||||
metadata,
|
||||
clock.UtcNow);
|
||||
|
||||
private static bool IsFallbackAllowed(GamePolicy policy, NetworkEndpoint? fallback) =>
|
||||
fallback is null || policy.FallbackPolicy == FallbackPolicyMode.DedicatedEndpointAllowed;
|
||||
|
||||
private static RendezvousErrorCode ValidateRegistration(RegisterSessionRequest request)
|
||||
{
|
||||
RendezvousErrorCode version = ContractValidation.ValidateContractVersion(request.ContractVersion);
|
||||
@@ -359,6 +373,8 @@ internal sealed class SessionLeaseService(
|
||||
|| !Enum.IsDefined(request.Visibility)
|
||||
|| !ContractValidation.IsCapacityValid(request.Capacity)
|
||||
|| !ContractValidation.IsMetadataValid(request.Metadata)
|
||||
|| request.DedicatedFallback is not null
|
||||
&& !ContractValidation.IsNetworkEndpointValid(request.DedicatedFallback)
|
||||
? RendezvousErrorCode.InvalidRequest
|
||||
: RendezvousErrorCode.None;
|
||||
}
|
||||
@@ -375,6 +391,8 @@ internal sealed class SessionLeaseService(
|
||||
|| !ContractValidation.IsDisplayNameValid(request.DisplayName)
|
||||
|| !ContractValidation.IsCapacityValid(request.Capacity)
|
||||
|| !ContractValidation.IsMetadataValid(request.Metadata)
|
||||
|| request.DedicatedFallback is not null
|
||||
&& !ContractValidation.IsNetworkEndpointValid(request.DedicatedFallback)
|
||||
? RendezvousErrorCode.InvalidRequest
|
||||
: RendezvousErrorCode.None;
|
||||
}
|
||||
@@ -424,6 +442,14 @@ internal sealed class SessionLeaseService(
|
||||
Metadata = request.Metadata
|
||||
.OrderBy(static item => item.Key, StringComparer.Ordinal)
|
||||
.ToDictionary(static item => item.Key, static item => item.Value, StringComparer.Ordinal),
|
||||
DedicatedFallback = request.DedicatedFallback is null
|
||||
? null
|
||||
: new NetworkEndpoint
|
||||
{
|
||||
AddressFamily = request.DedicatedFallback.AddressFamily,
|
||||
Address = request.DedicatedFallback.Address,
|
||||
Port = request.DedicatedFallback.Port,
|
||||
},
|
||||
};
|
||||
byte[] encoded = JsonSerializer.SerializeToUtf8Bytes(canonical, ContractJson.Options);
|
||||
byte[] digest = SHA256.HashData(encoded);
|
||||
|
||||
@@ -29,6 +29,7 @@ internal sealed record EphemeralStoreOptions
|
||||
public int MaxListings { get; init; } = 25_000;
|
||||
public int MaxPresenceBindings { get; init; } = 25_000;
|
||||
public int MaxJoinAttempts { get; init; } = 10_000;
|
||||
public int MaxOutcomeReports { get; init; } = 35_000;
|
||||
public int MaxReplayEntries { get; init; } = 30_000;
|
||||
public int MaxRevocations { get; init; } = 10_000;
|
||||
public int MaxIdempotencyEntries { get; init; } = 35_000;
|
||||
@@ -45,6 +46,7 @@ internal sealed record EphemeralStoreOptions
|
||||
RequirePositive(MaxListings, nameof(MaxListings));
|
||||
RequirePositive(MaxPresenceBindings, nameof(MaxPresenceBindings));
|
||||
RequirePositive(MaxJoinAttempts, nameof(MaxJoinAttempts));
|
||||
RequirePositive(MaxOutcomeReports, nameof(MaxOutcomeReports));
|
||||
RequirePositive(MaxReplayEntries, nameof(MaxReplayEntries));
|
||||
RequirePositive(MaxRevocations, nameof(MaxRevocations));
|
||||
RequirePositive(MaxIdempotencyEntries, nameof(MaxIdempotencyEntries));
|
||||
@@ -173,6 +175,7 @@ internal sealed record ListingDefinition
|
||||
public required int CurrentPlayers { get; init; }
|
||||
public required int MaximumPlayers { get; init; }
|
||||
public required IReadOnlyDictionary<string, string> Metadata { get; init; }
|
||||
public NetworkEndpoint? DedicatedFallback { get; init; }
|
||||
public required SecretFingerprint LeaseFingerprint { get; init; }
|
||||
public required MediationHandle HostPresenceHandle { get; init; }
|
||||
public required SecretFingerprint HostPresenceFingerprint { get; init; }
|
||||
@@ -189,7 +192,17 @@ internal sealed record StoredListing
|
||||
public static ListingDefinition Freeze(ListingDefinition source) => source with
|
||||
{
|
||||
Metadata = source.Metadata.ToFrozenDictionary(StringComparer.Ordinal),
|
||||
DedicatedFallback = CopyEndpoint(source.DedicatedFallback),
|
||||
};
|
||||
|
||||
internal static NetworkEndpoint? CopyEndpoint(NetworkEndpoint? endpoint) => endpoint is null
|
||||
? null
|
||||
: new NetworkEndpoint
|
||||
{
|
||||
AddressFamily = endpoint.AddressFamily,
|
||||
Address = endpoint.Address,
|
||||
Port = endpoint.Port,
|
||||
};
|
||||
}
|
||||
|
||||
internal sealed record CreateListingCommand(
|
||||
@@ -214,7 +227,8 @@ internal sealed record UpdateListingCommand(
|
||||
string DisplayName,
|
||||
int CurrentPlayers,
|
||||
int MaximumPlayers,
|
||||
IReadOnlyDictionary<string, string> Metadata);
|
||||
IReadOnlyDictionary<string, string> Metadata,
|
||||
NetworkEndpoint? DedicatedFallback);
|
||||
|
||||
internal sealed record DeleteListingCommand(
|
||||
SessionListingId ListingId,
|
||||
@@ -257,6 +271,7 @@ internal sealed record CreateJoinAttemptCommand
|
||||
public required SecretFingerprint ClientCapabilityFingerprint { get; init; }
|
||||
public required SecretFingerprint ConnectionTicketFingerprint { get; init; }
|
||||
public required string CapabilityDerivationSalt { get; init; }
|
||||
public NetworkEndpoint? DedicatedFallback { get; init; }
|
||||
public int ScopeAttemptLimit { get; init; } = int.MaxValue;
|
||||
|
||||
public override string ToString() => "[CreateJoinAttemptCommand: credentials redacted]";
|
||||
@@ -280,6 +295,7 @@ internal sealed record StoredJoinAttempt
|
||||
public required SecretFingerprint HostCapabilityFingerprint { get; init; }
|
||||
public required SecretFingerprint ClientCapabilityFingerprint { get; init; }
|
||||
public required SecretFingerprint ConnectionTicketFingerprint { get; init; }
|
||||
public NetworkEndpoint? DedicatedFallback { get; init; }
|
||||
public required DateTimeOffset ExpiresAt { get; init; }
|
||||
public required DateTimeOffset ConnectionTicketExpiresAt { get; init; }
|
||||
public AttemptEndpointBinding? HostEndpoint { get; init; }
|
||||
@@ -316,6 +332,16 @@ internal sealed record CancelJoinAttemptCommand(
|
||||
JoinAttemptId AttemptId,
|
||||
SecretFingerprint ClientCapabilityFingerprint);
|
||||
|
||||
internal sealed record ReportConnectionOutcomeCommand(
|
||||
JoinAttemptId AttemptId,
|
||||
SecretFingerprint ClientCapabilityFingerprint,
|
||||
ConnectionOutcomeKind Outcome,
|
||||
ConnectionElapsedBucket ElapsedBucket);
|
||||
|
||||
internal sealed record StoredConnectionOutcome(
|
||||
ConnectionOutcomeKind Outcome,
|
||||
ConnectionElapsedBucket ElapsedBucket);
|
||||
|
||||
internal sealed record ConsumeConnectionTicketCommand(
|
||||
JoinAttemptId AttemptId,
|
||||
SecretFingerprint ConnectionTicketFingerprint);
|
||||
@@ -336,6 +362,8 @@ internal enum StoreResultCode
|
||||
Draining = 6,
|
||||
ReplayRejected = 7,
|
||||
ServiceUnavailable = 8,
|
||||
StaleHost = 9,
|
||||
IncompatibleProtocol = 10,
|
||||
}
|
||||
|
||||
internal sealed record StoreResult<T>(StoreResultCode Code, T? Value = default, bool IsIdempotentReplay = false)
|
||||
@@ -359,6 +387,7 @@ internal interface IEphemeralRendezvousStore
|
||||
StoreResult<StoredJoinAttempt> CreateJoinAttempt(CreateJoinAttemptCommand command, CancellationToken cancellationToken = default);
|
||||
StoreResult<IReadOnlyList<StoredJoinAttempt>> BrowseHostJoinAttempts(HostJoinAttemptQuery query, CancellationToken cancellationToken = default);
|
||||
StoreResult<bool> CancelJoinAttempt(CancelJoinAttemptCommand command, CancellationToken cancellationToken = default);
|
||||
StoreResult<StoredConnectionOutcome> ReportConnectionOutcome(ReportConnectionOutcomeCommand command, CancellationToken cancellationToken = default);
|
||||
StoreResult<StoredJoinAttempt> BindAttemptEndpoint(BindAttemptEndpointCommand command, CancellationToken cancellationToken = default);
|
||||
StoreResult<IntroductionEndpoints> ConsumeIntroduction(MediationHandle handle, CancellationToken cancellationToken = default);
|
||||
StoreResult<bool> ConsumeConnectionTicket(ConsumeConnectionTicketCommand command, CancellationToken cancellationToken = default);
|
||||
|
||||
@@ -15,6 +15,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
private readonly Dictionary<MediationHandle, SessionListingId> _presenceHandles = [];
|
||||
private readonly Dictionary<MediationHandle, PresenceEntry> _presence = [];
|
||||
private readonly Dictionary<JoinAttemptId, AttemptEntry> _attempts = [];
|
||||
private readonly Dictionary<JoinAttemptId, OutcomeReportEntry> _outcomeReports = [];
|
||||
private readonly Dictionary<MediationHandle, JoinAttemptId> _attemptHandles = [];
|
||||
private readonly Dictionary<string, IdempotencyEntry> _idempotency = new(StringComparer.Ordinal);
|
||||
private readonly Dictionary<string, TimeSpan> _replay = new(StringComparer.Ordinal);
|
||||
@@ -183,7 +184,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
|| command.MaximumPlayers is <= 0 or > ContractLimits.SessionCapacityMaxPlayers
|
||||
|| command.CurrentPlayers < 0
|
||||
|| command.CurrentPlayers > command.MaximumPlayers
|
||||
|| !ContractValidation.IsMetadataValid(command.Metadata))
|
||||
|| !ContractValidation.IsMetadataValid(command.Metadata)
|
||||
|| command.DedicatedFallback is not null
|
||||
&& !ContractValidation.IsNetworkEndpointValid(command.DedicatedFallback))
|
||||
{
|
||||
throw new ArgumentException("Listing update invariants are invalid.", nameof(command));
|
||||
}
|
||||
@@ -213,6 +216,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
CurrentPlayers = command.CurrentPlayers,
|
||||
MaximumPlayers = command.MaximumPlayers,
|
||||
Metadata = command.Metadata,
|
||||
DedicatedFallback = command.DedicatedFallback,
|
||||
});
|
||||
entry.Version++;
|
||||
return new(StoreResultCode.Success, Snapshot(entry));
|
||||
@@ -362,14 +366,26 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
}
|
||||
|
||||
if (!_listings.TryGetValue(command.ListingId, out ListingEntry? listing)
|
||||
|| listing.Definition.Scope != command.Scope
|
||||
|| listing.Definition.ProtocolVersion != command.ProtocolVersion
|
||||
|| !_presence.ContainsKey(listing.Definition.HostPresenceHandle))
|
||||
|| listing.Definition.Scope != command.Scope)
|
||||
{
|
||||
return new(StoreResultCode.NotFound);
|
||||
}
|
||||
if (listing.Definition.ProtocolVersion != command.ProtocolVersion)
|
||||
{
|
||||
return new(StoreResultCode.IncompatibleProtocol);
|
||||
}
|
||||
if (!_presence.ContainsKey(listing.Definition.HostPresenceHandle))
|
||||
{
|
||||
return new(StoreResultCode.StaleHost);
|
||||
}
|
||||
|
||||
command = command with
|
||||
{
|
||||
DedicatedFallback = StoredListing.CopyEndpoint(listing.Definition.DedicatedFallback),
|
||||
};
|
||||
|
||||
if (_attempts.Count >= _options.MaxJoinAttempts
|
||||
|| _outcomeReports.Count >= _options.MaxOutcomeReports
|
||||
|| _idempotency.Count >= _options.MaxIdempotencyEntries
|
||||
|| _attempts.Values.Count(entry => entry.Command.Scope == command.Scope)
|
||||
>= command.ScopeAttemptLimit)
|
||||
@@ -387,6 +403,11 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
now + _options.JoinAttemptLifetime,
|
||||
WallDeadline(now, _options.JoinAttemptLifetime));
|
||||
_attempts.Add(command.AttemptId, attempt);
|
||||
_outcomeReports.Add(command.AttemptId, new(
|
||||
command.ListingId,
|
||||
command.ClientSubject,
|
||||
command.ClientCapabilityFingerprint,
|
||||
now + _options.JoinAttemptLifetime + _options.IdempotencyLifetime));
|
||||
_attemptHandles.Add(command.MediationHandle, command.AttemptId);
|
||||
_idempotency.Add(idempotencyKey, new(
|
||||
command.RequestFingerprint,
|
||||
@@ -460,6 +481,42 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
return new(StoreResultCode.Success, true);
|
||||
}, cancellationToken);
|
||||
|
||||
public StoreResult<StoredConnectionOutcome> ReportConnectionOutcome(
|
||||
ReportConnectionOutcomeCommand command,
|
||||
CancellationToken cancellationToken = default) => Atomic<StoredConnectionOutcome>(_ =>
|
||||
{
|
||||
ArgumentNullException.ThrowIfNull(command);
|
||||
if (command.AttemptId.Value == Guid.Empty
|
||||
|| !command.ClientCapabilityFingerprint.IsValid
|
||||
|| !ContractValidation.IsReportableConnectionOutcome(command.Outcome)
|
||||
|| !Enum.IsDefined(command.ElapsedBucket))
|
||||
{
|
||||
throw new ArgumentException("Connection outcome invariants are invalid.", nameof(command));
|
||||
}
|
||||
|
||||
if (!_available)
|
||||
{
|
||||
return new(StoreResultCode.ServiceUnavailable);
|
||||
}
|
||||
|
||||
if (!_outcomeReports.TryGetValue(command.AttemptId, out OutcomeReportEntry? entry)
|
||||
|| entry.ClientCapabilityFingerprint != command.ClientCapabilityFingerprint)
|
||||
{
|
||||
return new(StoreResultCode.NotFound);
|
||||
}
|
||||
|
||||
StoredConnectionOutcome reported = new(command.Outcome, command.ElapsedBucket);
|
||||
if (entry.Outcome is not null)
|
||||
{
|
||||
return entry.Outcome == reported
|
||||
? new(StoreResultCode.Success, entry.Outcome, true)
|
||||
: new(StoreResultCode.ReplayRejected);
|
||||
}
|
||||
|
||||
entry.Outcome = reported;
|
||||
return new(StoreResultCode.Success, reported);
|
||||
}, cancellationToken);
|
||||
|
||||
public StoreResult<StoredJoinAttempt> BindAttemptEndpoint(
|
||||
BindAttemptEndpointCommand command,
|
||||
CancellationToken cancellationToken = default) => Atomic<StoredJoinAttempt>(now =>
|
||||
@@ -681,6 +738,10 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
.Where(item => string.Equals(item.Value.Command.ClientSubject, subject, StringComparison.Ordinal))
|
||||
.Select(static item => item.Key)
|
||||
.ToArray();
|
||||
JoinAttemptId[] outcomeReports = _outcomeReports
|
||||
.Where(item => string.Equals(item.Value.ClientSubject, subject, StringComparison.Ordinal))
|
||||
.Select(static item => item.Key)
|
||||
.ToArray();
|
||||
foreach (SessionListingId listingId in listings)
|
||||
{
|
||||
RemoveListing(listingId);
|
||||
@@ -690,6 +751,10 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
{
|
||||
RemoveAttempt(attemptId);
|
||||
}
|
||||
foreach (JoinAttemptId attemptId in outcomeReports)
|
||||
{
|
||||
_outcomeReports.Remove(attemptId);
|
||||
}
|
||||
|
||||
return new(StoreResultCode.Success, listings.Length + attempts.Length);
|
||||
}, cancellationToken);
|
||||
@@ -799,6 +864,14 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
RemoveAttempt(attemptId);
|
||||
}
|
||||
|
||||
foreach (JoinAttemptId attemptId in _outcomeReports
|
||||
.Where(item => item.Value.Deadline <= now)
|
||||
.Select(static item => item.Key)
|
||||
.ToArray())
|
||||
{
|
||||
_outcomeReports.Remove(attemptId);
|
||||
}
|
||||
|
||||
foreach (SessionListingId listingId in _listings
|
||||
.Where(item => item.Value.LeaseDeadline <= now)
|
||||
.Select(static item => item.Key)
|
||||
@@ -815,6 +888,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
_presenceHandles.Clear();
|
||||
_presence.Clear();
|
||||
_attempts.Clear();
|
||||
_outcomeReports.Clear();
|
||||
_attemptHandles.Clear();
|
||||
_idempotency.Clear();
|
||||
_replay.Clear();
|
||||
@@ -837,6 +911,15 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
{
|
||||
RemoveAttempt(attemptId);
|
||||
}
|
||||
|
||||
|
||||
foreach (JoinAttemptId attemptId in _outcomeReports
|
||||
.Where(item => item.Value.ListingId == listingId)
|
||||
.Select(static item => item.Key)
|
||||
.ToArray())
|
||||
{
|
||||
_outcomeReports.Remove(attemptId);
|
||||
}
|
||||
}
|
||||
|
||||
private void RemoveAttempt(JoinAttemptId attemptId)
|
||||
@@ -875,6 +958,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
HostCapabilityFingerprint = entry.Command.HostCapabilityFingerprint,
|
||||
ClientCapabilityFingerprint = entry.Command.ClientCapabilityFingerprint,
|
||||
ConnectionTicketFingerprint = entry.Command.ConnectionTicketFingerprint,
|
||||
DedicatedFallback = StoredListing.CopyEndpoint(entry.Command.DedicatedFallback),
|
||||
ExpiresAt = entry.WallExpiresAt,
|
||||
ConnectionTicketExpiresAt = entry.TicketWallExpiresAt ?? default,
|
||||
HostEndpoint = entry.HostEndpoint,
|
||||
@@ -914,6 +998,8 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
|| listing.CurrentPlayers < 0
|
||||
|| listing.CurrentPlayers > listing.MaximumPlayers
|
||||
|| !ContractValidation.IsMetadataValid(listing.Metadata)
|
||||
|| listing.DedicatedFallback is not null
|
||||
&& !ContractValidation.IsNetworkEndpointValid(listing.DedicatedFallback)
|
||||
|| !listing.LeaseFingerprint.IsValid
|
||||
|| !listing.HostPresenceFingerprint.IsValid
|
||||
|| !IsDerivationSaltValid(listing.CapabilityDerivationSalt))
|
||||
@@ -954,6 +1040,8 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
|| !command.HostCapabilityFingerprint.IsValid
|
||||
|| !command.ClientCapabilityFingerprint.IsValid
|
||||
|| !command.ConnectionTicketFingerprint.IsValid
|
||||
|| command.DedicatedFallback is not null
|
||||
&& !ContractValidation.IsNetworkEndpointValid(command.DedicatedFallback)
|
||||
|| !IsDerivationSaltValid(command.CapabilityDerivationSalt)
|
||||
|| command.ScopeAttemptLimit <= 0)
|
||||
{
|
||||
@@ -1031,6 +1119,19 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
public bool IsCancelled { get; set; }
|
||||
}
|
||||
|
||||
private sealed class OutcomeReportEntry(
|
||||
SessionListingId listingId,
|
||||
string clientSubject,
|
||||
SecretFingerprint clientCapabilityFingerprint,
|
||||
TimeSpan deadline)
|
||||
{
|
||||
public SessionListingId ListingId { get; } = listingId;
|
||||
public string ClientSubject { get; } = clientSubject;
|
||||
public SecretFingerprint ClientCapabilityFingerprint { get; } = clientCapabilityFingerprint;
|
||||
public TimeSpan Deadline { get; } = deadline;
|
||||
public StoredConnectionOutcome? Outcome { get; set; }
|
||||
}
|
||||
|
||||
private sealed record IdempotencyEntry(
|
||||
string RequestFingerprint,
|
||||
object ResourceId,
|
||||
|
||||
@@ -13,6 +13,8 @@ internal static class StoreResultMapping
|
||||
StoreResultCode.Conflict => RendezvousErrorCode.Conflict,
|
||||
StoreResultCode.CapacityExceeded => RendezvousErrorCode.CapacityExceeded,
|
||||
StoreResultCode.ReplayRejected => RendezvousErrorCode.ReplayRejected,
|
||||
StoreResultCode.StaleHost => RendezvousErrorCode.StaleHost,
|
||||
StoreResultCode.IncompatibleProtocol => RendezvousErrorCode.IncompatibleProtocol,
|
||||
StoreResultCode.Draining or StoreResultCode.ServiceUnavailable =>
|
||||
RendezvousErrorCode.ServiceUnavailable,
|
||||
_ => RendezvousErrorCode.InternalError,
|
||||
|
||||
Reference in New Issue
Block a user