feat: add atomic ephemeral state (#6)
quality-gate / quality (push) Successful in 57s

Closes #6
This commit is contained in:
KyuubiYoru
2026-07-16 05:32:48 +02:00
parent 47382ddadc
commit 02ca502a76
7 changed files with 1754 additions and 2 deletions
@@ -0,0 +1,279 @@
using System.Collections.Frozen;
using System.Diagnostics;
using System.Net;
using FinalFactory.Rendezvous.Contracts;
namespace FinalFactory.Rendezvous.Server.State;
internal interface IWallClock
{
DateTimeOffset UtcNow { get; }
}
internal interface IMonotonicClock
{
TimeSpan Elapsed { get; }
}
internal sealed class SystemRendezvousClock : IWallClock, IMonotonicClock
{
private readonly long _origin = Stopwatch.GetTimestamp();
public DateTimeOffset UtcNow => DateTimeOffset.UtcNow;
public TimeSpan Elapsed => Stopwatch.GetElapsedTime(_origin);
}
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 MaxReplayEntries { get; init; } = 30_000;
public int MaxRevocations { get; init; } = 10_000;
public int MaxIdempotencyEntries { get; init; } = 35_000;
public TimeSpan LeaseLifetime { get; init; } = TimeSpan.FromSeconds(60);
public TimeSpan PresenceLifetime { get; init; } = TimeSpan.FromSeconds(20);
public TimeSpan JoinAttemptLifetime { get; init; } = TimeSpan.FromSeconds(30);
public TimeSpan ReplayLifetime { get; init; } = TimeSpan.FromSeconds(30);
public TimeSpan IdempotencyLifetime { get; init; } = TimeSpan.FromMinutes(2);
public TimeSpan GracefulDrainLifetime { get; init; } = TimeSpan.FromSeconds(30);
public void Validate()
{
RequirePositive(MaxListings, nameof(MaxListings));
RequirePositive(MaxPresenceBindings, nameof(MaxPresenceBindings));
RequirePositive(MaxJoinAttempts, nameof(MaxJoinAttempts));
RequirePositive(MaxReplayEntries, nameof(MaxReplayEntries));
RequirePositive(MaxRevocations, nameof(MaxRevocations));
RequirePositive(MaxIdempotencyEntries, nameof(MaxIdempotencyEntries));
RequireDuration(LeaseLifetime, TimeSpan.FromSeconds(60), nameof(LeaseLifetime));
RequireDuration(PresenceLifetime, TimeSpan.FromSeconds(20), nameof(PresenceLifetime));
RequireDuration(JoinAttemptLifetime, TimeSpan.FromSeconds(30), nameof(JoinAttemptLifetime));
RequireDuration(ReplayLifetime, TimeSpan.FromSeconds(30), nameof(ReplayLifetime));
RequireDuration(IdempotencyLifetime, TimeSpan.FromMinutes(10), nameof(IdempotencyLifetime));
RequireDuration(GracefulDrainLifetime, TimeSpan.FromSeconds(30), nameof(GracefulDrainLifetime));
}
private static void RequirePositive(int value, string name)
{
if (value <= 0)
{
throw new ArgumentOutOfRangeException(name, "Store capacity must be positive.");
}
}
private static void RequireDuration(TimeSpan value, TimeSpan maximum, string name)
{
if (value <= TimeSpan.Zero || value > maximum)
{
throw new ArgumentOutOfRangeException(name, $"Duration must be positive and no greater than {maximum}.");
}
}
}
internal readonly record struct TenantScope(GameId GameId, EnvironmentId EnvironmentId);
internal readonly record struct SecretFingerprint
{
public SecretFingerprint(string value)
{
if (string.IsNullOrWhiteSpace(value) || value.Length > 128)
{
throw new ArgumentException("Secret fingerprints must contain 1-128 characters.", nameof(value));
}
Value = value;
}
public string Value { get; }
public bool IsValid => !string.IsNullOrWhiteSpace(Value) && Value.Length <= 128;
public override string ToString() => "[REDACTED]";
}
internal readonly record struct ObservedEndpoint
{
public ObservedEndpoint(AddressFamilyKind addressFamily, string address, int port)
{
if (!IPAddress.TryParse(address, out IPAddress? parsed)
|| (addressFamily == AddressFamilyKind.Ipv4 && parsed.AddressFamily != System.Net.Sockets.AddressFamily.InterNetwork)
|| (addressFamily == AddressFamilyKind.Ipv6 && parsed.AddressFamily != System.Net.Sockets.AddressFamily.InterNetworkV6))
{
throw new ArgumentException("The address must match the declared address family.", nameof(address));
}
if (port is < 1 or > 65_535)
{
throw new ArgumentOutOfRangeException(nameof(port));
}
AddressFamily = addressFamily;
Address = parsed.ToString();
Port = port;
}
public AddressFamilyKind AddressFamily { get; }
public string Address { get; }
public int Port { get; }
public bool IsValid => !string.IsNullOrEmpty(Address)
&& Port is >= 1 and <= 65_535
&& AddressFamily is AddressFamilyKind.Ipv4 or AddressFamilyKind.Ipv6;
}
internal sealed record ListingDefinition
{
public required SessionListingId ListingId { get; init; }
public required LeaseId LeaseId { get; init; }
public required TenantScope Scope { get; init; }
public required string OwnerSubject { get; init; }
public required RegionId RegionId { get; init; }
public required uint ProtocolVersion { get; init; }
public required string BuildVersion { get; init; }
public required string DisplayName { get; init; }
public required ListingVisibility Visibility { get; init; }
public required PublisherTrustMode TrustMode { get; init; }
public required int CurrentPlayers { get; init; }
public required int MaximumPlayers { get; init; }
public required IReadOnlyDictionary<string, string> Metadata { get; init; }
public required SecretFingerprint LeaseFingerprint { get; init; }
public required MediationHandle HostPresenceHandle { get; init; }
public required SecretFingerprint HostPresenceFingerprint { get; init; }
}
internal sealed record StoredListing
{
public required ListingDefinition Definition { get; init; }
public required DateTimeOffset LeaseExpiresAt { get; init; }
public required long Version { get; init; }
public required bool HasFreshPresence { get; init; }
public static ListingDefinition Freeze(ListingDefinition source) => source with
{
Metadata = source.Metadata.ToFrozenDictionary(StringComparer.Ordinal),
};
}
internal sealed record CreateListingCommand(
string IdempotencyKey,
string RequestFingerprint,
ListingDefinition Listing,
int OwnerListingLimit = int.MaxValue);
internal sealed record RenewLeaseCommand(
SessionListingId ListingId,
LeaseId LeaseId,
SecretFingerprint LeaseFingerprint,
long ExpectedVersion);
internal sealed record DeleteListingCommand(
SessionListingId ListingId,
LeaseId LeaseId,
SecretFingerprint LeaseFingerprint);
internal sealed record BindHostPresenceCommand(
MediationHandle Handle,
SecretFingerprint CapabilityFingerprint,
ObservedEndpoint PublicEndpoint,
ObservedEndpoint? LocalEndpoint);
internal sealed record VisibleListingQuery(
TenantScope Scope,
uint ProtocolVersion,
RegionId? RegionId,
int MaximumResults = ContractLimits.BrowserPageMaxItems);
internal enum AttemptPeerRole
{
Host = 1,
Client = 2,
}
internal sealed record CreateJoinAttemptCommand
{
public required string IdempotencyOwner { get; init; }
public required string IdempotencyKey { get; init; }
public required string RequestFingerprint { get; init; }
public required string ClientSubject { get; init; }
public required JoinAttemptId AttemptId { get; init; }
public required MediationHandle MediationHandle { get; init; }
public required TenantScope Scope { get; init; }
public required SessionListingId ListingId { get; init; }
public required uint ProtocolVersion { get; init; }
public required SecretFingerprint HostCapabilityFingerprint { get; init; }
public required SecretFingerprint ClientCapabilityFingerprint { get; init; }
public int ScopeAttemptLimit { get; init; } = int.MaxValue;
}
internal sealed record AttemptEndpointBinding(
ObservedEndpoint PublicEndpoint,
ObservedEndpoint? LocalEndpoint);
internal sealed record StoredJoinAttempt
{
public required JoinAttemptId AttemptId { get; init; }
public required MediationHandle MediationHandle { get; init; }
public required TenantScope Scope { get; init; }
public required SessionListingId ListingId { get; init; }
public required string ClientSubject { get; init; }
public required uint ProtocolVersion { get; init; }
public required DateTimeOffset ExpiresAt { get; init; }
public AttemptEndpointBinding? HostEndpoint { get; init; }
public AttemptEndpointBinding? ClientEndpoint { get; init; }
public required bool IntroductionConsumed { get; init; }
}
internal sealed record BindAttemptEndpointCommand(
MediationHandle Handle,
AttemptPeerRole Role,
SecretFingerprint CapabilityFingerprint,
ObservedEndpoint PublicEndpoint,
ObservedEndpoint? LocalEndpoint);
internal sealed record IntroductionEndpoints(
JoinAttemptId AttemptId,
AttemptEndpointBinding Host,
AttemptEndpointBinding Client);
internal sealed record ReplayConsumption(
string Namespace,
string Key,
TimeSpan? Lifetime = null);
internal enum StoreResultCode
{
Success = 0,
NotFound = 1,
Expired = 2,
Revoked = 3,
Conflict = 4,
CapacityExceeded = 5,
Draining = 6,
ReplayRejected = 7,
ServiceUnavailable = 8,
}
internal sealed record StoreResult<T>(StoreResultCode Code, T? Value = default, bool IsIdempotentReplay = false)
{
public bool Succeeded => Code == StoreResultCode.Success;
}
internal interface IEphemeralRendezvousStore
{
Guid InstanceId { get; }
bool IsAvailable { get; }
bool IsDraining { get; }
StoreResult<StoredListing> CreateListing(CreateListingCommand command, CancellationToken cancellationToken = default);
StoreResult<StoredListing> RenewLease(RenewLeaseCommand command, CancellationToken cancellationToken = default);
StoreResult<bool> DeleteListing(DeleteListingCommand command, CancellationToken cancellationToken = default);
StoreResult<StoredListing> GetListing(SessionListingId listingId, bool requireFreshPresence, CancellationToken cancellationToken = default);
StoreResult<IReadOnlyList<StoredListing>> BrowseVisibleListings(VisibleListingQuery query, CancellationToken cancellationToken = default);
StoreResult<StoredListing> BindHostPresence(BindHostPresenceCommand command, CancellationToken cancellationToken = default);
StoreResult<StoredJoinAttempt> CreateJoinAttempt(CreateJoinAttemptCommand command, CancellationToken cancellationToken = default);
StoreResult<StoredJoinAttempt> BindAttemptEndpoint(BindAttemptEndpointCommand command, CancellationToken cancellationToken = default);
StoreResult<IntroductionEndpoints> ConsumeIntroduction(MediationHandle handle, CancellationToken cancellationToken = default);
StoreResult<bool> ConsumeReplay(ReplayConsumption consumption, CancellationToken cancellationToken = default);
StoreResult<bool> RevokeListing(SessionListingId listingId, CancellationToken cancellationToken = default);
StoreResult<int> RevokePrincipal(string subject, TimeSpan lifetime, CancellationToken cancellationToken = default);
void BeginDrain(CancellationToken cancellationToken = default);
}
@@ -0,0 +1,805 @@
using FinalFactory.Rendezvous.Contracts;
namespace FinalFactory.Rendezvous.Server.State;
internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousStore
{
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<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 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; }
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.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)
{
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<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)
{
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.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);
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)
{
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)
&& _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.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
|| listing.Definition.ProtocolVersion != command.ProtocolVersion
|| !_presence.ContainsKey(listing.Definition.HostPresenceHandle))
{
return new(StoreResultCode.NotFound);
}
if (_attempts.Count >= _options.MaxJoinAttempts
|| _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);
_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<StoredJoinAttempt> BindAttemptEndpoint(
BindAttemptEndpointCommand command,
CancellationToken cancellationToken = default) => Atomic<StoredJoinAttempt>(_ =>
{
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))
{
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);
public StoreResult<IntroductionEndpoints> ConsumeIntroduction(
MediationHandle handle,
CancellationToken cancellationToken = default) => Atomic<IntroductionEndpoints>(_ =>
{
if (!_available)
{
return new(StoreResultCode.ServiceUnavailable);
}
if (!_attemptHandles.TryGetValue(handle, out JoinAttemptId attemptId)
|| !_attempts.TryGetValue(attemptId, out AttemptEntry? attempt))
{
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;
return new(StoreResultCode.Success, new(
attempt.Command.AttemptId,
attempt.HostEndpoint,
attempt.ClientEndpoint));
}, 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();
foreach (SessionListingId listingId in listings)
{
RemoveListing(listingId);
}
foreach (JoinAttemptId attemptId in attempts)
{
RemoveAttempt(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)
{
cancellationToken.ThrowIfCancellationRequested();
lock (_gate)
{
cancellationToken.ThrowIfCancellationRequested();
TimeSpan now = _monotonicClock.Elapsed;
Cleanup(now);
return operation(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)
{
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 (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();
_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);
}
}
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,
ExpiresAt = entry.WallExpiresAt,
HostEndpoint = entry.HostEndpoint,
ClientEndpoint = entry.ClientEndpoint,
IntroductionConsumed = entry.IntroductionConsumed,
};
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.LeaseFingerprint.IsValid
|| !listing.HostPresenceFingerprint.IsValid)
{
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.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 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; } = 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 TimeSpan Deadline { get; } = deadline;
public DateTimeOffset WallExpiresAt { get; } = wallExpiresAt;
public AttemptEndpointBinding? HostEndpoint { get; set; }
public AttemptEndpointBinding? ClientEndpoint { get; set; }
public bool IntroductionConsumed { get; set; }
}
private sealed record IdempotencyEntry(
string RequestFingerprint,
object ResourceId,
TimeSpan Deadline);
}