From 02ca502a768f89f8aa30259c6152ebbf9eb7a33c Mon Sep 17 00:00:00 2001 From: KyuubiYoru Date: Thu, 16 Jul 2026 05:32:48 +0200 Subject: [PATCH] feat: add atomic ephemeral state (#6) Closes #6 --- .../0004-atomic-ephemeral-state.md | 100 +++ docs/architecture/README.md | 1 + src/FinalFactory.Rendezvous.Server/Program.cs | 19 +- .../State/EphemeralStateContracts.cs | 279 ++++++ .../State/InMemoryEphemeralRendezvousStore.cs | 805 ++++++++++++++++++ .../State/EphemeralStateTestData.cs | 105 +++ .../InMemoryEphemeralRendezvousStoreTests.cs | 447 ++++++++++ 7 files changed, 1754 insertions(+), 2 deletions(-) create mode 100644 docs/architecture/0004-atomic-ephemeral-state.md create mode 100644 src/FinalFactory.Rendezvous.Server/State/EphemeralStateContracts.cs create mode 100644 src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs create mode 100644 tests/FinalFactory.Rendezvous.Tests/State/EphemeralStateTestData.cs create mode 100644 tests/FinalFactory.Rendezvous.Tests/State/InMemoryEphemeralRendezvousStoreTests.cs diff --git a/docs/architecture/0004-atomic-ephemeral-state.md b/docs/architecture/0004-atomic-ephemeral-state.md new file mode 100644 index 0000000..b382370 --- /dev/null +++ b/docs/architecture/0004-atomic-ephemeral-state.md @@ -0,0 +1,100 @@ +# ADR 0004: atomic ephemeral state and single-active availability + +- Status: Accepted +- Date: 2026-07-16 +- Tracking: #6 + +## Context + +Listings, leases, endpoint observations, join attempts, and replay decisions must +move together. A partially committed authorization can expose an expired listing, +reuse a capability, or introduce an endpoint that was never authorized. V1 is a +single-active service, so it needs honest bounded in-memory behavior rather than +a database-shaped abstraction that implies unavailable durability or scale. + +## Decision + +`IEphemeralRendezvousStore` is the atomic boundary for directory, lease, presence, +attempt, endpoint, replay, revocation, and drain transitions. The v1 implementation +serializes each transition under one process-local lock. This deliberately favors +simple, auditable correctness at the initial 25,000-listing/10,000-attempt ceiling. +It retains only immutable listing data, opaque credential fingerprints, observed +endpoints, monotonic deadlines, and bounded idempotency/replay records. + +Every collection has an independent configured ceiling. An operation checks all +of the capacity it needs before changing any collection. Exhaustion returns +`CapacityExceeded`; it does not evict live state, partially insert an operation, +or grow a fallback queue. Policy-provided per-owner listing and per-tenant active +attempt quotas are evaluated inside the same creation transition, so concurrent +requests cannot pass a check performed outside the store. New join authorization returns `ServiceUnavailable` +when the atomic store is unavailable and `Draining` once drain starts. + +### Time and cleanup + +Expiry uses an injected monotonic clock. Wall time is used only to return an +informational `ExpiresAt` value. Moving the wall clock forward or backward cannot +expire or prolong authority. Cleanup runs deterministically at the start of every +store operation and removes presence, attempts, listings, replay entries, +idempotency records, and revocations at their deadline. Removal of a listing also +removes its presence handle and every linked attempt before another caller can +observe the store. + +### Concurrency and idempotency + +- Listing registration and join-attempt creation use an owner-scoped idempotency + key plus a canonical request fingerprint. An exact duplicate returns the + original live result; reuse with different input returns `Conflict`; replay + after the resource has expired returns `Expired` until the bounded idempotency + record itself expires. +- Lease renewal is compare-and-swap by version. A stale renewal returns the latest + version as `Conflict`. Renew/delete races are serialized: renewal either commits + before deletion or observes the listing as absent. +- Host presence refresh is an atomic whole-endpoint replacement because NAT + mappings can legitimately change. Attempt capabilities are different: the + first endpoint bound for each role wins, an identical datagram is idempotent, + and a different replay is rejected. Introduction is consumed once atomically. +- Cancellation is checked before waiting for the lock and again after acquiring + it. A cancellation observed at either point makes no change. Once a synchronous + transition starts, it completes atomically and does not expose partial state. + +### Visibility and revocation + +A listing is visible or joinable only when its lease and authenticated UDP host +presence are both fresh. Public browsing is tenant/protocol scoped, excludes +unlisted sessions, and uses a stable listing-ID order with the contract page +ceiling. Revoking a listing or principal removes every listing, presence, and +attempt path in the same transition. A revocation is inserted before removal; +if the bounded revocation pool is full, the operation rejects without deleting +anything. + +### Restart and graceful drain + +A process restart creates a new store instance ID and starts empty. Old listing, +lease, attempt, endpoint, idempotency, and consumption state is not recovered. +Publishers must re-register; old callers receive typed `NotFound`, `Expired`, or +`ServiceUnavailable` outcomes rather than an ambiguous success. No database is +required or supported for the single-active MVP. + +Drain is idempotent. It immediately rejects new registrations, attempts, and +lease extensions, while already-created attempts may bind endpoints and consume +their introduction during the configured window (at most 30 seconds). At the +deadline all active state is cleared atomically. Readiness is false while draining +or unavailable, and application shutdown starts drain before teardown. + +## Future shared-store mapping + +The interface uses explicit typed outcomes, TTLs, compare-and-swap versions, +idempotency records, and all-or-nothing multi-record transitions. A future Redis +implementation therefore requires authenticated transport, tenant-prefixed keys, +server-side scripts or transactions for each transition, TTLs based on the store's +authoritative time, and deterministic mediator routing. It must preserve these +semantics and pass the same contract tests before issue #18 may enable more than +one active instance. + +## Consequences + +- V1 has deterministic failure and restart behavior without durable gameplay state. +- A single lock is a measured capacity constraint, not a claim of horizontal scale. +- Transport and HTTP modules cannot bypass the store for authorization decisions. +- Operational code must treat `CapacityExceeded`, `Draining`, and + `ServiceUnavailable` as normal typed overload/availability outcomes. diff --git a/docs/architecture/README.md b/docs/architecture/README.md index ad87566..cce10b0 100644 --- a/docs/architecture/README.md +++ b/docs/architecture/README.md @@ -6,6 +6,7 @@ decision requires a superseding ADR and corresponding contract/test updates. - [ADR 0001: v1 control-plane boundaries and domain](0001-v1-control-plane-boundaries.md) - [ADR 0002: publisher trust, discovery, compatibility, and fallback](0002-publisher-trust-and-connection-policy.md) - [ADR 0003: state, privacy, availability, and safety budgets](0003-state-privacy-availability-and-budgets.md) +- [ADR 0004: atomic ephemeral state and single-active availability](0004-atomic-ephemeral-state.md) - [Threat model](../security/threat-model.md) - [Security promise and test matrix](../security/control-matrix.md) - [Versioned HTTP and UDP contracts](../contracts/README.md) diff --git a/src/FinalFactory.Rendezvous.Server/Program.cs b/src/FinalFactory.Rendezvous.Server/Program.cs index 8dbea94..ff55247 100644 --- a/src/FinalFactory.Rendezvous.Server/Program.cs +++ b/src/FinalFactory.Rendezvous.Server/Program.cs @@ -2,6 +2,7 @@ using System.Net; using FinalFactory.Rendezvous.Contracts; using FinalFactory.Rendezvous.Server.Http; using FinalFactory.Rendezvous.Server.Provisioning; +using FinalFactory.Rendezvous.Server.State; using FinalFactory.Rendezvous.Server.Transport; using Microsoft.OpenApi; @@ -35,6 +36,13 @@ builder.Services.AddOpenApi("v1", static options => builder.Services.ConfigureHttpJsonOptions(static options => ContractJson.Configure(options.SerializerOptions)); +SystemRendezvousClock rendezvousClock = new(); +InMemoryEphemeralRendezvousStore stateStore = new( + new EphemeralStoreOptions(), + rendezvousClock, + rendezvousClock); +builder.Services.AddSingleton(stateStore); + if (isOpenApiGeneration) { builder.Services.AddSingleton(new ProvisioningReadiness(false)); @@ -74,6 +82,7 @@ if (!isOpenApiGeneration) } WebApplication app = builder.Build(); +app.Lifetime.ApplicationStopping.Register(() => stateStore.BeginDrain()); app.MapOpenApi(); app.MapRendezvousContractEndpoints(); @@ -85,8 +94,14 @@ app.MapGet( .WithTags("Health"); app.MapGet( "/health/ready", - static (UdpMediatorService mediator, ProvisioningReadiness provisioning) => - mediator.LocalEndpoint is null || !provisioning.IsReady + static ( + UdpMediatorService mediator, + ProvisioningReadiness provisioning, + IEphemeralRendezvousStore state) => + mediator.LocalEndpoint is null + || !provisioning.IsReady + || !state.IsAvailable + || state.IsDraining ? Results.StatusCode(StatusCodes.Status503ServiceUnavailable) : Results.Ok(new HealthResponse { Status = "ready" })) .Produces() diff --git a/src/FinalFactory.Rendezvous.Server/State/EphemeralStateContracts.cs b/src/FinalFactory.Rendezvous.Server/State/EphemeralStateContracts.cs new file mode 100644 index 0000000..645733e --- /dev/null +++ b/src/FinalFactory.Rendezvous.Server/State/EphemeralStateContracts.cs @@ -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 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(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 CreateListing(CreateListingCommand command, CancellationToken cancellationToken = default); + StoreResult RenewLease(RenewLeaseCommand command, CancellationToken cancellationToken = default); + StoreResult DeleteListing(DeleteListingCommand command, CancellationToken cancellationToken = default); + StoreResult GetListing(SessionListingId listingId, bool requireFreshPresence, CancellationToken cancellationToken = default); + StoreResult> BrowseVisibleListings(VisibleListingQuery query, CancellationToken cancellationToken = default); + StoreResult BindHostPresence(BindHostPresenceCommand command, CancellationToken cancellationToken = default); + StoreResult CreateJoinAttempt(CreateJoinAttemptCommand command, CancellationToken cancellationToken = default); + StoreResult BindAttemptEndpoint(BindAttemptEndpointCommand command, CancellationToken cancellationToken = default); + StoreResult ConsumeIntroduction(MediationHandle handle, CancellationToken cancellationToken = default); + StoreResult ConsumeReplay(ReplayConsumption consumption, CancellationToken cancellationToken = default); + StoreResult RevokeListing(SessionListingId listingId, CancellationToken cancellationToken = default); + StoreResult RevokePrincipal(string subject, TimeSpan lifetime, CancellationToken cancellationToken = default); + void BeginDrain(CancellationToken cancellationToken = default); +} diff --git a/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs b/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs new file mode 100644 index 0000000..f558118 --- /dev/null +++ b/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs @@ -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 _listings = []; + private readonly Dictionary _leases = []; + private readonly Dictionary _presenceHandles = []; + private readonly Dictionary _presence = []; + private readonly Dictionary _attempts = []; + private readonly Dictionary _attemptHandles = []; + private readonly Dictionary _idempotency = new(StringComparer.Ordinal); + private readonly Dictionary _replay = new(StringComparer.Ordinal); + private readonly Dictionary _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 CreateListing( + CreateListingCommand command, + CancellationToken cancellationToken = default) => Atomic(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? admission = CheckNewWorkAdmission(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 RenewLease( + RenewLeaseCommand command, + CancellationToken cancellationToken = default) => Atomic(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 DeleteListing( + DeleteListingCommand command, + CancellationToken cancellationToken = default) => Atomic(_ => + { + 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 GetListing( + SessionListingId listingId, + bool requireFreshPresence, + CancellationToken cancellationToken = default) => Atomic(_ => + { + 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 BindHostPresence( + BindHostPresenceCommand command, + CancellationToken cancellationToken = default) => Atomic(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> BrowseVisibleListings( + VisibleListingQuery query, + CancellationToken cancellationToken = default) => Atomic>(_ => + { + 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 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 CreateJoinAttempt( + CreateJoinAttemptCommand command, + CancellationToken cancellationToken = default) => Atomic(now => + { + ArgumentNullException.ThrowIfNull(command); + ValidateAttempt(command); + ValidateIdempotency(command.IdempotencyKey, command.RequestFingerprint); + ValidateSubject(command.IdempotencyOwner, nameof(command.IdempotencyOwner)); + ValidateSubject(command.ClientSubject, nameof(command.ClientSubject)); + + StoreResult? admission = CheckNewWorkAdmission(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 BindAttemptEndpoint( + BindAttemptEndpointCommand command, + CancellationToken cancellationToken = default) => Atomic(_ => + { + 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 ConsumeIntroduction( + MediationHandle handle, + CancellationToken cancellationToken = default) => Atomic(_ => + { + 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 ConsumeReplay( + ReplayConsumption consumption, + CancellationToken cancellationToken = default) => Atomic(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 RevokeListing( + SessionListingId listingId, + CancellationToken cancellationToken = default) => Atomic(_ => + { + if (!_listings.ContainsKey(listingId)) + { + return new(StoreResultCode.NotFound); + } + + RemoveListing(listingId); + return new(StoreResultCode.Success, true); + }, cancellationToken); + + public StoreResult RevokePrincipal( + string subject, + TimeSpan lifetime, + CancellationToken cancellationToken = default) => Atomic(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 Atomic(Func> operation, CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + lock (_gate) + { + cancellationToken.ThrowIfCancellationRequested(); + TimeSpan now = _monotonicClock.Elapsed; + Cleanup(now); + return operation(now); + } + } + + private StoreResult? CheckNewWorkAdmission(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 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); +} diff --git a/tests/FinalFactory.Rendezvous.Tests/State/EphemeralStateTestData.cs b/tests/FinalFactory.Rendezvous.Tests/State/EphemeralStateTestData.cs new file mode 100644 index 0000000..e13a805 --- /dev/null +++ b/tests/FinalFactory.Rendezvous.Tests/State/EphemeralStateTestData.cs @@ -0,0 +1,105 @@ +using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.State; + +namespace FinalFactory.Rendezvous.Tests.State; + +internal sealed class ManualRendezvousClock : IWallClock, IMonotonicClock +{ + public DateTimeOffset UtcNow { get; private set; } = new(2026, 7, 16, 0, 0, 0, TimeSpan.Zero); + public TimeSpan Elapsed { get; private set; } + + public void Advance(TimeSpan duration) + { + Elapsed += duration; + UtcNow += duration; + } + + public void MoveWall(TimeSpan duration) => UtcNow += duration; +} + +internal sealed class EphemeralStateFixture +{ + private int _sequence; + + public EphemeralStateFixture(EphemeralStoreOptions? options = null) + { + Clock = new(); + Store = new(options ?? new EphemeralStoreOptions(), Clock, Clock); + } + + public ManualRendezvousClock Clock { get; } + public InMemoryEphemeralRendezvousStore Store { get; } + public TenantScope Scope { get; } = new(new GameId("space-game"), new EnvironmentId("test")); + + public CreateListingCommand ListingCommand( + string owner = "publisher-1", + string? idempotencyKey = null, + string? requestFingerprint = null) + { + int sequence = Interlocked.Increment(ref _sequence); + return new( + idempotencyKey ?? $"register-{sequence}", + requestFingerprint ?? $"request-{sequence}", + new ListingDefinition + { + ListingId = NewListingId(), + LeaseId = NewLeaseId(), + Scope = Scope, + OwnerSubject = owner, + RegionId = new RegionId("eu-central"), + ProtocolVersion = 7, + BuildVersion = "1.2.3", + DisplayName = "Test host", + Visibility = ListingVisibility.Public, + TrustMode = PublisherTrustMode.ManagedDedicated, + CurrentPlayers = 1, + MaximumPlayers = 8, + Metadata = new Dictionary(StringComparer.Ordinal) { ["mode"] = "coop" }, + LeaseFingerprint = Fingerprint($"lease-{sequence}"), + HostPresenceHandle = NewHandle(), + HostPresenceFingerprint = Fingerprint($"presence-{sequence}"), + }); + } + + public StoredListing CreateVisibleListing(out CreateListingCommand command) + { + command = ListingCommand(); + StoreResult created = Store.CreateListing(command); + Assert.True(created.Succeeded); + StoreResult bound = Store.BindHostPresence(new( + command.Listing.HostPresenceHandle, + command.Listing.HostPresenceFingerprint, + PublicEndpoint(40_000), + LocalEndpoint(40_000))); + Assert.True(bound.Succeeded); + return bound.Value!; + } + + public CreateJoinAttemptCommand AttemptCommand(StoredListing listing, string owner = "client-1") + { + int sequence = Interlocked.Increment(ref _sequence); + return new() + { + IdempotencyOwner = owner, + IdempotencyKey = $"join-{sequence}", + RequestFingerprint = $"join-request-{sequence}", + ClientSubject = owner, + AttemptId = NewAttemptId(), + MediationHandle = NewHandle(), + Scope = listing.Definition.Scope, + ListingId = listing.Definition.ListingId, + ProtocolVersion = listing.Definition.ProtocolVersion, + HostCapabilityFingerprint = Fingerprint($"host-{sequence}"), + ClientCapabilityFingerprint = Fingerprint($"client-{sequence}"), + }; + } + + public static SecretFingerprint Fingerprint(string value) => new(value); + public static ObservedEndpoint PublicEndpoint(int port) => new(AddressFamilyKind.Ipv4, "203.0.113.10", port); + public static ObservedEndpoint OtherPublicEndpoint(int port) => new(AddressFamilyKind.Ipv4, "198.51.100.20", port); + public static ObservedEndpoint LocalEndpoint(int port) => new(AddressFamilyKind.Ipv4, "192.168.1.20", port); + public static SessionListingId NewListingId() => new(Guid.NewGuid()); + public static LeaseId NewLeaseId() => new(Guid.NewGuid()); + public static JoinAttemptId NewAttemptId() => new(Guid.NewGuid()); + public static MediationHandle NewHandle() => new(Guid.NewGuid()); +} diff --git a/tests/FinalFactory.Rendezvous.Tests/State/InMemoryEphemeralRendezvousStoreTests.cs b/tests/FinalFactory.Rendezvous.Tests/State/InMemoryEphemeralRendezvousStoreTests.cs new file mode 100644 index 0000000..aa4754a --- /dev/null +++ b/tests/FinalFactory.Rendezvous.Tests/State/InMemoryEphemeralRendezvousStoreTests.cs @@ -0,0 +1,447 @@ +using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.State; + +namespace FinalFactory.Rendezvous.Tests.State; + +public sealed class InMemoryEphemeralRendezvousStoreTests +{ + [Fact] + public void DuplicateRegistrationIsIdempotentButChangedRequestConflicts() + { + EphemeralStateFixture fixture = new(); + CreateListingCommand command = fixture.ListingCommand(); + + StoreResult first = fixture.Store.CreateListing(command); + StoreResult duplicate = fixture.Store.CreateListing(command); + StoreResult changed = fixture.Store.CreateListing(command with { RequestFingerprint = "different" }); + + Assert.True(first.Succeeded); + Assert.True(duplicate.Succeeded); + Assert.True(duplicate.IsIdempotentReplay); + Assert.Equal(first.Value!.Definition.ListingId, duplicate.Value!.Definition.ListingId); + Assert.Equal(StoreResultCode.Conflict, changed.Code); + } + + [Fact] + public void LeaseAndPresenceExpiryUseMonotonicTime() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out _); + + fixture.Clock.Advance(TimeSpan.FromSeconds(20)); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(listing.Definition.ListingId, true).Code); + Assert.True(fixture.Store.GetListing(listing.Definition.ListingId, false).Succeeded); + + fixture.Clock.Advance(TimeSpan.FromSeconds(40)); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(listing.Definition.ListingId, false).Code); + } + + [Fact] + public void WallClockMovementDoesNotExpireOrExtendLease() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out CreateListingCommand command); + + fixture.Clock.MoveWall(TimeSpan.FromDays(30)); + Assert.True(fixture.Store.GetListing(listing.Definition.ListingId, false).Succeeded); + fixture.Clock.Advance(TimeSpan.FromSeconds(1)); + StoreResult renewed = fixture.Store.RenewLease(new( + listing.Definition.ListingId, + listing.Definition.LeaseId, + listing.Definition.LeaseFingerprint, + listing.Version)); + Assert.Equal(new DateTimeOffset(2026, 7, 16, 0, 1, 1, TimeSpan.Zero), renewed.Value!.LeaseExpiresAt); + fixture.Clock.MoveWall(TimeSpan.FromDays(-60)); + fixture.Clock.Advance(TimeSpan.FromSeconds(59)); + Assert.True(fixture.Store.GetListing(listing.Definition.ListingId, false).Succeeded); + fixture.Clock.Advance(TimeSpan.FromSeconds(1)); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(command.Listing.ListingId, false).Code); + } + + [Fact] + public async Task RenewDeleteRaceIsAtomicAndDeleteAlwaysWinsEventually() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out CreateListingCommand command); + using ManualResetEventSlim start = new(false); + + Task> renew = Task.Run(() => + { + start.Wait(); + return fixture.Store.RenewLease(new( + listing.Definition.ListingId, + listing.Definition.LeaseId, + listing.Definition.LeaseFingerprint, + listing.Version)); + }); + Task> delete = Task.Run(() => + { + start.Wait(); + return fixture.Store.DeleteListing(new( + command.Listing.ListingId, + command.Listing.LeaseId, + command.Listing.LeaseFingerprint)); + }); + + start.Set(); + await Task.WhenAll(renew, delete); + StoreResult renewResult = await renew; + StoreResult deleteResult = await delete; + + Assert.True(deleteResult.Succeeded); + Assert.Contains(renewResult.Code, new[] { StoreResultCode.Success, StoreResultCode.NotFound }); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(command.Listing.ListingId, false).Code); + } + + [Fact] + public void CompareAndSwapPreventsStaleRenewal() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out _); + RenewLeaseCommand command = new( + listing.Definition.ListingId, + listing.Definition.LeaseId, + listing.Definition.LeaseFingerprint, + listing.Version); + + StoreResult first = fixture.Store.RenewLease(command); + StoreResult stale = fixture.Store.RenewLease(command); + + Assert.Equal(2, first.Value!.Version); + Assert.Equal(StoreResultCode.Conflict, stale.Code); + Assert.Equal(2, stale.Value!.Version); + } + + [Fact] + public void JoinRequiresExactScopeProtocolAndFreshHostPresence() + { + EphemeralStateFixture fixture = new(); + CreateListingCommand listingCommand = fixture.ListingCommand(); + StoredListing listing = fixture.Store.CreateListing(listingCommand).Value!; + CreateJoinAttemptCommand attempt = fixture.AttemptCommand(listing); + + Assert.Equal(StoreResultCode.NotFound, fixture.Store.CreateJoinAttempt(attempt).Code); + fixture.Store.BindHostPresence(new( + listingCommand.Listing.HostPresenceHandle, + listingCommand.Listing.HostPresenceFingerprint, + EphemeralStateFixture.PublicEndpoint(40_000), + null)); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.CreateJoinAttempt(attempt with { ProtocolVersion = 8 }).Code); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.CreateJoinAttempt(attempt with + { + Scope = new(new("other-game"), new("test")), + }).Code); + Assert.True(fixture.Store.CreateJoinAttempt(attempt).Succeeded); + } + + [Fact] + public void DuplicateJoinIsIdempotentAndDoesNotAllocateTwice() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out _); + CreateJoinAttemptCommand command = fixture.AttemptCommand(listing); + + StoreResult first = fixture.Store.CreateJoinAttempt(command); + StoreResult duplicate = fixture.Store.CreateJoinAttempt(command); + + Assert.True(first.Succeeded); + Assert.True(duplicate.Succeeded); + Assert.True(duplicate.IsIdempotentReplay); + Assert.Equal(first.Value!.AttemptId, duplicate.Value!.AttemptId); + } + + [Fact] + public void BrowseReturnsOnlyFreshPublicCompatibleListingsInStableOrder() + { + EphemeralStateFixture fixture = new(); + StoredListing visible = fixture.CreateVisibleListing(out _); + CreateListingCommand staleCommand = fixture.ListingCommand(); + fixture.Store.CreateListing(staleCommand); + CreateListingCommand unlistedCommand = fixture.ListingCommand(); + unlistedCommand = unlistedCommand with + { + Listing = unlistedCommand.Listing with { Visibility = ListingVisibility.Unlisted }, + }; + fixture.Store.CreateListing(unlistedCommand); + fixture.Store.BindHostPresence(new( + unlistedCommand.Listing.HostPresenceHandle, + unlistedCommand.Listing.HostPresenceFingerprint, + EphemeralStateFixture.PublicEndpoint(40_099), + null)); + + StoreResult> result = fixture.Store.BrowseVisibleListings(new( + fixture.Scope, + visible.Definition.ProtocolVersion, + visible.Definition.RegionId)); + + Assert.True(result.Succeeded); + Assert.Collection(result.Value!, item => Assert.Equal(visible.Definition.ListingId, item.Definition.ListingId)); + } + + [Fact] + public async Task ConcurrentEndpointBindingAcceptsOneCompleteEndpointOnly() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out _); + CreateJoinAttemptCommand command = fixture.AttemptCommand(listing); + fixture.Store.CreateJoinAttempt(command); + BindAttemptEndpointCommand first = new( + command.MediationHandle, + AttemptPeerRole.Client, + command.ClientCapabilityFingerprint, + EphemeralStateFixture.PublicEndpoint(40_001), + EphemeralStateFixture.LocalEndpoint(40_001)); + BindAttemptEndpointCommand second = first with + { + PublicEndpoint = EphemeralStateFixture.OtherPublicEndpoint(50_001), + LocalEndpoint = null, + }; + using ManualResetEventSlim start = new(false); + + Task> left = Task.Run(() => { start.Wait(); return fixture.Store.BindAttemptEndpoint(first); }); + Task> right = Task.Run(() => { start.Wait(); return fixture.Store.BindAttemptEndpoint(second); }); + start.Set(); + await Task.WhenAll(left, right); + StoreResult leftResult = await left; + StoreResult rightResult = await right; + + Assert.Equal(1, new[] { leftResult, rightResult }.Count(static result => result.Succeeded)); + Assert.Equal(1, new[] { leftResult, rightResult }.Count(static result => result.Code == StoreResultCode.ReplayRejected)); + } + + [Fact] + public void AttemptCapabilitiesAndIntroductionAreOneTime() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out _); + CreateJoinAttemptCommand command = fixture.AttemptCommand(listing); + fixture.Store.CreateJoinAttempt(command); + BindAttemptEndpointCommand host = new( + command.MediationHandle, + AttemptPeerRole.Host, + command.HostCapabilityFingerprint, + EphemeralStateFixture.PublicEndpoint(40_010), + null); + BindAttemptEndpointCommand client = new( + command.MediationHandle, + AttemptPeerRole.Client, + command.ClientCapabilityFingerprint, + EphemeralStateFixture.OtherPublicEndpoint(40_020), + null); + + Assert.True(fixture.Store.BindAttemptEndpoint(host).Succeeded); + Assert.True(fixture.Store.BindAttemptEndpoint(client).Succeeded); + Assert.True(fixture.Store.BindAttemptEndpoint(client).IsIdempotentReplay); + Assert.True(fixture.Store.ConsumeIntroduction(command.MediationHandle).Succeeded); + Assert.Equal(StoreResultCode.ReplayRejected, fixture.Store.ConsumeIntroduction(command.MediationHandle).Code); + Assert.Equal(StoreResultCode.ReplayRejected, fixture.Store.BindAttemptEndpoint(client with + { + PublicEndpoint = EphemeralStateFixture.OtherPublicEndpoint(40_021), + }).Code); + } + + [Fact] + public void ExpiredAttemptCannotBeObservedBoundOrConsumed() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out _); + CreateJoinAttemptCommand command = fixture.AttemptCommand(listing); + fixture.Store.CreateJoinAttempt(command); + + fixture.Clock.Advance(TimeSpan.FromSeconds(30)); + + Assert.Equal(StoreResultCode.NotFound, fixture.Store.BindAttemptEndpoint(new( + command.MediationHandle, + AttemptPeerRole.Client, + command.ClientCapabilityFingerprint, + EphemeralStateFixture.PublicEndpoint(40_050), + null)).Code); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.ConsumeIntroduction(command.MediationHandle).Code); + } + + [Fact] + public void GenericReplayConsumptionIsBoundedAndExpires() + { + EphemeralStoreOptions options = new() { MaxReplayEntries = 1 }; + EphemeralStateFixture fixture = new(options); + + Assert.True(fixture.Store.ConsumeReplay(new("ticket", "one")).Succeeded); + Assert.Equal(StoreResultCode.ReplayRejected, fixture.Store.ConsumeReplay(new("ticket", "one")).Code); + Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.ConsumeReplay(new("ticket", "two")).Code); + fixture.Clock.Advance(options.ReplayLifetime); + Assert.True(fixture.Store.ConsumeReplay(new("ticket", "two")).Succeeded); + } + + [Fact] + public void PresenceAttemptAndRevocationPoolsShedWithoutPartialMutation() + { + EphemeralStoreOptions options = new() + { + MaxPresenceBindings = 1, + MaxJoinAttempts = 1, + MaxRevocations = 1, + }; + EphemeralStateFixture fixture = new(options); + StoredListing first = fixture.CreateVisibleListing(out CreateListingCommand firstCommand); + CreateListingCommand secondCommand = fixture.ListingCommand(owner: "publisher-2"); + fixture.Store.CreateListing(secondCommand); + + Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.BindHostPresence(new( + secondCommand.Listing.HostPresenceHandle, + secondCommand.Listing.HostPresenceFingerprint, + EphemeralStateFixture.OtherPublicEndpoint(42_000), + null)).Code); + Assert.True(fixture.Store.CreateJoinAttempt(fixture.AttemptCommand(first)).Succeeded); + Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.CreateJoinAttempt(fixture.AttemptCommand(first, "client-2")).Code); + Assert.True(fixture.Store.RevokePrincipal("unrelated", TimeSpan.FromMinutes(1)).Succeeded); + Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.RevokePrincipal(firstCommand.Listing.OwnerSubject, TimeSpan.FromMinutes(1)).Code); + Assert.True(fixture.Store.GetListing(first.Definition.ListingId, true).Succeeded); + Assert.True(fixture.Store.GetListing(secondCommand.Listing.ListingId, false).Succeeded); + } + + [Fact] + public void ExhaustionShedsNewListingWithoutMutatingExistingState() + { + EphemeralStoreOptions options = new() { MaxListings = 1 }; + EphemeralStateFixture fixture = new(options); + CreateListingCommand first = fixture.ListingCommand(); + CreateListingCommand second = fixture.ListingCommand(); + + Assert.True(fixture.Store.CreateListing(first).Succeeded); + Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.CreateListing(second).Code); + Assert.True(fixture.Store.GetListing(first.Listing.ListingId, false).Succeeded); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(second.Listing.ListingId, false).Code); + } + + [Fact] + public void IdempotencyPoolExhaustionDoesNotCreateUntrackedResource() + { + EphemeralStoreOptions options = new() { MaxListings = 2, MaxIdempotencyEntries = 1 }; + EphemeralStateFixture fixture = new(options); + CreateListingCommand first = fixture.ListingCommand(); + CreateListingCommand second = fixture.ListingCommand(); + + Assert.True(fixture.Store.CreateListing(first).Succeeded); + Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.CreateListing(second).Code); + Assert.True(fixture.Store.GetListing(first.Listing.ListingId, false).Succeeded); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(second.Listing.ListingId, false).Code); + } + + [Fact] + public void PolicyQuotasAreCheckedInsideAtomicCreation() + { + EphemeralStateFixture fixture = new(); + CreateListingCommand first = fixture.ListingCommand(owner: "publisher-quota") with { OwnerListingLimit = 1 }; + CreateListingCommand second = fixture.ListingCommand(owner: "publisher-quota") with { OwnerListingLimit = 1 }; + Assert.True(fixture.Store.CreateListing(first).Succeeded); + Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.CreateListing(second).Code); + + fixture.Store.BindHostPresence(new( + first.Listing.HostPresenceHandle, + first.Listing.HostPresenceFingerprint, + EphemeralStateFixture.PublicEndpoint(42_100), + null)); + StoredListing listing = fixture.Store.GetListing(first.Listing.ListingId, true).Value!; + CreateJoinAttemptCommand attempt = fixture.AttemptCommand(listing) with { ScopeAttemptLimit = 1 }; + Assert.True(fixture.Store.CreateJoinAttempt(attempt).Succeeded); + Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.CreateJoinAttempt( + fixture.AttemptCommand(listing, "client-quota-2") with { ScopeAttemptLimit = 1 }).Code); + } + + [Fact] + public void InvalidDefaultSecurityValuesCannotEnterStore() + { + EphemeralStateFixture fixture = new(); + CreateListingCommand command = fixture.ListingCommand(); + + Assert.Throws(() => fixture.Store.CreateListing(command with + { + Listing = command.Listing with { LeaseFingerprint = default }, + })); + } + + [Fact] + public void RevocationRemovesEveryPathAndBlocksNewWorkAtomically() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out CreateListingCommand command); + CreateJoinAttemptCommand attempt = fixture.AttemptCommand(listing, command.Listing.OwnerSubject); + fixture.Store.CreateJoinAttempt(attempt); + + StoreResult revoked = fixture.Store.RevokePrincipal(command.Listing.OwnerSubject, TimeSpan.FromMinutes(1)); + + Assert.True(revoked.Succeeded); + Assert.Equal(2, revoked.Value); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(listing.Definition.ListingId, false).Code); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.BindAttemptEndpoint(new( + attempt.MediationHandle, + AttemptPeerRole.Client, + attempt.ClientCapabilityFingerprint, + EphemeralStateFixture.PublicEndpoint(40_030), + null)).Code); + Assert.Equal(StoreResultCode.Revoked, fixture.Store.CreateListing(fixture.ListingCommand(owner: command.Listing.OwnerSubject)).Code); + } + + [Fact] + public void RestartHasNewGenerationAndNoEphemeralState() + { + EphemeralStateFixture before = new(); + StoredListing listing = before.CreateVisibleListing(out _); + EphemeralStateFixture after = new(); + + Assert.NotEqual(before.Store.InstanceId, after.Store.InstanceId); + Assert.Equal(StoreResultCode.NotFound, after.Store.GetListing(listing.Definition.ListingId, false).Code); + } + + [Fact] + public void DrainRejectsNewWorkAllowsInflightCompletionThenClearsState() + { + EphemeralStoreOptions options = new() { GracefulDrainLifetime = TimeSpan.FromSeconds(5) }; + EphemeralStateFixture fixture = new(options); + StoredListing listing = fixture.CreateVisibleListing(out _); + CreateJoinAttemptCommand attempt = fixture.AttemptCommand(listing); + fixture.Store.CreateJoinAttempt(attempt); + fixture.Store.BeginDrain(); + + Assert.Equal(StoreResultCode.Draining, fixture.Store.CreateListing(fixture.ListingCommand()).Code); + Assert.Equal(StoreResultCode.Draining, fixture.Store.CreateJoinAttempt(fixture.AttemptCommand(listing)).Code); + Assert.True(fixture.Store.BindAttemptEndpoint(new( + attempt.MediationHandle, + AttemptPeerRole.Client, + attempt.ClientCapabilityFingerprint, + EphemeralStateFixture.PublicEndpoint(41_000), + null)).Succeeded); + + fixture.Clock.Advance(options.GracefulDrainLifetime); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(listing.Definition.ListingId, false).Code); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.BindAttemptEndpoint(new( + attempt.MediationHandle, + AttemptPeerRole.Host, + attempt.HostCapabilityFingerprint, + EphemeralStateFixture.PublicEndpoint(41_001), + null)).Code); + } + + [Fact] + public void PrecancelledOperationHasNoPartialEffect() + { + EphemeralStateFixture fixture = new(); + CreateListingCommand command = fixture.ListingCommand(); + using CancellationTokenSource cancellation = new(); + cancellation.Cancel(); + + Assert.Throws(() => fixture.Store.CreateListing(command, cancellation.Token)); + Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(command.Listing.ListingId, false).Code); + } + + [Fact] + public void UnavailableStoreFailsNewAuthorizationClosedAndErasesActiveState() + { + EphemeralStateFixture fixture = new(); + StoredListing listing = fixture.CreateVisibleListing(out _); + fixture.Store.MarkUnavailable(); + + Assert.Equal(StoreResultCode.ServiceUnavailable, fixture.Store.GetListing(listing.Definition.ListingId, false).Code); + Assert.Equal(StoreResultCode.ServiceUnavailable, fixture.Store.CreateJoinAttempt(fixture.AttemptCommand(listing)).Code); + } +}