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 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 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 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)
|
- [Threat model](../security/threat-model.md)
|
||||||
- [Security promise and test matrix](../security/control-matrix.md)
|
- [Security promise and test matrix](../security/control-matrix.md)
|
||||||
- [Versioned HTTP and UDP contracts](../contracts/README.md)
|
- [Versioned HTTP and UDP contracts](../contracts/README.md)
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ using System.Net;
|
|||||||
using FinalFactory.Rendezvous.Contracts;
|
using FinalFactory.Rendezvous.Contracts;
|
||||||
using FinalFactory.Rendezvous.Server.Http;
|
using FinalFactory.Rendezvous.Server.Http;
|
||||||
using FinalFactory.Rendezvous.Server.Provisioning;
|
using FinalFactory.Rendezvous.Server.Provisioning;
|
||||||
|
using FinalFactory.Rendezvous.Server.State;
|
||||||
using FinalFactory.Rendezvous.Server.Transport;
|
using FinalFactory.Rendezvous.Server.Transport;
|
||||||
using Microsoft.OpenApi;
|
using Microsoft.OpenApi;
|
||||||
|
|
||||||
@@ -35,6 +36,13 @@ builder.Services.AddOpenApi("v1", static options =>
|
|||||||
builder.Services.ConfigureHttpJsonOptions(static options =>
|
builder.Services.ConfigureHttpJsonOptions(static options =>
|
||||||
ContractJson.Configure(options.SerializerOptions));
|
ContractJson.Configure(options.SerializerOptions));
|
||||||
|
|
||||||
|
SystemRendezvousClock rendezvousClock = new();
|
||||||
|
InMemoryEphemeralRendezvousStore stateStore = new(
|
||||||
|
new EphemeralStoreOptions(),
|
||||||
|
rendezvousClock,
|
||||||
|
rendezvousClock);
|
||||||
|
builder.Services.AddSingleton<IEphemeralRendezvousStore>(stateStore);
|
||||||
|
|
||||||
if (isOpenApiGeneration)
|
if (isOpenApiGeneration)
|
||||||
{
|
{
|
||||||
builder.Services.AddSingleton(new ProvisioningReadiness(false));
|
builder.Services.AddSingleton(new ProvisioningReadiness(false));
|
||||||
@@ -74,6 +82,7 @@ if (!isOpenApiGeneration)
|
|||||||
}
|
}
|
||||||
|
|
||||||
WebApplication app = builder.Build();
|
WebApplication app = builder.Build();
|
||||||
|
app.Lifetime.ApplicationStopping.Register(() => stateStore.BeginDrain());
|
||||||
|
|
||||||
app.MapOpenApi();
|
app.MapOpenApi();
|
||||||
app.MapRendezvousContractEndpoints();
|
app.MapRendezvousContractEndpoints();
|
||||||
@@ -85,8 +94,14 @@ app.MapGet(
|
|||||||
.WithTags("Health");
|
.WithTags("Health");
|
||||||
app.MapGet(
|
app.MapGet(
|
||||||
"/health/ready",
|
"/health/ready",
|
||||||
static (UdpMediatorService mediator, ProvisioningReadiness provisioning) =>
|
static (
|
||||||
mediator.LocalEndpoint is null || !provisioning.IsReady
|
UdpMediatorService mediator,
|
||||||
|
ProvisioningReadiness provisioning,
|
||||||
|
IEphemeralRendezvousStore state) =>
|
||||||
|
mediator.LocalEndpoint is null
|
||||||
|
|| !provisioning.IsReady
|
||||||
|
|| !state.IsAvailable
|
||||||
|
|| state.IsDraining
|
||||||
? Results.StatusCode(StatusCodes.Status503ServiceUnavailable)
|
? Results.StatusCode(StatusCodes.Status503ServiceUnavailable)
|
||||||
: Results.Ok(new HealthResponse { Status = "ready" }))
|
: Results.Ok(new HealthResponse { Status = "ready" }))
|
||||||
.Produces<HealthResponse>()
|
.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