Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 02ca502a76 |
@@ -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.
|
||||
@@ -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)
|
||||
|
||||
@@ -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<IEphemeralRendezvousStore>(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<HealthResponse>()
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
@@ -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<string, string>(StringComparer.Ordinal) { ["mode"] = "coop" },
|
||||
LeaseFingerprint = Fingerprint($"lease-{sequence}"),
|
||||
HostPresenceHandle = NewHandle(),
|
||||
HostPresenceFingerprint = Fingerprint($"presence-{sequence}"),
|
||||
});
|
||||
}
|
||||
|
||||
public StoredListing CreateVisibleListing(out CreateListingCommand command)
|
||||
{
|
||||
command = ListingCommand();
|
||||
StoreResult<StoredListing> created = Store.CreateListing(command);
|
||||
Assert.True(created.Succeeded);
|
||||
StoreResult<StoredListing> 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());
|
||||
}
|
||||
@@ -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<StoredListing> first = fixture.Store.CreateListing(command);
|
||||
StoreResult<StoredListing> duplicate = fixture.Store.CreateListing(command);
|
||||
StoreResult<StoredListing> 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<StoredListing> 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<StoreResult<StoredListing>> renew = Task.Run(() =>
|
||||
{
|
||||
start.Wait();
|
||||
return fixture.Store.RenewLease(new(
|
||||
listing.Definition.ListingId,
|
||||
listing.Definition.LeaseId,
|
||||
listing.Definition.LeaseFingerprint,
|
||||
listing.Version));
|
||||
});
|
||||
Task<StoreResult<bool>> 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<StoredListing> renewResult = await renew;
|
||||
StoreResult<bool> 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<StoredListing> first = fixture.Store.RenewLease(command);
|
||||
StoreResult<StoredListing> 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<StoredJoinAttempt> first = fixture.Store.CreateJoinAttempt(command);
|
||||
StoreResult<StoredJoinAttempt> 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<IReadOnlyList<StoredListing>> 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<StoreResult<StoredJoinAttempt>> left = Task.Run(() => { start.Wait(); return fixture.Store.BindAttemptEndpoint(first); });
|
||||
Task<StoreResult<StoredJoinAttempt>> right = Task.Run(() => { start.Wait(); return fixture.Store.BindAttemptEndpoint(second); });
|
||||
start.Set();
|
||||
await Task.WhenAll(left, right);
|
||||
StoreResult<StoredJoinAttempt> leftResult = await left;
|
||||
StoreResult<StoredJoinAttempt> 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<ArgumentException>(() => 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<int> 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<OperationCanceledException>(() => 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);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user