diff --git a/README.md b/README.md index bf6a5af..b361ca7 100644 --- a/README.md +++ b/README.md @@ -113,6 +113,9 @@ signing, staged promotion, rollback, and migration are defined in [releases and compatibility](docs/releases/README.md). The scriptable host/browser/join diagnostic and its stable automation contract are documented in the [TestClient integration guide](docs/integration/test-client.md). +Optional bounded SSE deltas, reconnect/reset semantics, proxy requirements, and +polling fallback are documented in +[live session-list updates](docs/integration/live-session-updates.md). The package, gameplay-socket, host-admission, provisioning, metadata, key rotation, versioning, and secure rollout seams are in the [game integration guide](docs/integration/sdk-seams.md). diff --git a/docs/api/rendezvous-v1.json b/docs/api/rendezvous-v1.json index 63d8c8b..dc41e3b 100644 --- a/docs/api/rendezvous-v1.json +++ b/docs/api/rendezvous-v1.json @@ -1177,6 +1177,159 @@ } } }, + "/v1/sessions/stream": { + "get": { + "tags": [ + "Sessions" + ], + "operationId": "StreamSessions", + "parameters": [ + { + "name": "contractVersion", + "in": "query", + "required": true, + "schema": { + "type": "integer", + "format": "int32" + } + }, + { + "name": "gameId", + "in": "query", + "required": true, + "schema": { + "type": "string" + } + }, + { + "name": "environmentId", + "in": "query", + "required": true, + "schema": { + "type": "string" + } + }, + { + "name": "protocolVersion", + "in": "query", + "required": true, + "schema": { + "type": "integer", + "format": "uint32" + } + }, + { + "name": "regionId", + "in": "query", + "schema": { + "type": "string" + } + }, + { + "name": "excludeFull", + "in": "query", + "schema": { + "type": "boolean" + } + }, + { + "name": "streamCursor", + "in": "query", + "schema": { + "type": "string" + } + }, + { + "name": "Last-Event-ID", + "in": "header", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "OK", + "headers": { + "X-Rendezvous-Correlation-ID": { + "description": "Safe request correlation identifier generated by the service.", + "schema": { + "type": "string" + } + } + }, + "content": { + "text/event-stream": { + "schema": { + "$ref": "#/components/schemas/SessionStreamEvent" + } + } + } + }, + "400": { + "description": "Bad Request", + "headers": { + "X-Rendezvous-Correlation-ID": { + "description": "Safe request correlation identifier generated by the service.", + "schema": { + "type": "string" + } + } + }, + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ApiError" + } + } + } + }, + "429": { + "description": "Too Many Requests", + "headers": { + "X-Rendezvous-Correlation-ID": { + "description": "Safe request correlation identifier generated by the service.", + "schema": { + "type": "string" + } + }, + "Retry-After": { + "description": "Whole seconds before the caller should retry (1-60).", + "schema": { + "type": "integer", + "format": "int32" + } + } + }, + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ApiError" + } + } + } + }, + "503": { + "description": "Service Unavailable", + "headers": { + "X-Rendezvous-Correlation-ID": { + "description": "Safe request correlation identifier generated by the service.", + "schema": { + "type": "string" + } + } + }, + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ApiError" + } + } + } + } + } + } + }, "/v1/sessions/{listingId}/join-attempts": { "get": { "tags": [ @@ -2657,7 +2810,8 @@ "BrowseSessionsResponse": { "required": [ "contractVersion", - "items" + "items", + "streamCursor" ], "type": "object", "properties": { @@ -2676,6 +2830,9 @@ "null", "string" ] + }, + "streamCursor": { + "type": "string" } } }, @@ -3545,6 +3702,54 @@ "type": "string", "format": "uuid" }, + "SessionStreamEvent": { + "required": [ + "contractVersion", + "kind", + "cursor" + ], + "type": "object", + "properties": { + "contractVersion": { + "type": "integer", + "format": "int32" + }, + "kind": { + "$ref": "#/components/schemas/SessionStreamEventKind" + }, + "cursor": { + "type": "string" + }, + "session": { + "oneOf": [ + { + "type": "null" + }, + { + "$ref": "#/components/schemas/SessionListing" + } + ] + }, + "listingId": { + "oneOf": [ + { + "type": "null" + }, + { + "$ref": "#/components/schemas/SessionListingId" + } + ] + } + } + }, + "SessionStreamEventKind": { + "enum": [ + "sessionUpsert", + "sessionRemove", + "reset", + "keepalive" + ] + }, "UpdateSessionRequest": { "required": [ "contractVersion", @@ -3563,6 +3768,33 @@ "leaseToken": { "type": "string" }, + "regionId": { + "oneOf": [ + { + "type": "null" + }, + { + "$ref": "#/components/schemas/RegionId" + } + ] + }, + "protocolVersion": { + "type": [ + "null", + "integer" + ], + "format": "uint32" + }, + "visibility": { + "oneOf": [ + { + "type": "null" + }, + { + "$ref": "#/components/schemas/ListingVisibility" + } + ] + }, "buildVersion": { "type": "string" }, diff --git a/docs/integration/live-session-updates.md b/docs/integration/live-session-updates.md new file mode 100644 index 0000000..b114532 --- /dev/null +++ b/docs/integration/live-session-updates.md @@ -0,0 +1,114 @@ +# Live session-list updates + +Tracking: #26 + +Live updates are an optional acceleration for an open server browser. The +bounded `GET /v1/sessions` snapshot remains the source of truth, and join +authorization still revalidates current capacity, presence, policy, and +compatibility. A displayed player count is advisory, never an admission promise. + +## Snapshot, stream, reset + +Every `BrowseSessionsResponse` includes `streamCursor` in addition to its normal +pagination cursor. Connect to `GET /v1/sessions/stream` with the same game, +environment, protocol, optional region, and `excludeFull` filter. Send the most +recent stream cursor as `Last-Event-ID`. + +| SSE event | Contract kind | UI action | +| --- | --- | --- | +| `session_upsert` | `sessionUpsert` | Add or replace the complete public projection by listing ID. | +| `session_remove` | `sessionRemove` | Remove the listing ID. | +| `reset` | `reset` | Discard local state, fetch a fresh snapshot, then reconnect with its cursor. | +| `keepalive` | `keepalive` | Preserve the cursor and connection; do not change UI state. | + +Each SSE `id` equals the opaque cursor inside its JSON event. Cursors are signed, +short-lived, monotonically ordered, and bound to the complete filter. A missing, +expired, corrupted, foreign, future, or replay-gapped cursor produces `reset` +instead of a potentially incomplete view. Do not parse or retain it as a stable +identifier. + +Updates cover creation after fresh UDP presence, public-field/capacity changes, +presence staleness and recovery, lease expiry, deregistration, operator or +principal revocation, and visibility/region/protocol changes. Events contain the +same bounded public `SessionListing` as snapshots. They never contain raw peer +endpoints, lease tokens, punch capabilities, tickets, publisher subjects, or +internal store identifiers. + +## SDK and polling fallback + +```csharp +BrowseSessionsRequest filter = new() +{ + GameId = new("space-game"), + EnvironmentId = new("production"), + ProtocolVersion = 7, + RegionId = new("eu-central"), + ExcludeFull = true, +}; +RendezvousClientResult snapshot = + await browser.BrowseAsync(filter, cancellationToken); + +await foreach (RendezvousClientResult update in + browser.StreamAsync(filter, snapshot.Value!.StreamCursor, cancellationToken)) +{ + if (!update.IsSuccess) + { + // Switch to bounded polling with jittered backoff. + break; + } + // Apply upsert/remove by listing ID. On reset, discard and browse again. +} +``` + +Cancellation or enumerator disposal closes the response and releases the server +subscription. A normal connection-duration close is a reconnect signal: use the +last applied event cursor. Repeated failures, unsupported platform HTTP stacks, +and restrictive proxies fall back to snapshots with exponential jittered +backoff, a capped interval, and `Retry-After`. Never open parallel streams to +compensate for a slow UI. + +## TestClient + +```bash +dotnet run --project src/FinalFactory.Rendezvous.TestClient \ + --configuration Release --no-build -- \ + watch --service https://rendezvous.example.invalid/ \ + --game space-game --environment production --region eu-central --protocol 7 \ + --run-seconds 60 --json +``` + +`watch.snapshot`, `watch.session-upsert`, `watch.session-remove`, +`watch.keepalive`, and `watch.reconnect` are stable diagnostics. Add +`--exercise-reset --script` to corrupt the snapshot cursor deliberately and +verify a typed reset plus snapshot refresh. Use `--exercise-reconnect --script` +while producing one update to close the first stream deliberately, reconnect +from its prior cursor, and verify that the same ordered event is replayed. +Polished list diffing, selection retention, animation, and accessibility remain +in each game. + +## Bounds and slow consumers + +The v1 journal retains at most 4,096 public-only changes. It admits at most 256 +subscribers total and 64 per tenant, reads at most 128 changes per batch, +waits a configurable 50 milliseconds after a live change and coalesces the +resulting batch to the final change per listing, sends a keepalive every 15 +seconds, and closes a connection after five minutes. A consumer behind the +replay window receives `reset`; it never acquires an unbounded queue. + +Normal optional-work concurrency and per-source/tenant rate controls apply for +the stream lifetime. Exhaustion returns typed HTTP `429` before streaming. +Shutdown cancels streams; reconnect only after readiness returns and expect a +reset after a single-active restart because listings and replay are ephemeral. + +## Reverse proxy + +- Disable response buffering (`X-Accel-Buffering: no` is also emitted), + compression, transformation, and caching for `text/event-stream`. +- Preserve `Last-Event-ID`; set upstream/read timeouts above the 15-second + keepalive and around six minutes for the five-minute connection ceiling. +- Flush events promptly and use HTTP/2 only when streaming semantics survive. +- Preserve the source-IP trust boundary and abuse controls; do not add a bypass. + +Verify the deployed proxy with an idle keepalive, update, reconnect, invalid +cursor reset, slow reader, and graceful shutdown. An in-process pass does not +prove that a production proxy is non-buffering. diff --git a/docs/integration/test-client.md b/docs/integration/test-client.md index a807fee..e58e8a9 100644 --- a/docs/integration/test-client.md +++ b/docs/integration/test-client.md @@ -146,7 +146,14 @@ Never use an unbounded sleep to orchestrate processes. Wait for versioned events such as `host.ready` and apply a deadline. Useful success events are `host.registered`, `host.ready`, `host.direct-traffic`, `host.deregistered`, `browse.completed`, `browse.session`, `join.connected`, `join.direct-traffic`, -and `join.outcome-report`. +`join.outcome-report`, `watch.snapshot`, `watch.session-upsert`, +`watch.session-remove`, `watch.reset`, `watch.reconnect`, and `watch.complete`. + +For a bounded live-directory diagnostic, use `watch --run-seconds 60`. Add +`--exercise-reset --script` to prove fail-closed cursor recovery, or +`--exercise-reconnect --script` while changing one listing to prove ordered +`Last-Event-ID` replay after a deliberate disconnect. The full event and proxy +contract is in [live session-list updates](live-session-updates.md). | Exit | Meaning | | ---: | --- | diff --git a/src/FinalFactory.Rendezvous.Client/RendezvousClientAbstractions.cs b/src/FinalFactory.Rendezvous.Client/RendezvousClientAbstractions.cs index 7579856..a8b6b23 100644 --- a/src/FinalFactory.Rendezvous.Client/RendezvousClientAbstractions.cs +++ b/src/FinalFactory.Rendezvous.Client/RendezvousClientAbstractions.cs @@ -144,6 +144,11 @@ public interface IRendezvousSessionBrowserClient EnvironmentId environmentId, uint protocolVersion, CancellationToken cancellationToken = default); + + IAsyncEnumerable> StreamAsync( + BrowseSessionsRequest request, + string streamCursor, + CancellationToken cancellationToken = default); } public interface IRendezvousJoinClient diff --git a/src/FinalFactory.Rendezvous.Client/RendezvousHttpTransport.cs b/src/FinalFactory.Rendezvous.Client/RendezvousHttpTransport.cs index 577acf2..d4a87d7 100644 --- a/src/FinalFactory.Rendezvous.Client/RendezvousHttpTransport.cs +++ b/src/FinalFactory.Rendezvous.Client/RendezvousHttpTransport.cs @@ -114,6 +114,49 @@ internal sealed class RendezvousHttpTransport } } + internal async Task> OpenStreamAsync( + Func requestFactory, + CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + using CancellationTokenSource requestTimeout = + CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + requestTimeout.CancelAfter(_options.RequestTimeout); + CancellationToken requestCancellation = requestTimeout.Token; + using HttpRequestMessage request = requestFactory(); + HttpResponseMessage? response = null; + try + { + response = await _httpClient.SendAsync( + request, + HttpCompletionOption.ResponseHeadersRead, + requestCancellation).ConfigureAwait(false); + if (response.IsSuccessStatusCode) + { + HttpResponseMessage ownedResponse = response; + response = null; + return RendezvousClientResult.Success(ownedResponse); + } + + ApiError error = await ReadErrorAsync(response, requestCancellation).ConfigureAwait(false); + int? retryAfter = error.RetryAfterSeconds ?? GetRetryAfterSeconds(response.Headers.RetryAfter); + return RendezvousClientResult.Failure( + error.Code, + error.Message, + retryAfter); + } + catch (Exception exception) when (IsTransientTransportFailure(exception, cancellationToken)) + { + return RendezvousClientResult.Failure( + RendezvousErrorCode.ServiceUnavailable, + "The Rendezvous event stream could not be opened."); + } + finally + { + response?.Dispose(); + } + } + internal static HttpRequestMessage JsonRequest( HttpMethod method, string uri, diff --git a/src/FinalFactory.Rendezvous.Client/RendezvousPublisherClient.cs b/src/FinalFactory.Rendezvous.Client/RendezvousPublisherClient.cs index 16c40a0..306da63 100644 --- a/src/FinalFactory.Rendezvous.Client/RendezvousPublisherClient.cs +++ b/src/FinalFactory.Rendezvous.Client/RendezvousPublisherClient.cs @@ -84,6 +84,9 @@ public sealed class RendezvousPublisherClient : IRendezvousPublisherClient { ContractVersion = request.ContractVersion, LeaseToken = session.LeaseToken, + RegionId = request.RegionId, + ProtocolVersion = request.ProtocolVersion, + Visibility = request.Visibility, BuildVersion = request.BuildVersion, DisplayName = request.DisplayName, Capacity = CopyCapacity(request.Capacity), diff --git a/src/FinalFactory.Rendezvous.Client/RendezvousSessionBrowserClient.cs b/src/FinalFactory.Rendezvous.Client/RendezvousSessionBrowserClient.cs index 7b7f185..6a2e5da 100644 --- a/src/FinalFactory.Rendezvous.Client/RendezvousSessionBrowserClient.cs +++ b/src/FinalFactory.Rendezvous.Client/RendezvousSessionBrowserClient.cs @@ -1,3 +1,7 @@ +using System.Net.Http.Headers; +using System.Runtime.CompilerServices; +using System.Text; +using System.Text.Json; using FinalFactory.Rendezvous.Contracts; namespace FinalFactory.Rendezvous.Client; @@ -106,5 +110,264 @@ public sealed class RendezvousSessionBrowserClient : IRendezvousSessionBrowserCl cancellationToken); } + public async IAsyncEnumerable> StreamAsync( + BrowseSessionsRequest request, + string streamCursor, + [EnumeratorCancellation] CancellationToken cancellationToken = default) + { + if (request is null) + { + throw new ArgumentNullException(nameof(request)); + } + if (string.IsNullOrWhiteSpace(streamCursor) + || !ContractValidation.IsCursorValid(streamCursor)) + { + throw new ArgumentException("A valid snapshot stream cursor is required.", nameof(streamCursor)); + } + + string query = $"v1/sessions/stream?contractVersion={request.ContractVersion}" + + $"&gameId={Escape(request.GameId.Value)}" + + $"&environmentId={Escape(request.EnvironmentId.Value)}" + + $"&protocolVersion={request.ProtocolVersion}" + + $"&excludeFull={request.ExcludeFull.ToString().ToLowerInvariant()}" + + (request.RegionId.HasValue ? $"®ionId={Escape(request.RegionId.Value.Value)}" : string.Empty); + RendezvousClientResult opened = await _transport.OpenStreamAsync( + () => + { + HttpRequestMessage message = new(HttpMethod.Get, query); + message.Headers.Accept.Add(new MediaTypeWithQualityHeaderValue("text/event-stream")); + message.Headers.TryAddWithoutValidation("Last-Event-ID", streamCursor); + return message; + }, + cancellationToken).ConfigureAwait(false); + if (!opened.IsSuccess || opened.Value is null) + { + yield return RendezvousClientResult.Failure( + opened.Error, + opened.Message, + opened.RetryAfterSeconds); + yield break; + } + + using HttpResponseMessage response = opened.Value; + if (!string.Equals( + response.Content.Headers.ContentType?.MediaType, + "text/event-stream", + StringComparison.OrdinalIgnoreCase)) + { + yield return RendezvousClientResult.Failure( + RendezvousErrorCode.InternalError, + "The service returned an invalid event-stream content type."); + yield break; + } + + using Stream source = await response.Content.ReadAsStreamAsync().ConfigureAwait(false); + using SseLineReader reader = new(source); + while (true) + { + SseReadResult? read = null; + RendezvousClientResult? readFailure = null; + bool cancelled = false; + try + { + read = await ReadEventAsync(reader, cancellationToken).ConfigureAwait(false); + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + cancelled = true; + } + catch (Exception exception) when (exception is IOException or JsonException or InvalidDataException) + { + readFailure = RendezvousClientResult.Failure( + RendezvousErrorCode.InternalError, + "The service returned an invalid or oversized event stream."); + } + if (cancelled) + { + yield break; + } + if (readFailure is not null) + { + yield return readFailure; + yield break; + } + + if (read!.EndOfStream) + { + yield break; + } + yield return read.Result!; + if (!read.Result!.IsSuccess) + { + yield break; + } + } + } + + private static async Task ReadEventAsync( + SseLineReader reader, + CancellationToken cancellationToken) + { + string? eventName = null; + string? id = null; + string? data = null; + int bytes = 0; + while (true) + { + string? line = await reader.ReadLineAsync(cancellationToken).ConfigureAwait(false); + if (line is null) + { + return eventName is null && id is null && data is null + ? SseReadResult.End + : throw new InvalidDataException("The final SSE event was incomplete."); + } + bytes += Encoding.UTF8.GetByteCount(line) + 1; + if (bytes > ContractLimits.SessionStreamEventMaxBytes) + { + throw new InvalidDataException("The SSE event exceeded the contract limit."); + } + if (line.Length == 0) + { + break; + } + if (line.StartsWith("event: ", StringComparison.Ordinal)) + { + eventName = line[7..]; + } + else if (line.StartsWith("id: ", StringComparison.Ordinal)) + { + id = line[4..]; + } + else if (line.StartsWith("data: ", StringComparison.Ordinal)) + { + data = line[6..]; + } + } + + SessionStreamEvent? item = data is null + ? null + : JsonSerializer.Deserialize(data, ContractJson.Options); + if (item is null + || !string.Equals(item.Cursor, id, StringComparison.Ordinal) + || !string.Equals(eventName, EventName(item.Kind), StringComparison.Ordinal) + || !IsValidShape(item)) + { + return new(false, RendezvousClientResult.Failure( + RendezvousErrorCode.InternalError, + "The service returned an invalid event envelope.")); + } + return new(false, RendezvousClientResult.Success(item)); + } + + private static string EventName(SessionStreamEventKind kind) => kind switch + { + SessionStreamEventKind.SessionUpsert => "session_upsert", + SessionStreamEventKind.SessionRemove => "session_remove", + SessionStreamEventKind.Reset => "reset", + SessionStreamEventKind.Keepalive => "keepalive", + _ => string.Empty, + }; + + private static bool IsValidShape(SessionStreamEvent item) => + item.ContractVersion == ContractLimits.ContractVersion + && !string.IsNullOrWhiteSpace(item.Cursor) + && ContractValidation.IsCursorValid(item.Cursor) + && (item.Kind == SessionStreamEventKind.SessionUpsert + && item.Session is not null + && IsValidListing(item.Session) + && item.ListingId is null + || item.Kind == SessionStreamEventKind.SessionRemove + && item.Session is null + && item.ListingId.HasValue + && item.ListingId.Value.Value != Guid.Empty + || item.Kind is SessionStreamEventKind.Reset or SessionStreamEventKind.Keepalive + && item.Session is null + && item.ListingId is null); + + private static bool IsValidListing(SessionListing listing) => + listing.ContractVersion == ContractLimits.ContractVersion + && listing.ListingId.Value != Guid.Empty + && !string.IsNullOrWhiteSpace(listing.GameId.Value) + && !string.IsNullOrWhiteSpace(listing.EnvironmentId.Value) + && !string.IsNullOrWhiteSpace(listing.RegionId.Value) + && listing.ProtocolVersion != 0 + && ContractValidation.IsBuildVersionValid(listing.BuildVersion) + && ContractValidation.IsDisplayNameValid(listing.DisplayName) + && listing.Visibility == ListingVisibility.Public + && Enum.IsDefined(typeof(PublisherTrustMode), listing.PublisherTrustMode) + && ContractValidation.IsCapacityValid(listing.Capacity) + && ContractValidation.IsMetadataValid(listing.Metadata) + && (listing.DedicatedFallback is null + || ContractValidation.IsNetworkEndpointValid(listing.DedicatedFallback)); + private static string Escape(string value) => Uri.EscapeDataString(value ?? string.Empty); + + private sealed class SseReadResult + { + public SseReadResult( + bool endOfStream, + RendezvousClientResult? result) + { + EndOfStream = endOfStream; + Result = result; + } + + public bool EndOfStream { get; } + public RendezvousClientResult? Result { get; } + public static SseReadResult End { get; } = new(true, null); + } + + private sealed class SseLineReader(Stream source) : IDisposable + { + private static readonly UTF8Encoding Utf8 = new(false, true); + private readonly byte[] _buffer = new byte[4096]; + private readonly MemoryStream _line = new(); + private int _offset; + private int _count; + + public async Task ReadLineAsync(CancellationToken cancellationToken) + { + while (true) + { + if (_offset >= _count) + { + _count = await source.ReadAsync( + _buffer.AsMemory(), + cancellationToken).ConfigureAwait(false); + _offset = 0; + if (_count == 0) + { + if (_line.Length == 0) + { + return null; + } + return TakeLine(); + } + } + + byte value = _buffer[_offset++]; + if (value == (byte)'\n') + { + return TakeLine(); + } + if (_line.Length >= ContractLimits.SessionStreamEventMaxBytes) + { + throw new InvalidDataException("An SSE line exceeded the contract limit."); + } + _line.WriteByte(value); + } + } + + public void Dispose() => _line.Dispose(); + + private string TakeLine() + { + byte[] bytes = _line.ToArray(); + _line.SetLength(0); + int length = bytes.Length > 0 && bytes[^1] == (byte)'\r' + ? bytes.Length - 1 + : bytes.Length; + return Utf8.GetString(bytes, 0, length); + } + } } diff --git a/src/FinalFactory.Rendezvous.Contracts/ContractLimits.cs b/src/FinalFactory.Rendezvous.Contracts/ContractLimits.cs index f29d8d4..4f6ff50 100644 --- a/src/FinalFactory.Rendezvous.Contracts/ContractLimits.cs +++ b/src/FinalFactory.Rendezvous.Contracts/ContractLimits.cs @@ -5,6 +5,7 @@ public static class ContractLimits public const int ContractVersion = 1; public const int HttpRequestMaxBytes = 16 * 1024; public const int BrowserResponseMaxBytes = 256 * 1024; + public const int SessionStreamEventMaxBytes = 32 * 1024; public const int UdpDatagramMaxBytes = 1_200; public const int MetadataMaxBytes = 4 * 1024; public const int MetadataMaxKeys = 32; diff --git a/src/FinalFactory.Rendezvous.Contracts/Http/SessionContracts.cs b/src/FinalFactory.Rendezvous.Contracts/Http/SessionContracts.cs index a200f39..f1d3b07 100644 --- a/src/FinalFactory.Rendezvous.Contracts/Http/SessionContracts.cs +++ b/src/FinalFactory.Rendezvous.Contracts/Http/SessionContracts.cs @@ -140,6 +140,10 @@ public sealed class UpdateSessionRequest [JsonRequired] public string LeaseToken { get; set; } = string.Empty; + public RegionId? RegionId { get; set; } + public uint? ProtocolVersion { get; set; } + public ListingVisibility? Visibility { get; set; } + [JsonRequired] public string BuildVersion { get; set; } = string.Empty; @@ -194,6 +198,32 @@ public sealed class BrowseSessionsResponse public List Items { get; set; } = []; public string? NextCursor { get; set; } + + [JsonRequired] + public string StreamCursor { get; set; } = string.Empty; +} + +public enum SessionStreamEventKind +{ + SessionUpsert = 1, + SessionRemove = 2, + Reset = 3, + Keepalive = 4, +} + +public sealed class SessionStreamEvent +{ + [JsonRequired] + public int ContractVersion { get; set; } = ContractLimits.ContractVersion; + + [JsonRequired] + public SessionStreamEventKind Kind { get; set; } + + [JsonRequired] + public string Cursor { get; set; } = string.Empty; + + public SessionListing? Session { get; set; } + public SessionListingId? ListingId { get; set; } } public sealed class GetSessionResponse diff --git a/src/FinalFactory.Rendezvous.Server/Browser/SessionBrowserService.cs b/src/FinalFactory.Rendezvous.Server/Browser/SessionBrowserService.cs index aea0deb..ed38af2 100644 --- a/src/FinalFactory.Rendezvous.Server/Browser/SessionBrowserService.cs +++ b/src/FinalFactory.Rendezvous.Server/Browser/SessionBrowserService.cs @@ -12,6 +12,8 @@ internal sealed record BrowserServiceResult(RendezvousErrorCode Error, T? Val internal sealed class SessionBrowserService( IEphemeralRendezvousStore store, SessionBrowserCursorCodec cursors, + SessionStreamCursorCodec streamCursors, + SessionChangeJournal changes, IWallClock clock) { public BrowserServiceResult Browse( @@ -45,9 +47,25 @@ internal sealed class SessionBrowserService( request.PageSize + 1, after, request.ExcludeFull); - StoreResult> found = store.BrowseVisibleListings( - query, - cancellationToken); + StoreResult> found = default!; + long streamRevision = 0; + bool stableSnapshot = false; + for (int attempt = 0; attempt < 3; attempt++) + { + long before = changes.CurrentRevision; + found = store.BrowseVisibleListings(query, cancellationToken); + long afterRevision = changes.CurrentRevision; + if (before == afterRevision) + { + streamRevision = afterRevision; + stableSnapshot = true; + break; + } + } + if (!stableSnapshot) + { + return new(RendezvousErrorCode.ServiceUnavailable); + } if (!found.Succeeded || found.Value is null) { return new(found.Code == StoreResultCode.ServiceUnavailable @@ -65,7 +83,12 @@ internal sealed class SessionBrowserService( string? nextCursor = hasMore ? cursors.Encode(query, items[^1].ListingId, clock.UtcNow) : null; - BrowseSessionsResponse response = new() { Items = items, NextCursor = nextCursor }; + BrowseSessionsResponse response = new() + { + Items = items, + NextCursor = nextCursor, + StreamCursor = streamCursors.Encode(query, streamRevision, clock.UtcNow), + }; if (JsonSerializer.SerializeToUtf8Bytes(response, ContractJson.Options).Length <= ContractLimits.BrowserResponseMaxBytes) { @@ -76,7 +99,10 @@ internal sealed class SessionBrowserService( hasMore = true; } - return new(RendezvousErrorCode.None, new BrowseSessionsResponse()); + return new(RendezvousErrorCode.None, new BrowseSessionsResponse + { + StreamCursor = streamCursors.Encode(query, streamRevision, clock.UtcNow), + }); } public BrowserServiceResult Get( diff --git a/src/FinalFactory.Rendezvous.Server/Browser/SessionChangeJournal.cs b/src/FinalFactory.Rendezvous.Server/Browser/SessionChangeJournal.cs new file mode 100644 index 0000000..a80e841 --- /dev/null +++ b/src/FinalFactory.Rendezvous.Server/Browser/SessionChangeJournal.cs @@ -0,0 +1,276 @@ +using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.State; + +namespace FinalFactory.Rendezvous.Server.Browser; + +internal sealed record SessionChangeJournalOptions +{ + public int ReplayCapacity { get; init; } = 4096; + public int MaximumSubscribers { get; init; } = 256; + public int MaximumSubscribersPerTenant { get; init; } = 64; + public int MaximumBatchSize { get; init; } = 128; + public TimeSpan CoalesceInterval { get; init; } = TimeSpan.FromMilliseconds(50); + public TimeSpan KeepaliveInterval { get; init; } = TimeSpan.FromSeconds(15); + public TimeSpan MaximumConnectionDuration { get; init; } = TimeSpan.FromMinutes(5); + + public void Validate() + { + if (ReplayCapacity is < 64 or > 65_536 + || MaximumSubscribers is < 1 or > 4096 + || MaximumSubscribersPerTenant < 1 + || MaximumSubscribersPerTenant > MaximumSubscribers + || MaximumBatchSize is < 1 or > 1024 + || CoalesceInterval < TimeSpan.Zero + || CoalesceInterval > TimeSpan.FromSeconds(1) + || KeepaliveInterval < TimeSpan.FromSeconds(1) + || KeepaliveInterval > TimeSpan.FromMinutes(1) + || MaximumConnectionDuration < KeepaliveInterval + || MaximumConnectionDuration > TimeSpan.FromMinutes(30)) + { + throw new ArgumentOutOfRangeException(nameof(SessionChangeJournalOptions)); + } + } +} + +internal sealed record SessionChange( + long Revision, + SessionListingProjection? Before, + SessionListingProjection? After) +{ + public SessionListingId ListingId => (After ?? Before)!.Listing.ListingId; +} + +internal sealed record SessionChangeBatch( + long CurrentRevision, + bool RequiresReset, + IReadOnlyList Changes); + +internal sealed class SessionListingProjection +{ + private SessionListingProjection(SessionListing listing, bool visible) + { + Listing = listing; + Visible = visible; + } + + public SessionListing Listing { get; } + public bool Visible { get; } + + public static SessionListingProjection From(StoredListing stored) => new( + new SessionListing + { + ListingId = stored.Definition.ListingId, + GameId = stored.Definition.Scope.GameId, + EnvironmentId = stored.Definition.Scope.EnvironmentId, + RegionId = stored.Definition.RegionId, + ProtocolVersion = stored.Definition.ProtocolVersion, + BuildVersion = stored.Definition.BuildVersion, + DisplayName = stored.Definition.DisplayName, + Visibility = stored.Definition.Visibility, + PublisherTrustMode = stored.Definition.TrustMode, + Capacity = new SessionCapacity + { + CurrentPlayers = stored.Definition.CurrentPlayers, + MaximumPlayers = stored.Definition.MaximumPlayers, + }, + Metadata = new Dictionary(stored.Definition.Metadata, StringComparer.Ordinal), + DedicatedFallback = StoredListing.CopyEndpoint(stored.Definition.DedicatedFallback), + }, + stored.HasFreshPresence && stored.Definition.Visibility == ListingVisibility.Public); + + public bool Matches(VisibleListingQuery query) => Visible + && Listing.GameId == query.Scope.GameId + && Listing.EnvironmentId == query.Scope.EnvironmentId + && Listing.ProtocolVersion == query.ProtocolVersion + && (!query.RegionId.HasValue || Listing.RegionId == query.RegionId.Value) + && (!query.ExcludeFull + || Listing.Capacity.CurrentPlayers < Listing.Capacity.MaximumPlayers); + + public static bool Equivalent(SessionListingProjection? left, SessionListingProjection? right) + { + if (ReferenceEquals(left, right)) + { + return true; + } + if (left is null || right is null || left.Visible != right.Visible) + { + return false; + } + + SessionListing a = left.Listing; + SessionListing b = right.Listing; + return a.ListingId == b.ListingId + && a.GameId == b.GameId + && a.EnvironmentId == b.EnvironmentId + && a.RegionId == b.RegionId + && a.ProtocolVersion == b.ProtocolVersion + && string.Equals(a.BuildVersion, b.BuildVersion, StringComparison.Ordinal) + && string.Equals(a.DisplayName, b.DisplayName, StringComparison.Ordinal) + && a.Visibility == b.Visibility + && a.PublisherTrustMode == b.PublisherTrustMode + && a.Capacity.CurrentPlayers == b.Capacity.CurrentPlayers + && a.Capacity.MaximumPlayers == b.Capacity.MaximumPlayers + && a.Metadata.Count == b.Metadata.Count + && a.Metadata.All(item => b.Metadata.TryGetValue(item.Key, out string? value) + && string.Equals(item.Value, value, StringComparison.Ordinal)) + && EndpointEquals(a.DedicatedFallback, b.DedicatedFallback); + } + + private static bool EndpointEquals(NetworkEndpoint? left, NetworkEndpoint? right) => + left is null && right is null + || left is not null && right is not null + && left.AddressFamily == right.AddressFamily + && string.Equals(left.Address, right.Address, StringComparison.Ordinal) + && left.Port == right.Port; +} + +internal sealed class SessionChangeJournal +{ + private readonly object _gate = new(); + private readonly SessionChangeJournalOptions _options; + private readonly Queue _changes = []; + private TaskCompletionSource _changed = NewSignal(); + private long _revision; + private int _subscribers; + private readonly Dictionary _subscribersByTenant = []; + + public SessionChangeJournal(SessionChangeJournalOptions options) + { + ArgumentNullException.ThrowIfNull(options); + options.Validate(); + _options = options; + } + + public SessionChangeJournalOptions Options => _options; + + public long CurrentRevision + { + get + { + lock (_gate) + { + return _revision; + } + } + } + + public void Publish(StoredListing? before, StoredListing? after) + { + SessionListingProjection? previous = before is null ? null : SessionListingProjection.From(before); + SessionListingProjection? current = after is null ? null : SessionListingProjection.From(after); + if (SessionListingProjection.Equivalent(previous, current) + || previous is { Visible: false } && current is null + || previous is null && current is { Visible: false }) + { + return; + } + + TaskCompletionSource signal; + long revision; + lock (_gate) + { + revision = ++_revision; + _changes.Enqueue(new SessionChange(revision, previous, current)); + while (_changes.Count > _options.ReplayCapacity) + { + _changes.Dequeue(); + } + signal = _changed; + _changed = NewSignal(); + } + signal.TrySetResult(revision); + } + + public SessionChangeBatch ReadAfter(long revision) + { + lock (_gate) + { + long oldest = _changes.TryPeek(out SessionChange? first) + ? first.Revision + : _revision + 1; + if (revision < oldest - 1 || revision > _revision) + { + return new(_revision, true, []); + } + + SessionChange[] changes = _changes + .Where(change => change.Revision > revision) + .Take(_options.MaximumBatchSize) + .ToArray(); + return new(_revision, false, changes); + } + } + + public async Task WaitForChangeAsync( + long revision, + TimeSpan timeout, + CancellationToken cancellationToken) + { + Task signal; + lock (_gate) + { + if (_revision > revision) + { + return true; + } + signal = _changed.Task; + } + + try + { + await signal.WaitAsync(timeout, cancellationToken).ConfigureAwait(false); + return true; + } + catch (TimeoutException) + { + return false; + } + } + + public bool TrySubscribe(TenantScope scope, out IDisposable? lease) + { + lock (_gate) + { + if (_subscribers >= _options.MaximumSubscribers + || _subscribersByTenant.GetValueOrDefault(scope) + >= _options.MaximumSubscribersPerTenant) + { + lease = null; + return false; + } + _subscribers++; + _subscribersByTenant[scope] = _subscribersByTenant.GetValueOrDefault(scope) + 1; + lease = new Subscription(this, scope); + return true; + } + } + + private void Release(TenantScope scope) + { + lock (_gate) + { + _subscribers--; + int remaining = _subscribersByTenant[scope] - 1; + if (remaining == 0) + { + _subscribersByTenant.Remove(scope); + } + else + { + _subscribersByTenant[scope] = remaining; + } + } + } + + private static TaskCompletionSource NewSignal() => new( + TaskCreationOptions.RunContinuationsAsynchronously); + + private sealed class Subscription( + SessionChangeJournal owner, + TenantScope scope) : IDisposable + { + private SessionChangeJournal? _owner = owner; + + public void Dispose() => Interlocked.Exchange(ref _owner, null)?.Release(scope); + } +} diff --git a/src/FinalFactory.Rendezvous.Server/Browser/SessionStreamCursorCodec.cs b/src/FinalFactory.Rendezvous.Server/Browser/SessionStreamCursorCodec.cs new file mode 100644 index 0000000..bbc3100 --- /dev/null +++ b/src/FinalFactory.Rendezvous.Server/Browser/SessionStreamCursorCodec.cs @@ -0,0 +1,108 @@ +using System.Security.Cryptography; +using System.Text.Json; +using System.Text.Json.Serialization; +using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.State; + +namespace FinalFactory.Rendezvous.Server.Browser; + +internal sealed class SessionStreamCursorCodec : IDisposable +{ + private const string Prefix = "rvs1"; + private readonly EphemeralCursorProtector _protector = new(); + + public string Encode(VisibleListingQuery query, long revision, DateTimeOffset now) + { + ArgumentOutOfRangeException.ThrowIfNegative(revision); + SessionStreamCursorPayload payload = new() + { + GameId = query.Scope.GameId.Value, + EnvironmentId = query.Scope.EnvironmentId.Value, + ProtocolVersion = query.ProtocolVersion, + RegionId = query.RegionId?.Value, + ExcludeFull = query.ExcludeFull, + Revision = revision, + ExpiresAtUnixSeconds = now.AddMinutes(10).ToUnixTimeSeconds(), + }; + byte[] encoded = JsonSerializer.SerializeToUtf8Bytes(payload, ContractJson.Options); + try + { + return _protector.Protect(Prefix, encoded); + } + finally + { + CryptographicOperations.ZeroMemory(encoded); + } + } + + public bool TryDecode( + string? cursor, + VisibleListingQuery query, + DateTimeOffset now, + out long revision) + { + revision = 0; + if (cursor is null || !_protector.TryUnprotect(Prefix, cursor, out byte[] encodedPayload)) + { + return false; + } + + SessionStreamCursorPayload? payload; + try + { + payload = JsonSerializer.Deserialize( + encodedPayload, + ContractJson.Options); + } + catch (JsonException) + { + payload = null; + } + finally + { + CryptographicOperations.ZeroMemory(encodedPayload); + } + + if (payload is null + || payload.Revision < 0 + || payload.ExpiresAtUnixSeconds <= now.ToUnixTimeSeconds() + || !string.Equals(payload.GameId, query.Scope.GameId.Value, StringComparison.Ordinal) + || !string.Equals(payload.EnvironmentId, query.Scope.EnvironmentId.Value, StringComparison.Ordinal) + || payload.ProtocolVersion != query.ProtocolVersion + || !string.Equals(payload.RegionId, query.RegionId?.Value, StringComparison.Ordinal) + || payload.ExcludeFull != query.ExcludeFull) + { + return false; + } + + revision = payload.Revision; + return true; + } + + public void Dispose() => _protector.Dispose(); + + public override string ToString() => "[SessionStreamCursorCodec: key and cursors redacted]"; +} + +internal sealed class SessionStreamCursorPayload +{ + [JsonRequired] + public string GameId { get; set; } = string.Empty; + + [JsonRequired] + public string EnvironmentId { get; set; } = string.Empty; + + [JsonRequired] + public uint ProtocolVersion { get; set; } + + public string? RegionId { get; set; } + + [JsonRequired] + public bool ExcludeFull { get; set; } + + [JsonRequired] + public long Revision { get; set; } + + [JsonRequired] + public long ExpiresAtUnixSeconds { get; set; } +} diff --git a/src/FinalFactory.Rendezvous.Server/Browser/SessionStreamService.cs b/src/FinalFactory.Rendezvous.Server/Browser/SessionStreamService.cs new file mode 100644 index 0000000..b2c513b --- /dev/null +++ b/src/FinalFactory.Rendezvous.Server/Browser/SessionStreamService.cs @@ -0,0 +1,172 @@ +using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.State; + +namespace FinalFactory.Rendezvous.Server.Browser; + +internal sealed record SessionStreamReadResult( + bool RequiresReset, + IReadOnlyList Events); + +internal sealed class SessionStreamSubscription : IDisposable +{ + private IDisposable? _lease; + + public SessionStreamSubscription( + VisibleListingQuery query, + long revision, + bool requiresReset, + IDisposable lease) + { + Query = query; + Revision = revision; + RequiresReset = requiresReset; + _lease = lease; + } + + public VisibleListingQuery Query { get; } + public long Revision { get; set; } + public bool RequiresReset { get; set; } + + public void Dispose() => Interlocked.Exchange(ref _lease, null)?.Dispose(); +} + +internal sealed class SessionStreamService( + SessionChangeJournal changes, + SessionStreamCursorCodec cursors, + IWallClock clock) +{ + public BrowserServiceResult Subscribe( + BrowseSessionsRequest request, + string? cursor) + { + ArgumentNullException.ThrowIfNull(request); + RendezvousErrorCode validation = Validate(request); + if (validation != RendezvousErrorCode.None) + { + return new(validation); + } + VisibleListingQuery query = new( + new TenantScope(request.GameId, request.EnvironmentId), + request.ProtocolVersion, + request.RegionId, + ContractLimits.BrowserPageMaxItems, + ExcludeFull: request.ExcludeFull); + if (!changes.TrySubscribe(query.Scope, out IDisposable? lease) || lease is null) + { + return new(RendezvousErrorCode.CapacityExceeded); + } + bool validCursor = cursors.TryDecode(cursor, query, clock.UtcNow, out long revision); + if (!validCursor) + { + revision = changes.CurrentRevision; + } + return new(RendezvousErrorCode.None, new SessionStreamSubscription( + query, + revision, + requiresReset: !validCursor, + lease)); + } + + public SessionStreamReadResult Read(SessionStreamSubscription subscription) + { + ArgumentNullException.ThrowIfNull(subscription); + if (subscription.RequiresReset) + { + subscription.RequiresReset = false; + return new(true, []); + } + + SessionChangeBatch batch = changes.ReadAfter(subscription.Revision); + if (batch.RequiresReset) + { + subscription.Revision = batch.CurrentRevision; + return new(true, []); + } + if (batch.Changes.Count == 0) + { + return new(false, []); + } + + Dictionary coalesced = []; + foreach (SessionChange change in batch.Changes) + { + bool beforeMatches = change.Before?.Matches(subscription.Query) == true; + bool afterMatches = change.After?.Matches(subscription.Query) == true; + if (!beforeMatches && !afterMatches) + { + continue; + } + coalesced[change.ListingId] = afterMatches + ? new(change.Revision, SessionStreamEventKind.SessionUpsert, change.After!.Listing) + : new(change.Revision, SessionStreamEventKind.SessionRemove, null); + } + + subscription.Revision = batch.Changes[^1].Revision; + SessionStreamEvent[] events = coalesced + .OrderBy(static item => item.Value.Revision) + .Select(item => ToEvent(item.Key, item.Value, subscription.Query)) + .ToArray(); + return new(false, events); + } + + public async Task WaitForChangeAsync( + SessionStreamSubscription subscription, + CancellationToken cancellationToken) + { + bool changed = await changes.WaitForChangeAsync( + subscription.Revision, + changes.Options.KeepaliveInterval, + cancellationToken).ConfigureAwait(false); + if (changed && changes.Options.CoalesceInterval > TimeSpan.Zero) + { + await Task.Delay(changes.Options.CoalesceInterval, cancellationToken) + .ConfigureAwait(false); + } + return changed; + } + + public SessionStreamEvent ResetEvent(SessionStreamSubscription subscription) => new() + { + Kind = SessionStreamEventKind.Reset, + Cursor = cursors.Encode(subscription.Query, subscription.Revision, clock.UtcNow), + }; + + public SessionStreamEvent KeepaliveEvent(SessionStreamSubscription subscription) => new() + { + Kind = SessionStreamEventKind.Keepalive, + Cursor = cursors.Encode(subscription.Query, subscription.Revision, clock.UtcNow), + }; + + public TimeSpan MaximumConnectionDuration => changes.Options.MaximumConnectionDuration; + + private SessionStreamEvent ToEvent( + SessionListingId listingId, + PendingDelta delta, + VisibleListingQuery query) => new() + { + Kind = delta.Kind, + Cursor = cursors.Encode(query, delta.Revision, clock.UtcNow), + Session = delta.Session, + ListingId = delta.Kind == SessionStreamEventKind.SessionRemove ? listingId : null, + }; + + private static RendezvousErrorCode Validate(BrowseSessionsRequest request) + { + RendezvousErrorCode version = ContractValidation.ValidateContractVersion(request.ContractVersion); + if (version != RendezvousErrorCode.None) + { + return version; + } + return string.IsNullOrEmpty(request.GameId.Value) + || string.IsNullOrEmpty(request.EnvironmentId.Value) + || request.ProtocolVersion == 0 + || request.RegionId.HasValue && string.IsNullOrEmpty(request.RegionId.Value.Value) + ? RendezvousErrorCode.InvalidRequest + : RendezvousErrorCode.None; + } + + private sealed record PendingDelta( + long Revision, + SessionStreamEventKind Kind, + SessionListing? Session); +} diff --git a/src/FinalFactory.Rendezvous.Server/FinalFactory.Rendezvous.Server.csproj b/src/FinalFactory.Rendezvous.Server/FinalFactory.Rendezvous.Server.csproj index 918f605..ddb7ca0 100644 --- a/src/FinalFactory.Rendezvous.Server/FinalFactory.Rendezvous.Server.csproj +++ b/src/FinalFactory.Rendezvous.Server/FinalFactory.Rendezvous.Server.csproj @@ -47,6 +47,9 @@ + + + diff --git a/src/FinalFactory.Rendezvous.Server/Http/ContractEndpoints.cs b/src/FinalFactory.Rendezvous.Server/Http/ContractEndpoints.cs index fb24fcb..086d99c 100644 --- a/src/FinalFactory.Rendezvous.Server/Http/ContractEndpoints.cs +++ b/src/FinalFactory.Rendezvous.Server/Http/ContractEndpoints.cs @@ -1,4 +1,5 @@ using System.Net; +using System.Text.Json; using FinalFactory.Rendezvous.Contracts; using FinalFactory.Rendezvous.Server.Abuse; using FinalFactory.Rendezvous.Server.Browser; @@ -69,6 +70,12 @@ internal static class ContractEndpoints .Produces(StatusCodes.Status429TooManyRequests) .Produces(StatusCodes.Status503ServiceUnavailable) .WithName("BrowseSessions"); + sessions.MapGet("/stream", StreamSessions) + .Produces(StatusCodes.Status200OK, contentType: "text/event-stream") + .Produces(StatusCodes.Status400BadRequest) + .Produces(StatusCodes.Status429TooManyRequests) + .Produces(StatusCodes.Status503ServiceUnavailable) + .WithName("StreamSessions"); sessions.MapGet("/{listingId}", GetSession) .Produces() .Produces(StatusCodes.Status400BadRequest) @@ -349,6 +356,123 @@ internal static class ContractEndpoints } } + private static async Task StreamSessions( + [FromQuery] int contractVersion, + [FromQuery] string gameId, + [FromQuery] string environmentId, + [FromQuery] uint protocolVersion, + [FromQuery] string? regionId, + [FromQuery] bool? excludeFull, + [FromQuery] string? streamCursor, + [FromHeader(Name = "Last-Event-ID")] string? lastEventId, + [FromServices] SessionStreamService streams, + [FromServices] AbuseProtectionService abuseProtection, + HttpContext httpContext, + CancellationToken cancellationToken) + { + if (!GameId.TryParse(gameId, out GameId parsedGameId) + || !EnvironmentId.TryParse(environmentId, out EnvironmentId parsedEnvironmentId) + || regionId is not null && !RegionId.TryParse(regionId, out _)) + { + return Error(RendezvousErrorCode.InvalidRequest); + } + if (!TryAcquireIdentity( + abuseProtection, + httpContext, + "StreamSessions", + Tenant(parsedGameId, parsedEnvironmentId), + null, + null, + out AbuseProtectionService.AbuseLease? abuseLease)) + { + return RateLimited(httpContext); + } + + using (abuseLease) + { + BrowserServiceResult subscribed = streams.Subscribe(new() + { + ContractVersion = contractVersion, + GameId = parsedGameId, + EnvironmentId = parsedEnvironmentId, + ProtocolVersion = protocolVersion, + RegionId = regionId is null ? null : new RegionId(regionId), + ExcludeFull = excludeFull ?? false, + }, string.IsNullOrEmpty(lastEventId) ? streamCursor : lastEventId); + if (!subscribed.Succeeded || subscribed.Value is null) + { + return Error(subscribed.Error); + } + + using SessionStreamSubscription subscription = subscribed.Value; + using CancellationTokenSource duration = CancellationTokenSource.CreateLinkedTokenSource( + cancellationToken); + duration.CancelAfter(streams.MaximumConnectionDuration); + HttpResponse response = httpContext.Response; + response.StatusCode = StatusCodes.Status200OK; + response.ContentType = "text/event-stream"; + response.Headers.CacheControl = "no-cache, no-store"; + response.Headers["X-Accel-Buffering"] = "no"; + await response.StartAsync(duration.Token).ConfigureAwait(false); + try + { + while (!duration.IsCancellationRequested) + { + SessionStreamReadResult read = streams.Read(subscription); + if (read.RequiresReset) + { + await WriteSseAsync(response, streams.ResetEvent(subscription), duration.Token) + .ConfigureAwait(false); + await response.Body.FlushAsync(duration.Token).ConfigureAwait(false); + break; + } + if (read.Events.Count > 0) + { + foreach (SessionStreamEvent item in read.Events) + { + await WriteSseAsync(response, item, duration.Token).ConfigureAwait(false); + } + await response.Body.FlushAsync(duration.Token).ConfigureAwait(false); + continue; + } + + bool changed = await streams.WaitForChangeAsync(subscription, duration.Token) + .ConfigureAwait(false); + if (!changed) + { + await WriteSseAsync( + response, + streams.KeepaliveEvent(subscription), + duration.Token).ConfigureAwait(false); + await response.Body.FlushAsync(duration.Token).ConfigureAwait(false); + } + } + } + catch (OperationCanceledException) when (duration.IsCancellationRequested) + { + } + return Results.Empty; + } + } + + private static async Task WriteSseAsync( + HttpResponse response, + SessionStreamEvent item, + CancellationToken cancellationToken) + { + string eventName = item.Kind switch + { + SessionStreamEventKind.SessionUpsert => "session_upsert", + SessionStreamEventKind.SessionRemove => "session_remove", + SessionStreamEventKind.Reset => "reset", + _ => "keepalive", + }; + string data = JsonSerializer.Serialize(item, ContractJson.Options); + await response.WriteAsync( + $"id: {item.Cursor}\nevent: {eventName}\ndata: {data}\n\n", + cancellationToken).ConfigureAwait(false); + } + private static IResult GetSession( SessionListingId listingId, [FromQuery] int contractVersion, diff --git a/src/FinalFactory.Rendezvous.Server/Program.cs b/src/FinalFactory.Rendezvous.Server/Program.cs index c594fb0..de25148 100644 --- a/src/FinalFactory.Rendezvous.Server/Program.cs +++ b/src/FinalFactory.Rendezvous.Server/Program.cs @@ -255,6 +255,7 @@ builder.Services.Configure(options => options.ShutdownTimeout = TimeSpan.FromSeconds(deploymentOptions.DrainDeadlineSeconds + 10)); SystemRendezvousClock rendezvousClock = new(); +SessionChangeJournal sessionChanges = new(new SessionChangeJournalOptions()); EphemeralStoreOptions stateOptions = new() { GracefulDrainLifetime = TimeSpan.FromSeconds(deploymentOptions.DrainDeadlineSeconds), @@ -262,8 +263,10 @@ EphemeralStoreOptions stateOptions = new() InMemoryEphemeralRendezvousStore stateStore = new( stateOptions, rendezvousClock, - rendezvousClock); + rendezvousClock, + sessionChanges); builder.Services.AddSingleton(stateStore); +builder.Services.AddSingleton(sessionChanges); builder.Services.AddSingleton(stateStore); builder.Services.AddSingleton(rendezvousClock); builder.Services.AddSingleton(rendezvousClock); @@ -297,7 +300,9 @@ else builder.Services.AddSingleton(SessionLeaseTiming.From(stateOptions)); builder.Services.AddSingleton(); builder.Services.AddSingleton(); + builder.Services.AddSingleton(); builder.Services.AddSingleton(); + builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); diff --git a/src/FinalFactory.Rendezvous.Server/Sessions/SessionLeaseService.cs b/src/FinalFactory.Rendezvous.Server/Sessions/SessionLeaseService.cs index 23addbe..57fe498 100644 --- a/src/FinalFactory.Rendezvous.Server/Sessions/SessionLeaseService.cs +++ b/src/FinalFactory.Rendezvous.Server/Sessions/SessionLeaseService.cs @@ -240,7 +240,15 @@ internal sealed class SessionLeaseService( } StoredListing ownedListing = listing!; - PublisherAuthorizationResult authorized = AuthorizeExisting(principal, ownedListing, request.Metadata); + PublisherAuthorizationResult authorized = authorization.Authorize( + principal, + ownedListing.Definition.Scope.GameId, + ownedListing.Definition.Scope.EnvironmentId, + request.RegionId ?? ownedListing.Definition.RegionId, + request.ProtocolVersion ?? ownedListing.Definition.ProtocolVersion, + request.Visibility ?? ownedListing.Definition.Visibility, + request.Metadata, + clock.UtcNow); if (!authorized.IsAllowed || authorized.Context is null) { return new(MapAuthorization(authorized.Error)); @@ -261,7 +269,10 @@ internal sealed class SessionLeaseService( request.Capacity.CurrentPlayers, request.Capacity.MaximumPlayers, request.Metadata, - request.DedicatedFallback), cancellationToken); + request.DedicatedFallback, + request.RegionId, + request.ProtocolVersion, + request.Visibility), cancellationToken); return updated.Succeeded ? new(RendezvousErrorCode.None, true) : new(updated.Code.ToContractError()); @@ -391,6 +402,9 @@ internal sealed class SessionLeaseService( || !ContractValidation.IsDisplayNameValid(request.DisplayName) || !ContractValidation.IsCapacityValid(request.Capacity) || !ContractValidation.IsMetadataValid(request.Metadata) + || request.RegionId.HasValue && string.IsNullOrEmpty(request.RegionId.Value.Value) + || request.ProtocolVersion.HasValue && request.ProtocolVersion.Value == 0 + || request.Visibility.HasValue && !Enum.IsDefined(request.Visibility.Value) || request.DedicatedFallback is not null && !ContractValidation.IsNetworkEndpointValid(request.DedicatedFallback) ? RendezvousErrorCode.InvalidRequest diff --git a/src/FinalFactory.Rendezvous.Server/State/EphemeralStateContracts.cs b/src/FinalFactory.Rendezvous.Server/State/EphemeralStateContracts.cs index cd15893..c7fab4d 100644 --- a/src/FinalFactory.Rendezvous.Server/State/EphemeralStateContracts.cs +++ b/src/FinalFactory.Rendezvous.Server/State/EphemeralStateContracts.cs @@ -228,7 +228,10 @@ internal sealed record UpdateListingCommand( int CurrentPlayers, int MaximumPlayers, IReadOnlyDictionary Metadata, - NetworkEndpoint? DedicatedFallback); + NetworkEndpoint? DedicatedFallback, + RegionId? RegionId = null, + uint? ProtocolVersion = null, + ListingVisibility? Visibility = null); internal sealed record DeleteListingCommand( SessionListingId ListingId, diff --git a/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs b/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs index f046a5b..9c5db33 100644 --- a/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs +++ b/src/FinalFactory.Rendezvous.Server/State/InMemoryEphemeralRendezvousStore.cs @@ -1,4 +1,5 @@ using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.Browser; namespace FinalFactory.Rendezvous.Server.State; @@ -8,6 +9,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto private readonly object _gate = new(); private readonly EphemeralStoreOptions _options; private readonly IMonotonicClock _monotonicClock; + private readonly SessionChangeJournal? _sessionChanges; private readonly DateTimeOffset _wallOrigin; private readonly TimeSpan _monotonicOrigin; private readonly Dictionary _listings = []; @@ -45,7 +47,8 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto public InMemoryEphemeralRendezvousStore( EphemeralStoreOptions options, IWallClock wallClock, - IMonotonicClock monotonicClock) + IMonotonicClock monotonicClock, + SessionChangeJournal? sessionChanges = null) { ArgumentNullException.ThrowIfNull(options); ArgumentNullException.ThrowIfNull(wallClock); @@ -53,6 +56,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto options.Validate(); _options = options; _monotonicClock = monotonicClock; + _sessionChanges = sessionChanges; _wallOrigin = wallClock.UtcNow; _monotonicOrigin = monotonicClock.Elapsed; InstanceId = Guid.NewGuid(); @@ -225,7 +229,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto now + _options.IdempotencyLifetime); _idempotency.Add(idempotencyKey, idempotency); EnqueueDeadline(_idempotencyExpiries, idempotencyKey, idempotency.Deadline); - return new(StoreResultCode.Success, Snapshot(entry)); + StoredListing created = Snapshot(entry); + _sessionChanges?.Publish(null, created); + return new(StoreResultCode.Success, created); }, cancellationToken); public StoreResult RenewLease( @@ -277,6 +283,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto || command.MaximumPlayers is <= 0 or > ContractLimits.SessionCapacityMaxPlayers || command.CurrentPlayers < 0 || command.CurrentPlayers > command.MaximumPlayers + || command.RegionId.HasValue && string.IsNullOrEmpty(command.RegionId.Value.Value) + || command.ProtocolVersion.HasValue && command.ProtocolVersion.Value == 0 + || command.Visibility.HasValue && !Enum.IsDefined(command.Visibility.Value) || !ContractValidation.IsMetadataValid(command.Metadata) || command.DedicatedFallback is not null && !ContractValidation.IsNetworkEndpointValid(command.DedicatedFallback)) @@ -302,8 +311,12 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto return new(StoreResultCode.NotFound); } + StoredListing before = Snapshot(entry); entry.Definition = StoredListing.Freeze(entry.Definition with { + RegionId = command.RegionId ?? entry.Definition.RegionId, + ProtocolVersion = command.ProtocolVersion ?? entry.Definition.ProtocolVersion, + Visibility = command.Visibility ?? entry.Definition.Visibility, BuildVersion = command.BuildVersion, DisplayName = command.DisplayName, CurrentPlayers = command.CurrentPlayers, @@ -312,7 +325,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto DedicatedFallback = command.DedicatedFallback, }); entry.Version++; - return new(StoreResultCode.Success, Snapshot(entry)); + StoredListing after = Snapshot(entry); + _sessionChanges?.Publish(before, after); + return new(StoreResultCode.Success, after); }, cancellationToken); public StoreResult DeleteListing( @@ -382,6 +397,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto return new(StoreResultCode.CapacityExceeded); } + StoredListing before = Snapshot(entry); bool isNewPresence = !_presence.ContainsKey(command.Handle); PresenceEntry presence = new( command.PublicEndpoint, @@ -396,7 +412,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto command.Handle, presence.Deadline); } - return new(StoreResultCode.Success, Snapshot(entry)); + StoredListing after = Snapshot(entry); + _sessionChanges?.Publish(before, after); + return new(StoreResultCode.Success, after); }, cancellationToken, eagerCleanup: false); public StoreResult> BrowseVisibleListings( @@ -996,6 +1014,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto private void ClearActiveState() { + StoredListing[] removedListings = _listings.Values.Select(Snapshot).ToArray(); _listings.Clear(); _listingCountsByOwner.Clear(); _listingExpiries.Clear(); @@ -1017,6 +1036,10 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto _idempotencyExpiries.Clear(); _replay.Clear(); _replayExpiries.Clear(); + foreach (StoredListing listing in removedListings) + { + _sessionChanges?.Publish(listing, null); + } } private void RemoveListing(SessionListingId listingId) @@ -1026,6 +1049,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto return; } + StoredListing removed = Snapshot(listing); _leases.Remove(listing.Definition.LeaseId); DecrementCount(_listingCountsByOwner, listing.Definition.OwnerSubject); _presenceHandles.Remove(listing.Definition.HostPresenceHandle); @@ -1045,6 +1069,8 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto RemoveOutcome(attemptId); } } + + _sessionChanges?.Publish(removed, null); } private void RemoveAttempt(JoinAttemptId attemptId) @@ -1186,6 +1212,12 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto } else if (_presence.Remove(candidate.Key)) { + if (_presenceHandles.TryGetValue(candidate.Key, out SessionListingId listingId) + && _listings.TryGetValue(listingId, out ListingEntry? listing)) + { + StoredListing after = Snapshot(listing); + _sessionChanges?.Publish(after with { HasFreshPresence = true }, after); + } removed++; } } diff --git a/src/FinalFactory.Rendezvous.TestClient/RendezvousCommandRunner.cs b/src/FinalFactory.Rendezvous.TestClient/RendezvousCommandRunner.cs index 6410358..da2cbee 100644 --- a/src/FinalFactory.Rendezvous.TestClient/RendezvousCommandRunner.cs +++ b/src/FinalFactory.Rendezvous.TestClient/RendezvousCommandRunner.cs @@ -21,6 +21,7 @@ internal sealed class RendezvousCommandRunner : ITestClientCommandRunner { TestClientMode.Host => RunHostAsync(options, output, cancellationToken), TestClientMode.Browse => RunBrowseAsync(options, output, cancellationToken), + TestClientMode.Watch => RunWatchAsync(options, output, cancellationToken), TestClientMode.Join => RunJoinAsync(options, output, input, cancellationToken), _ => Task.FromResult(TestClientExitCode.Usage), }; @@ -553,6 +554,178 @@ internal sealed class RendezvousCommandRunner : ITestClientCommandRunner } } + private static async Task RunWatchAsync( + TestClientOptions options, + TestClientOutput output, + CancellationToken cancellationToken) + { + using CancellationTokenSource watch = CancellationTokenSource.CreateLinkedTokenSource( + cancellationToken); + watch.CancelAfter(options.RunDuration ?? options.OperationTimeout); + using HttpClient http = CreateHttpClient(options); + RendezvousSessionBrowserClient browser = new(http, ClientOptions(options)); + BrowseSessionsRequest request = BrowseRequest(options); + RendezvousClientResult snapshot = await browser.BrowseAsync( + request, + watch.Token).ConfigureAwait(false); + if (!snapshot.IsSuccess || snapshot.Value is null) + { + WriteServiceFailure(output, "watch.snapshot", "directory", snapshot); + return TestClientExitCode.ServiceFailure; + } + output.Write( + "watch.snapshot", + "complete", + phase: "directory", + count: snapshot.Value.Items.Count); + string cursor = options.ExerciseReset + ? CorruptCursor(snapshot.Value.StreamCursor) + : snapshot.Value.StreamCursor; + output.Write("watch.stream", "started", phase: "live-directory"); + SessionStreamEvent? expectedReplay = null; + bool reconnectExerciseCompleted = false; + int emptyConnections = 0; + try + { + while (true) + { + string connectionCursor = cursor; + bool receivedEvent = false; + bool deliberateReconnect = false; + await foreach (RendezvousClientResult result in browser + .StreamAsync(request, cursor, watch.Token) + .ConfigureAwait(false)) + { + if (!result.IsSuccess || result.Value is null) + { + WriteServiceFailure(output, "watch.stream", "live-directory", result); + return TestClientExitCode.ServiceFailure; + } + receivedEvent = true; + SessionStreamEvent item = result.Value; + if (expectedReplay is not null) + { + if (!SameStreamEvent(expectedReplay, item)) + { + output.WriteError( + "watch.reconnect", + "failed", + "The reconnect did not replay the expected ordered event.", + phase: "live-directory"); + return TestClientExitCode.ServiceFailure; + } + output.Write("watch.reconnect", "verified", phase: "live-directory"); + expectedReplay = null; + reconnectExerciseCompleted = true; + cursor = item.Cursor; + if (options.Script) + { + return TestClientExitCode.Success; + } + } + else + { + cursor = item.Cursor; + } + + switch (item.Kind) + { + case SessionStreamEventKind.SessionUpsert when item.Session is not null: + output.Write( + "watch.session-upsert", + "available", + phase: "live-directory", + listingId: item.Session.ListingId.ToString(), + displayName: item.Session.DisplayName); + break; + case SessionStreamEventKind.SessionRemove when item.ListingId.HasValue: + output.Write( + "watch.session-remove", + "removed", + phase: "live-directory", + listingId: item.ListingId.Value.ToString()); + break; + case SessionStreamEventKind.Reset: + output.Write("watch.reset", "required", phase: "live-directory"); + RendezvousClientResult refreshed = await browser.BrowseAsync( + request, + watch.Token).ConfigureAwait(false); + if (!refreshed.IsSuccess || refreshed.Value is null) + { + WriteServiceFailure(output, "watch.snapshot", "directory", refreshed); + return TestClientExitCode.ServiceFailure; + } + output.Write( + "watch.snapshot", + "refreshed", + phase: "directory", + count: refreshed.Value.Items.Count); + return TestClientExitCode.Success; + case SessionStreamEventKind.Keepalive: + output.Write("watch.keepalive", "alive", phase: "live-directory"); + break; + } + + bool listingDelta = item.Kind is SessionStreamEventKind.SessionUpsert + or SessionStreamEventKind.SessionRemove; + if (options.ExerciseReconnect + && !reconnectExerciseCompleted + && listingDelta + && expectedReplay is null) + { + expectedReplay = item; + cursor = connectionCursor; + deliberateReconnect = true; + output.Write("watch.reconnect", "started", phase: "live-directory"); + break; + } + if (options.Script && listingDelta) + { + return TestClientExitCode.Success; + } + } + + if (deliberateReconnect) + { + continue; + } + emptyConnections = receivedEvent ? 0 : emptyConnections + 1; + if (emptyConnections >= 3) + { + output.WriteError( + "watch.reconnect", + "failed", + "The stream closed repeatedly without an event; use bounded polling fallback.", + phase: "live-directory"); + return TestClientExitCode.ServiceFailure; + } + output.Write("watch.reconnect", "required", phase: "live-directory"); + await Task.Delay(TimeSpan.FromMilliseconds(250), watch.Token).ConfigureAwait(false); + } + } + catch (OperationCanceledException) when (!cancellationToken.IsCancellationRequested) + { + output.Write("watch.complete", "complete", phase: "lifecycle"); + return TestClientExitCode.Success; + } + } + + private static bool SameStreamEvent(SessionStreamEvent expected, SessionStreamEvent actual) => + expected.Kind == actual.Kind + && string.Equals(expected.Cursor, actual.Cursor, StringComparison.Ordinal) + && expected.ListingId == actual.ListingId + && expected.Session?.ListingId == actual.Session?.ListingId; + + private static string CorruptCursor(string cursor) + { + if (string.IsNullOrEmpty(cursor)) + { + return "invalid-stream-cursor"; + } + char replacement = cursor[^1] == 'a' ? 'b' : 'a'; + return cursor[..^1] + replacement; + } + private static async Task SelectListingAsync( TestClientOptions options, TestClientOutput output, diff --git a/src/FinalFactory.Rendezvous.TestClient/TestClientOptions.cs b/src/FinalFactory.Rendezvous.TestClient/TestClientOptions.cs index 76127dc..ec3bd2c 100644 --- a/src/FinalFactory.Rendezvous.TestClient/TestClientOptions.cs +++ b/src/FinalFactory.Rendezvous.TestClient/TestClientOptions.cs @@ -8,6 +8,7 @@ internal enum TestClientMode { Host, Browse, + Watch, Join, } @@ -34,6 +35,8 @@ internal sealed class TestClientOptions internal bool Script { get; init; } internal bool Json { get; init; } internal bool ExitAfterEcho { get; init; } + internal bool ExerciseReconnect { get; init; } + internal bool ExerciseReset { get; init; } } internal sealed class TestClientParseResult @@ -63,6 +66,7 @@ internal static class TestClientOptionParser Usage: rendezvous-test-client host [options] rendezvous-test-client browse [options] + rendezvous-test-client watch [options] rendezvous-test-client join [options] Common options: @@ -88,6 +92,11 @@ internal static class TestClientOptionParser --run-seconds NUMBER Stop after 1-86400 seconds --exit-after-echo Stop after an authenticated ping/echo/ack exchange + Watch options: + --run-seconds NUMBER Stop after 1-86400 seconds + --exercise-reconnect Disconnect after an update and verify ordered replay + --exercise-reset Corrupt the snapshot cursor and verify reset/refresh + Join options: --listing UUID Join an exact listing; otherwise browse/select @@ -129,6 +138,8 @@ internal static class TestClientOptionParser bool script = false; bool json = false; bool exitAfterEcho = false; + bool exerciseReconnect = false; + bool exerciseReset = false; HashSet seen = new(StringComparer.Ordinal); for (int index = 1; index < args.Length; index++) @@ -138,7 +149,8 @@ internal static class TestClientOptionParser { return TestClientParseResult.Help(); } - if (option is "--script" or "--json" or "--exit-after-echo") + if (option is "--script" or "--json" or "--exit-after-echo" + or "--exercise-reconnect" or "--exercise-reset") { if (!seen.Add(option)) { @@ -147,6 +159,8 @@ internal static class TestClientOptionParser script |= option == "--script"; json |= option == "--json"; exitAfterEcho |= option == "--exit-after-echo"; + exerciseReconnect |= option == "--exercise-reconnect"; + exerciseReset |= option == "--exercise-reset"; continue; } if (!option.StartsWith("--", StringComparison.Ordinal) @@ -279,8 +293,11 @@ internal static class TestClientOptionParser return TestClientParseResult.Failure("One or more game, environment, region, build, or display values violate v1 limits."); } if (listingId.HasValue && mode != TestClientMode.Join - || runSeconds.HasValue && mode != TestClientMode.Host + || runSeconds.HasValue && mode is not (TestClientMode.Host or TestClientMode.Watch) || exitAfterEcho && mode != TestClientMode.Host + || exerciseReconnect && mode != TestClientMode.Watch + || exerciseReset && mode != TestClientMode.Watch + || exerciseReconnect && exerciseReset || metadata.Count > 0 && mode != TestClientMode.Host || dedicatedFallback is not null && mode != TestClientMode.Host || seen.Contains("--publisher-credential-env") && mode != TestClientMode.Host @@ -316,6 +333,8 @@ internal static class TestClientOptionParser Script = script, Json = json, ExitAfterEcho = exitAfterEcho, + ExerciseReconnect = exerciseReconnect, + ExerciseReset = exerciseReset, }); } diff --git a/tests/FinalFactory.Rendezvous.Tests/Browser/SessionBrowserTestData.cs b/tests/FinalFactory.Rendezvous.Tests/Browser/SessionBrowserTestData.cs index 82de9e0..a9a9e4b 100644 --- a/tests/FinalFactory.Rendezvous.Tests/Browser/SessionBrowserTestData.cs +++ b/tests/FinalFactory.Rendezvous.Tests/Browser/SessionBrowserTestData.cs @@ -7,18 +7,25 @@ namespace FinalFactory.Rendezvous.Tests.Browser; internal sealed class SessionBrowserFixture : IDisposable { - private readonly EphemeralStateFixture _state = new(); + private readonly EphemeralStateFixture _state; public SessionBrowserFixture() { + Changes = new(new SessionChangeJournalOptions()); + _state = new(changes: Changes); Cursors = new(); - Browser = new(_state.Store, Cursors, _state.Clock); + StreamCursors = new(); + Browser = new(_state.Store, Cursors, StreamCursors, Changes, _state.Clock); + Streams = new(Changes, StreamCursors, _state.Clock); } public InMemoryEphemeralRendezvousStore Store => _state.Store; public ManualRendezvousClock Clock => _state.Clock; public SessionBrowserCursorCodec Cursors { get; } + public SessionStreamCursorCodec StreamCursors { get; } + public SessionChangeJournal Changes { get; } public SessionBrowserService Browser { get; } + public SessionStreamService Streams { get; } public TenantScope Scope => _state.Scope; public StoredListing Add( @@ -67,5 +74,9 @@ internal sealed class SessionBrowserFixture : IDisposable PageSize = pageSize, }; - public void Dispose() => Cursors.Dispose(); + public void Dispose() + { + Cursors.Dispose(); + StreamCursors.Dispose(); + } } diff --git a/tests/FinalFactory.Rendezvous.Tests/Browser/SessionStreamServiceTests.cs b/tests/FinalFactory.Rendezvous.Tests/Browser/SessionStreamServiceTests.cs new file mode 100644 index 0000000..d69df83 --- /dev/null +++ b/tests/FinalFactory.Rendezvous.Tests/Browser/SessionStreamServiceTests.cs @@ -0,0 +1,325 @@ +using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.Browser; +using FinalFactory.Rendezvous.Server.State; +using FinalFactory.Rendezvous.Tests.State; + +namespace FinalFactory.Rendezvous.Tests.Browser; + +public sealed class SessionStreamServiceTests +{ + [Fact] + public void SnapshotPlusUpdateMatchesFreshProjection() + { + using SessionBrowserFixture fixture = new(); + StoredListing listing = fixture.Add(); + BrowseSessionsRequest request = fixture.Request(); + BrowseSessionsResponse snapshot = AssertSuccess(fixture.Browser.Browse(request)); + using SessionStreamSubscription subscription = AssertSuccess( + fixture.Streams.Subscribe(request, snapshot.StreamCursor)); + + StoreResult updated = fixture.Store.UpdateListing(Update( + listing, + displayName: "Updated host", + currentPlayers: 4)); + Assert.True(updated.Succeeded); + SessionStreamEvent delta = Assert.Single(fixture.Streams.Read(subscription).Events); + Assert.Equal(SessionStreamEventKind.SessionUpsert, delta.Kind); + Assert.Equal("Updated host", delta.Session!.DisplayName); + Assert.Equal(4, delta.Session.Capacity.CurrentPlayers); + + BrowseSessionsResponse fresh = AssertSuccess(fixture.Browser.Browse(request)); + SessionListing expected = Assert.Single(fresh.Items); + Assert.Equal(expected.DisplayName, delta.Session.DisplayName); + Assert.Equal(expected.Capacity.CurrentPlayers, delta.Session.Capacity.CurrentPlayers); + Assert.DoesNotContain("lease", System.Text.Json.JsonSerializer.Serialize(delta, ContractJson.Options), StringComparison.OrdinalIgnoreCase); + } + + [Fact] + public void PresenceStalenessRecoveryAndRevocationProduceRemoveUpsertRemove() + { + using SessionBrowserFixture fixture = new(); + StoredListing listing = fixture.Add(); + BrowseSessionsRequest request = fixture.Request(); + BrowseSessionsResponse snapshot = AssertSuccess(fixture.Browser.Browse(request)); + using SessionStreamSubscription subscription = AssertSuccess( + fixture.Streams.Subscribe(request, snapshot.StreamCursor)); + + fixture.Clock.Advance(TimeSpan.FromSeconds(21)); + AssertSuccess(fixture.Browser.Browse(request)); + SessionStreamEvent stale = Assert.Single(fixture.Streams.Read(subscription).Events); + Assert.Equal(SessionStreamEventKind.SessionRemove, stale.Kind); + Assert.Equal(listing.Definition.ListingId, stale.ListingId); + + StoreResult rebound = fixture.Store.BindHostPresence(new( + listing.Definition.HostPresenceHandle, + listing.Definition.HostPresenceFingerprint, + new ObservedEndpoint(AddressFamilyKind.Ipv4, "203.0.113.20", 40_020), + null)); + Assert.True(rebound.Succeeded); + Assert.Equal( + SessionStreamEventKind.SessionUpsert, + Assert.Single(fixture.Streams.Read(subscription).Events).Kind); + + Assert.True(fixture.Store.RevokeListing(listing.Definition.ListingId).Succeeded); + Assert.Equal( + SessionStreamEventKind.SessionRemove, + Assert.Single(fixture.Streams.Read(subscription).Events).Kind); + } + + [Fact] + public void CreationAndLeaseExpiryProduceUpsertThenRemove() + { + using SessionBrowserFixture fixture = new(); + BrowseSessionsRequest request = fixture.Request(); + BrowseSessionsResponse snapshot = AssertSuccess(fixture.Browser.Browse(request)); + using SessionStreamSubscription subscription = AssertSuccess( + fixture.Streams.Subscribe(request, snapshot.StreamCursor)); + + StoredListing listing = fixture.Add(); + SessionStreamEvent created = Assert.Single(fixture.Streams.Read(subscription).Events); + Assert.Equal(SessionStreamEventKind.SessionUpsert, created.Kind); + Assert.Equal(listing.Definition.ListingId, created.Session!.ListingId); + + for (int refresh = 0; refresh < 3; refresh++) + { + fixture.Clock.Advance(TimeSpan.FromSeconds(19)); + Assert.True(fixture.Store.BindHostPresence(new( + listing.Definition.HostPresenceHandle, + listing.Definition.HostPresenceFingerprint, + new ObservedEndpoint(AddressFamilyKind.Ipv4, "203.0.113.20", 40_020), + null)).Succeeded); + } + fixture.Clock.Advance(TimeSpan.FromSeconds(4)); + AssertSuccess(fixture.Browser.Browse(request)); + + SessionStreamEvent expired = Assert.Single(fixture.Streams.Read(subscription).Events); + Assert.Equal(SessionStreamEventKind.SessionRemove, expired.Kind); + Assert.Equal(listing.Definition.ListingId, expired.ListingId); + Assert.Equal( + StoreResultCode.NotFound, + fixture.Store.GetListing(listing.Definition.ListingId, requireFreshPresence: false).Code); + } + + [Fact] + public void ScopeProtocolRegionAndFullFiltersNeverLeak() + { + using SessionBrowserFixture fixture = new(); + BrowseSessionsRequest request = fixture.Request(); + request.ExcludeFull = true; + BrowseSessionsResponse snapshot = AssertSuccess(fixture.Browser.Browse(request)); + using SessionStreamSubscription subscription = AssertSuccess( + fixture.Streams.Subscribe(request, snapshot.StreamCursor)); + + fixture.Add(scope: new(new("other-game"), fixture.Scope.EnvironmentId)); + fixture.Add(protocolVersion: 99); + fixture.Add(regionId: new("other-region")); + fixture.Add(currentPlayers: 8, maximumPlayers: 8); + Assert.Empty(fixture.Streams.Read(subscription).Events); + } + + [Fact] + public void VisibilityCompatibilityAndRegionChangesEnterAndLeaveTheFilter() + { + using SessionBrowserFixture fixture = new(); + StoredListing listing = fixture.Add(); + BrowseSessionsRequest request = fixture.Request(); + BrowseSessionsResponse snapshot = AssertSuccess(fixture.Browser.Browse(request)); + using SessionStreamSubscription subscription = AssertSuccess( + fixture.Streams.Subscribe(request, snapshot.StreamCursor)); + + listing = fixture.Store.UpdateListing(Update( + listing, + listing.Definition.DisplayName, + 1, + visibility: ListingVisibility.Unlisted)).Value!; + Assert.Equal(SessionStreamEventKind.SessionRemove, SingleKind(fixture, subscription)); + listing = fixture.Store.UpdateListing(Update( + listing, + listing.Definition.DisplayName, + 1, + visibility: ListingVisibility.Public)).Value!; + Assert.Equal(SessionStreamEventKind.SessionUpsert, SingleKind(fixture, subscription)); + listing = fixture.Store.UpdateListing(Update( + listing, + listing.Definition.DisplayName, + 1, + protocolVersion: 99)).Value!; + Assert.Equal(SessionStreamEventKind.SessionRemove, SingleKind(fixture, subscription)); + listing = fixture.Store.UpdateListing(Update( + listing, + listing.Definition.DisplayName, + 1, + protocolVersion: 7)).Value!; + Assert.Equal(SessionStreamEventKind.SessionUpsert, SingleKind(fixture, subscription)); + listing = fixture.Store.UpdateListing(Update( + listing, + listing.Definition.DisplayName, + 1, + regionId: new RegionId("other-region"))).Value!; + Assert.Equal(SessionStreamEventKind.SessionRemove, SingleKind(fixture, subscription)); + } + + [Fact] + public void ReplayGapAndForeignCursorForceResetAndSubscriberLimitFailsClosed() + { + ManualRendezvousClock clock = new(); + SessionChangeJournal changes = new(new SessionChangeJournalOptions + { + ReplayCapacity = 64, + MaximumSubscribers = 1, + MaximumSubscribersPerTenant = 1, + }); + using SessionStreamCursorCodec cursors = new(); + SessionStreamService streams = new(changes, cursors, clock); + EphemeralStateFixture state = new(changes: changes); + StoredListing listing = state.CreateVisibleListing(out _); + BrowseSessionsRequest request = new() + { + GameId = state.Scope.GameId, + EnvironmentId = state.Scope.EnvironmentId, + ProtocolVersion = listing.Definition.ProtocolVersion, + }; + VisibleListingQuery query = new(state.Scope, listing.Definition.ProtocolVersion, null); + string initial = cursors.Encode(query, changes.CurrentRevision, clock.UtcNow); + using SessionStreamSubscription subscription = AssertSuccess(streams.Subscribe(request, initial)); + Assert.Equal( + RendezvousErrorCode.CapacityExceeded, + streams.Subscribe(request, initial).Error); + + for (int index = 0; index < 65; index++) + { + listing = state.Store.UpdateListing(Update( + listing, + displayName: $"Host {index}", + currentPlayers: index % 8)).Value!; + } + Assert.True(streams.Read(subscription).RequiresReset); + subscription.Dispose(); + + using SessionStreamSubscription foreign = AssertSuccess(streams.Subscribe( + request, + "not-a-valid-cursor")); + Assert.True(streams.Read(foreign).RequiresReset); + } + + [Fact] + public void BurstIsBoundedAndCoalescedWithoutLosingFinalState() + { + ManualRendezvousClock clock = new(); + SessionChangeJournal changes = new(new SessionChangeJournalOptions + { + ReplayCapacity = 1024, + MaximumBatchSize = 128, + }); + using SessionStreamCursorCodec cursors = new(); + SessionStreamService streams = new(changes, cursors, clock); + EphemeralStateFixture state = new(changes: changes); + StoredListing listing = state.CreateVisibleListing(out _); + BrowseSessionsRequest request = new() + { + GameId = state.Scope.GameId, + EnvironmentId = state.Scope.EnvironmentId, + ProtocolVersion = listing.Definition.ProtocolVersion, + }; + VisibleListingQuery query = new(state.Scope, listing.Definition.ProtocolVersion, null); + using SessionStreamSubscription subscription = AssertSuccess(streams.Subscribe( + request, + cursors.Encode(query, changes.CurrentRevision, clock.UtcNow))); + + for (int index = 0; index < 1000; index++) + { + listing = state.Store.UpdateListing(Update( + listing, + displayName: $"Host {index}", + currentPlayers: index % 8)).Value!; + } + + List emitted = []; + while (subscription.Revision < changes.CurrentRevision) + { + SessionStreamReadResult read = streams.Read(subscription); + Assert.False(read.RequiresReset); + emitted.AddRange(read.Events); + } + + Assert.Equal(8, emitted.Count); + SessionStreamEvent final = emitted[^1]; + Assert.Equal(SessionStreamEventKind.SessionUpsert, final.Kind); + Assert.Equal("Host 999", final.Session!.DisplayName); + Assert.Equal(7, final.Session.Capacity.CurrentPlayers); + } + + [Fact] + public void SubscriberLimitIsEnforcedPerTenantAndReleasedOnDispose() + { + ManualRendezvousClock clock = new(); + SessionChangeJournal changes = new(new SessionChangeJournalOptions + { + MaximumSubscribers = 2, + MaximumSubscribersPerTenant = 1, + }); + using SessionStreamCursorCodec cursors = new(); + SessionStreamService streams = new(changes, cursors, clock); + TenantScope firstScope = new(new("first-game"), new("production")); + TenantScope secondScope = new(new("second-game"), new("production")); + BrowseSessionsRequest firstRequest = Request(firstScope); + BrowseSessionsRequest secondRequest = Request(secondScope); + string firstCursor = cursors.Encode( + new VisibleListingQuery(firstScope, 7, null), + changes.CurrentRevision, + clock.UtcNow); + string secondCursor = cursors.Encode( + new VisibleListingQuery(secondScope, 7, null), + changes.CurrentRevision, + clock.UtcNow); + + SessionStreamSubscription first = AssertSuccess(streams.Subscribe(firstRequest, firstCursor)); + Assert.Equal( + RendezvousErrorCode.CapacityExceeded, + streams.Subscribe(firstRequest, firstCursor).Error); + using SessionStreamSubscription second = AssertSuccess( + streams.Subscribe(secondRequest, secondCursor)); + first.Dispose(); + using SessionStreamSubscription replacement = AssertSuccess( + streams.Subscribe(firstRequest, firstCursor)); + } + + private static BrowseSessionsRequest Request(TenantScope scope) => new() + { + GameId = scope.GameId, + EnvironmentId = scope.EnvironmentId, + ProtocolVersion = 7, + }; + + private static UpdateListingCommand Update( + StoredListing listing, + string displayName, + int currentPlayers, + RegionId? regionId = null, + uint? protocolVersion = null, + ListingVisibility? visibility = null) => new( + listing.Definition.ListingId, + listing.Definition.LeaseId, + listing.Definition.LeaseFingerprint, + listing.Definition.OwnerSubject, + listing.Definition.BuildVersion, + displayName, + currentPlayers, + listing.Definition.MaximumPlayers, + listing.Definition.Metadata, + listing.Definition.DedicatedFallback, + regionId, + protocolVersion, + visibility); + + private static SessionStreamEventKind SingleKind( + SessionBrowserFixture fixture, + SessionStreamSubscription subscription) => + Assert.Single(fixture.Streams.Read(subscription).Events).Kind; + + private static T AssertSuccess(BrowserServiceResult result) + { + Assert.True(result.Succeeded, result.Error.ToString()); + return Assert.IsType(result.Value); + } +} diff --git a/tests/FinalFactory.Rendezvous.Tests/Client/RendezvousClientBehaviorTests.cs b/tests/FinalFactory.Rendezvous.Tests/Client/RendezvousClientBehaviorTests.cs index 92aa29a..2185bc2 100644 --- a/tests/FinalFactory.Rendezvous.Tests/Client/RendezvousClientBehaviorTests.cs +++ b/tests/FinalFactory.Rendezvous.Tests/Client/RendezvousClientBehaviorTests.cs @@ -167,6 +167,88 @@ public sealed class RendezvousClientBehaviorTests Assert.Contains("gameId=space-game", handler.RequestUris[0].Query, StringComparison.Ordinal); } + [Fact] + public async Task StreamRejectsMalformedAndOversizedEventEnvelopes() + { + string[] bodies = + [ + "event: session_upsert\nid: valid-cursor\ndata: {}\n\n", + "data: " + new string('x', ContractLimits.SessionStreamEventMaxBytes + 1) + "\n\n", + ]; + foreach (string body in bodies) + { + StringContent content = new(body, Encoding.UTF8, "text/event-stream"); + ScriptedHandler handler = new(Response(HttpStatusCode.OK, content)); + using HttpClient httpClient = new(handler) + { + BaseAddress = new("http://rendezvous.test/"), + }; + RendezvousSessionBrowserClient browser = new(httpClient); + await using IAsyncEnumerator> events = browser + .StreamAsync(new BrowseSessionsRequest + { + GameId = new("space-game"), + EnvironmentId = new("production"), + ProtocolVersion = 7, + }, "valid-stream-cursor") + .GetAsyncEnumerator(); + + Assert.True(await events.MoveNextAsync()); + Assert.False(events.Current.IsSuccess); + Assert.Equal(RendezvousErrorCode.InternalError, events.Current.Error); + } + } + + [Fact] + public async Task StreamRequiresANonEmptySnapshotCursor() + { + using HttpClient httpClient = new(new ScriptedHandler()) + { + BaseAddress = new("http://rendezvous.test/"), + }; + RendezvousSessionBrowserClient browser = new(httpClient); + await Assert.ThrowsAsync(async () => + { + await foreach (RendezvousClientResult _ in browser.StreamAsync( + new BrowseSessionsRequest + { + GameId = new("space-game"), + EnvironmentId = new("production"), + ProtocolVersion = 7, + }, + string.Empty)) + { + } + }); + } + + [Fact] + public async Task StreamOpeningIsBoundedByTheConfiguredRequestTimeout() + { + using HttpClient httpClient = new(new SilentHandler()) + { + BaseAddress = new("http://rendezvous.test/"), + }; + RendezvousSessionBrowserClient browser = new( + httpClient, + new RendezvousClientOptions + { + MaximumSafeRetries = 0, + RequestTimeout = TimeSpan.FromMilliseconds(20), + }); + await using IAsyncEnumerator> events = browser + .StreamAsync(new BrowseSessionsRequest + { + GameId = new("space-game"), + EnvironmentId = new("production"), + ProtocolVersion = 7, + }, "valid-stream-cursor") + .GetAsyncEnumerator(); + + Assert.True(await events.MoveNextAsync().AsTask().WaitAsync(TimeSpan.FromSeconds(2))); + Assert.Equal(RendezvousErrorCode.ServiceUnavailable, events.Current.Error); + } + [Fact] public async Task LeaseMaintainerReportsLeaseLoss() { diff --git a/tests/FinalFactory.Rendezvous.Tests/Client/RendezvousClientIntegrationTests.cs b/tests/FinalFactory.Rendezvous.Tests/Client/RendezvousClientIntegrationTests.cs index 445f5dd..87515f7 100644 --- a/tests/FinalFactory.Rendezvous.Tests/Client/RendezvousClientIntegrationTests.cs +++ b/tests/FinalFactory.Rendezvous.Tests/Client/RendezvousClientIntegrationTests.cs @@ -142,6 +142,77 @@ public sealed class RendezvousClientIntegrationTests } } + [Fact] + public async Task BrowserStreamResetsInvalidCursorReplaysReconnectAndReleasesConnections() + { + await using ClientTestHost host = await ClientTestHost.StartAsync(); + RendezvousPublisherClient publisher = new(host.HttpClient); + RendezvousSessionBrowserClient browser = new(host.HttpClient); + PublishedSession session = AssertSuccess(await publisher.RegisterAsync( + CreateRegistration(200), + host.PublisherCredential)); + BindPresence(host, session, 41_200); + BrowseSessionsRequest request = BrowseRequest(); + BrowseSessionsResponse snapshot = AssertSuccess(await browser.BrowseAsync(request)); + Assert.False(string.IsNullOrWhiteSpace(snapshot.StreamCursor)); + + using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(10)); + await using (IAsyncEnumerator> invalid = browser + .StreamAsync(request, CorruptCursor(snapshot.StreamCursor), timeout.Token) + .GetAsyncEnumerator(timeout.Token)) + { + Assert.True(await invalid.MoveNextAsync()); + Assert.Equal(SessionStreamEventKind.Reset, AssertSuccess(invalid.Current).Kind); + Assert.False(await invalid.MoveNextAsync()); + } + await using IAsyncEnumerator> events = browser + .StreamAsync(request, snapshot.StreamCursor, timeout.Token) + .GetAsyncEnumerator(timeout.Token); + Task upsertPending = events.MoveNextAsync().AsTask(); + Assert.True((await publisher.UpdateAsync( + session, + new UpdateSessionRequest + { + BuildVersion = "2.0.0", + DisplayName = "Live update", + Capacity = new() { CurrentPlayers = 3, MaximumPlayers = 8 }, + Metadata = new() { ["mode"] = "online-coop" }, + }, + host.PublisherCredential, + timeout.Token)).IsSuccess); + Assert.True(await upsertPending); + SessionStreamEvent upsert = AssertSuccess(events.Current); + Assert.Equal(SessionStreamEventKind.SessionUpsert, upsert.Kind); + Assert.Equal("Live update", upsert.Session!.DisplayName); + + await using (IAsyncEnumerator> replay = browser + .StreamAsync(request, snapshot.StreamCursor, timeout.Token) + .GetAsyncEnumerator(timeout.Token)) + { + Assert.True(await replay.MoveNextAsync()); + SessionStreamEvent replayed = AssertSuccess(replay.Current); + Assert.Equal(SessionStreamEventKind.SessionUpsert, replayed.Kind); + Assert.Equal(upsert.Cursor, replayed.Cursor); + Assert.Equal("Live update", replayed.Session!.DisplayName); + } + + Task removePending = events.MoveNextAsync().AsTask(); + Assert.True((await publisher.DeregisterAsync( + session, + host.PublisherCredential, + timeout.Token)).IsSuccess); + Assert.True(await removePending); + SessionStreamEvent remove = AssertSuccess(events.Current); + Assert.Equal(SessionStreamEventKind.SessionRemove, remove.Kind); + Assert.Equal(session.ListingId, remove.ListingId); + } + + private static string CorruptCursor(string cursor) + { + char replacement = cursor[^1] == 'a' ? 'b' : 'a'; + return cursor[..^1] + replacement; + } + private static T AssertSuccess(RendezvousClientResult result) { Assert.True(result.IsSuccess, result.Message); @@ -210,7 +281,8 @@ public sealed class RendezvousClientIntegrationTests { ManualRendezvousClock clock = new(ProvisioningTestData.Now); EphemeralStoreOptions stateOptions = new(); - InMemoryEphemeralRendezvousStore store = new(stateOptions, clock, clock); + SessionChangeJournal changes = new(new SessionChangeJournalOptions()); + InMemoryEphemeralRendezvousStore store = new(stateOptions, clock, clock, changes); EphemeralCapabilityIssuer capabilities = new(); ProvisioningRuntime provisioning = ProvisioningRuntime.Create( ProvisioningTestData.CreateOptions(), @@ -239,7 +311,10 @@ public sealed class RendezvousClientIntegrationTests builder.Services.AddSingleton(SessionLeaseTiming.From(stateOptions)); builder.Services.AddSingleton(); builder.Services.AddSingleton(); + builder.Services.AddSingleton(); + builder.Services.AddSingleton(changes); builder.Services.AddSingleton(); + builder.Services.AddSingleton(); WebApplication app = builder.Build(); app.UseExceptionHandler(); diff --git a/tests/FinalFactory.Rendezvous.Tests/Contracts/OpenApiCompatibilityTests.cs b/tests/FinalFactory.Rendezvous.Tests/Contracts/OpenApiCompatibilityTests.cs index 787a565..9480677 100644 --- a/tests/FinalFactory.Rendezvous.Tests/Contracts/OpenApiCompatibilityTests.cs +++ b/tests/FinalFactory.Rendezvous.Tests/Contracts/OpenApiCompatibilityTests.cs @@ -17,6 +17,7 @@ public sealed class OpenApiCompatibilityTests "/v1/operator/principals/revoke", "/v1/operator/status", "/v1/sessions", + "/v1/sessions/stream", "/v1/sessions/{listingId}", "/v1/sessions/{listingId}/join-attempts", "/v1/sessions/{listingId}/renew", @@ -68,6 +69,25 @@ public sealed class OpenApiCompatibilityTests Assert.DoesNotContain(listingProperties, static property => property.Contains("token", StringComparison.OrdinalIgnoreCase) || property.Contains("playerId", StringComparison.OrdinalIgnoreCase)); + JsonElement streamProperties = schemas.GetProperty("SessionStreamEvent") + .GetProperty("properties"); + Assert.True(streamProperties.TryGetProperty("contractVersion", out _)); + Assert.True(streamProperties.TryGetProperty("kind", out _)); + Assert.True(streamProperties.TryGetProperty("cursor", out _)); + Assert.True(streamProperties.TryGetProperty("session", out _)); + Assert.True(streamProperties.TryGetProperty("listingId", out _)); + Assert.DoesNotContain(streamProperties.EnumerateObject(), static property => + property.Name.Contains("token", StringComparison.OrdinalIgnoreCase) + || property.Name.Contains("capability", StringComparison.OrdinalIgnoreCase) + || property.Name.Contains("ticket", StringComparison.OrdinalIgnoreCase) + || property.Name.Contains("endpoint", StringComparison.OrdinalIgnoreCase)); + Assert.True(root.GetProperty("paths") + .GetProperty("/v1/sessions/stream") + .GetProperty("get") + .GetProperty("responses") + .GetProperty("200") + .GetProperty("content") + .TryGetProperty("text/event-stream", out _)); JsonElement dedicatedFallback = schemas.GetProperty("SessionListing") .GetProperty("properties") .GetProperty("dedicatedFallback"); @@ -205,7 +225,7 @@ public sealed class OpenApiCompatibilityTests } } - Assert.Equal(17, overloadContracts); + Assert.Equal(18, overloadContracts); (string Path, string Method)[] bodyOperations = [ ("/v1/sessions", "post"), diff --git a/tests/FinalFactory.Rendezvous.Tests/JoinAttempts/JoinAttemptHttpEndpointTests.cs b/tests/FinalFactory.Rendezvous.Tests/JoinAttempts/JoinAttemptHttpEndpointTests.cs index 18ad84c..f0772d6 100644 --- a/tests/FinalFactory.Rendezvous.Tests/JoinAttempts/JoinAttemptHttpEndpointTests.cs +++ b/tests/FinalFactory.Rendezvous.Tests/JoinAttempts/JoinAttemptHttpEndpointTests.cs @@ -288,7 +288,8 @@ public sealed class JoinAttemptHttpEndpointTests { ManualRendezvousClock clock = new(ProvisioningTestData.Now); EphemeralStoreOptions stateOptions = new(); - InMemoryEphemeralRendezvousStore store = new(stateOptions, clock, clock); + SessionChangeJournal changes = new(new SessionChangeJournalOptions()); + InMemoryEphemeralRendezvousStore store = new(stateOptions, clock, clock, changes); EphemeralCapabilityIssuer capabilities = new(); ProvisioningRuntime provisioning = ProvisioningRuntime.Create( ProvisioningTestData.CreateOptions(), @@ -318,7 +319,10 @@ public sealed class JoinAttemptHttpEndpointTests builder.Services.AddSingleton(SessionLeaseTiming.From(stateOptions)); builder.Services.AddSingleton(); builder.Services.AddSingleton(); + builder.Services.AddSingleton(); + builder.Services.AddSingleton(changes); builder.Services.AddSingleton(); + builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); ConnectionOutcomeMetrics outcomeMetrics = new(); diff --git a/tests/FinalFactory.Rendezvous.Tests/Sessions/SessionHttpEndpointTests.cs b/tests/FinalFactory.Rendezvous.Tests/Sessions/SessionHttpEndpointTests.cs index 51baa91..57dd3e5 100644 --- a/tests/FinalFactory.Rendezvous.Tests/Sessions/SessionHttpEndpointTests.cs +++ b/tests/FinalFactory.Rendezvous.Tests/Sessions/SessionHttpEndpointTests.cs @@ -28,7 +28,8 @@ public sealed class SessionHttpEndpointTests { ManualRendezvousClock clock = new(ProvisioningTestData.Now); EphemeralStoreOptions stateOptions = new(); - InMemoryEphemeralRendezvousStore store = new(stateOptions, clock, clock); + SessionChangeJournal changes = new(new SessionChangeJournalOptions()); + InMemoryEphemeralRendezvousStore store = new(stateOptions, clock, clock, changes); EphemeralCapabilityIssuer capabilities = new(); ProvisioningRuntime provisioning = ProvisioningRuntime.Create( ProvisioningTestData.CreateOptions(), @@ -57,7 +58,10 @@ public sealed class SessionHttpEndpointTests builder.Services.AddSingleton(SessionLeaseTiming.From(stateOptions)); builder.Services.AddSingleton(); builder.Services.AddSingleton(); + builder.Services.AddSingleton(); + builder.Services.AddSingleton(changes); builder.Services.AddSingleton(); + builder.Services.AddSingleton(); await using WebApplication app = builder.Build(); app.UseExceptionHandler(); app.UseMiddleware(); diff --git a/tests/FinalFactory.Rendezvous.Tests/State/EphemeralStateTestData.cs b/tests/FinalFactory.Rendezvous.Tests/State/EphemeralStateTestData.cs index 523edf6..9a4f509 100644 --- a/tests/FinalFactory.Rendezvous.Tests/State/EphemeralStateTestData.cs +++ b/tests/FinalFactory.Rendezvous.Tests/State/EphemeralStateTestData.cs @@ -1,4 +1,5 @@ using FinalFactory.Rendezvous.Contracts; +using FinalFactory.Rendezvous.Server.Browser; using FinalFactory.Rendezvous.Server.State; namespace FinalFactory.Rendezvous.Tests.State; @@ -24,10 +25,12 @@ internal sealed class EphemeralStateFixture { private int _sequence; - public EphemeralStateFixture(EphemeralStoreOptions? options = null) + public EphemeralStateFixture( + EphemeralStoreOptions? options = null, + SessionChangeJournal? changes = null) { Clock = new(); - Store = new(options ?? new EphemeralStoreOptions(), Clock, Clock); + Store = new(options ?? new EphemeralStoreOptions(), Clock, Clock, changes); } public ManualRendezvousClock Clock { get; } diff --git a/tests/FinalFactory.Rendezvous.Tests/TestClient/TestClientCommandTests.cs b/tests/FinalFactory.Rendezvous.Tests/TestClient/TestClientCommandTests.cs index 44be390..ea805ee 100644 --- a/tests/FinalFactory.Rendezvous.Tests/TestClient/TestClientCommandTests.cs +++ b/tests/FinalFactory.Rendezvous.Tests/TestClient/TestClientCommandTests.cs @@ -59,6 +59,29 @@ public sealed class TestClientCommandTests Assert.Equal(130, (int)TestClientExitCode.Cancelled); } + [Fact] + public void WatchModeSupportsBoundedRuntimeAndDeliberateRecoveryExercisesOnly() + { + TestClientParseResult reset = TestClientOptionParser.Parse( + ["watch", "--run-seconds", "30", "--exercise-reset", "--script", "--json"]); + Assert.True(reset.Succeeded, reset.Error); + TestClientOptions resetOptions = Assert.IsType(reset.Options); + Assert.Equal(TestClientMode.Watch, resetOptions.Mode); + Assert.Equal(TimeSpan.FromSeconds(30), resetOptions.RunDuration); + Assert.True(resetOptions.ExerciseReset); + + TestClientParseResult reconnect = TestClientOptionParser.Parse( + ["watch", "--exercise-reconnect", "--script"]); + Assert.True(reconnect.Succeeded, reconnect.Error); + Assert.True(Assert.IsType(reconnect.Options).ExerciseReconnect); + + Assert.False(TestClientOptionParser.Parse(["browse", "--exercise-reconnect"]).Succeeded); + Assert.False(TestClientOptionParser.Parse(["browse", "--exercise-reset"]).Succeeded); + Assert.False(TestClientOptionParser.Parse( + ["watch", "--exercise-reconnect", "--exercise-reset"]).Succeeded); + Assert.False(TestClientOptionParser.Parse(["join", "--run-seconds", "30"]).Succeeded); + } + [Fact] public void HostFailureBudgetStopsAuthorityLossAndBoundsTransientRetries() { diff --git a/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/browse-sessions.json b/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/browse-sessions.json index 2b6055b..ced832a 100644 --- a/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/browse-sessions.json +++ b/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/browse-sessions.json @@ -1 +1 @@ -{"contractVersion":1,"items":[{"contractVersion":1,"listingId":"00112233-4455-6677-8899-aabbccddeeff","gameId":"space-game","environmentId":"production","regionId":"eu-central","protocolVersion":7,"buildVersion":"1.4.2","displayName":"Europa Relay","visibility":"public","publisherTrustMode":"managedDedicated","capacity":{"currentPlayers":2,"maximumPlayers":8},"metadata":{"mode":"co-op","map":"europa"}}],"nextCursor":"cursor-002"} +{"contractVersion":1,"items":[{"contractVersion":1,"listingId":"00112233-4455-6677-8899-aabbccddeeff","gameId":"space-game","environmentId":"production","regionId":"eu-central","protocolVersion":7,"buildVersion":"1.4.2","displayName":"Europa Relay","visibility":"public","publisherTrustMode":"managedDedicated","capacity":{"currentPlayers":2,"maximumPlayers":8},"metadata":{"mode":"co-op","map":"europa"}}],"nextCursor":"cursor-002","streamCursor":"stream-cursor-002"} diff --git a/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/client-public-api.txt b/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/client-public-api.txt index f896627..cd24c2d 100644 --- a/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/client-public-api.txt +++ b/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/client-public-api.txt @@ -40,6 +40,7 @@ TYPE FinalFactory.Rendezvous.Client.IRendezvousSessionBrowserClient METHOD System.Threading.Tasks.Task>> BrowseAllAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.Int32 maximumPages, System.Threading.CancellationToken cancellationToken) METHOD System.Threading.Tasks.Task> BrowseAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.Threading.CancellationToken cancellationToken) METHOD System.Threading.Tasks.Task> GetAsync(FinalFactory.Rendezvous.Contracts.SessionListingId listingId, FinalFactory.Rendezvous.Contracts.GameId gameId, FinalFactory.Rendezvous.Contracts.EnvironmentId environmentId, System.UInt32 protocolVersion, System.Threading.CancellationToken cancellationToken) + METHOD System.Collections.Generic.IAsyncEnumerable> StreamAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.String streamCursor, System.Threading.CancellationToken cancellationToken) TYPE FinalFactory.Rendezvous.Client.LeaseMaintenanceResult PROP FinalFactory.Rendezvous.Contracts.RendezvousErrorCode Error {get;} PROP FinalFactory.Rendezvous.Client.LeaseMaintenanceStopReason Reason {get;} @@ -212,6 +213,7 @@ TYPE FinalFactory.Rendezvous.Client.RendezvousSessionBrowserClient METHOD System.Threading.Tasks.Task>> BrowseAllAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.Int32 maximumPages, System.Threading.CancellationToken cancellationToken) METHOD System.Threading.Tasks.Task> BrowseAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.Threading.CancellationToken cancellationToken) METHOD System.Threading.Tasks.Task> GetAsync(FinalFactory.Rendezvous.Contracts.SessionListingId listingId, FinalFactory.Rendezvous.Contracts.GameId gameId, FinalFactory.Rendezvous.Contracts.EnvironmentId environmentId, System.UInt32 protocolVersion, System.Threading.CancellationToken cancellationToken) + METHOD System.Collections.Generic.IAsyncEnumerable> StreamAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.String streamCursor, System.Threading.CancellationToken cancellationToken) TYPE FinalFactory.Rendezvous.Client.SessionLeaseMaintainer EVENT System.EventHandler LeaseLost METHOD System.Threading.Tasks.ValueTask DisposeAsync() diff --git a/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/contracts-public-api.txt b/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/contracts-public-api.txt index 2936bc8..8d7d548 100644 --- a/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/contracts-public-api.txt +++ b/tests/FinalFactory.Rendezvous.Tests/TestData/Contracts/v1/contracts-public-api.txt @@ -28,6 +28,7 @@ TYPE FinalFactory.Rendezvous.Contracts.BrowseSessionsResponse PROP System.Int32 ContractVersion {get;set;} PROP System.Collections.Generic.List Items {get;set;} PROP System.String NextCursor {get;set;} + PROP System.String StreamCursor {get;set;} TYPE FinalFactory.Rendezvous.Contracts.ConnectionElapsedBucket ENUM UnderOneSecond=1 ENUM OneToFiveSeconds=2 @@ -84,6 +85,7 @@ TYPE FinalFactory.Rendezvous.Contracts.ContractLimits FIELD System.Int32 OpaqueHttpCredentialMaxCharacters=1024 FIELD System.Int32 RegionIdMaxCharacters=32 FIELD System.Int32 SessionCapacityMaxPlayers=10000 + FIELD System.Int32 SessionStreamEventMaxBytes=32768 FIELD System.Int32 UdpCapabilityMaxCharacters=192 FIELD System.Int32 UdpDatagramMaxBytes=1200 TYPE FinalFactory.Rendezvous.Contracts.ContractValidation @@ -345,6 +347,18 @@ TYPE FinalFactory.Rendezvous.Contracts.SessionListingId METHOD System.Boolean TryParse(System.String value, FinalFactory.Rendezvous.Contracts.SessionListingId& id) METHOD System.Boolean op_Equality(FinalFactory.Rendezvous.Contracts.SessionListingId left, FinalFactory.Rendezvous.Contracts.SessionListingId right) METHOD System.Boolean op_Inequality(FinalFactory.Rendezvous.Contracts.SessionListingId left, FinalFactory.Rendezvous.Contracts.SessionListingId right) +TYPE FinalFactory.Rendezvous.Contracts.SessionStreamEvent + CTOR () + PROP System.Int32 ContractVersion {get;set;} + PROP System.String Cursor {get;set;} + PROP FinalFactory.Rendezvous.Contracts.SessionStreamEventKind Kind {get;set;} + PROP System.Nullable ListingId {get;set;} + PROP FinalFactory.Rendezvous.Contracts.SessionListing Session {get;set;} +TYPE FinalFactory.Rendezvous.Contracts.SessionStreamEventKind + ENUM SessionUpsert=1 + ENUM SessionRemove=2 + ENUM Reset=3 + ENUM Keepalive=4 TYPE FinalFactory.Rendezvous.Contracts.UdpDecodeError ENUM None=0 ENUM DatagramTooLarge=1 @@ -371,3 +385,6 @@ TYPE FinalFactory.Rendezvous.Contracts.UpdateSessionRequest PROP System.String DisplayName {get;set;} PROP System.String LeaseToken {get;set;} PROP System.Collections.Generic.Dictionary Metadata {get;set;} + PROP System.Nullable ProtocolVersion {get;set;} + PROP System.Nullable RegionId {get;set;} + PROP System.Nullable Visibility {get;set;}