Files
Rendezvous/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs
T
KyuubiYoru 94aba8a3bb
quality-gate / quality (push) Successful in 59s
feat(client): standardize connection outcomes (#13)
2026-07-16 10:18:41 +02:00

1140 lines
43 KiB
C#

using FinalFactory.Rendezvous.Contracts;
namespace FinalFactory.Rendezvous.Server.State;
internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousStore
{
private static readonly TimeSpan UdpMaintenanceInterval = TimeSpan.FromSeconds(1);
private readonly object _gate = new();
private readonly EphemeralStoreOptions _options;
private readonly IMonotonicClock _monotonicClock;
private readonly DateTimeOffset _wallOrigin;
private readonly TimeSpan _monotonicOrigin;
private readonly Dictionary<SessionListingId, ListingEntry> _listings = [];
private readonly Dictionary<LeaseId, SessionListingId> _leases = [];
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);
private readonly Dictionary<string, TimeSpan> _revocations = new(StringComparer.Ordinal);
private TimeSpan? _drainDeadline;
private TimeSpan _nextUdpMaintenance;
private long _maintenanceSweepCount;
private bool _available = true;
public InMemoryEphemeralRendezvousStore(
EphemeralStoreOptions options,
IWallClock wallClock,
IMonotonicClock monotonicClock)
{
ArgumentNullException.ThrowIfNull(options);
ArgumentNullException.ThrowIfNull(wallClock);
ArgumentNullException.ThrowIfNull(monotonicClock);
options.Validate();
_options = options;
_monotonicClock = monotonicClock;
_wallOrigin = wallClock.UtcNow;
_monotonicOrigin = monotonicClock.Elapsed;
InstanceId = Guid.NewGuid();
}
public Guid InstanceId { get; }
internal long MaintenanceSweepCount => Interlocked.Read(ref _maintenanceSweepCount);
public bool IsAvailable
{
get
{
lock (_gate)
{
return _available;
}
}
}
public bool IsDraining
{
get
{
lock (_gate)
{
return _drainDeadline.HasValue;
}
}
}
public StoreResult<StoredListing> CreateListing(
CreateListingCommand command,
CancellationToken cancellationToken = default) => Atomic<StoredListing>(now =>
{
ArgumentNullException.ThrowIfNull(command);
ValidateListing(command.Listing);
if (command.OwnerListingLimit <= 0)
{
throw new ArgumentOutOfRangeException(nameof(command), "Owner listing limit must be positive.");
}
ValidateIdempotency(command.IdempotencyKey, command.RequestFingerprint);
StoreResult<StoredListing>? admission = CheckNewWorkAdmission<StoredListing>(command.Listing.OwnerSubject);
if (admission is not null)
{
return admission;
}
string idempotencyKey = $"listing:{command.Listing.Scope.GameId}:{command.Listing.Scope.EnvironmentId}:{command.Listing.OwnerSubject}:{command.IdempotencyKey}";
if (_idempotency.TryGetValue(idempotencyKey, out IdempotencyEntry? previous))
{
if (!string.Equals(previous.RequestFingerprint, command.RequestFingerprint, StringComparison.Ordinal))
{
return new(StoreResultCode.Conflict);
}
if (previous.ResourceId is SessionListingId listingId
&& _listings.TryGetValue(listingId, out ListingEntry? existing))
{
return new(StoreResultCode.Success, Snapshot(existing), true);
}
return new(StoreResultCode.Expired);
}
if (_listings.Count >= _options.MaxListings
|| _idempotency.Count >= _options.MaxIdempotencyEntries
|| _listings.Values.Count(entry => string.Equals(
entry.Definition.OwnerSubject,
command.Listing.OwnerSubject,
StringComparison.Ordinal)) >= command.OwnerListingLimit)
{
return new(StoreResultCode.CapacityExceeded);
}
ListingDefinition frozen = StoredListing.Freeze(command.Listing);
if (_listings.ContainsKey(frozen.ListingId)
|| _leases.ContainsKey(frozen.LeaseId)
|| HandleExists(frozen.HostPresenceHandle))
{
return new(StoreResultCode.Conflict);
}
ListingEntry entry = new(
frozen,
now + _options.LeaseLifetime,
WallDeadline(now, _options.LeaseLifetime),
version: 1);
_listings.Add(frozen.ListingId, entry);
_leases.Add(frozen.LeaseId, frozen.ListingId);
_presenceHandles.Add(frozen.HostPresenceHandle, frozen.ListingId);
_idempotency.Add(idempotencyKey, new(
command.RequestFingerprint,
frozen.ListingId,
now + _options.IdempotencyLifetime));
return new(StoreResultCode.Success, Snapshot(entry));
}, cancellationToken);
public StoreResult<StoredListing> RenewLease(
RenewLeaseCommand command,
CancellationToken cancellationToken = default) => Atomic<StoredListing>(now =>
{
ArgumentNullException.ThrowIfNull(command);
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (_drainDeadline.HasValue)
{
return new(StoreResultCode.Draining);
}
if (!_listings.TryGetValue(command.ListingId, out ListingEntry? entry))
{
return new(StoreResultCode.NotFound);
}
if (entry.Definition.LeaseId != command.LeaseId
|| entry.Definition.LeaseFingerprint != command.LeaseFingerprint
|| !string.Equals(entry.Definition.OwnerSubject, command.OwnerSubject, StringComparison.Ordinal))
{
return new(StoreResultCode.NotFound);
}
if (entry.Version != command.ExpectedVersion)
{
return new(StoreResultCode.Conflict, Snapshot(entry));
}
entry.LeaseDeadline = now + _options.LeaseLifetime;
entry.WallExpiresAt = WallDeadline(now, _options.LeaseLifetime);
entry.Version++;
return new(StoreResultCode.Success, Snapshot(entry));
}, cancellationToken);
public StoreResult<StoredListing> UpdateListing(
UpdateListingCommand command,
CancellationToken cancellationToken = default) => Atomic<StoredListing>(_ =>
{
ArgumentNullException.ThrowIfNull(command);
ValidateSubject(command.OwnerSubject, nameof(command.OwnerSubject));
if (!ContractValidation.IsBuildVersionValid(command.BuildVersion)
|| !ContractValidation.IsDisplayNameValid(command.DisplayName)
|| command.MaximumPlayers is <= 0 or > ContractLimits.SessionCapacityMaxPlayers
|| command.CurrentPlayers < 0
|| command.CurrentPlayers > command.MaximumPlayers
|| !ContractValidation.IsMetadataValid(command.Metadata)
|| command.DedicatedFallback is not null
&& !ContractValidation.IsNetworkEndpointValid(command.DedicatedFallback))
{
throw new ArgumentException("Listing update invariants are invalid.", nameof(command));
}
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (_drainDeadline.HasValue)
{
return new(StoreResultCode.Draining);
}
if (!_listings.TryGetValue(command.ListingId, out ListingEntry? entry)
|| entry.Definition.LeaseId != command.LeaseId
|| entry.Definition.LeaseFingerprint != command.LeaseFingerprint
|| !string.Equals(entry.Definition.OwnerSubject, command.OwnerSubject, StringComparison.Ordinal))
{
return new(StoreResultCode.NotFound);
}
entry.Definition = StoredListing.Freeze(entry.Definition with
{
BuildVersion = command.BuildVersion,
DisplayName = command.DisplayName,
CurrentPlayers = command.CurrentPlayers,
MaximumPlayers = command.MaximumPlayers,
Metadata = command.Metadata,
DedicatedFallback = command.DedicatedFallback,
});
entry.Version++;
return new(StoreResultCode.Success, Snapshot(entry));
}, cancellationToken);
public StoreResult<bool> DeleteListing(
DeleteListingCommand command,
CancellationToken cancellationToken = default) => Atomic<bool>(_ =>
{
ArgumentNullException.ThrowIfNull(command);
if (!_listings.TryGetValue(command.ListingId, out ListingEntry? entry)
|| entry.Definition.LeaseId != command.LeaseId
|| entry.Definition.LeaseFingerprint != command.LeaseFingerprint
|| !string.Equals(entry.Definition.OwnerSubject, command.OwnerSubject, StringComparison.Ordinal))
{
return new(StoreResultCode.NotFound);
}
RemoveListing(command.ListingId);
return new(StoreResultCode.Success, true);
}, cancellationToken);
public StoreResult<StoredListing> GetListing(
SessionListingId listingId,
bool requireFreshPresence,
CancellationToken cancellationToken = default) => Atomic<StoredListing>(_ =>
{
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (!_listings.TryGetValue(listingId, out ListingEntry? entry))
{
return new(StoreResultCode.NotFound);
}
bool fresh = _presence.ContainsKey(entry.Definition.HostPresenceHandle);
return requireFreshPresence && !fresh
? new(StoreResultCode.NotFound)
: new(StoreResultCode.Success, Snapshot(entry));
}, cancellationToken);
public StoreResult<StoredListing> BindHostPresence(
BindHostPresenceCommand command,
CancellationToken cancellationToken = default) => Atomic<StoredListing>(now =>
{
ArgumentNullException.ThrowIfNull(command);
if (command.Handle.Value == Guid.Empty || !command.CapabilityFingerprint.IsValid)
{
throw new ArgumentException("Host presence binding is invalid.", nameof(command));
}
ValidateEndpoint(command.PublicEndpoint, command.LocalEndpoint);
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (!_presenceHandles.TryGetValue(command.Handle, out SessionListingId listingId)
|| !_listings.TryGetValue(listingId, out ListingEntry? entry)
|| entry.LeaseDeadline <= now
|| entry.Definition.HostPresenceFingerprint != command.CapabilityFingerprint)
{
return new(StoreResultCode.NotFound);
}
if (!_presence.ContainsKey(command.Handle) && _presence.Count >= _options.MaxPresenceBindings)
{
return new(StoreResultCode.CapacityExceeded);
}
_presence[command.Handle] = new(
command.PublicEndpoint,
command.LocalEndpoint,
now + _options.PresenceLifetime);
return new(StoreResultCode.Success, Snapshot(entry));
}, cancellationToken, eagerCleanup: false);
public StoreResult<IReadOnlyList<StoredListing>> BrowseVisibleListings(
VisibleListingQuery query,
CancellationToken cancellationToken = default) => Atomic<IReadOnlyList<StoredListing>>(_ =>
{
ArgumentNullException.ThrowIfNull(query);
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (!IsScopeValid(query.Scope)
|| query.ProtocolVersion == 0
|| (query.RegionId.HasValue && string.IsNullOrEmpty(query.RegionId.Value.Value))
|| query.MaximumResults <= 0
|| query.MaximumResults > ContractLimits.BrowserPageMaxItems + 1)
{
throw new ArgumentOutOfRangeException(nameof(query));
}
IReadOnlyList<StoredListing> visible = _listings.Values
.Where(entry => entry.Definition.Scope == query.Scope
&& entry.Definition.ProtocolVersion == query.ProtocolVersion
&& entry.Definition.Visibility == ListingVisibility.Public
&& (!query.RegionId.HasValue || entry.Definition.RegionId == query.RegionId.Value)
&& (!query.AfterListingId.HasValue
|| entry.Definition.ListingId.Value.CompareTo(query.AfterListingId.Value.Value) > 0)
&& (!query.ExcludeFull
|| entry.Definition.CurrentPlayers < entry.Definition.MaximumPlayers)
&& _presence.ContainsKey(entry.Definition.HostPresenceHandle))
.OrderBy(static entry => entry.Definition.ListingId.Value)
.Take(query.MaximumResults)
.Select(Snapshot)
.ToArray();
return new(StoreResultCode.Success, visible);
}, cancellationToken);
public StoreResult<StoredJoinAttempt> CreateJoinAttempt(
CreateJoinAttemptCommand command,
CancellationToken cancellationToken = default) => Atomic<StoredJoinAttempt>(now =>
{
ArgumentNullException.ThrowIfNull(command);
ValidateAttempt(command);
ValidateIdempotency(command.IdempotencyKey, command.RequestFingerprint);
ValidateSubject(command.IdempotencyOwner, nameof(command.IdempotencyOwner));
ValidateSubject(command.ClientSubject, nameof(command.ClientSubject));
StoreResult<StoredJoinAttempt>? admission = CheckNewWorkAdmission<StoredJoinAttempt>(command.ClientSubject);
if (admission is not null)
{
return admission;
}
string idempotencyKey = $"attempt:{command.Scope.GameId}:{command.Scope.EnvironmentId}:{command.IdempotencyOwner}:{command.IdempotencyKey}";
if (_idempotency.TryGetValue(idempotencyKey, out IdempotencyEntry? previous))
{
if (!string.Equals(previous.RequestFingerprint, command.RequestFingerprint, StringComparison.Ordinal))
{
return new(StoreResultCode.Conflict);
}
if (previous.ResourceId is JoinAttemptId attemptId
&& _attempts.TryGetValue(attemptId, out AttemptEntry? priorAttempt))
{
return new(StoreResultCode.Success, Snapshot(priorAttempt), true);
}
return new(StoreResultCode.Expired);
}
if (!_listings.TryGetValue(command.ListingId, out ListingEntry? listing)
|| 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)
{
return new(StoreResultCode.CapacityExceeded);
}
if (_attempts.ContainsKey(command.AttemptId) || HandleExists(command.MediationHandle))
{
return new(StoreResultCode.Conflict);
}
AttemptEntry attempt = new(
command,
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,
command.AttemptId,
now + _options.IdempotencyLifetime));
return new(StoreResultCode.Success, Snapshot(attempt));
}, cancellationToken);
public StoreResult<IReadOnlyList<StoredJoinAttempt>> BrowseHostJoinAttempts(
HostJoinAttemptQuery query,
CancellationToken cancellationToken = default) => Atomic<IReadOnlyList<StoredJoinAttempt>>(_ =>
{
ArgumentNullException.ThrowIfNull(query);
if (query.ListingId.Value == Guid.Empty
|| !query.LeaseFingerprint.IsValid
|| query.MaximumResults is < 1 or > ContractLimits.BrowserPageMaxItems + 1)
{
throw new ArgumentException("Host attempt query invariants are invalid.", nameof(query));
}
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (!_listings.TryGetValue(query.ListingId, out ListingEntry? listing)
|| listing.Definition.LeaseFingerprint != query.LeaseFingerprint)
{
return new(StoreResultCode.NotFound);
}
IReadOnlyList<StoredJoinAttempt> attempts = _attempts.Values
.Where(entry => entry.Command.ListingId == query.ListingId
&& (!entry.IntroductionConsumed || entry.IsCancelled)
&& (!query.AfterAttemptId.HasValue
|| entry.Command.AttemptId.Value.CompareTo(query.AfterAttemptId.Value.Value) > 0))
.OrderBy(static entry => entry.Command.AttemptId.Value)
.Take(query.MaximumResults)
.Select(Snapshot)
.ToArray();
return new(StoreResultCode.Success, attempts);
}, cancellationToken);
public StoreResult<bool> CancelJoinAttempt(
CancelJoinAttemptCommand command,
CancellationToken cancellationToken = default) => Atomic<bool>(_ =>
{
ArgumentNullException.ThrowIfNull(command);
if (command.AttemptId.Value == Guid.Empty || !command.ClientCapabilityFingerprint.IsValid)
{
throw new ArgumentException("Join cancellation invariants are invalid.", nameof(command));
}
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (!_attempts.TryGetValue(command.AttemptId, out AttemptEntry? attempt)
|| attempt.ClientCapabilityFingerprint != command.ClientCapabilityFingerprint)
{
return new(StoreResultCode.NotFound);
}
if (attempt.IsCancelled)
{
return new(StoreResultCode.Success, true, true);
}
attempt.IsCancelled = true;
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 =>
{
ArgumentNullException.ThrowIfNull(command);
if (command.Handle.Value == Guid.Empty
|| command.Role is not (AttemptPeerRole.Host or AttemptPeerRole.Client)
|| !command.CapabilityFingerprint.IsValid)
{
throw new ArgumentException("Attempt endpoint binding is invalid.", nameof(command));
}
ValidateEndpoint(command.PublicEndpoint, command.LocalEndpoint);
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (!_attemptHandles.TryGetValue(command.Handle, out JoinAttemptId attemptId)
|| !_attempts.TryGetValue(attemptId, out AttemptEntry? attempt)
|| attempt.Deadline <= now
|| attempt.IsCancelled)
{
return new(StoreResultCode.NotFound);
}
if (!HasFreshHostPresence(attempt, now))
{
return new(StoreResultCode.NotFound);
}
SecretFingerprint expected = command.Role == AttemptPeerRole.Host
? attempt.HostCapabilityFingerprint
: attempt.ClientCapabilityFingerprint;
if (expected != command.CapabilityFingerprint)
{
return new(StoreResultCode.NotFound);
}
AttemptEndpointBinding binding = new(command.PublicEndpoint, command.LocalEndpoint);
AttemptEndpointBinding? current = command.Role == AttemptPeerRole.Host
? attempt.HostEndpoint
: attempt.ClientEndpoint;
if (current is not null)
{
return current == binding
? new(StoreResultCode.Success, Snapshot(attempt), true)
: new(StoreResultCode.ReplayRejected);
}
if (command.Role == AttemptPeerRole.Host)
{
attempt.HostEndpoint = binding;
}
else
{
attempt.ClientEndpoint = binding;
}
return new(StoreResultCode.Success, Snapshot(attempt));
}, cancellationToken, eagerCleanup: false);
public StoreResult<IntroductionEndpoints> ConsumeIntroduction(
MediationHandle handle,
CancellationToken cancellationToken = default) => Atomic<IntroductionEndpoints>(now =>
{
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (!_attemptHandles.TryGetValue(handle, out JoinAttemptId attemptId)
|| !_attempts.TryGetValue(attemptId, out AttemptEntry? attempt)
|| attempt.Deadline <= now
|| attempt.IsCancelled)
{
return new(StoreResultCode.NotFound);
}
if (!HasFreshHostPresence(attempt, now))
{
return new(StoreResultCode.NotFound);
}
if (attempt.IntroductionConsumed)
{
return new(StoreResultCode.ReplayRejected);
}
if (attempt.HostEndpoint is null || attempt.ClientEndpoint is null)
{
return new(StoreResultCode.Conflict);
}
attempt.IntroductionConsumed = true;
TimeSpan ticketLifetime = TimeSpan.FromTicks(Math.Min(
_options.ConnectionTicketLifetime.Ticks,
(attempt.Deadline - now).Ticks));
attempt.TicketDeadline = now + ticketLifetime;
attempt.TicketWallExpiresAt = WallDeadline(now, ticketLifetime);
return new(StoreResultCode.Success, new(
Snapshot(attempt),
attempt.HostEndpoint,
attempt.ClientEndpoint));
}, cancellationToken, eagerCleanup: false);
public StoreResult<bool> ConsumeConnectionTicket(
ConsumeConnectionTicketCommand command,
CancellationToken cancellationToken = default) => Atomic<bool>(now =>
{
ArgumentNullException.ThrowIfNull(command);
if (command.AttemptId.Value == Guid.Empty || !command.ConnectionTicketFingerprint.IsValid)
{
throw new ArgumentException("Connection ticket invariants are invalid.", nameof(command));
}
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (!_attempts.TryGetValue(command.AttemptId, out AttemptEntry? attempt)
|| attempt.ConnectionTicketFingerprint != command.ConnectionTicketFingerprint)
{
return new(StoreResultCode.NotFound);
}
if (!attempt.IntroductionConsumed)
{
return new(StoreResultCode.Conflict);
}
if (attempt.IsCancelled)
{
return new(StoreResultCode.Conflict);
}
if (!attempt.TicketDeadline.HasValue || attempt.TicketDeadline.Value <= now)
{
return new(StoreResultCode.Expired);
}
if (attempt.ConnectionTicketConsumed)
{
return new(StoreResultCode.ReplayRejected);
}
attempt.ConnectionTicketConsumed = true;
return new(StoreResultCode.Success, true);
}, cancellationToken);
public StoreResult<bool> ConsumeReplay(
ReplayConsumption consumption,
CancellationToken cancellationToken = default) => Atomic<bool>(now =>
{
ArgumentNullException.ThrowIfNull(consumption);
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
ValidateReplay(consumption);
string key = $"{consumption.Namespace}:{consumption.Key}";
if (_replay.ContainsKey(key))
{
return new(StoreResultCode.ReplayRejected);
}
if (_replay.Count >= _options.MaxReplayEntries)
{
return new(StoreResultCode.CapacityExceeded);
}
TimeSpan lifetime = consumption.Lifetime ?? _options.ReplayLifetime;
if (lifetime <= TimeSpan.Zero || lifetime > _options.ReplayLifetime)
{
throw new ArgumentOutOfRangeException(nameof(consumption), "Replay lifetime exceeds the configured ceiling.");
}
_replay.Add(key, now + lifetime);
return new(StoreResultCode.Success, true);
}, cancellationToken);
public StoreResult<bool> RevokeListing(
SessionListingId listingId,
CancellationToken cancellationToken = default) => Atomic<bool>(_ =>
{
if (!_listings.ContainsKey(listingId))
{
return new(StoreResultCode.NotFound);
}
RemoveListing(listingId);
return new(StoreResultCode.Success, true);
}, cancellationToken);
public StoreResult<int> RevokePrincipal(
string subject,
TimeSpan lifetime,
CancellationToken cancellationToken = default) => Atomic<int>(now =>
{
ValidateSubject(subject, nameof(subject));
if (lifetime <= TimeSpan.Zero || lifetime > TimeSpan.FromMinutes(10))
{
throw new ArgumentOutOfRangeException(nameof(lifetime));
}
if (!_revocations.ContainsKey(subject) && _revocations.Count >= _options.MaxRevocations)
{
return new(StoreResultCode.CapacityExceeded);
}
_revocations[subject] = now + lifetime;
SessionListingId[] listings = _listings
.Where(item => string.Equals(item.Value.Definition.OwnerSubject, subject, StringComparison.Ordinal))
.Select(static item => item.Key)
.ToArray();
JoinAttemptId[] attempts = _attempts
.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);
}
foreach (JoinAttemptId attemptId in attempts)
{
RemoveAttempt(attemptId);
}
foreach (JoinAttemptId attemptId in outcomeReports)
{
_outcomeReports.Remove(attemptId);
}
return new(StoreResultCode.Success, listings.Length + attempts.Length);
}, cancellationToken);
public void BeginDrain(CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
lock (_gate)
{
cancellationToken.ThrowIfCancellationRequested();
if (!_drainDeadline.HasValue)
{
_drainDeadline = _monotonicClock.Elapsed + _options.GracefulDrainLifetime;
}
}
}
internal void MarkUnavailable()
{
lock (_gate)
{
_available = false;
ClearActiveState();
}
}
private StoreResult<T> Atomic<T>(
Func<TimeSpan, StoreResult<T>> operation,
CancellationToken cancellationToken,
bool eagerCleanup = true)
{
cancellationToken.ThrowIfCancellationRequested();
lock (_gate)
{
cancellationToken.ThrowIfCancellationRequested();
TimeSpan now = _monotonicClock.Elapsed;
// Authenticated UDP duplicates need O(1) store work. Their operations
// check exact resource deadlines and amortize physical expiry removal.
bool drainExpired = _drainDeadline is TimeSpan drainDeadline
&& now >= drainDeadline;
if (drainExpired || eagerCleanup || now >= _nextUdpMaintenance)
{
Cleanup(now);
_nextUdpMaintenance = now + UdpMaintenanceInterval;
}
return operation(now);
}
}
private bool HasFreshHostPresence(AttemptEntry attempt, TimeSpan now) =>
_listings.TryGetValue(attempt.Command.ListingId, out ListingEntry? listing)
&& listing.LeaseDeadline > now
&& _presence.TryGetValue(
listing.Definition.HostPresenceHandle,
out PresenceEntry? presence)
&& presence.Deadline > now;
private StoreResult<T>? CheckNewWorkAdmission<T>(string subject)
{
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (_drainDeadline.HasValue)
{
return new(StoreResultCode.Draining);
}
return _revocations.ContainsKey(subject)
? new(StoreResultCode.Revoked)
: null;
}
private void Cleanup(TimeSpan now)
{
_maintenanceSweepCount++;
if (_drainDeadline is TimeSpan drainDeadline && now >= drainDeadline)
{
ClearActiveState();
}
RemoveExpired(_revocations, now);
RemoveExpired(_replay, now);
foreach (string key in _idempotency
.Where(item => item.Value.Deadline <= now)
.Select(static item => item.Key)
.ToArray())
{
_idempotency.Remove(key);
}
foreach (MediationHandle handle in _presence
.Where(item => item.Value.Deadline <= now)
.Select(static item => item.Key)
.ToArray())
{
_presence.Remove(handle);
}
foreach (JoinAttemptId attemptId in _attempts
.Where(item => item.Value.Deadline <= now)
.Select(static item => item.Key)
.ToArray())
{
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)
.ToArray())
{
RemoveListing(listingId);
}
}
private void ClearActiveState()
{
_listings.Clear();
_leases.Clear();
_presenceHandles.Clear();
_presence.Clear();
_attempts.Clear();
_outcomeReports.Clear();
_attemptHandles.Clear();
_idempotency.Clear();
_replay.Clear();
}
private void RemoveListing(SessionListingId listingId)
{
if (!_listings.Remove(listingId, out ListingEntry? listing))
{
return;
}
_leases.Remove(listing.Definition.LeaseId);
_presenceHandles.Remove(listing.Definition.HostPresenceHandle);
_presence.Remove(listing.Definition.HostPresenceHandle);
foreach (JoinAttemptId attemptId in _attempts
.Where(item => item.Value.Command.ListingId == listingId)
.Select(static item => item.Key)
.ToArray())
{
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)
{
if (_attempts.Remove(attemptId, out AttemptEntry? attempt))
{
_attemptHandles.Remove(attempt.Command.MediationHandle);
}
}
private bool HandleExists(MediationHandle handle) =>
_presenceHandles.ContainsKey(handle) || _attemptHandles.ContainsKey(handle);
private DateTimeOffset WallDeadline(TimeSpan now, TimeSpan lifetime) =>
_wallOrigin + (now - _monotonicOrigin) + lifetime;
private StoredListing Snapshot(ListingEntry entry) => new()
{
Definition = entry.Definition,
LeaseExpiresAt = entry.WallExpiresAt,
Version = entry.Version,
HasFreshPresence = _presence.ContainsKey(entry.Definition.HostPresenceHandle),
};
private static StoredJoinAttempt Snapshot(AttemptEntry entry) => new()
{
AttemptId = entry.Command.AttemptId,
MediationHandle = entry.Command.MediationHandle,
Scope = entry.Command.Scope,
ListingId = entry.Command.ListingId,
ClientSubject = entry.Command.ClientSubject,
ProtocolVersion = entry.Command.ProtocolVersion,
IdempotencyKey = entry.Command.IdempotencyKey,
RequestFingerprint = entry.Command.RequestFingerprint,
CapabilityDerivationSalt = entry.Command.CapabilityDerivationSalt,
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,
ClientEndpoint = entry.ClientEndpoint,
IntroductionConsumed = entry.IntroductionConsumed,
ConnectionTicketConsumed = entry.ConnectionTicketConsumed,
IsCancelled = entry.IsCancelled,
};
private static void RemoveExpired(Dictionary<string, TimeSpan> entries, TimeSpan now)
{
foreach (string key in entries
.Where(item => item.Value <= now)
.Select(static item => item.Key)
.ToArray())
{
entries.Remove(key);
}
}
private static void ValidateListing(ListingDefinition listing)
{
ArgumentNullException.ThrowIfNull(listing);
ValidateSubject(listing.OwnerSubject, nameof(listing.OwnerSubject));
ArgumentNullException.ThrowIfNull(listing.Metadata);
if (listing.ListingId.Value == Guid.Empty
|| listing.LeaseId.Value == Guid.Empty
|| listing.HostPresenceHandle.Value == Guid.Empty
|| !IsScopeValid(listing.Scope)
|| string.IsNullOrEmpty(listing.RegionId.Value)
|| listing.ProtocolVersion == 0
|| !ContractValidation.IsBuildVersionValid(listing.BuildVersion)
|| !ContractValidation.IsDisplayNameValid(listing.DisplayName)
|| !Enum.IsDefined(listing.Visibility)
|| !Enum.IsDefined(listing.TrustMode)
|| listing.MaximumPlayers is <= 0 or > ContractLimits.SessionCapacityMaxPlayers
|| 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))
{
throw new ArgumentException("Listing invariants are invalid.", nameof(listing));
}
}
private static void ValidateIdempotency(string key, string requestFingerprint)
{
if (string.IsNullOrWhiteSpace(key) || key.Length > 128)
{
throw new ArgumentException("Idempotency keys must contain 1-128 characters.", nameof(key));
}
if (string.IsNullOrWhiteSpace(requestFingerprint) || requestFingerprint.Length > 128)
{
throw new ArgumentException("Request fingerprints must contain 1-128 characters.", nameof(requestFingerprint));
}
}
private static void ValidateReplay(ReplayConsumption consumption)
{
if (string.IsNullOrWhiteSpace(consumption.Namespace) || consumption.Namespace.Length > 64
|| string.IsNullOrWhiteSpace(consumption.Key) || consumption.Key.Length > 128)
{
throw new ArgumentException("Replay namespace/key is invalid.", nameof(consumption));
}
}
private static void ValidateAttempt(CreateJoinAttemptCommand command)
{
if (command.AttemptId.Value == Guid.Empty
|| command.MediationHandle.Value == Guid.Empty
|| command.ListingId.Value == Guid.Empty
|| !IsScopeValid(command.Scope)
|| command.ProtocolVersion == 0
|| !command.HostCapabilityFingerprint.IsValid
|| !command.ClientCapabilityFingerprint.IsValid
|| !command.ConnectionTicketFingerprint.IsValid
|| command.DedicatedFallback is not null
&& !ContractValidation.IsNetworkEndpointValid(command.DedicatedFallback)
|| !IsDerivationSaltValid(command.CapabilityDerivationSalt)
|| command.ScopeAttemptLimit <= 0)
{
throw new ArgumentException("Join attempt invariants are invalid.", nameof(command));
}
}
private static void ValidateEndpoint(ObservedEndpoint publicEndpoint, ObservedEndpoint? localEndpoint)
{
if (!publicEndpoint.IsValid || (localEndpoint.HasValue && !localEndpoint.Value.IsValid))
{
throw new ArgumentException("Observed endpoints must be valid immutable endpoint values.", nameof(publicEndpoint));
}
}
private static bool IsScopeValid(TenantScope scope) =>
!string.IsNullOrEmpty(scope.GameId.Value) && !string.IsNullOrEmpty(scope.EnvironmentId.Value);
private static bool IsDerivationSaltValid(string? value) => value is not null
&& value.Length == 43
&& value.All(static character =>
character is >= 'A' and <= 'Z'
or >= 'a' and <= 'z'
or >= '0' and <= '9'
or '-'
or '_');
private static void ValidateSubject(string subject, string parameterName)
{
if (string.IsNullOrWhiteSpace(subject) || subject.Length > 256)
{
throw new ArgumentException("Subjects must contain 1-256 characters.", parameterName);
}
}
private sealed class ListingEntry(
ListingDefinition definition,
TimeSpan leaseDeadline,
DateTimeOffset wallExpiresAt,
long version)
{
public ListingDefinition Definition { get; set; } = definition;
public TimeSpan LeaseDeadline { get; set; } = leaseDeadline;
public DateTimeOffset WallExpiresAt { get; set; } = wallExpiresAt;
public long Version { get; set; } = version;
}
private sealed class PresenceEntry(
ObservedEndpoint publicEndpoint,
ObservedEndpoint? localEndpoint,
TimeSpan deadline)
{
public ObservedEndpoint PublicEndpoint { get; } = publicEndpoint;
public ObservedEndpoint? LocalEndpoint { get; } = localEndpoint;
public TimeSpan Deadline { get; } = deadline;
}
private sealed class AttemptEntry(
CreateJoinAttemptCommand command,
TimeSpan deadline,
DateTimeOffset wallExpiresAt)
{
public CreateJoinAttemptCommand Command { get; } = command;
public SecretFingerprint HostCapabilityFingerprint { get; } = command.HostCapabilityFingerprint;
public SecretFingerprint ClientCapabilityFingerprint { get; } = command.ClientCapabilityFingerprint;
public SecretFingerprint ConnectionTicketFingerprint { get; } = command.ConnectionTicketFingerprint;
public TimeSpan Deadline { get; } = deadline;
public DateTimeOffset WallExpiresAt { get; } = wallExpiresAt;
public TimeSpan? TicketDeadline { get; set; }
public DateTimeOffset? TicketWallExpiresAt { get; set; }
public AttemptEndpointBinding? HostEndpoint { get; set; }
public AttemptEndpointBinding? ClientEndpoint { get; set; }
public bool IntroductionConsumed { get; set; }
public bool ConnectionTicketConsumed { get; set; }
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,
TimeSpan Deadline);
}