Compare commits

...

1 Commits

Author SHA1 Message Date
KyuubiYoru 06c4ecf8f3 feat(browser): stream bounded live session updates (#26)
quality-gate / quality (push) Failing after 1m47s
quality-gate / container (push) Has been skipped
2026-07-16 23:25:48 +02:00
34 changed files with 2249 additions and 27 deletions
+3
View File
@@ -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).
+233 -1
View File
@@ -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"
},
+114
View File
@@ -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<BrowseSessionsResponse> snapshot =
await browser.BrowseAsync(filter, cancellationToken);
await foreach (RendezvousClientResult<SessionStreamEvent> 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.
+8 -1
View File
@@ -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 |
| ---: | --- |
@@ -144,6 +144,11 @@ public interface IRendezvousSessionBrowserClient
EnvironmentId environmentId,
uint protocolVersion,
CancellationToken cancellationToken = default);
IAsyncEnumerable<RendezvousClientResult<SessionStreamEvent>> StreamAsync(
BrowseSessionsRequest request,
string streamCursor,
CancellationToken cancellationToken = default);
}
public interface IRendezvousJoinClient
@@ -114,6 +114,49 @@ internal sealed class RendezvousHttpTransport
}
}
internal async Task<RendezvousClientResult<HttpResponseMessage>> OpenStreamAsync(
Func<HttpRequestMessage> 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<HttpResponseMessage>(
error.Code,
error.Message,
retryAfter);
}
catch (Exception exception) when (IsTransientTransportFailure(exception, cancellationToken))
{
return RendezvousClientResult.Failure<HttpResponseMessage>(
RendezvousErrorCode.ServiceUnavailable,
"The Rendezvous event stream could not be opened.");
}
finally
{
response?.Dispose();
}
}
internal static HttpRequestMessage JsonRequest<T>(
HttpMethod method,
string uri,
@@ -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),
@@ -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<RendezvousClientResult<SessionStreamEvent>> 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 ? $"&regionId={Escape(request.RegionId.Value.Value)}" : string.Empty);
RendezvousClientResult<HttpResponseMessage> 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<SessionStreamEvent>(
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<SessionStreamEvent>(
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<SessionStreamEvent>? 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<SessionStreamEvent>(
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<SseReadResult> 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<SessionStreamEvent>(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<SessionStreamEvent>(
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<SessionStreamEvent>? result)
{
EndOfStream = endOfStream;
Result = result;
}
public bool EndOfStream { get; }
public RendezvousClientResult<SessionStreamEvent>? 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<string?> 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);
}
}
}
@@ -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;
@@ -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<SessionListing> 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
@@ -12,6 +12,8 @@ internal sealed record BrowserServiceResult<T>(RendezvousErrorCode Error, T? Val
internal sealed class SessionBrowserService(
IEphemeralRendezvousStore store,
SessionBrowserCursorCodec cursors,
SessionStreamCursorCodec streamCursors,
SessionChangeJournal changes,
IWallClock clock)
{
public BrowserServiceResult<BrowseSessionsResponse> Browse(
@@ -45,9 +47,25 @@ internal sealed class SessionBrowserService(
request.PageSize + 1,
after,
request.ExcludeFull);
StoreResult<IReadOnlyList<StoredListing>> found = store.BrowseVisibleListings(
query,
cancellationToken);
StoreResult<IReadOnlyList<StoredListing>> 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<GetSessionResponse> Get(
@@ -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<SessionChange> 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<string, string>(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<SessionChange> _changes = [];
private TaskCompletionSource<long> _changed = NewSignal();
private long _revision;
private int _subscribers;
private readonly Dictionary<TenantScope, int> _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<long> 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<bool> WaitForChangeAsync(
long revision,
TimeSpan timeout,
CancellationToken cancellationToken)
{
Task<long> 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<long> 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);
}
}
@@ -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<SessionStreamCursorPayload>(
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; }
}
@@ -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<SessionStreamEvent> 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<SessionStreamSubscription> 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<SessionListingId, PendingDelta> 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<bool> 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);
}
@@ -47,6 +47,9 @@
<Compile Include="Browser/EphemeralCursorProtector.cs" />
<Compile Include="Browser/SessionBrowserCursorCodec.cs" />
<Compile Include="Browser/SessionBrowserService.cs" />
<Compile Include="Browser/SessionChangeJournal.cs" />
<Compile Include="Browser/SessionStreamCursorCodec.cs" />
<Compile Include="Browser/SessionStreamService.cs" />
<Compile Include="ConnectionOutcomes/ConnectionOutcomeService.cs" />
<Compile Include="Deployment/DeploymentOptions.cs" />
<Compile Include="Deployment/GracefulDrainService.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<ApiError>(StatusCodes.Status429TooManyRequests)
.Produces<ApiError>(StatusCodes.Status503ServiceUnavailable)
.WithName("BrowseSessions");
sessions.MapGet("/stream", StreamSessions)
.Produces<SessionStreamEvent>(StatusCodes.Status200OK, contentType: "text/event-stream")
.Produces<ApiError>(StatusCodes.Status400BadRequest)
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.Produces<ApiError>(StatusCodes.Status503ServiceUnavailable)
.WithName("StreamSessions");
sessions.MapGet("/{listingId}", GetSession)
.Produces<GetSessionResponse>()
.Produces<ApiError>(StatusCodes.Status400BadRequest)
@@ -349,6 +356,123 @@ internal static class ContractEndpoints
}
}
private static async Task<IResult> 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<SessionStreamSubscription> 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,
@@ -255,6 +255,7 @@ builder.Services.Configure<HostOptions>(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<IEphemeralRendezvousStore>(stateStore);
builder.Services.AddSingleton<IWallClock>(rendezvousClock);
builder.Services.AddSingleton<IMonotonicClock>(rendezvousClock);
@@ -297,7 +300,9 @@ else
builder.Services.AddSingleton(SessionLeaseTiming.From(stateOptions));
builder.Services.AddSingleton<SessionLeaseService>();
builder.Services.AddSingleton<SessionBrowserCursorCodec>();
builder.Services.AddSingleton<SessionStreamCursorCodec>();
builder.Services.AddSingleton<SessionBrowserService>();
builder.Services.AddSingleton<SessionStreamService>();
builder.Services.AddSingleton<JoinAttemptCursorCodec>();
builder.Services.AddSingleton<JoinAttemptService>();
builder.Services.AddSingleton<ConnectionOutcomeMetrics>();
@@ -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
@@ -228,7 +228,10 @@ internal sealed record UpdateListingCommand(
int CurrentPlayers,
int MaximumPlayers,
IReadOnlyDictionary<string, string> Metadata,
NetworkEndpoint? DedicatedFallback);
NetworkEndpoint? DedicatedFallback,
RegionId? RegionId = null,
uint? ProtocolVersion = null,
ListingVisibility? Visibility = null);
internal sealed record DeleteListingCommand(
SessionListingId ListingId,
@@ -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<SessionListingId, ListingEntry> _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<StoredListing> 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<bool> 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<IReadOnlyList<StoredListing>> 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++;
}
}
@@ -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<TestClientExitCode> 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<BrowseSessionsResponse> 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<SessionStreamEvent> 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<BrowseSessionsResponse> 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<SessionSelection> SelectListingAsync(
TestClientOptions options,
TestClientOutput output,
@@ -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<string> 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,
});
}
@@ -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();
}
}
@@ -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<StoredListing> 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<StoredListing> 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<SessionStreamEvent> 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<T>(BrowserServiceResult<T> result)
{
Assert.True(result.Succeeded, result.Error.ToString());
return Assert.IsType<T>(result.Value);
}
}
@@ -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<RendezvousClientResult<SessionStreamEvent>> 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<ArgumentException>(async () =>
{
await foreach (RendezvousClientResult<SessionStreamEvent> _ 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<RendezvousClientResult<SessionStreamEvent>> 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()
{
@@ -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<RendezvousClientResult<SessionStreamEvent>> 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<RendezvousClientResult<SessionStreamEvent>> events = browser
.StreamAsync(request, snapshot.StreamCursor, timeout.Token)
.GetAsyncEnumerator(timeout.Token);
Task<bool> 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<RendezvousClientResult<SessionStreamEvent>> 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<bool> 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<T>(RendezvousClientResult<T> 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<SessionLeaseService>();
builder.Services.AddSingleton<SessionBrowserCursorCodec>();
builder.Services.AddSingleton<SessionStreamCursorCodec>();
builder.Services.AddSingleton(changes);
builder.Services.AddSingleton<SessionBrowserService>();
builder.Services.AddSingleton<SessionStreamService>();
WebApplication app = builder.Build();
app.UseExceptionHandler();
@@ -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"),
@@ -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<SessionLeaseService>();
builder.Services.AddSingleton<SessionBrowserCursorCodec>();
builder.Services.AddSingleton<SessionStreamCursorCodec>();
builder.Services.AddSingleton(changes);
builder.Services.AddSingleton<SessionBrowserService>();
builder.Services.AddSingleton<SessionStreamService>();
builder.Services.AddSingleton<JoinAttemptCursorCodec>();
builder.Services.AddSingleton<JoinAttemptService>();
ConnectionOutcomeMetrics outcomeMetrics = new();
@@ -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<SessionLeaseService>();
builder.Services.AddSingleton<SessionBrowserCursorCodec>();
builder.Services.AddSingleton<SessionStreamCursorCodec>();
builder.Services.AddSingleton(changes);
builder.Services.AddSingleton<SessionBrowserService>();
builder.Services.AddSingleton<SessionStreamService>();
await using WebApplication app = builder.Build();
app.UseExceptionHandler();
app.UseMiddleware<HttpAbuseProtectionMiddleware>();
@@ -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; }
@@ -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<TestClientOptions>(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<TestClientOptions>(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()
{
@@ -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"}
@@ -40,6 +40,7 @@ TYPE FinalFactory.Rendezvous.Client.IRendezvousSessionBrowserClient
METHOD System.Threading.Tasks.Task<FinalFactory.Rendezvous.Client.RendezvousClientResult<System.Collections.Generic.IReadOnlyList<FinalFactory.Rendezvous.Contracts.SessionListing>>> BrowseAllAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.Int32 maximumPages, System.Threading.CancellationToken cancellationToken)
METHOD System.Threading.Tasks.Task<FinalFactory.Rendezvous.Client.RendezvousClientResult<FinalFactory.Rendezvous.Contracts.BrowseSessionsResponse>> BrowseAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.Threading.CancellationToken cancellationToken)
METHOD System.Threading.Tasks.Task<FinalFactory.Rendezvous.Client.RendezvousClientResult<FinalFactory.Rendezvous.Contracts.GetSessionResponse>> 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<FinalFactory.Rendezvous.Client.RendezvousClientResult<FinalFactory.Rendezvous.Contracts.SessionStreamEvent>> 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<FinalFactory.Rendezvous.Client.RendezvousClientResult<System.Collections.Generic.IReadOnlyList<FinalFactory.Rendezvous.Contracts.SessionListing>>> BrowseAllAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.Int32 maximumPages, System.Threading.CancellationToken cancellationToken)
METHOD System.Threading.Tasks.Task<FinalFactory.Rendezvous.Client.RendezvousClientResult<FinalFactory.Rendezvous.Contracts.BrowseSessionsResponse>> BrowseAsync(FinalFactory.Rendezvous.Contracts.BrowseSessionsRequest request, System.Threading.CancellationToken cancellationToken)
METHOD System.Threading.Tasks.Task<FinalFactory.Rendezvous.Client.RendezvousClientResult<FinalFactory.Rendezvous.Contracts.GetSessionResponse>> 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<FinalFactory.Rendezvous.Client.RendezvousClientResult<FinalFactory.Rendezvous.Contracts.SessionStreamEvent>> 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()
@@ -28,6 +28,7 @@ TYPE FinalFactory.Rendezvous.Contracts.BrowseSessionsResponse
PROP System.Int32 ContractVersion {get;set;}
PROP System.Collections.Generic.List<FinalFactory.Rendezvous.Contracts.SessionListing> 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<FinalFactory.Rendezvous.Contracts.SessionListingId> 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<System.String,System.String> Metadata {get;set;}
PROP System.Nullable<System.UInt32> ProtocolVersion {get;set;}
PROP System.Nullable<FinalFactory.Rendezvous.Contracts.RegionId> RegionId {get;set;}
PROP System.Nullable<FinalFactory.Rendezvous.Contracts.ListingVisibility> Visibility {get;set;}