Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| be732de7c9 |
@@ -91,6 +91,9 @@ Tenant policy, publisher/operator principals, and production key custody are
|
||||
defined in [game provisioning and signing-key lifecycle](docs/security/provisioning.md).
|
||||
Layered HTTP/UDP budgets, overload behavior, and safe operational tuning are
|
||||
defined in [hostile-input and overload protection](docs/security/abuse-protection.md).
|
||||
Health semantics, bounded telemetry, alerting, audit privacy, and the authenticated
|
||||
operator controls are defined in the
|
||||
[observability and operator runbook](docs/operations/observability-and-operator-runbook.md).
|
||||
The scriptable host/browser/join diagnostic and its stable automation contract are
|
||||
documented in the [TestClient integration guide](docs/integration/test-client.md).
|
||||
The always-on three-party scenarios, optional Linux namespace topology, and
|
||||
@@ -112,7 +115,8 @@ dotnet test Rendezvous.slnx --configuration Release --no-build
|
||||
Run the bootstrap server with
|
||||
`dotnet run --project src/FinalFactory.Rendezvous.Server`. It serves HTTP health endpoints and binds
|
||||
the configured UDP mediator port; both stop through normal host cancellation.
|
||||
The launch profile uses an ephemeral development-only signing key. Production
|
||||
The launch profile uses separate ephemeral development-only publisher and operator
|
||||
signing keys. Production
|
||||
startup fails closed until externally supplied game policies and `env:` signing
|
||||
key references resolve to valid key material; no reusable game secret is stored
|
||||
in this repository or the public Client package.
|
||||
|
||||
+1619
-4
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,142 @@
|
||||
# Observability and operator runbook
|
||||
|
||||
This runbook defines the production signals and privileged controls for the
|
||||
Rendezvous service. The service emits `System.Diagnostics.Metrics` instruments
|
||||
from the `FinalFactory.Rendezvous` meter and distributed-tracing activities from
|
||||
`FinalFactory.Rendezvous.Server`. Connect those sources to the deployment's
|
||||
OpenTelemetry or equivalent collector. Do not add identifiers to metric labels.
|
||||
|
||||
## Health and readiness
|
||||
|
||||
- `GET /health/live` proves that the HTTP process can answer. It deliberately
|
||||
remains independent of provisioning, the state store, drain state, and optional
|
||||
listeners so an orchestrator does not restart a recoverable dependency failure.
|
||||
- `GET /health/ready` returns success only after the HTTP path is answering, the
|
||||
required IPv4 UDP socket is bound, any configured IPv6 UDP socket is bound,
|
||||
provisioning loaded successfully, the store is available, and drain has not
|
||||
started. A failed check returns `503` and removes the instance from new work.
|
||||
- A graceful drain immediately makes readiness fail while liveness remains healthy.
|
||||
Existing work may complete until the bounded store drain deadline.
|
||||
|
||||
## Metrics and traces
|
||||
|
||||
| Instrument | Purpose | Bounded dimensions |
|
||||
| --- | --- | --- |
|
||||
| `rendezvous.http.requests` / `rendezvous.http.duration` | HTTP volume and latency | operation, status code |
|
||||
| `rendezvous.udp.results` / `rendezvous.udp.duration` | UDP mediation volume and processing latency | frozen/litenet operation, result |
|
||||
| `rendezvous.limiter.drops` | Requests shed by admission controls | transport, fixed partition class |
|
||||
| `rendezvous.operator.authentication` | Accepted, forbidden, and rejected operator authentication | result |
|
||||
| `rendezvous.audit.events` | Privileged action outcomes | fixed action, result |
|
||||
| `rendezvous.connection.outcomes` | Client-reported direct-connect outcomes | normalized outcome, elapsed bucket |
|
||||
| `rendezvous.pairing.latency` | Time from attempt creation to successful peer introduction | none |
|
||||
| `rendezvous.queue.depth` | Active join-attempt queue depth | none |
|
||||
| `rendezvous.store.active_listings` / `active_leases` / `active_attempts` / `replay_markers` | Current ephemeral load | none |
|
||||
| `rendezvous.store.expiry_churn` | Cumulative natural expiry activity | none |
|
||||
| `rendezvous.store.available` | Store health (`1` available, `0` unavailable) | none |
|
||||
|
||||
HTTP responses include `X-Rendezvous-Correlation-ID`. It is a generated trace ID
|
||||
or random value, never a caller-supplied session or player identifier. UDP and
|
||||
HTTP activities contain operation-level data only. Logs and traces must not add
|
||||
tokens, capabilities, session/listing IDs, player subjects, metadata, raw IP
|
||||
addresses, or endpoint values.
|
||||
|
||||
Recommended dashboard panels are request rate and p50/p95/p99 latency by fixed
|
||||
operation, UDP result ratio, direct connection success ratio, pairing latency,
|
||||
active listings/attempts, expiry churn, limiter drops, store availability,
|
||||
operator authentication results, audit action results, and signing-key windows.
|
||||
|
||||
## Alerts
|
||||
|
||||
Tune thresholds from the normal production baseline, then keep these conditions
|
||||
as distinct actionable alerts:
|
||||
|
||||
- **Signing key expiry:** page when any required signing key has less than seven
|
||||
days before `signUntil`; escalate at 24 hours. Confirm a replacement is signing
|
||||
and the previous key remains verify-only for the maximum credential lifetime.
|
||||
- **Authentication spike:** warn when rejected or forbidden operator authentication
|
||||
exceeds five attempts in five minutes. Treat unexpected publisher-authentication
|
||||
growth as a possible credential or integration incident.
|
||||
- **Direct success regression:** warn when the connected outcome ratio falls more
|
||||
than 20% below its seven-day same-region baseline for 15 minutes, with a minimum
|
||||
sample floor. Break down only by bounded outcome and time bucket.
|
||||
- **Saturation:** warn when queue depth remains above 70% of the configured attempt
|
||||
limit, limiter drops are sustained, or p95 latency exceeds the service objective;
|
||||
page at 90% or when lease-critical traffic is shed.
|
||||
- **Store degradation:** page immediately when `rendezvous.store.available` is zero
|
||||
or readiness fails for the store. Rising expiry churn without corresponding new
|
||||
work is a warning for stalled clients or clock/configuration mistakes.
|
||||
- **Listener/config readiness:** page when no ready instances remain. Investigate
|
||||
UDP bind failures, a configured-but-unbound IPv6 listener, provisioning errors,
|
||||
and unintended drain state separately.
|
||||
|
||||
## Operator authentication and controls
|
||||
|
||||
Operator credentials use a signing key configured with `CredentialKinds:
|
||||
["Operator"]`. Operator keys cannot be scoped to a game/environment or used for
|
||||
publisher credentials. Mint short-lived operator credentials through the trusted
|
||||
provisioning process, outside the public Rendezvous HTTP service, and grant only
|
||||
the required permission. Never place credentials in command history, URLs, logs,
|
||||
or support tickets.
|
||||
|
||||
The application also enforces a default-deny source boundary. Configure at most
|
||||
32 exact operator source IPs in
|
||||
`Rendezvous:AbuseProtection:OperatorAllowedAddresses`; an empty list disables all
|
||||
operator HTTP access. Development permits loopback only. Production must place
|
||||
`/v1/operator/*` behind a private management listener or reverse-proxy ACL, list
|
||||
only the resulting trusted management source addresses, and block that path on
|
||||
the public edge. If forwarded headers are enabled, keep the existing exact-proxy,
|
||||
single-hop trust policy and allowlist the post-forwarding operator source. Verify
|
||||
from both an allowed management host and a denied public host before deployment.
|
||||
Denied sources are charged to the bounded general HTTP partition before credential
|
||||
or request-body processing, then receive `404`; sustained denied traffic receives
|
||||
the same typed `429` overload response as other public traffic.
|
||||
|
||||
Operator traffic has a dedicated, bounded rate/concurrency partition and critical
|
||||
tracker-key reserve. Public browse/join saturation therefore cannot consume the
|
||||
operator control budget, while compromised management sources remain rate-limited.
|
||||
|
||||
The OpenAPI document defines the separate `OperatorBearer` scheme. All endpoints
|
||||
are under `/v1/operator`:
|
||||
|
||||
| Endpoint | Permission | Confirmation |
|
||||
| --- | --- | --- |
|
||||
| `GET /status` | `ReadPolicy` | none; returns aggregates, tenant status, safe key status, and audit counts |
|
||||
| `POST /listings/revoke` | `RevokePublisher` | repeat the exact listing ID in `confirmListingId` |
|
||||
| `POST /principals/revoke` | `RevokePublisher` | repeat the exact subject and choose a 1-600 second revocation lifetime |
|
||||
| `POST /keys/revoke` | `RotateKeys` | repeat the exact key ID; runtime revocation is immediate |
|
||||
| `POST /drain` | `ManagePolicy` | send the exact value `DRAIN` |
|
||||
|
||||
Publisher credentials are rejected on this surface even if their subject resembles
|
||||
an operator. Destructive responses do not echo identifiers. The status response
|
||||
does not expose player identities, raw endpoints, session metadata, capabilities,
|
||||
or tokens. Every authenticated operator action, rejected confirmation, and
|
||||
permission denial is audited with actor and target fingerprints.
|
||||
|
||||
Key revocation is process-local in the current single-instance store. Apply the
|
||||
same revocation to every instance, then replace configuration before restarting;
|
||||
a restart reconstructs the configured key ring. Principal revocation is bounded
|
||||
to ten minutes and removes that principal's active listings and attempts. Use
|
||||
listing revocation for one targeted session and drain before planned shutdown.
|
||||
|
||||
## Audit retention and incident handling
|
||||
|
||||
The in-process audit trail defaults to 10,000 entries and 30 days. It evicts the
|
||||
oldest record at capacity and purges expired records on the next write. Configure
|
||||
`Rendezvous:Audit:MaxEntries` and `RetentionDays` within their validated bounds.
|
||||
Export the structured `AuditTrail` log events through the deployment's protected
|
||||
logging pipeline when durable retention is required; the in-memory trail is not a
|
||||
durable compliance archive. Those events include only timestamps, fixed action
|
||||
fields, correlation IDs, and actor/target fingerprints.
|
||||
|
||||
Audit records retain timestamp, fixed action/result, target kind, correlation ID,
|
||||
and 96-bit SHA-256 fingerprints of actor and target. Routine logs contain only the
|
||||
fixed action/result/target kind and correlation ID. Restrict audit access to the
|
||||
operator role, retain aggregates only as long as operationally necessary, and
|
||||
delete raw exported audit data according to the 30-day policy unless an incident
|
||||
hold is approved.
|
||||
|
||||
During an incident: confirm readiness and store health; capture aggregate graphs
|
||||
and correlation IDs; revoke the narrowest listing, principal, or key; drain only
|
||||
when isolation is required; record the action in the incident timeline; and verify
|
||||
that direct success, limiter drops, and authentication rates return to baseline.
|
||||
Do not copy player data, endpoints, or credentials into the incident record.
|
||||
@@ -20,6 +20,10 @@ traffic is also isolated by its authenticated scope.
|
||||
Health probes use their own source-prefix budget so public API overload cannot
|
||||
make a healthy instance fail its orchestrator probes, while health traffic is
|
||||
still bounded.
|
||||
Operator endpoints likewise use a separate bounded rate/concurrency partition
|
||||
backed by the critical tracker reserve. They first require an exact source IP
|
||||
from the default-deny `OperatorAllowedAddresses` policy, so public traffic
|
||||
cannot spend the incident-response budget.
|
||||
3. Once an endpoint has safely derived identities, it also acquires applicable
|
||||
tenant, principal or capability, and listing/attempt budgets. Secret
|
||||
capabilities are represented only by bounded SHA-256 fingerprints.
|
||||
@@ -41,7 +45,8 @@ load shedding.
|
||||
|
||||
`Rendezvous:AbuseProtection:MaxTrackedKeys` is a hard combined ceiling for rate
|
||||
and active-concurrency keys. General HTTP and UDP traffic cannot consume the
|
||||
configured `CriticalTrackedKeyReserve`; lease operations and health probes may
|
||||
configured `CriticalTrackedKeyReserve`; lease operations, health probes, and
|
||||
allowlisted operator controls may
|
||||
use that reserve but never exceed the hard ceiling. A request that would exceed
|
||||
its applicable ceiling fails closed without adding state. Fixed-window rate keys
|
||||
are cleared at the next window boundary; concurrency keys are removed as their
|
||||
|
||||
@@ -20,6 +20,8 @@ internal sealed class AbuseProtectionOptions
|
||||
|
||||
public string[] TrustedProxyAddresses { get; set; } = [];
|
||||
|
||||
public string[] OperatorAllowedAddresses { get; set; } = [];
|
||||
|
||||
[Range(1, 100_000)]
|
||||
public int HealthGlobalRequestsPerWindow { get; set; } = 1_000;
|
||||
|
||||
@@ -32,6 +34,18 @@ internal sealed class AbuseProtectionOptions
|
||||
[Range(1, 1_000)]
|
||||
public int HealthIpPrefixConcurrency { get; set; } = 8;
|
||||
|
||||
[Range(1, 100_000)]
|
||||
public int OperatorGlobalRequestsPerWindow { get; set; } = 1_000;
|
||||
|
||||
[Range(1, 10_000)]
|
||||
public int OperatorGlobalConcurrency { get; set; } = 32;
|
||||
|
||||
[Range(1, 100_000)]
|
||||
public int OperatorIpPrefixRequestsPerWindow { get; set; } = 120;
|
||||
|
||||
[Range(1, 1_000)]
|
||||
public int OperatorIpPrefixConcurrency { get; set; } = 8;
|
||||
|
||||
[Range(1, 1_000_000)]
|
||||
public int HttpGlobalRequestsPerWindow { get; set; } = 20_000;
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@ using System.Buffers;
|
||||
using System.Net;
|
||||
using System.Security.Cryptography;
|
||||
using System.Text;
|
||||
using FinalFactory.Rendezvous.Server.Observability;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Abuse;
|
||||
@@ -12,13 +13,23 @@ internal sealed class AbuseProtectionService
|
||||
private readonly TimeProvider _timeProvider;
|
||||
private readonly TrackerState _httpTracker;
|
||||
private readonly TrackerState _udpTracker;
|
||||
private readonly RendezvousTelemetry? _telemetry;
|
||||
private readonly HashSet<string> _operatorAllowedAddresses;
|
||||
|
||||
public AbuseProtectionService(
|
||||
IOptions<AbuseProtectionOptions> options,
|
||||
TimeProvider? timeProvider = null)
|
||||
TimeProvider? timeProvider = null,
|
||||
RendezvousTelemetry? telemetry = null)
|
||||
{
|
||||
_options = options.Value;
|
||||
_timeProvider = timeProvider ?? TimeProvider.System;
|
||||
_telemetry = telemetry;
|
||||
_operatorAllowedAddresses = options.Value.OperatorAllowedAddresses
|
||||
.Select(static value => IPAddress.TryParse(value, out IPAddress? address)
|
||||
? NormalizeAddress(address).ToString()
|
||||
: string.Empty)
|
||||
.Where(static value => value.Length > 0)
|
||||
.ToHashSet(StringComparer.Ordinal);
|
||||
DateTimeOffset now = _timeProvider.GetUtcNow();
|
||||
_httpTracker = new(now);
|
||||
_udpTracker = new(now);
|
||||
@@ -87,6 +98,35 @@ internal sealed class AbuseProtectionService
|
||||
out retryAfterSeconds);
|
||||
}
|
||||
|
||||
public bool IsOperatorSourceAllowed(IPAddress? remoteAddress) =>
|
||||
remoteAddress is not null
|
||||
&& _operatorAllowedAddresses.Contains(NormalizeAddress(remoteAddress).ToString());
|
||||
|
||||
public bool TryAcquireOperatorIngress(
|
||||
IPAddress? remoteAddress,
|
||||
out AbuseLease? lease,
|
||||
out int retryAfterSeconds)
|
||||
{
|
||||
string prefix = GetNetworkPrefix(remoteAddress);
|
||||
RateDimension[] rates =
|
||||
[
|
||||
new("operator:rate:global", _options.OperatorGlobalRequestsPerWindow),
|
||||
new($"operator:rate:ip:{prefix}", _options.OperatorIpPrefixRequestsPerWindow),
|
||||
];
|
||||
RateDimension[] concurrency =
|
||||
[
|
||||
new("operator:concurrency:global", _options.OperatorGlobalConcurrency),
|
||||
new($"operator:concurrency:ip:{prefix}", _options.OperatorIpPrefixConcurrency),
|
||||
];
|
||||
return TryAcquire(
|
||||
rates,
|
||||
concurrency,
|
||||
TrackerDomain.Http,
|
||||
true,
|
||||
out lease,
|
||||
out retryAfterSeconds);
|
||||
}
|
||||
|
||||
public bool TryAcquireHttpIdentity(
|
||||
string operation,
|
||||
string? tenant,
|
||||
@@ -266,7 +306,37 @@ internal sealed class AbuseProtectionService
|
||||
out int retryAfterSeconds)
|
||||
{
|
||||
TrackerState tracker = domain == TrackerDomain.Udp ? _udpTracker : _httpTracker;
|
||||
bool accepted;
|
||||
lock (tracker.Gate)
|
||||
{
|
||||
accepted = TryAcquireLocked(
|
||||
tracker,
|
||||
rates,
|
||||
concurrency,
|
||||
domain,
|
||||
canUseCriticalReserve,
|
||||
out lease,
|
||||
out retryAfterSeconds);
|
||||
}
|
||||
|
||||
if (!accepted)
|
||||
{
|
||||
_telemetry?.RecordLimiterDrop(
|
||||
domain == TrackerDomain.Udp ? "udp" : "http",
|
||||
"rate-or-concurrency");
|
||||
}
|
||||
|
||||
return accepted;
|
||||
}
|
||||
|
||||
private bool TryAcquireLocked(
|
||||
TrackerState tracker,
|
||||
ReadOnlySpan<RateDimension> rates,
|
||||
ReadOnlySpan<RateDimension> concurrency,
|
||||
TrackerDomain domain,
|
||||
bool canUseCriticalReserve,
|
||||
out AbuseLease? lease,
|
||||
out int retryAfterSeconds)
|
||||
{
|
||||
DateTimeOffset now = _timeProvider.GetUtcNow();
|
||||
TimeSpan window = TimeSpan.FromSeconds(_options.WindowSeconds);
|
||||
@@ -327,7 +397,6 @@ internal sealed class AbuseProtectionService
|
||||
lease = new AbuseLease(this, tracker, acquiredConcurrency);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
private static bool CanAcquireAll(
|
||||
TrackerState tracker,
|
||||
@@ -415,6 +484,9 @@ internal sealed class AbuseProtectionService
|
||||
return "unknown";
|
||||
}
|
||||
|
||||
private static IPAddress NormalizeAddress(IPAddress address) =>
|
||||
address.IsIPv4MappedToIPv6 ? address.MapToIPv4() : address;
|
||||
|
||||
private readonly record struct RateDimension(string Key, int Limit);
|
||||
|
||||
private enum TrackerDomain
|
||||
|
||||
@@ -19,16 +19,71 @@ internal sealed class HttpAbuseProtectionMiddleware(
|
||||
string operation = context.GetEndpoint()?.Metadata.GetMetadata<IEndpointNameMetadata>()
|
||||
?.EndpointName ?? "Unmatched";
|
||||
bool healthEndpoint = operation is "GetLiveness" or "GetReadiness";
|
||||
bool acquired = healthEndpoint
|
||||
? protection.TryAcquireHealthIngress(
|
||||
bool operatorEndpoint = operation is
|
||||
"GetOperatorStatus"
|
||||
or "RevokeOperatorListing"
|
||||
or "RevokeOperatorPrincipal"
|
||||
or "RevokeOperatorSigningKey"
|
||||
or "BeginOperatorDrain";
|
||||
if (operatorEndpoint
|
||||
&& !protection.IsOperatorSourceAllowed(context.Connection.RemoteIpAddress))
|
||||
{
|
||||
bool deniedSourceAdmitted = protection.TryAcquireHttpIngress(
|
||||
context.Connection.RemoteIpAddress,
|
||||
out AbuseProtectionService.AbuseLease? lease,
|
||||
out int retryAfterSeconds)
|
||||
: protection.TryAcquireHttpIngress(
|
||||
"Unmatched",
|
||||
out AbuseProtectionService.AbuseLease? deniedSourceLease,
|
||||
out int deniedRetryAfterSeconds);
|
||||
using (deniedSourceLease)
|
||||
{
|
||||
if (!deniedSourceAdmitted)
|
||||
{
|
||||
context.Response.Headers.RetryAfter = deniedRetryAfterSeconds.ToString(
|
||||
System.Globalization.CultureInfo.InvariantCulture);
|
||||
await WriteErrorAsync(
|
||||
context,
|
||||
StatusCodes.Status429TooManyRequests,
|
||||
RendezvousErrorCode.RateLimited,
|
||||
"The request rate limit was exceeded.",
|
||||
deniedRetryAfterSeconds).ConfigureAwait(false);
|
||||
return;
|
||||
}
|
||||
|
||||
await WriteErrorAsync(
|
||||
context,
|
||||
StatusCodes.Status404NotFound,
|
||||
RendezvousErrorCode.NotFound,
|
||||
"The requested resource was not found.").ConfigureAwait(false);
|
||||
}
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
AbuseProtectionService.AbuseLease? lease;
|
||||
int retryAfterSeconds;
|
||||
bool acquired;
|
||||
if (healthEndpoint)
|
||||
{
|
||||
acquired = protection.TryAcquireHealthIngress(
|
||||
context.Connection.RemoteIpAddress,
|
||||
out lease,
|
||||
out retryAfterSeconds);
|
||||
}
|
||||
else if (operatorEndpoint)
|
||||
{
|
||||
acquired = protection.TryAcquireOperatorIngress(
|
||||
context.Connection.RemoteIpAddress,
|
||||
out lease,
|
||||
out retryAfterSeconds);
|
||||
}
|
||||
else
|
||||
{
|
||||
acquired = protection.TryAcquireHttpIngress(
|
||||
context.Connection.RemoteIpAddress,
|
||||
operation,
|
||||
out lease,
|
||||
out retryAfterSeconds);
|
||||
}
|
||||
|
||||
if (!acquired)
|
||||
{
|
||||
context.Response.Headers.RetryAfter = retryAfterSeconds.ToString(
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using FinalFactory.Rendezvous.Server.Observability;
|
||||
using FinalFactory.Rendezvous.Server.Sessions;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
|
||||
@@ -15,6 +16,10 @@ internal sealed class ConnectionOutcomeMetrics
|
||||
{
|
||||
private readonly object _gate = new();
|
||||
private readonly Dictionary<(ConnectionOutcomeKind, ConnectionElapsedBucket), long> _counts = [];
|
||||
private readonly RendezvousTelemetry? _telemetry;
|
||||
|
||||
public ConnectionOutcomeMetrics(RendezvousTelemetry? telemetry = null) =>
|
||||
_telemetry = telemetry;
|
||||
|
||||
internal void Record(ConnectionOutcomeKind outcome, ConnectionElapsedBucket elapsedBucket)
|
||||
{
|
||||
@@ -24,6 +29,8 @@ internal sealed class ConnectionOutcomeMetrics
|
||||
_counts.TryGetValue(key, out long count);
|
||||
_counts[key] = count + 1;
|
||||
}
|
||||
|
||||
_telemetry?.RecordConnectionOutcome(outcome.ToString(), elapsedBucket.ToString());
|
||||
}
|
||||
|
||||
internal long GetCount(ConnectionOutcomeKind outcome, ConnectionElapsedBucket elapsedBucket)
|
||||
|
||||
@@ -4,7 +4,8 @@ using Microsoft.AspNetCore.Diagnostics;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Http;
|
||||
|
||||
internal sealed class RendezvousExceptionHandler : IExceptionHandler
|
||||
internal sealed partial class RendezvousExceptionHandler(
|
||||
ILogger<RendezvousExceptionHandler> logger) : IExceptionHandler
|
||||
{
|
||||
public async ValueTask<bool> TryHandleAsync(
|
||||
HttpContext httpContext,
|
||||
@@ -26,6 +27,13 @@ internal sealed class RendezvousExceptionHandler : IExceptionHandler
|
||||
: invalidRequest
|
||||
? StatusCodes.Status400BadRequest
|
||||
: StatusCodes.Status500InternalServerError;
|
||||
LogRequestFailure(
|
||||
logger,
|
||||
payloadTooLarge ? "payload-too-large" : invalidRequest ? "invalid-request" : "internal-error",
|
||||
httpContext.Response.StatusCode,
|
||||
httpContext.Response.Headers["X-Rendezvous-Correlation-ID"].ToString() is { Length: > 0 } value
|
||||
? value
|
||||
: "unavailable");
|
||||
await httpContext.Response.WriteAsJsonAsync(
|
||||
new ApiError
|
||||
{
|
||||
@@ -42,4 +50,14 @@ internal sealed class RendezvousExceptionHandler : IExceptionHandler
|
||||
cancellationToken).ConfigureAwait(false);
|
||||
return true;
|
||||
}
|
||||
|
||||
[LoggerMessage(
|
||||
EventId = 200,
|
||||
Level = LogLevel.Warning,
|
||||
Message = "Request failed with {FailureKind} and HTTP status {StatusCode}; correlation {CorrelationId}")]
|
||||
private static partial void LogRequestFailure(
|
||||
ILogger logger,
|
||||
string failureKind,
|
||||
int statusCode,
|
||||
string correlationId);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
using System.ComponentModel.DataAnnotations;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Observability;
|
||||
|
||||
internal sealed class AuditOptions
|
||||
{
|
||||
public const string SectionName = "Rendezvous:Audit";
|
||||
|
||||
[Range(100, 100_000)]
|
||||
public int MaxEntries { get; set; } = 10_000;
|
||||
|
||||
[Range(1, 30)]
|
||||
public int RetentionDays { get; set; } = 30;
|
||||
}
|
||||
@@ -0,0 +1,134 @@
|
||||
using System.Security.Cryptography;
|
||||
using System.Text;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Observability;
|
||||
|
||||
internal sealed partial class AuditTrail
|
||||
{
|
||||
private readonly object _gate = new();
|
||||
private readonly LinkedList<AuditEntry> _entries = [];
|
||||
private readonly AuditOptions _options;
|
||||
private readonly TimeProvider _timeProvider;
|
||||
private readonly ILogger<AuditTrail> _logger;
|
||||
private readonly RendezvousTelemetry _telemetry;
|
||||
|
||||
public AuditTrail(
|
||||
IOptions<AuditOptions> options,
|
||||
ILogger<AuditTrail> logger,
|
||||
RendezvousTelemetry telemetry,
|
||||
TimeProvider? timeProvider = null)
|
||||
{
|
||||
_options = options.Value;
|
||||
_logger = logger;
|
||||
_telemetry = telemetry;
|
||||
_timeProvider = timeProvider ?? TimeProvider.System;
|
||||
}
|
||||
|
||||
public void Record(
|
||||
string actorSubject,
|
||||
string action,
|
||||
string result,
|
||||
string targetKind,
|
||||
string targetIdentifier,
|
||||
string correlationId)
|
||||
{
|
||||
DateTimeOffset now = _timeProvider.GetUtcNow();
|
||||
AuditEntry entry = new(
|
||||
now,
|
||||
Fingerprint(actorSubject),
|
||||
action,
|
||||
result,
|
||||
targetKind,
|
||||
Fingerprint(targetIdentifier),
|
||||
correlationId);
|
||||
lock (_gate)
|
||||
{
|
||||
PurgeExpired(now);
|
||||
|
||||
while (_entries.Count >= _options.MaxEntries)
|
||||
{
|
||||
_entries.RemoveFirst();
|
||||
}
|
||||
|
||||
_entries.AddLast(entry);
|
||||
}
|
||||
|
||||
_telemetry.RecordAudit(action, result);
|
||||
LogOperatorAction(
|
||||
_logger,
|
||||
entry.Timestamp,
|
||||
entry.ActorFingerprint,
|
||||
action,
|
||||
result,
|
||||
targetKind,
|
||||
entry.TargetFingerprint,
|
||||
correlationId);
|
||||
}
|
||||
|
||||
public IReadOnlyDictionary<string, long> GetAggregateCounts()
|
||||
{
|
||||
lock (_gate)
|
||||
{
|
||||
PurgeExpired(_timeProvider.GetUtcNow());
|
||||
return _entries
|
||||
.GroupBy(static entry => $"{entry.Action}:{entry.Result}", StringComparer.Ordinal)
|
||||
.ToDictionary(
|
||||
static group => group.Key,
|
||||
static group => (long)group.Count(),
|
||||
StringComparer.Ordinal);
|
||||
}
|
||||
}
|
||||
|
||||
internal IReadOnlyList<AuditEntry> GetEntriesForTests()
|
||||
{
|
||||
lock (_gate)
|
||||
{
|
||||
PurgeExpired(_timeProvider.GetUtcNow());
|
||||
return _entries.ToArray();
|
||||
}
|
||||
}
|
||||
|
||||
private void PurgeExpired(DateTimeOffset now)
|
||||
{
|
||||
DateTimeOffset oldest = now.AddDays(-_options.RetentionDays);
|
||||
while (_entries.First is { Value.Timestamp: var timestamp }
|
||||
&& timestamp < oldest)
|
||||
{
|
||||
_entries.RemoveFirst();
|
||||
}
|
||||
}
|
||||
|
||||
private static string Fingerprint(string value)
|
||||
{
|
||||
byte[] digest = SHA256.HashData(Encoding.UTF8.GetBytes(value));
|
||||
return Convert.ToHexString(digest.AsSpan(0, 12));
|
||||
}
|
||||
|
||||
[LoggerMessage(
|
||||
EventId = 100,
|
||||
Level = LogLevel.Information,
|
||||
Message = "Operator audit at {Timestamp}: actor {ActorFingerprint} action {Action} completed with {Result} for {TargetKind} target {TargetFingerprint}; correlation {CorrelationId}")]
|
||||
private static partial void LogOperatorAction(
|
||||
ILogger logger,
|
||||
DateTimeOffset timestamp,
|
||||
string actorFingerprint,
|
||||
string action,
|
||||
string result,
|
||||
string targetKind,
|
||||
string targetFingerprint,
|
||||
string correlationId);
|
||||
}
|
||||
|
||||
internal sealed record AuditEntry(
|
||||
DateTimeOffset Timestamp,
|
||||
string ActorFingerprint,
|
||||
string Action,
|
||||
string Result,
|
||||
string TargetKind,
|
||||
string TargetFingerprint,
|
||||
string CorrelationId)
|
||||
{
|
||||
public override string ToString() =>
|
||||
$"[AuditEntry {Action}/{Result}; actor and target fingerprinted]";
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Observability;
|
||||
|
||||
internal static class HealthEndpoints
|
||||
{
|
||||
public static IEndpointRouteBuilder MapRendezvousHealthEndpoints(
|
||||
this IEndpointRouteBuilder endpoints)
|
||||
{
|
||||
endpoints.MapGet(
|
||||
"/health/live",
|
||||
static () => Results.Ok(new HealthResponse { Status = "live" }))
|
||||
.Produces<HealthResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.WithName("GetLiveness")
|
||||
.WithTags("Health");
|
||||
endpoints.MapGet(
|
||||
"/health/ready",
|
||||
static (RendezvousReadiness readiness) =>
|
||||
!readiness.GetSnapshot().IsReady
|
||||
? Results.StatusCode(StatusCodes.Status503ServiceUnavailable)
|
||||
: Results.Ok(new HealthResponse { Status = "ready" }))
|
||||
.Produces<HealthResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.Produces(StatusCodes.Status503ServiceUnavailable)
|
||||
.WithName("GetReadiness")
|
||||
.WithTags("Health");
|
||||
return endpoints;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
using FinalFactory.Rendezvous.Server.Provisioning;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
using FinalFactory.Rendezvous.Server.Transport;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Observability;
|
||||
|
||||
internal sealed class RendezvousReadiness(
|
||||
UdpMediatorService mediator,
|
||||
ProvisioningReadiness provisioning,
|
||||
IEphemeralRendezvousStore state,
|
||||
IOptions<UdpMediatorOptions> udpOptions)
|
||||
{
|
||||
public ReadinessSnapshot GetSnapshot()
|
||||
{
|
||||
bool ipv6Required = !string.IsNullOrWhiteSpace(udpOptions.Value.Ipv6ListenAddress);
|
||||
return new ReadinessSnapshot(
|
||||
HttpListenerReady: true,
|
||||
UdpIpv4ListenerReady: mediator.LocalEndpoint is not null,
|
||||
UdpIpv6ListenerReady: !ipv6Required || mediator.LocalIpv6Endpoint is not null,
|
||||
ProvisioningReady: provisioning.IsReady,
|
||||
StoreAvailable: state.IsAvailable,
|
||||
Draining: state.IsDraining);
|
||||
}
|
||||
}
|
||||
|
||||
internal sealed record ReadinessSnapshot(
|
||||
bool HttpListenerReady,
|
||||
bool UdpIpv4ListenerReady,
|
||||
bool UdpIpv6ListenerReady,
|
||||
bool ProvisioningReady,
|
||||
bool StoreAvailable,
|
||||
bool Draining)
|
||||
{
|
||||
public bool IsReady => HttpListenerReady
|
||||
&& UdpIpv4ListenerReady
|
||||
&& UdpIpv6ListenerReady
|
||||
&& ProvisioningReady
|
||||
&& StoreAvailable
|
||||
&& !Draining;
|
||||
}
|
||||
@@ -0,0 +1,126 @@
|
||||
using System.Diagnostics;
|
||||
using System.Diagnostics.Metrics;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Observability;
|
||||
|
||||
internal sealed class RendezvousTelemetry : IDisposable
|
||||
{
|
||||
public const string MeterName = "FinalFactory.Rendezvous";
|
||||
public const string ActivitySourceName = "FinalFactory.Rendezvous.Server";
|
||||
|
||||
private readonly InMemoryEphemeralRendezvousStore _store;
|
||||
private readonly Meter _meter = new(MeterName, "1.0.0");
|
||||
private readonly ActivitySource _activities = new(ActivitySourceName, "1.0.0");
|
||||
private readonly Counter<long> _httpRequests;
|
||||
private readonly Histogram<double> _httpDuration;
|
||||
private readonly Counter<long> _udpResults;
|
||||
private readonly Histogram<double> _udpDuration;
|
||||
private readonly Counter<long> _limiterDrops;
|
||||
private readonly Counter<long> _auditEvents;
|
||||
private readonly Counter<long> _connectionOutcomes;
|
||||
private readonly Counter<long> _operatorAuthentication;
|
||||
private readonly Histogram<double> _pairingLatency;
|
||||
|
||||
public RendezvousTelemetry(InMemoryEphemeralRendezvousStore store)
|
||||
{
|
||||
_store = store;
|
||||
_httpRequests = _meter.CreateCounter<long>("rendezvous.http.requests");
|
||||
_httpDuration = _meter.CreateHistogram<double>(
|
||||
"rendezvous.http.duration",
|
||||
"ms");
|
||||
_udpResults = _meter.CreateCounter<long>("rendezvous.udp.results");
|
||||
_udpDuration = _meter.CreateHistogram<double>(
|
||||
"rendezvous.udp.duration",
|
||||
"ms");
|
||||
_limiterDrops = _meter.CreateCounter<long>("rendezvous.limiter.drops");
|
||||
_auditEvents = _meter.CreateCounter<long>("rendezvous.audit.events");
|
||||
_connectionOutcomes = _meter.CreateCounter<long>("rendezvous.connection.outcomes");
|
||||
_operatorAuthentication = _meter.CreateCounter<long>("rendezvous.operator.authentication");
|
||||
_pairingLatency = _meter.CreateHistogram<double>(
|
||||
"rendezvous.pairing.latency",
|
||||
"ms");
|
||||
_meter.CreateObservableGauge(
|
||||
"rendezvous.store.active_listings",
|
||||
() => _store.GetMetricsSnapshot().ActiveListings);
|
||||
_meter.CreateObservableGauge(
|
||||
"rendezvous.store.active_leases",
|
||||
() => _store.GetMetricsSnapshot().ActiveListings);
|
||||
_meter.CreateObservableGauge(
|
||||
"rendezvous.store.active_attempts",
|
||||
() => _store.GetMetricsSnapshot().ActiveJoinAttempts);
|
||||
_meter.CreateObservableGauge(
|
||||
"rendezvous.queue.depth",
|
||||
() => _store.GetMetricsSnapshot().ActiveJoinAttempts);
|
||||
_meter.CreateObservableGauge(
|
||||
"rendezvous.store.replay_markers",
|
||||
() => _store.GetMetricsSnapshot().ReplayMarkers);
|
||||
_meter.CreateObservableGauge(
|
||||
"rendezvous.store.available",
|
||||
() => _store.GetMetricsSnapshot().IsAvailable ? 1 : 0);
|
||||
_meter.CreateObservableCounter(
|
||||
"rendezvous.store.expiry_churn",
|
||||
() => _store.GetMetricsSnapshot().ExpiryChurn);
|
||||
}
|
||||
|
||||
public Activity? StartActivity(string name, ActivityKind kind = ActivityKind.Internal) =>
|
||||
_activities.StartActivity(name, kind);
|
||||
|
||||
public void RecordHttp(string operation, int statusCode, double elapsedMilliseconds)
|
||||
{
|
||||
TagList tags = new()
|
||||
{
|
||||
{ "operation", operation },
|
||||
{ "status_code", statusCode },
|
||||
};
|
||||
_httpRequests.Add(1, tags);
|
||||
_httpDuration.Record(elapsedMilliseconds, tags);
|
||||
}
|
||||
|
||||
public void RecordUdp(string operation, string result, double elapsedMilliseconds)
|
||||
{
|
||||
TagList tags = new()
|
||||
{
|
||||
{ "operation", operation },
|
||||
{ "result", result },
|
||||
};
|
||||
_udpResults.Add(1, tags);
|
||||
_udpDuration.Record(elapsedMilliseconds, tags);
|
||||
}
|
||||
|
||||
public void RecordLimiterDrop(string transport, string partition) =>
|
||||
_limiterDrops.Add(1, new TagList
|
||||
{
|
||||
{ "transport", transport },
|
||||
{ "partition", partition },
|
||||
});
|
||||
|
||||
public void RecordAudit(string action, string result) =>
|
||||
_auditEvents.Add(1, new TagList
|
||||
{
|
||||
{ "action", action },
|
||||
{ "result", result },
|
||||
});
|
||||
|
||||
public void RecordConnectionOutcome(string outcome, string elapsedBucket) =>
|
||||
_connectionOutcomes.Add(1, new TagList
|
||||
{
|
||||
{ "outcome", outcome },
|
||||
{ "elapsed_bucket", elapsedBucket },
|
||||
});
|
||||
|
||||
public void RecordOperatorAuthentication(string result) =>
|
||||
_operatorAuthentication.Add(1, new TagList
|
||||
{
|
||||
{ "result", result },
|
||||
});
|
||||
|
||||
public void RecordPairingLatency(double elapsedMilliseconds) =>
|
||||
_pairingLatency.Record(elapsedMilliseconds);
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
_activities.Dispose();
|
||||
_meter.Dispose();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
using System.Diagnostics;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Observability;
|
||||
|
||||
internal sealed class TelemetryMiddleware(
|
||||
RequestDelegate next,
|
||||
RendezvousTelemetry telemetry)
|
||||
{
|
||||
public async Task InvokeAsync(HttpContext context)
|
||||
{
|
||||
string operation = context.GetEndpoint()?.Metadata.GetMetadata<IEndpointNameMetadata>()
|
||||
?.EndpointName ?? "Unmatched";
|
||||
long started = Stopwatch.GetTimestamp();
|
||||
using Activity? activity = telemetry.StartActivity(
|
||||
$"HTTP {operation}",
|
||||
ActivityKind.Server);
|
||||
string correlationId = activity?.TraceId.ToString() ?? Guid.NewGuid().ToString("N");
|
||||
context.Response.Headers["X-Rendezvous-Correlation-ID"] = correlationId;
|
||||
activity?.SetTag("rendezvous.operation", operation);
|
||||
try
|
||||
{
|
||||
await next(context).ConfigureAwait(false);
|
||||
}
|
||||
finally
|
||||
{
|
||||
telemetry.RecordHttp(
|
||||
operation,
|
||||
context.Response.StatusCode,
|
||||
Stopwatch.GetElapsedTime(started).TotalMilliseconds);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,428 @@
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using FinalFactory.Rendezvous.Server.Observability;
|
||||
using FinalFactory.Rendezvous.Server.Provisioning;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Operations;
|
||||
|
||||
internal static class OperatorEndpoints
|
||||
{
|
||||
private const string CorrelationHeader = "X-Rendezvous-Correlation-ID";
|
||||
|
||||
public static IEndpointRouteBuilder MapOperatorEndpoints(this IEndpointRouteBuilder endpoints)
|
||||
{
|
||||
RouteGroupBuilder group = endpoints.MapGroup("/v1/operator").WithTags("Operator");
|
||||
group.MapGet("/status", GetStatus)
|
||||
.Produces<OperatorStatusResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
|
||||
.Produces<ApiError>(StatusCodes.Status403Forbidden)
|
||||
.Produces<ApiError>(StatusCodes.Status404NotFound)
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.WithName("GetOperatorStatus");
|
||||
group.MapPost("/listings/revoke", RevokeListing)
|
||||
.Accepts<RevokeListingRequest>("application/json")
|
||||
.Produces<OperatorActionResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status400BadRequest)
|
||||
.Produces<ApiError>(StatusCodes.Status413PayloadTooLarge)
|
||||
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
|
||||
.Produces<ApiError>(StatusCodes.Status403Forbidden)
|
||||
.Produces<ApiError>(StatusCodes.Status404NotFound)
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.Produces<ApiError>(StatusCodes.Status503ServiceUnavailable)
|
||||
.WithName("RevokeOperatorListing");
|
||||
group.MapPost("/principals/revoke", RevokePrincipal)
|
||||
.Accepts<RevokePrincipalRequest>("application/json")
|
||||
.Produces<OperatorActionResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status400BadRequest)
|
||||
.Produces<ApiError>(StatusCodes.Status413PayloadTooLarge)
|
||||
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
|
||||
.Produces<ApiError>(StatusCodes.Status403Forbidden)
|
||||
.Produces<ApiError>(StatusCodes.Status404NotFound)
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.Produces<ApiError>(StatusCodes.Status503ServiceUnavailable)
|
||||
.WithName("RevokeOperatorPrincipal");
|
||||
group.MapPost("/keys/revoke", RevokeSigningKey)
|
||||
.Accepts<RevokeSigningKeyRequest>("application/json")
|
||||
.Produces<OperatorActionResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status400BadRequest)
|
||||
.Produces<ApiError>(StatusCodes.Status413PayloadTooLarge)
|
||||
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
|
||||
.Produces<ApiError>(StatusCodes.Status403Forbidden)
|
||||
.Produces<ApiError>(StatusCodes.Status404NotFound)
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.WithName("RevokeOperatorSigningKey");
|
||||
group.MapPost("/drain", BeginDrain)
|
||||
.Accepts<BeginDrainRequest>("application/json")
|
||||
.Produces<OperatorActionResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status400BadRequest)
|
||||
.Produces<ApiError>(StatusCodes.Status413PayloadTooLarge)
|
||||
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
|
||||
.Produces<ApiError>(StatusCodes.Status403Forbidden)
|
||||
.Produces<ApiError>(StatusCodes.Status404NotFound)
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.WithName("BeginOperatorDrain");
|
||||
return endpoints;
|
||||
}
|
||||
|
||||
private static IResult GetStatus(
|
||||
[FromHeader(Name = "Authorization")] string? authorization,
|
||||
[FromServices] PrincipalCredentialService credentials,
|
||||
[FromServices] IWallClock clock,
|
||||
[FromServices] OperatorService service,
|
||||
[FromServices] AuditTrail audit,
|
||||
[FromServices] RendezvousTelemetry telemetry,
|
||||
HttpContext context)
|
||||
{
|
||||
if (!TryAuthorize(
|
||||
authorization,
|
||||
OperatorPermission.ReadPolicy,
|
||||
"inspect-status",
|
||||
credentials,
|
||||
clock,
|
||||
audit,
|
||||
telemetry,
|
||||
context,
|
||||
out OperatorPrincipal? principal,
|
||||
out IResult? failure))
|
||||
{
|
||||
return failure!;
|
||||
}
|
||||
|
||||
OperatorStatusResponse response = service.GetStatus();
|
||||
audit.Record(
|
||||
principal!.Subject,
|
||||
"inspect-status",
|
||||
"succeeded",
|
||||
"service",
|
||||
"rendezvous",
|
||||
Correlation(context));
|
||||
return Results.Ok(response);
|
||||
}
|
||||
|
||||
private static IResult RevokeListing(
|
||||
[FromBody] RevokeListingRequest request,
|
||||
[FromHeader(Name = "Authorization")] string? authorization,
|
||||
[FromServices] PrincipalCredentialService credentials,
|
||||
[FromServices] IWallClock clock,
|
||||
[FromServices] OperatorService service,
|
||||
[FromServices] AuditTrail audit,
|
||||
[FromServices] RendezvousTelemetry telemetry,
|
||||
HttpContext context,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
if (!TryAuthorize(
|
||||
authorization,
|
||||
OperatorPermission.RevokePublisher,
|
||||
"revoke-listing",
|
||||
credentials,
|
||||
clock,
|
||||
audit,
|
||||
telemetry,
|
||||
context,
|
||||
out OperatorPrincipal? principal,
|
||||
out IResult? failure))
|
||||
{
|
||||
return failure!;
|
||||
}
|
||||
|
||||
bool valid = SessionListingId.TryParse(request.ListingId, out SessionListingId listingId);
|
||||
if (!valid || !string.Equals(request.ListingId, request.ConfirmListingId, StringComparison.Ordinal))
|
||||
{
|
||||
AuditRejected(audit, principal!, "revoke-listing", "listing", request.ListingId, context);
|
||||
return BadRequest("A valid listing ID and an exact repeated confirmation are required.");
|
||||
}
|
||||
|
||||
StoreResult<bool> result = service.RevokeListing(listingId, cancellationToken);
|
||||
return StoreActionResult(
|
||||
result.Code,
|
||||
audit,
|
||||
principal!,
|
||||
"revoke-listing",
|
||||
"listing",
|
||||
request.ListingId,
|
||||
context,
|
||||
affectedResources: result.Succeeded ? 1 : null);
|
||||
}
|
||||
|
||||
private static IResult RevokePrincipal(
|
||||
[FromBody] RevokePrincipalRequest request,
|
||||
[FromHeader(Name = "Authorization")] string? authorization,
|
||||
[FromServices] PrincipalCredentialService credentials,
|
||||
[FromServices] IWallClock clock,
|
||||
[FromServices] OperatorService service,
|
||||
[FromServices] AuditTrail audit,
|
||||
[FromServices] RendezvousTelemetry telemetry,
|
||||
HttpContext context,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
if (!TryAuthorize(
|
||||
authorization,
|
||||
OperatorPermission.RevokePublisher,
|
||||
"revoke-principal",
|
||||
credentials,
|
||||
clock,
|
||||
audit,
|
||||
telemetry,
|
||||
context,
|
||||
out OperatorPrincipal? principal,
|
||||
out IResult? failure))
|
||||
{
|
||||
return failure!;
|
||||
}
|
||||
|
||||
bool safeSubject = request.Subject is { Length: > 0 and <= 128 }
|
||||
&& request.Subject.All(static character => character is >= '!' and <= '~');
|
||||
if (!safeSubject
|
||||
|| !string.Equals(request.Subject, request.ConfirmSubject, StringComparison.Ordinal)
|
||||
|| request.LifetimeSeconds is < 1 or > 600)
|
||||
{
|
||||
AuditRejected(audit, principal!, "revoke-principal", "principal", request.Subject, context);
|
||||
return BadRequest("A valid subject, exact repeated confirmation, and 1-600 second lifetime are required.");
|
||||
}
|
||||
|
||||
StoreResult<int> result = service.RevokePrincipal(
|
||||
request.Subject,
|
||||
TimeSpan.FromSeconds(request.LifetimeSeconds),
|
||||
cancellationToken);
|
||||
return StoreActionResult(
|
||||
result.Code,
|
||||
audit,
|
||||
principal!,
|
||||
"revoke-principal",
|
||||
"principal",
|
||||
request.Subject,
|
||||
context,
|
||||
result.Value);
|
||||
}
|
||||
|
||||
private static IResult RevokeSigningKey(
|
||||
[FromBody] RevokeSigningKeyRequest request,
|
||||
[FromHeader(Name = "Authorization")] string? authorization,
|
||||
[FromServices] PrincipalCredentialService credentials,
|
||||
[FromServices] IWallClock clock,
|
||||
[FromServices] OperatorService service,
|
||||
[FromServices] AuditTrail audit,
|
||||
[FromServices] RendezvousTelemetry telemetry,
|
||||
HttpContext context)
|
||||
{
|
||||
if (!TryAuthorize(
|
||||
authorization,
|
||||
OperatorPermission.RotateKeys,
|
||||
"revoke-signing-key",
|
||||
credentials,
|
||||
clock,
|
||||
audit,
|
||||
telemetry,
|
||||
context,
|
||||
out OperatorPrincipal? principal,
|
||||
out IResult? failure))
|
||||
{
|
||||
return failure!;
|
||||
}
|
||||
|
||||
bool safeKeyId = request.KeyId is { Length: > 0 and <= 64 }
|
||||
&& request.KeyId.All(static character => character is
|
||||
>= 'A' and <= 'Z'
|
||||
or >= 'a' and <= 'z'
|
||||
or >= '0' and <= '9'
|
||||
or '-'
|
||||
or '_');
|
||||
if (!safeKeyId || !string.Equals(request.KeyId, request.ConfirmKeyId, StringComparison.Ordinal))
|
||||
{
|
||||
AuditRejected(audit, principal!, "revoke-signing-key", "signing-key", request.KeyId, context);
|
||||
return BadRequest("A valid key ID and an exact repeated confirmation are required.");
|
||||
}
|
||||
|
||||
bool revoked = service.RevokeSigningKey(request.KeyId);
|
||||
string result = revoked ? "succeeded" : "not-found";
|
||||
audit.Record(
|
||||
principal!.Subject,
|
||||
"revoke-signing-key",
|
||||
result,
|
||||
"signing-key",
|
||||
request.KeyId,
|
||||
Correlation(context));
|
||||
return revoked
|
||||
? Results.Ok(new OperatorActionResponse { Status = "completed" })
|
||||
: Error(RendezvousErrorCode.NotFound, "The requested resource was not found.");
|
||||
}
|
||||
|
||||
private static IResult BeginDrain(
|
||||
[FromBody] BeginDrainRequest request,
|
||||
[FromHeader(Name = "Authorization")] string? authorization,
|
||||
[FromServices] PrincipalCredentialService credentials,
|
||||
[FromServices] IWallClock clock,
|
||||
[FromServices] OperatorService service,
|
||||
[FromServices] AuditTrail audit,
|
||||
[FromServices] RendezvousTelemetry telemetry,
|
||||
HttpContext context,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
if (!TryAuthorize(
|
||||
authorization,
|
||||
OperatorPermission.ManagePolicy,
|
||||
"begin-drain",
|
||||
credentials,
|
||||
clock,
|
||||
audit,
|
||||
telemetry,
|
||||
context,
|
||||
out OperatorPrincipal? principal,
|
||||
out IResult? failure))
|
||||
{
|
||||
return failure!;
|
||||
}
|
||||
|
||||
if (!string.Equals(request.Confirmation, "DRAIN", StringComparison.Ordinal))
|
||||
{
|
||||
AuditRejected(audit, principal!, "begin-drain", "service", "rendezvous", context);
|
||||
return BadRequest("The confirmation value must be exactly 'DRAIN'.");
|
||||
}
|
||||
|
||||
service.BeginDrain(cancellationToken);
|
||||
audit.Record(
|
||||
principal!.Subject,
|
||||
"begin-drain",
|
||||
"succeeded",
|
||||
"service",
|
||||
"rendezvous",
|
||||
Correlation(context));
|
||||
return Results.Ok(new OperatorActionResponse { Status = "draining" });
|
||||
}
|
||||
|
||||
private static bool TryAuthorize(
|
||||
string? authorization,
|
||||
OperatorPermission requiredPermission,
|
||||
string operation,
|
||||
PrincipalCredentialService credentials,
|
||||
IWallClock clock,
|
||||
AuditTrail audit,
|
||||
RendezvousTelemetry telemetry,
|
||||
HttpContext context,
|
||||
out OperatorPrincipal? principal,
|
||||
out IResult? failure)
|
||||
{
|
||||
principal = null;
|
||||
failure = null;
|
||||
const string prefix = "Bearer ";
|
||||
if (authorization is null
|
||||
|| !authorization.StartsWith(prefix, StringComparison.OrdinalIgnoreCase))
|
||||
{
|
||||
telemetry.RecordOperatorAuthentication("rejected");
|
||||
failure = AuthenticationRequired(context);
|
||||
return false;
|
||||
}
|
||||
|
||||
CredentialValidationResult validation = credentials.Validate(
|
||||
authorization[prefix.Length..],
|
||||
clock.UtcNow);
|
||||
if (!validation.IsValid || validation.Principal is not OperatorPrincipal candidate)
|
||||
{
|
||||
telemetry.RecordOperatorAuthentication("rejected");
|
||||
failure = AuthenticationRequired(context);
|
||||
return false;
|
||||
}
|
||||
|
||||
if (!candidate.Permissions.Contains(requiredPermission))
|
||||
{
|
||||
telemetry.RecordOperatorAuthentication("forbidden");
|
||||
audit.Record(
|
||||
candidate.Subject,
|
||||
operation,
|
||||
"forbidden",
|
||||
"operator-operation",
|
||||
operation,
|
||||
Correlation(context));
|
||||
failure = Error(RendezvousErrorCode.Forbidden, "The operator is not authorized for this operation.");
|
||||
return false;
|
||||
}
|
||||
|
||||
telemetry.RecordOperatorAuthentication("accepted");
|
||||
principal = candidate;
|
||||
return true;
|
||||
}
|
||||
|
||||
private static IResult StoreActionResult(
|
||||
StoreResultCode code,
|
||||
AuditTrail audit,
|
||||
OperatorPrincipal principal,
|
||||
string action,
|
||||
string targetKind,
|
||||
string targetIdentifier,
|
||||
HttpContext context,
|
||||
int? affectedResources)
|
||||
{
|
||||
string auditResult = code == StoreResultCode.Success
|
||||
? "succeeded"
|
||||
: code.ToString().ToLowerInvariant();
|
||||
audit.Record(
|
||||
principal.Subject,
|
||||
action,
|
||||
auditResult,
|
||||
targetKind,
|
||||
targetIdentifier,
|
||||
Correlation(context));
|
||||
return code switch
|
||||
{
|
||||
StoreResultCode.Success => Results.Ok(new OperatorActionResponse
|
||||
{
|
||||
Status = "completed",
|
||||
AffectedResources = affectedResources,
|
||||
}),
|
||||
StoreResultCode.NotFound => Error(
|
||||
RendezvousErrorCode.NotFound,
|
||||
"The requested resource was not found."),
|
||||
StoreResultCode.CapacityExceeded => Error(
|
||||
RendezvousErrorCode.CapacityExceeded,
|
||||
"The operation could not be retained within the configured capacity."),
|
||||
StoreResultCode.ServiceUnavailable or StoreResultCode.Draining => Error(
|
||||
RendezvousErrorCode.ServiceUnavailable,
|
||||
"The service is not available for this operation."),
|
||||
_ => Error(RendezvousErrorCode.Conflict, "The operation could not be completed."),
|
||||
};
|
||||
}
|
||||
|
||||
private static void AuditRejected(
|
||||
AuditTrail audit,
|
||||
OperatorPrincipal principal,
|
||||
string action,
|
||||
string targetKind,
|
||||
string? targetIdentifier,
|
||||
HttpContext context) => audit.Record(
|
||||
principal.Subject,
|
||||
action,
|
||||
"rejected",
|
||||
targetKind,
|
||||
targetIdentifier ?? string.Empty,
|
||||
Correlation(context));
|
||||
|
||||
private static string Correlation(HttpContext context) =>
|
||||
context.Response.Headers[CorrelationHeader].ToString() is { Length: > 0 } value
|
||||
? value
|
||||
: "unavailable";
|
||||
|
||||
private static IResult AuthenticationRequired(HttpContext context)
|
||||
{
|
||||
context.Response.Headers.WWWAuthenticate = "Bearer realm=\"operator\"";
|
||||
return Error(
|
||||
RendezvousErrorCode.AuthenticationRequired,
|
||||
"A valid operator bearer credential is required.");
|
||||
}
|
||||
|
||||
private static IResult BadRequest(string message) => Error(RendezvousErrorCode.InvalidRequest, message);
|
||||
|
||||
private static IResult Error(RendezvousErrorCode code, string message) => Results.Json(
|
||||
new ApiError { Code = code, Message = message },
|
||||
ContractJson.Options,
|
||||
statusCode: code switch
|
||||
{
|
||||
RendezvousErrorCode.AuthenticationRequired => StatusCodes.Status401Unauthorized,
|
||||
RendezvousErrorCode.Forbidden => StatusCodes.Status403Forbidden,
|
||||
RendezvousErrorCode.NotFound => StatusCodes.Status404NotFound,
|
||||
RendezvousErrorCode.CapacityExceeded => StatusCodes.Status429TooManyRequests,
|
||||
RendezvousErrorCode.ServiceUnavailable => StatusCodes.Status503ServiceUnavailable,
|
||||
RendezvousErrorCode.Conflict => StatusCodes.Status409Conflict,
|
||||
_ => StatusCodes.Status400BadRequest,
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,82 @@
|
||||
namespace FinalFactory.Rendezvous.Server.Operations;
|
||||
|
||||
internal sealed record OperatorStatusResponse
|
||||
{
|
||||
public required string Status { get; init; }
|
||||
public required OperatorReadinessResponse Readiness { get; init; }
|
||||
public required OperatorStoreResponse Store { get; init; }
|
||||
public required IReadOnlyList<OperatorTenantResponse> Tenants { get; init; }
|
||||
public required IReadOnlyList<OperatorSigningKeyResponse> SigningKeys { get; init; }
|
||||
public required IReadOnlyDictionary<string, long> AuditCounts { get; init; }
|
||||
}
|
||||
|
||||
internal sealed record OperatorReadinessResponse
|
||||
{
|
||||
public required bool HttpListener { get; init; }
|
||||
public required bool UdpIpv4Listener { get; init; }
|
||||
public required bool UdpIpv6Listener { get; init; }
|
||||
public required bool Provisioning { get; init; }
|
||||
public required bool Store { get; init; }
|
||||
public required bool Draining { get; init; }
|
||||
}
|
||||
|
||||
internal sealed record OperatorStoreResponse
|
||||
{
|
||||
public required int ActiveListings { get; init; }
|
||||
public required int FreshPresenceBindings { get; init; }
|
||||
public required int ActiveJoinAttempts { get; init; }
|
||||
public required int RetainedOutcomeReports { get; init; }
|
||||
public required int ReplayMarkers { get; init; }
|
||||
public required int PrincipalRevocations { get; init; }
|
||||
public required int IdempotencyEntries { get; init; }
|
||||
public required long MaintenanceSweeps { get; init; }
|
||||
public required long ExpiryChurn { get; init; }
|
||||
}
|
||||
|
||||
internal sealed record OperatorTenantResponse
|
||||
{
|
||||
public required string GameId { get; init; }
|
||||
public required string EnvironmentId { get; init; }
|
||||
public required string Status { get; init; }
|
||||
}
|
||||
|
||||
internal sealed record OperatorSigningKeyResponse
|
||||
{
|
||||
public required string KeyId { get; init; }
|
||||
public required string Status { get; init; }
|
||||
public required DateTimeOffset SignUntil { get; init; }
|
||||
public required DateTimeOffset VerifyUntil { get; init; }
|
||||
public string? GameId { get; init; }
|
||||
public string? EnvironmentId { get; init; }
|
||||
public required IReadOnlyList<string> CredentialKinds { get; init; }
|
||||
}
|
||||
|
||||
internal sealed record OperatorActionResponse
|
||||
{
|
||||
public required string Status { get; init; }
|
||||
public int? AffectedResources { get; init; }
|
||||
}
|
||||
|
||||
internal sealed record RevokeListingRequest
|
||||
{
|
||||
public required string ListingId { get; init; }
|
||||
public required string ConfirmListingId { get; init; }
|
||||
}
|
||||
|
||||
internal sealed record RevokePrincipalRequest
|
||||
{
|
||||
public required string Subject { get; init; }
|
||||
public required string ConfirmSubject { get; init; }
|
||||
public required int LifetimeSeconds { get; init; }
|
||||
}
|
||||
|
||||
internal sealed record RevokeSigningKeyRequest
|
||||
{
|
||||
public required string KeyId { get; init; }
|
||||
public required string ConfirmKeyId { get; init; }
|
||||
}
|
||||
|
||||
internal sealed record BeginDrainRequest
|
||||
{
|
||||
public required string Confirmation { get; init; }
|
||||
}
|
||||
@@ -0,0 +1,81 @@
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using FinalFactory.Rendezvous.Server.Observability;
|
||||
using FinalFactory.Rendezvous.Server.Provisioning;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Server.Operations;
|
||||
|
||||
internal sealed class OperatorService(
|
||||
InMemoryEphemeralRendezvousStore store,
|
||||
ProvisioningRuntime provisioning,
|
||||
RendezvousReadiness readiness,
|
||||
AuditTrail audit,
|
||||
IWallClock clock)
|
||||
{
|
||||
public OperatorStatusResponse GetStatus()
|
||||
{
|
||||
ReadinessSnapshot readinessSnapshot = readiness.GetSnapshot();
|
||||
EphemeralStoreSnapshot storeSnapshot = store.GetSnapshot();
|
||||
return new OperatorStatusResponse
|
||||
{
|
||||
Status = readinessSnapshot.IsReady ? "ready" : "not-ready",
|
||||
Readiness = new OperatorReadinessResponse
|
||||
{
|
||||
HttpListener = readinessSnapshot.HttpListenerReady,
|
||||
UdpIpv4Listener = readinessSnapshot.UdpIpv4ListenerReady,
|
||||
UdpIpv6Listener = readinessSnapshot.UdpIpv6ListenerReady,
|
||||
Provisioning = readinessSnapshot.ProvisioningReady,
|
||||
Store = readinessSnapshot.StoreAvailable,
|
||||
Draining = readinessSnapshot.Draining,
|
||||
},
|
||||
Store = new OperatorStoreResponse
|
||||
{
|
||||
ActiveListings = storeSnapshot.ActiveListings,
|
||||
FreshPresenceBindings = storeSnapshot.FreshPresenceBindings,
|
||||
ActiveJoinAttempts = storeSnapshot.ActiveJoinAttempts,
|
||||
RetainedOutcomeReports = storeSnapshot.RetainedOutcomeReports,
|
||||
ReplayMarkers = storeSnapshot.ReplayMarkers,
|
||||
PrincipalRevocations = storeSnapshot.PrincipalRevocations,
|
||||
IdempotencyEntries = storeSnapshot.IdempotencyEntries,
|
||||
MaintenanceSweeps = storeSnapshot.MaintenanceSweeps,
|
||||
ExpiryChurn = storeSnapshot.ExpiryChurn,
|
||||
},
|
||||
Tenants = provisioning.Policies.EnabledPolicies
|
||||
.OrderBy(static policy => policy.GameId.Value, StringComparer.Ordinal)
|
||||
.ThenBy(static policy => policy.EnvironmentId.Value, StringComparer.Ordinal)
|
||||
.Select(static policy => new OperatorTenantResponse
|
||||
{
|
||||
GameId = policy.GameId.Value,
|
||||
EnvironmentId = policy.EnvironmentId.Value,
|
||||
Status = "enabled",
|
||||
})
|
||||
.ToArray(),
|
||||
SigningKeys = provisioning.SigningKeys.GetStatuses(clock.UtcNow)
|
||||
.Select(static key => new OperatorSigningKeyResponse
|
||||
{
|
||||
KeyId = key.KeyId,
|
||||
Status = key.Status,
|
||||
SignUntil = key.SignUntil,
|
||||
VerifyUntil = key.VerifyUntil,
|
||||
GameId = key.GameId,
|
||||
EnvironmentId = key.EnvironmentId,
|
||||
CredentialKinds = key.CredentialKinds,
|
||||
})
|
||||
.ToArray(),
|
||||
AuditCounts = audit.GetAggregateCounts(),
|
||||
};
|
||||
}
|
||||
|
||||
public StoreResult<bool> RevokeListing(
|
||||
SessionListingId listingId,
|
||||
CancellationToken cancellationToken) => store.RevokeListing(listingId, cancellationToken);
|
||||
|
||||
public StoreResult<int> RevokePrincipal(
|
||||
string subject,
|
||||
TimeSpan lifetime,
|
||||
CancellationToken cancellationToken) => store.RevokePrincipal(subject, lifetime, cancellationToken);
|
||||
|
||||
public bool RevokeSigningKey(string keyId) => provisioning.SigningKeys.Revoke(keyId);
|
||||
|
||||
public void BeginDrain(CancellationToken cancellationToken) => store.BeginDrain(cancellationToken);
|
||||
}
|
||||
@@ -5,6 +5,8 @@ using FinalFactory.Rendezvous.Server.Browser;
|
||||
using FinalFactory.Rendezvous.Server.ConnectionOutcomes;
|
||||
using FinalFactory.Rendezvous.Server.Http;
|
||||
using FinalFactory.Rendezvous.Server.JoinAttempts;
|
||||
using FinalFactory.Rendezvous.Server.Observability;
|
||||
using FinalFactory.Rendezvous.Server.Operations;
|
||||
using FinalFactory.Rendezvous.Server.Provisioning;
|
||||
using FinalFactory.Rendezvous.Server.Sessions;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
@@ -13,6 +15,9 @@ using Microsoft.AspNetCore.HttpOverrides;
|
||||
using Microsoft.OpenApi;
|
||||
|
||||
WebApplicationBuilder builder = WebApplication.CreateBuilder(args);
|
||||
builder.Logging.AddFilter(
|
||||
"Microsoft.AspNetCore.Diagnostics.ExceptionHandlerMiddleware",
|
||||
LogLevel.None);
|
||||
bool isOpenApiGeneration = string.Equals(
|
||||
System.Reflection.Assembly.GetEntryAssembly()?.GetName().Name,
|
||||
"GetDocument.Insider",
|
||||
@@ -61,6 +66,14 @@ builder.Services.AddOpenApi("v1", static options =>
|
||||
In = ParameterLocation.Header,
|
||||
Description = "Attempt-scoped client capability returned only to the joining caller.",
|
||||
};
|
||||
const string operatorSchemeName = "OperatorBearer";
|
||||
document.Components.SecuritySchemes[operatorSchemeName] = new OpenApiSecurityScheme
|
||||
{
|
||||
Type = SecuritySchemeType.Http,
|
||||
Scheme = "bearer",
|
||||
BearerFormat = "rv1 operator credential",
|
||||
Description = "Operator-only credential with an explicit permission set.",
|
||||
};
|
||||
|
||||
HashSet<string> securedOperations = new(StringComparer.Ordinal)
|
||||
{
|
||||
@@ -69,8 +82,17 @@ builder.Services.AddOpenApi("v1", static options =>
|
||||
"UpdateSession",
|
||||
"DeleteSession",
|
||||
};
|
||||
HashSet<string> operatorOperations = new(StringComparer.Ordinal)
|
||||
{
|
||||
"GetOperatorStatus",
|
||||
"RevokeOperatorListing",
|
||||
"RevokeOperatorPrincipal",
|
||||
"RevokeOperatorSigningKey",
|
||||
"BeginOperatorDrain",
|
||||
};
|
||||
OpenApiSecuritySchemeReference reference = new(schemeName, document, null);
|
||||
OpenApiSecuritySchemeReference attemptReference = new(attemptSchemeName, document, null);
|
||||
OpenApiSecuritySchemeReference operatorReference = new(operatorSchemeName, document, null);
|
||||
foreach (OpenApiPathItem path in document.Paths.Values)
|
||||
{
|
||||
if (path.Operations is null)
|
||||
@@ -99,20 +121,47 @@ builder.Services.AddOpenApi("v1", static options =>
|
||||
});
|
||||
}
|
||||
|
||||
foreach (OpenApiOperation operation in path.Operations.Values.Where(
|
||||
operation => operatorOperations.Contains(
|
||||
operation.OperationId ?? string.Empty)))
|
||||
{
|
||||
operation.Security ??= [];
|
||||
operation.Security.Add(new OpenApiSecurityRequirement
|
||||
{
|
||||
[operatorReference] = [],
|
||||
});
|
||||
}
|
||||
|
||||
foreach (OpenApiOperation operation in path.Operations.Values)
|
||||
{
|
||||
if (operation.Responses is null
|
||||
|| !operation.Responses.TryGetValue(
|
||||
StatusCodes.Status429TooManyRequests.ToString(
|
||||
System.Globalization.CultureInfo.InvariantCulture),
|
||||
out IOpenApiResponse? response)
|
||||
|| response is not OpenApiResponse concreteResponse)
|
||||
if (operation.Responses is null)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
foreach ((string status, IOpenApiResponse response) in operation.Responses)
|
||||
{
|
||||
if (response is not OpenApiResponse concreteResponse)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
concreteResponse.Headers ??=
|
||||
new Dictionary<string, IOpenApiHeader>(StringComparer.OrdinalIgnoreCase);
|
||||
concreteResponse.Headers["X-Rendezvous-Correlation-ID"] = new OpenApiHeader
|
||||
{
|
||||
Description = "Safe request correlation identifier generated by the service.",
|
||||
Schema = new OpenApiSchema
|
||||
{
|
||||
Type = JsonSchemaType.String,
|
||||
},
|
||||
};
|
||||
if (string.Equals(
|
||||
status,
|
||||
StatusCodes.Status429TooManyRequests.ToString(
|
||||
System.Globalization.CultureInfo.InvariantCulture),
|
||||
StringComparison.Ordinal))
|
||||
{
|
||||
concreteResponse.Headers["Retry-After"] = new OpenApiHeader
|
||||
{
|
||||
Description = "Whole seconds before the caller should retry (1-60).",
|
||||
@@ -124,6 +173,8 @@ builder.Services.AddOpenApi("v1", static options =>
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return Task.CompletedTask;
|
||||
});
|
||||
@@ -166,8 +217,18 @@ builder.Services
|
||||
&& addresses.All(
|
||||
static value => IPAddress.TryParse(value, out _)),
|
||||
"Trusted proxy addresses must contain at most 32 literal IP addresses.")
|
||||
.Validate(
|
||||
options => options.OperatorAllowedAddresses is { Length: <= 32 } addresses
|
||||
&& addresses.All(
|
||||
static value => IPAddress.TryParse(value, out _)),
|
||||
"Operator allowed addresses must contain at most 32 literal IP addresses.")
|
||||
.ValidateOnStart();
|
||||
builder.Services.AddSingleton<AbuseProtectionService>();
|
||||
builder.Services
|
||||
.AddOptions<AuditOptions>()
|
||||
.BindConfiguration(AuditOptions.SectionName)
|
||||
.ValidateDataAnnotations()
|
||||
.ValidateOnStart();
|
||||
AbuseProtectionOptions configuredAbuseProtection = builder.Configuration
|
||||
.GetSection(AbuseProtectionOptions.SectionName)
|
||||
.Get<AbuseProtectionOptions>() ?? new AbuseProtectionOptions();
|
||||
@@ -180,8 +241,13 @@ InMemoryEphemeralRendezvousStore stateStore = new(
|
||||
stateOptions,
|
||||
rendezvousClock,
|
||||
rendezvousClock);
|
||||
builder.Services.AddSingleton(stateStore);
|
||||
builder.Services.AddSingleton<IEphemeralRendezvousStore>(stateStore);
|
||||
builder.Services.AddSingleton<IWallClock>(rendezvousClock);
|
||||
builder.Services.AddSingleton<IMonotonicClock>(rendezvousClock);
|
||||
builder.Services.AddSingleton<RendezvousTelemetry>();
|
||||
builder.Services.AddSingleton<AuditTrail>();
|
||||
builder.Services.AddSingleton<RendezvousReadiness>();
|
||||
|
||||
if (isOpenApiGeneration)
|
||||
{
|
||||
@@ -214,6 +280,7 @@ else
|
||||
builder.Services.AddSingleton<JoinAttemptService>();
|
||||
builder.Services.AddSingleton<ConnectionOutcomeMetrics>();
|
||||
builder.Services.AddSingleton<ConnectionOutcomeService>();
|
||||
builder.Services.AddSingleton<OperatorService>();
|
||||
builder.Services.AddSingleton(new ProvisioningReadiness(true));
|
||||
}
|
||||
|
||||
@@ -246,34 +313,13 @@ if (TrustedProxyForwarding.IsEnabled(configuredAbuseProtection))
|
||||
{
|
||||
app.UseForwardedHeaders();
|
||||
}
|
||||
app.UseMiddleware<TelemetryMiddleware>();
|
||||
app.UseExceptionHandler();
|
||||
app.UseMiddleware<HttpAbuseProtectionMiddleware>();
|
||||
app.MapOpenApi();
|
||||
app.MapRendezvousContractEndpoints();
|
||||
app.MapGet(
|
||||
"/health/live",
|
||||
static () => Results.Ok(new HealthResponse { Status = "live" }))
|
||||
.Produces<HealthResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.WithName("GetLiveness")
|
||||
.WithTags("Health");
|
||||
app.MapGet(
|
||||
"/health/ready",
|
||||
static (
|
||||
UdpMediatorService mediator,
|
||||
ProvisioningReadiness provisioning,
|
||||
IEphemeralRendezvousStore state) =>
|
||||
mediator.LocalEndpoint is null
|
||||
|| !provisioning.IsReady
|
||||
|| !state.IsAvailable
|
||||
|| state.IsDraining
|
||||
? Results.StatusCode(StatusCodes.Status503ServiceUnavailable)
|
||||
: Results.Ok(new HealthResponse { Status = "ready" }))
|
||||
.Produces<HealthResponse>()
|
||||
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
|
||||
.Produces(StatusCodes.Status503ServiceUnavailable)
|
||||
.WithName("GetReadiness")
|
||||
.WithTags("Health");
|
||||
app.MapOperatorEndpoints();
|
||||
app.MapRendezvousHealthEndpoints();
|
||||
|
||||
await app.RunAsync();
|
||||
|
||||
|
||||
@@ -142,8 +142,36 @@ internal sealed class SigningKeyRing : IDisposable
|
||||
return VerificationKeyLookup.Available;
|
||||
}
|
||||
|
||||
public bool Revoke(string keyId) =>
|
||||
_keys.ContainsKey(keyId) && _runtimeRevocations.TryAdd(keyId, 0);
|
||||
public bool Revoke(string keyId)
|
||||
{
|
||||
if (!_keys.ContainsKey(keyId))
|
||||
{
|
||||
return false;
|
||||
}
|
||||
|
||||
_runtimeRevocations.TryAdd(keyId, 0);
|
||||
return true;
|
||||
}
|
||||
|
||||
public IReadOnlyList<SigningKeyStatus> GetStatuses(DateTimeOffset now) => _keys.Values
|
||||
.OrderBy(static key => key.KeyId, StringComparer.Ordinal)
|
||||
.Select(key => new SigningKeyStatus(
|
||||
key.KeyId,
|
||||
IsRevoked(key)
|
||||
? "revoked"
|
||||
: now < key.NotBefore
|
||||
? "not-yet-valid"
|
||||
: now < key.SignUntil
|
||||
? "signing"
|
||||
: now < key.VerifyUntil
|
||||
? "verify-only"
|
||||
: "retired",
|
||||
key.SignUntil,
|
||||
key.VerifyUntil,
|
||||
key.GameId,
|
||||
key.EnvironmentId,
|
||||
key.CredentialKinds.Select(static kind => kind.ToString()).Order().ToArray()))
|
||||
.ToArray();
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
@@ -210,6 +238,15 @@ internal sealed class SigningKeyRing : IDisposable
|
||||
}
|
||||
}
|
||||
|
||||
internal sealed record SigningKeyStatus(
|
||||
string KeyId,
|
||||
string Status,
|
||||
DateTimeOffset SignUntil,
|
||||
DateTimeOffset VerifyUntil,
|
||||
string? GameId,
|
||||
string? EnvironmentId,
|
||||
IReadOnlyList<string> CredentialKinds);
|
||||
|
||||
internal sealed class SigningKey : IDisposable
|
||||
{
|
||||
private byte[]? _material;
|
||||
|
||||
@@ -297,6 +297,7 @@ internal sealed record StoredJoinAttempt
|
||||
public required SecretFingerprint ConnectionTicketFingerprint { get; init; }
|
||||
public NetworkEndpoint? DedicatedFallback { get; init; }
|
||||
public required DateTimeOffset ExpiresAt { get; init; }
|
||||
public required TimeSpan CreatedAtMonotonic { get; init; }
|
||||
public required DateTimeOffset ConnectionTicketExpiresAt { get; init; }
|
||||
public AttemptEndpointBinding? HostEndpoint { get; init; }
|
||||
public AttemptEndpointBinding? ClientEndpoint { get; init; }
|
||||
|
||||
@@ -23,7 +23,10 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
private TimeSpan? _drainDeadline;
|
||||
private TimeSpan _nextUdpMaintenance;
|
||||
private long _maintenanceSweepCount;
|
||||
private long _expiryChurn;
|
||||
private bool _available = true;
|
||||
private EphemeralStoreSnapshot? _metricsSnapshot;
|
||||
private TimeSpan _metricsSnapshotAt = TimeSpan.MinValue;
|
||||
|
||||
public InMemoryEphemeralRendezvousStore(
|
||||
EphemeralStoreOptions options,
|
||||
@@ -66,6 +69,50 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
}
|
||||
}
|
||||
|
||||
internal EphemeralStoreSnapshot GetSnapshot()
|
||||
{
|
||||
lock (_gate)
|
||||
{
|
||||
TimeSpan now = _monotonicClock.Elapsed;
|
||||
Cleanup(now);
|
||||
EphemeralStoreSnapshot snapshot = CreateSnapshot();
|
||||
_metricsSnapshot = snapshot;
|
||||
_metricsSnapshotAt = now;
|
||||
return snapshot;
|
||||
}
|
||||
}
|
||||
|
||||
internal EphemeralStoreSnapshot GetMetricsSnapshot()
|
||||
{
|
||||
lock (_gate)
|
||||
{
|
||||
TimeSpan now = _monotonicClock.Elapsed;
|
||||
if (_metricsSnapshot is null
|
||||
|| now < _metricsSnapshotAt
|
||||
|| now - _metricsSnapshotAt >= TimeSpan.FromMilliseconds(100))
|
||||
{
|
||||
Cleanup(now);
|
||||
_metricsSnapshot = CreateSnapshot();
|
||||
_metricsSnapshotAt = now;
|
||||
}
|
||||
|
||||
return _metricsSnapshot;
|
||||
}
|
||||
}
|
||||
|
||||
private EphemeralStoreSnapshot CreateSnapshot() => new(
|
||||
_listings.Count,
|
||||
_presence.Count,
|
||||
_attempts.Count,
|
||||
_outcomeReports.Count,
|
||||
_replay.Count,
|
||||
_revocations.Count,
|
||||
_idempotency.Count,
|
||||
_maintenanceSweepCount,
|
||||
_expiryChurn,
|
||||
_available,
|
||||
_drainDeadline.HasValue);
|
||||
|
||||
public StoreResult<StoredListing> CreateListing(
|
||||
CreateListingCommand command,
|
||||
CancellationToken cancellationToken = default) => Atomic<StoredListing>(now =>
|
||||
@@ -400,6 +447,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
|
||||
AttemptEntry attempt = new(
|
||||
command,
|
||||
now,
|
||||
now + _options.JoinAttemptLifetime,
|
||||
WallDeadline(now, _options.JoinAttemptLifetime));
|
||||
_attempts.Add(command.AttemptId, attempt);
|
||||
@@ -729,6 +777,10 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
return new(StoreResultCode.CapacityExceeded);
|
||||
}
|
||||
|
||||
int activeResourcesBefore = _listings.Count
|
||||
+ _presence.Count
|
||||
+ _attempts.Count
|
||||
+ _outcomeReports.Count;
|
||||
_revocations[subject] = now + lifetime;
|
||||
SessionListingId[] listings = _listings
|
||||
.Where(item => string.Equals(item.Value.Definition.OwnerSubject, subject, StringComparison.Ordinal))
|
||||
@@ -756,7 +808,11 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
_outcomeReports.Remove(attemptId);
|
||||
}
|
||||
|
||||
return new(StoreResultCode.Success, listings.Length + attempts.Length);
|
||||
int activeResourcesAfter = _listings.Count
|
||||
+ _presence.Count
|
||||
+ _attempts.Count
|
||||
+ _outcomeReports.Count;
|
||||
return new(StoreResultCode.Success, activeResourcesBefore - activeResourcesAfter);
|
||||
}, cancellationToken);
|
||||
|
||||
public void BeginDrain(CancellationToken cancellationToken = default)
|
||||
@@ -768,6 +824,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
if (!_drainDeadline.HasValue)
|
||||
{
|
||||
_drainDeadline = _monotonicClock.Elapsed + _options.GracefulDrainLifetime;
|
||||
_metricsSnapshot = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -778,6 +835,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
{
|
||||
_available = false;
|
||||
ClearActiveState();
|
||||
_metricsSnapshot = null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -801,7 +859,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
_nextUdpMaintenance = now + UdpMaintenanceInterval;
|
||||
}
|
||||
|
||||
return operation(now);
|
||||
StoreResult<T> result = operation(now);
|
||||
_metricsSnapshot = null;
|
||||
return result;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -838,44 +898,54 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
ClearActiveState();
|
||||
}
|
||||
|
||||
RemoveExpired(_revocations, now);
|
||||
RemoveExpired(_replay, now);
|
||||
foreach (string key in _idempotency
|
||||
_expiryChurn += RemoveExpired(_revocations, now);
|
||||
_expiryChurn += RemoveExpired(_replay, now);
|
||||
string[] expiredIdempotency = _idempotency
|
||||
.Where(item => item.Value.Deadline <= now)
|
||||
.Select(static item => item.Key)
|
||||
.ToArray())
|
||||
.ToArray();
|
||||
_expiryChurn += expiredIdempotency.Length;
|
||||
foreach (string key in expiredIdempotency)
|
||||
{
|
||||
_idempotency.Remove(key);
|
||||
}
|
||||
|
||||
foreach (MediationHandle handle in _presence
|
||||
MediationHandle[] expiredPresence = _presence
|
||||
.Where(item => item.Value.Deadline <= now)
|
||||
.Select(static item => item.Key)
|
||||
.ToArray())
|
||||
.ToArray();
|
||||
_expiryChurn += expiredPresence.Length;
|
||||
foreach (MediationHandle handle in expiredPresence)
|
||||
{
|
||||
_presence.Remove(handle);
|
||||
}
|
||||
|
||||
foreach (JoinAttemptId attemptId in _attempts
|
||||
JoinAttemptId[] expiredAttempts = _attempts
|
||||
.Where(item => item.Value.Deadline <= now)
|
||||
.Select(static item => item.Key)
|
||||
.ToArray())
|
||||
.ToArray();
|
||||
_expiryChurn += expiredAttempts.Length;
|
||||
foreach (JoinAttemptId attemptId in expiredAttempts)
|
||||
{
|
||||
RemoveAttempt(attemptId);
|
||||
}
|
||||
|
||||
foreach (JoinAttemptId attemptId in _outcomeReports
|
||||
JoinAttemptId[] expiredOutcomes = _outcomeReports
|
||||
.Where(item => item.Value.Deadline <= now)
|
||||
.Select(static item => item.Key)
|
||||
.ToArray())
|
||||
.ToArray();
|
||||
_expiryChurn += expiredOutcomes.Length;
|
||||
foreach (JoinAttemptId attemptId in expiredOutcomes)
|
||||
{
|
||||
_outcomeReports.Remove(attemptId);
|
||||
}
|
||||
|
||||
foreach (SessionListingId listingId in _listings
|
||||
SessionListingId[] expiredListings = _listings
|
||||
.Where(item => item.Value.LeaseDeadline <= now)
|
||||
.Select(static item => item.Key)
|
||||
.ToArray())
|
||||
.ToArray();
|
||||
_expiryChurn += expiredListings.Length;
|
||||
foreach (SessionListingId listingId in expiredListings)
|
||||
{
|
||||
RemoveListing(listingId);
|
||||
}
|
||||
@@ -959,6 +1029,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
ClientCapabilityFingerprint = entry.Command.ClientCapabilityFingerprint,
|
||||
ConnectionTicketFingerprint = entry.Command.ConnectionTicketFingerprint,
|
||||
DedicatedFallback = StoredListing.CopyEndpoint(entry.Command.DedicatedFallback),
|
||||
CreatedAtMonotonic = entry.CreatedAtMonotonic,
|
||||
ExpiresAt = entry.WallExpiresAt,
|
||||
ConnectionTicketExpiresAt = entry.TicketWallExpiresAt ?? default,
|
||||
HostEndpoint = entry.HostEndpoint,
|
||||
@@ -968,15 +1039,18 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
IsCancelled = entry.IsCancelled,
|
||||
};
|
||||
|
||||
private static void RemoveExpired(Dictionary<string, TimeSpan> entries, TimeSpan now)
|
||||
private static int RemoveExpired(Dictionary<string, TimeSpan> entries, TimeSpan now)
|
||||
{
|
||||
foreach (string key in entries
|
||||
string[] expired = entries
|
||||
.Where(item => item.Value <= now)
|
||||
.Select(static item => item.Key)
|
||||
.ToArray())
|
||||
.ToArray();
|
||||
foreach (string key in expired)
|
||||
{
|
||||
entries.Remove(key);
|
||||
}
|
||||
|
||||
return expired.Length;
|
||||
}
|
||||
|
||||
private static void ValidateListing(ListingDefinition listing)
|
||||
@@ -1101,10 +1175,12 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
|
||||
private sealed class AttemptEntry(
|
||||
CreateJoinAttemptCommand command,
|
||||
TimeSpan createdAtMonotonic,
|
||||
TimeSpan deadline,
|
||||
DateTimeOffset wallExpiresAt)
|
||||
{
|
||||
public CreateJoinAttemptCommand Command { get; } = command;
|
||||
public TimeSpan CreatedAtMonotonic { get; } = createdAtMonotonic;
|
||||
public SecretFingerprint HostCapabilityFingerprint { get; } = command.HostCapabilityFingerprint;
|
||||
public SecretFingerprint ClientCapabilityFingerprint { get; } = command.ClientCapabilityFingerprint;
|
||||
public SecretFingerprint ConnectionTicketFingerprint { get; } = command.ConnectionTicketFingerprint;
|
||||
@@ -1137,3 +1213,16 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
||||
object ResourceId,
|
||||
TimeSpan Deadline);
|
||||
}
|
||||
|
||||
internal sealed record EphemeralStoreSnapshot(
|
||||
int ActiveListings,
|
||||
int FreshPresenceBindings,
|
||||
int ActiveJoinAttempts,
|
||||
int RetainedOutcomeReports,
|
||||
int ReplayMarkers,
|
||||
int PrincipalRevocations,
|
||||
int IdempotencyEntries,
|
||||
long MaintenanceSweeps,
|
||||
long ExpiryChurn,
|
||||
bool IsAvailable,
|
||||
bool IsDraining);
|
||||
|
||||
@@ -1,8 +1,10 @@
|
||||
using System.Diagnostics;
|
||||
using System.Net;
|
||||
using System.Net.Sockets;
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using FinalFactory.Rendezvous.Server.Abuse;
|
||||
using FinalFactory.Rendezvous.Server.JoinAttempts;
|
||||
using FinalFactory.Rendezvous.Server.Observability;
|
||||
using FinalFactory.Rendezvous.Server.Sessions;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
|
||||
@@ -38,7 +40,9 @@ internal sealed class NatMediationProcessor(
|
||||
IEphemeralRendezvousStore store,
|
||||
ISessionCapabilityService capabilities,
|
||||
JoinAttemptService joinAttempts,
|
||||
AbuseProtectionService? abuseProtection = null)
|
||||
AbuseProtectionService? abuseProtection = null,
|
||||
RendezvousTelemetry? telemetry = null,
|
||||
IMonotonicClock? monotonicClock = null)
|
||||
{
|
||||
public NatMediationResult ProcessDatagram(
|
||||
ReadOnlySpan<byte> encoded,
|
||||
@@ -68,6 +72,26 @@ internal sealed class NatMediationProcessor(
|
||||
INatIntroductionSink introductionSink,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
long started = Stopwatch.GetTimestamp();
|
||||
using Activity? activity = telemetry?.StartActivity("UDP frozen", ActivityKind.Server);
|
||||
NatMediationResult result = ProcessDatagramCore(
|
||||
encoded,
|
||||
observedPublicEndpoint,
|
||||
introductionSink,
|
||||
cancellationToken);
|
||||
telemetry?.RecordUdp(
|
||||
"frozen",
|
||||
result.ToString(),
|
||||
Stopwatch.GetElapsedTime(started).TotalMilliseconds);
|
||||
return result;
|
||||
}
|
||||
|
||||
private NatMediationResult ProcessDatagramCore(
|
||||
ReadOnlySpan<byte> encoded,
|
||||
IPEndPoint observedPublicEndpoint,
|
||||
INatIntroductionSink introductionSink,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
|
||||
if (!RendezvousUdpCodec.TryDecode(encoded, out PresenceDatagram? datagram, out _)
|
||||
|| datagram is null
|
||||
@@ -141,12 +165,22 @@ internal sealed class NatMediationProcessor(
|
||||
IPEndPoint observedPublicEndpoint,
|
||||
string token,
|
||||
INatIntroductionSink introductionSink,
|
||||
CancellationToken cancellationToken = default) => ProcessRequestCore(
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
long started = Stopwatch.GetTimestamp();
|
||||
using Activity? activity = telemetry?.StartActivity("UDP litenet", ActivityKind.Server);
|
||||
NatMediationResult result = ProcessRequestCore(
|
||||
claimedLocalEndpoint,
|
||||
observedPublicEndpoint,
|
||||
token,
|
||||
introductionSink,
|
||||
cancellationToken);
|
||||
telemetry?.RecordUdp(
|
||||
"litenet",
|
||||
result.ToString(),
|
||||
Stopwatch.GetElapsedTime(started).TotalMilliseconds);
|
||||
return result;
|
||||
}
|
||||
|
||||
private NatMediationResult ProcessRequestCore(
|
||||
IPEndPoint claimedLocalEndpoint,
|
||||
@@ -255,6 +289,14 @@ internal sealed class NatMediationProcessor(
|
||||
try
|
||||
{
|
||||
introductionSink.Introduce(CreatePlan(consumed.Value, ticket.Value));
|
||||
if (telemetry is not null && monotonicClock is not null)
|
||||
{
|
||||
telemetry.RecordPairingLatency(Math.Max(
|
||||
0,
|
||||
(monotonicClock.Elapsed - consumed.Value.Attempt.CreatedAtMonotonic)
|
||||
.TotalMilliseconds));
|
||||
}
|
||||
|
||||
return NatMediationResult.Introduced;
|
||||
}
|
||||
catch (Exception exception) when (exception is SocketException
|
||||
|
||||
@@ -1,5 +1,8 @@
|
||||
{
|
||||
"Rendezvous": {
|
||||
"AbuseProtection": {
|
||||
"OperatorAllowedAddresses": ["127.0.0.1", "::1"]
|
||||
},
|
||||
"Provisioning": {
|
||||
"Issuer": "final-factory-rendezvous-development",
|
||||
"Audience": "final-factory-rendezvous",
|
||||
@@ -14,6 +17,14 @@
|
||||
"NotBefore": "2025-01-01T00:00:00Z",
|
||||
"SignUntil": "2035-01-01T00:00:00Z",
|
||||
"VerifyUntil": "2035-01-02T00:00:00Z"
|
||||
},
|
||||
{
|
||||
"KeyId": "development-operator-1",
|
||||
"SecretReference": "development:ephemeral/rendezvous-operator-signing",
|
||||
"CredentialKinds": ["Operator"],
|
||||
"NotBefore": "2025-01-01T00:00:00Z",
|
||||
"SignUntil": "2035-01-01T00:00:00Z",
|
||||
"VerifyUntil": "2035-01-02T00:00:00Z"
|
||||
}
|
||||
],
|
||||
"Games": [
|
||||
|
||||
@@ -6,16 +6,25 @@
|
||||
"MaxDatagramsPerPoll": 256,
|
||||
"PollIntervalMilliseconds": 2
|
||||
},
|
||||
"Audit": {
|
||||
"MaxEntries": 10000,
|
||||
"RetentionDays": 30
|
||||
},
|
||||
"AbuseProtection": {
|
||||
"WindowSeconds": 1,
|
||||
"MaxTrackedKeys": 100000,
|
||||
"CriticalTrackedKeyReserve": 2048,
|
||||
"UdpTrackedKeyLimit": 70000,
|
||||
"TrustedProxyAddresses": [],
|
||||
"OperatorAllowedAddresses": [],
|
||||
"HealthGlobalRequestsPerWindow": 1000,
|
||||
"HealthGlobalConcurrency": 32,
|
||||
"HealthIpPrefixRequestsPerWindow": 120,
|
||||
"HealthIpPrefixConcurrency": 8,
|
||||
"OperatorGlobalRequestsPerWindow": 1000,
|
||||
"OperatorGlobalConcurrency": 32,
|
||||
"OperatorIpPrefixRequestsPerWindow": 120,
|
||||
"OperatorIpPrefixConcurrency": 8,
|
||||
"HttpGlobalRequestsPerWindow": 20000,
|
||||
"HttpOptionalRequestsPerWindow": 18000,
|
||||
"HttpIpPrefixRequestsPerWindow": 500,
|
||||
|
||||
@@ -11,6 +11,11 @@ public sealed class OpenApiCompatibilityTests
|
||||
"/v1/join-attempts",
|
||||
"/v1/join-attempts/{attemptId}",
|
||||
"/v1/join-attempts/{attemptId}/outcome",
|
||||
"/v1/operator/drain",
|
||||
"/v1/operator/keys/revoke",
|
||||
"/v1/operator/listings/revoke",
|
||||
"/v1/operator/principals/revoke",
|
||||
"/v1/operator/status",
|
||||
"/v1/sessions",
|
||||
"/v1/sessions/{listingId}",
|
||||
"/v1/sessions/{listingId}/join-attempts",
|
||||
@@ -104,6 +109,11 @@ public sealed class OpenApiCompatibilityTests
|
||||
Assert.Equal(
|
||||
"X-Rendezvous-Client-Punch-Capability",
|
||||
attemptCapability.GetProperty("name").GetString());
|
||||
JsonElement operatorBearer = root.GetProperty("components")
|
||||
.GetProperty("securitySchemes")
|
||||
.GetProperty("OperatorBearer");
|
||||
Assert.Equal("http", operatorBearer.GetProperty("type").GetString());
|
||||
Assert.Equal("bearer", operatorBearer.GetProperty("scheme").GetString());
|
||||
(string Path, string Method)[] publisherOperations =
|
||||
[
|
||||
("/v1/sessions", "post"),
|
||||
@@ -120,6 +130,23 @@ public sealed class OpenApiCompatibilityTests
|
||||
Assert.True(security[0].TryGetProperty("PublisherBearer", out _));
|
||||
}
|
||||
|
||||
(string Path, string Method)[] operatorOperations =
|
||||
[
|
||||
("/v1/operator/status", "get"),
|
||||
("/v1/operator/listings/revoke", "post"),
|
||||
("/v1/operator/principals/revoke", "post"),
|
||||
("/v1/operator/keys/revoke", "post"),
|
||||
("/v1/operator/drain", "post"),
|
||||
];
|
||||
foreach ((string operationPath, string method) in operatorOperations)
|
||||
{
|
||||
JsonElement security = root.GetProperty("paths")
|
||||
.GetProperty(operationPath)
|
||||
.GetProperty(method)
|
||||
.GetProperty("security");
|
||||
Assert.True(security[0].TryGetProperty("OperatorBearer", out _));
|
||||
}
|
||||
|
||||
JsonElement cancelParameters = root.GetProperty("paths")
|
||||
.GetProperty("/v1/join-attempts/{attemptId}")
|
||||
.GetProperty("delete")
|
||||
@@ -157,6 +184,15 @@ public sealed class OpenApiCompatibilityTests
|
||||
static item => item.Name is "get" or "post" or "put" or "delete"))
|
||||
{
|
||||
JsonElement responses = operation.Value.GetProperty("responses");
|
||||
foreach (JsonProperty response in responses.EnumerateObject())
|
||||
{
|
||||
JsonElement correlation = response.Value.GetProperty("headers")
|
||||
.GetProperty("X-Rendezvous-Correlation-ID");
|
||||
Assert.Equal(
|
||||
"string",
|
||||
correlation.GetProperty("schema").GetProperty("type").GetString());
|
||||
}
|
||||
|
||||
if (!responses.TryGetProperty("429", out JsonElement overloaded))
|
||||
{
|
||||
continue;
|
||||
@@ -171,7 +207,7 @@ public sealed class OpenApiCompatibilityTests
|
||||
}
|
||||
}
|
||||
|
||||
Assert.Equal(12, overloadContracts);
|
||||
Assert.Equal(17, overloadContracts);
|
||||
(string Path, string Method)[] bodyOperations =
|
||||
[
|
||||
("/v1/sessions", "post"),
|
||||
@@ -180,6 +216,10 @@ public sealed class OpenApiCompatibilityTests
|
||||
("/v1/sessions/{listingId}", "delete"),
|
||||
("/v1/join-attempts", "post"),
|
||||
("/v1/join-attempts/{attemptId}/outcome", "post"),
|
||||
("/v1/operator/listings/revoke", "post"),
|
||||
("/v1/operator/principals/revoke", "post"),
|
||||
("/v1/operator/keys/revoke", "post"),
|
||||
("/v1/operator/drain", "post"),
|
||||
];
|
||||
foreach ((string operationPath, string method) in bodyOperations)
|
||||
{
|
||||
|
||||
@@ -0,0 +1,253 @@
|
||||
using System.Collections.Concurrent;
|
||||
using System.Diagnostics;
|
||||
using System.Diagnostics.Metrics;
|
||||
using System.Net;
|
||||
using FinalFactory.Rendezvous.Server.Observability;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
using FinalFactory.Rendezvous.Server.Transport;
|
||||
using FinalFactory.Rendezvous.Tests.State;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Tests.Observability;
|
||||
|
||||
[CollectionDefinition(RendezvousTelemetryIsolation.Name, DisableParallelization = true)]
|
||||
public sealed class RendezvousTelemetryIsolation
|
||||
{
|
||||
public const string Name = "Rendezvous telemetry";
|
||||
}
|
||||
|
||||
[Collection(RendezvousTelemetryIsolation.Name)]
|
||||
public sealed class ObservabilityTests
|
||||
{
|
||||
private static readonly HashSet<string> AllowedTagKeys =
|
||||
[
|
||||
"operation",
|
||||
"status_code",
|
||||
"result",
|
||||
"transport",
|
||||
"partition",
|
||||
"action",
|
||||
"outcome",
|
||||
"elapsed_bucket",
|
||||
];
|
||||
|
||||
[Fact]
|
||||
public void MetricsAndTracesUseBoundedDimensionsWithoutSensitiveValues()
|
||||
{
|
||||
EphemeralStateFixture fixture = new();
|
||||
using RendezvousTelemetry telemetry = new(fixture.Store);
|
||||
List<Measurement> measurements = [];
|
||||
using MeterListener meterListener = new();
|
||||
meterListener.InstrumentPublished = (instrument, listener) =>
|
||||
{
|
||||
if (instrument.Meter.Name == RendezvousTelemetry.MeterName)
|
||||
{
|
||||
listener.EnableMeasurementEvents(instrument);
|
||||
}
|
||||
};
|
||||
meterListener.SetMeasurementEventCallback<long>((instrument, value, tags, _) =>
|
||||
measurements.Add(new(instrument.Name, value, Tags(tags))));
|
||||
meterListener.SetMeasurementEventCallback<int>((instrument, value, tags, _) =>
|
||||
measurements.Add(new(instrument.Name, value, Tags(tags))));
|
||||
meterListener.SetMeasurementEventCallback<double>((instrument, value, tags, _) =>
|
||||
measurements.Add(new(instrument.Name, value, Tags(tags))));
|
||||
meterListener.Start();
|
||||
|
||||
Activity? observed = null;
|
||||
List<string> activityData = [];
|
||||
using ActivityListener activityListener = new()
|
||||
{
|
||||
ShouldListenTo = source => source.Name == RendezvousTelemetry.ActivitySourceName,
|
||||
Sample = static (ref ActivityCreationOptions<ActivityContext> _) =>
|
||||
ActivitySamplingResult.AllData,
|
||||
ActivityStopped = activity =>
|
||||
{
|
||||
observed = activity;
|
||||
activityData.Add(activity.DisplayName);
|
||||
activityData.AddRange(activity.TagObjects.Select(static tag => $"{tag.Key}={tag.Value}"));
|
||||
},
|
||||
};
|
||||
ActivitySource.AddActivityListener(activityListener);
|
||||
|
||||
const string secret = "secret-player-token-canary";
|
||||
using (telemetry.StartActivity("HTTP GetOperatorStatus", ActivityKind.Server))
|
||||
{
|
||||
telemetry.RecordHttp("GetOperatorStatus", 200, 3.5);
|
||||
telemetry.RecordUdp("frozen", "Introduced", 1.25);
|
||||
telemetry.RecordLimiterDrop("udp", "rate-or-concurrency");
|
||||
telemetry.RecordAudit("revoke-listing", "succeeded");
|
||||
telemetry.RecordConnectionOutcome("Connected", "UnderOneSecond");
|
||||
telemetry.RecordOperatorAuthentication("accepted");
|
||||
telemetry.RecordPairingLatency(12.5);
|
||||
}
|
||||
NatMediationProcessor processor = new(null!, null!, null!, telemetry: telemetry);
|
||||
Assert.Equal(
|
||||
NatMediationResult.Dropped,
|
||||
processor.ProcessRequest(
|
||||
new IPEndPoint(IPAddress.Parse("10.0.0.8"), 9000),
|
||||
new IPEndPoint(IPAddress.Parse("203.0.113.8"), 50000),
|
||||
secret,
|
||||
NoopIntroductionSink.Instance));
|
||||
|
||||
meterListener.RecordObservableInstruments();
|
||||
Assert.NotNull(observed);
|
||||
string flattened = string.Join('|', measurements.Select(static item => item.ToString()));
|
||||
Assert.DoesNotContain(secret, flattened, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(secret, string.Join('|', activityData), StringComparison.Ordinal);
|
||||
Assert.Contains(measurements, static item => item.Name == "rendezvous.http.requests");
|
||||
Assert.Contains(measurements, static item => item.Name == "rendezvous.udp.results");
|
||||
Assert.Contains(measurements, static item => item.Name == "rendezvous.limiter.drops");
|
||||
Assert.Contains(measurements, static item => item.Name == "rendezvous.queue.depth");
|
||||
Assert.Contains(measurements, static item => item.Name == "rendezvous.store.active_leases");
|
||||
Assert.Contains(measurements, static item => item.Name == "rendezvous.store.expiry_churn");
|
||||
Assert.Contains(measurements, static item => item.Name == "rendezvous.pairing.latency");
|
||||
Assert.All(measurements.SelectMany(static item => item.Tags), static tag =>
|
||||
Assert.Contains(tag.Key, AllowedTagKeys));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void MetricScrapeExpiresIdleStateAndReportsExpiryChurn()
|
||||
{
|
||||
EphemeralStateFixture fixture = new();
|
||||
Assert.True(fixture.Store.CreateListing(fixture.ListingCommand()).Succeeded);
|
||||
using RendezvousTelemetry telemetry = new(fixture.Store);
|
||||
List<Measurement> measurements = [];
|
||||
using MeterListener listener = new();
|
||||
listener.InstrumentPublished = (instrument, meterListener) =>
|
||||
{
|
||||
if (instrument.Meter.Name == RendezvousTelemetry.MeterName)
|
||||
{
|
||||
meterListener.EnableMeasurementEvents(instrument);
|
||||
}
|
||||
};
|
||||
listener.SetMeasurementEventCallback<long>((instrument, value, tags, _) =>
|
||||
measurements.Add(new(instrument.Name, value, Tags(tags))));
|
||||
listener.SetMeasurementEventCallback<int>((instrument, value, tags, _) =>
|
||||
measurements.Add(new(instrument.Name, value, Tags(tags))));
|
||||
listener.Start();
|
||||
|
||||
listener.RecordObservableInstruments();
|
||||
Assert.Equal(
|
||||
1,
|
||||
Assert.Single(measurements, static item => item.Name == "rendezvous.store.active_leases").Value);
|
||||
|
||||
fixture.Clock.Advance(TimeSpan.FromSeconds(61));
|
||||
measurements.Clear();
|
||||
listener.RecordObservableInstruments();
|
||||
Assert.Equal(
|
||||
0,
|
||||
Assert.Single(measurements, static item => item.Name == "rendezvous.store.active_leases").Value);
|
||||
Assert.True(Assert.Single(
|
||||
measurements,
|
||||
static item => item.Name == "rendezvous.store.expiry_churn").Value >= 1);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void AuditTrailFingerprintsIdentifiersEnforcesRetentionAndBoundsCapacity()
|
||||
{
|
||||
EphemeralStateFixture fixture = new();
|
||||
using RendezvousTelemetry telemetry = new(fixture.Store);
|
||||
CapturingLogger<AuditTrail> logger = new();
|
||||
ManualTimeProvider time = new(new DateTimeOffset(2026, 7, 16, 0, 0, 0, TimeSpan.Zero));
|
||||
AuditTrail audit = new(
|
||||
Options.Create(new AuditOptions { MaxEntries = 100, RetentionDays = 1 }),
|
||||
logger,
|
||||
telemetry,
|
||||
time);
|
||||
const string actor = "operator-secret-subject";
|
||||
const string target = "player-secret-subject";
|
||||
|
||||
for (int index = 0; index < 101; index++)
|
||||
{
|
||||
audit.Record(actor, "revoke-principal", "succeeded", "principal", target, "safe-correlation");
|
||||
}
|
||||
|
||||
IReadOnlyList<AuditEntry> bounded = audit.GetEntriesForTests();
|
||||
Assert.Equal(100, bounded.Count);
|
||||
Assert.All(bounded, entry =>
|
||||
{
|
||||
Assert.DoesNotContain(actor, entry.ToString(), StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(target, entry.ToString(), StringComparison.Ordinal);
|
||||
Assert.NotEqual(actor, entry.ActorFingerprint);
|
||||
Assert.NotEqual(target, entry.TargetFingerprint);
|
||||
});
|
||||
Assert.DoesNotContain(actor, string.Join('|', logger.Messages), StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(target, string.Join('|', logger.Messages), StringComparison.Ordinal);
|
||||
|
||||
time.Advance(TimeSpan.FromDays(2));
|
||||
Assert.Empty(audit.GetEntriesForTests());
|
||||
Assert.Empty(audit.GetAggregateCounts());
|
||||
audit.Record(actor, "inspect-status", "succeeded", "service", "rendezvous", "safe-correlation");
|
||||
Assert.Single(audit.GetEntriesForTests());
|
||||
Assert.Equal(1, audit.GetAggregateCounts()["inspect-status:succeeded"]);
|
||||
}
|
||||
|
||||
private static KeyValuePair<string, object?>[] Tags(
|
||||
ReadOnlySpan<KeyValuePair<string, object?>> tags) => tags.ToArray();
|
||||
|
||||
private sealed record Measurement(
|
||||
string Name,
|
||||
double Value,
|
||||
KeyValuePair<string, object?>[] Tags);
|
||||
|
||||
private sealed class ManualTimeProvider(DateTimeOffset now) : TimeProvider
|
||||
{
|
||||
private DateTimeOffset _now = now;
|
||||
public override DateTimeOffset GetUtcNow() => _now;
|
||||
public void Advance(TimeSpan duration) => _now += duration;
|
||||
}
|
||||
|
||||
private sealed class NoopIntroductionSink : INatIntroductionSink
|
||||
{
|
||||
public static NoopIntroductionSink Instance { get; } = new();
|
||||
public void Introduce(NatIntroductionPlan plan)
|
||||
{
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
internal sealed class CapturingLogger<T> : ILogger<T>
|
||||
{
|
||||
public List<string> Messages { get; } = [];
|
||||
|
||||
public IDisposable? BeginScope<TState>(TState state) where TState : notnull => null;
|
||||
public bool IsEnabled(LogLevel logLevel) => true;
|
||||
|
||||
public void Log<TState>(
|
||||
LogLevel logLevel,
|
||||
EventId eventId,
|
||||
TState state,
|
||||
Exception? exception,
|
||||
Func<TState, Exception?, string> formatter) => Messages.Add(formatter(state, exception));
|
||||
}
|
||||
|
||||
internal sealed class CapturingLoggerProvider : ILoggerProvider
|
||||
{
|
||||
public ConcurrentQueue<string> Messages { get; } = new();
|
||||
|
||||
public ILogger CreateLogger(string categoryName) => new Sink(Messages);
|
||||
|
||||
public void Dispose() => GC.SuppressFinalize(this);
|
||||
|
||||
private sealed class Sink(ConcurrentQueue<string> messages) : ILogger
|
||||
{
|
||||
public IDisposable BeginScope<TState>(TState state) where TState : notnull => Scope.Instance;
|
||||
public bool IsEnabled(LogLevel logLevel) => true;
|
||||
|
||||
public void Log<TState>(
|
||||
LogLevel logLevel,
|
||||
EventId eventId,
|
||||
TState state,
|
||||
Exception? exception,
|
||||
Func<TState, Exception?, string> formatter) => messages.Enqueue(formatter(state, exception));
|
||||
}
|
||||
|
||||
private sealed class Scope : IDisposable
|
||||
{
|
||||
public static Scope Instance { get; } = new();
|
||||
public void Dispose()
|
||||
{
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,519 @@
|
||||
using System.Diagnostics;
|
||||
using System.Diagnostics.Metrics;
|
||||
using System.Net;
|
||||
using System.Net.Http.Headers;
|
||||
using System.Net.Http.Json;
|
||||
using System.Text.Json;
|
||||
using FinalFactory.Rendezvous.Contracts;
|
||||
using FinalFactory.Rendezvous.Server.Abuse;
|
||||
using FinalFactory.Rendezvous.Server.Http;
|
||||
using FinalFactory.Rendezvous.Server.JoinAttempts;
|
||||
using FinalFactory.Rendezvous.Server.Observability;
|
||||
using FinalFactory.Rendezvous.Server.Operations;
|
||||
using FinalFactory.Rendezvous.Server.Provisioning;
|
||||
using FinalFactory.Rendezvous.Server.Sessions;
|
||||
using FinalFactory.Rendezvous.Server.State;
|
||||
using FinalFactory.Rendezvous.Server.Transport;
|
||||
using FinalFactory.Rendezvous.Tests.Observability;
|
||||
using FinalFactory.Rendezvous.Tests.Provisioning;
|
||||
using FinalFactory.Rendezvous.Tests.State;
|
||||
using Microsoft.AspNetCore.Builder;
|
||||
using Microsoft.AspNetCore.Hosting;
|
||||
using Microsoft.AspNetCore.Hosting.Server;
|
||||
using Microsoft.AspNetCore.Hosting.Server.Features;
|
||||
using Microsoft.AspNetCore.Http;
|
||||
using Microsoft.AspNetCore.Routing;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Logging;
|
||||
|
||||
namespace FinalFactory.Rendezvous.Tests.Operations;
|
||||
|
||||
[Collection(RendezvousTelemetryIsolation.Name)]
|
||||
public sealed class OperatorEndpointTests
|
||||
{
|
||||
[Fact]
|
||||
public async Task OperatorSurfaceSeparatesAuthenticationConfirmsActionsAndRedactsInspection()
|
||||
{
|
||||
await using OperatorTestHost host = await OperatorTestHost.StartAsync();
|
||||
List<string> telemetryData = [];
|
||||
using MeterListener meterListener = new();
|
||||
meterListener.InstrumentPublished = (instrument, listener) =>
|
||||
{
|
||||
if (instrument.Meter.Name == RendezvousTelemetry.MeterName)
|
||||
{
|
||||
listener.EnableMeasurementEvents(instrument);
|
||||
}
|
||||
};
|
||||
meterListener.SetMeasurementEventCallback<long>((instrument, value, tags, _) =>
|
||||
CaptureMeasurement(telemetryData, instrument, value, tags));
|
||||
meterListener.SetMeasurementEventCallback<int>((instrument, value, tags, _) =>
|
||||
CaptureMeasurement(telemetryData, instrument, value, tags));
|
||||
meterListener.SetMeasurementEventCallback<double>((instrument, value, tags, _) =>
|
||||
CaptureMeasurement(telemetryData, instrument, value, tags));
|
||||
meterListener.Start();
|
||||
using ActivityListener activityListener = new()
|
||||
{
|
||||
ShouldListenTo = static source => source.Name == RendezvousTelemetry.ActivitySourceName,
|
||||
Sample = static (ref ActivityCreationOptions<ActivityContext> _) =>
|
||||
ActivitySamplingResult.AllData,
|
||||
ActivityStopped = activity =>
|
||||
{
|
||||
telemetryData.Add(activity.DisplayName);
|
||||
telemetryData.AddRange(activity.TagObjects.Select(static tag => $"{tag.Key}={tag.Value}"));
|
||||
},
|
||||
};
|
||||
ActivitySource.AddActivityListener(activityListener);
|
||||
|
||||
using HttpResponseMessage liveBeforeDependencies = await host.Client.GetAsync("/health/live");
|
||||
Assert.Equal(HttpStatusCode.OK, liveBeforeDependencies.StatusCode);
|
||||
using HttpResponseMessage readyBeforeUdp = await host.Client.GetAsync("/health/ready");
|
||||
Assert.Equal(HttpStatusCode.ServiceUnavailable, readyBeforeUdp.StatusCode);
|
||||
|
||||
using HttpResponseMessage unauthenticated = await host.Client.GetAsync("/v1/operator/status");
|
||||
Assert.Equal(HttpStatusCode.Unauthorized, unauthenticated.StatusCode);
|
||||
AuthenticationHeaderValue challenge = Assert.Single(
|
||||
unauthenticated.Headers.WwwAuthenticate);
|
||||
Assert.Equal("Bearer", challenge.Scheme);
|
||||
Assert.Equal("realm=\"operator\"", challenge.Parameter);
|
||||
|
||||
using HttpResponseMessage publisher = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Get,
|
||||
"/v1/operator/status",
|
||||
host.PublisherCredential);
|
||||
Assert.Equal(HttpStatusCode.Unauthorized, publisher.StatusCode);
|
||||
|
||||
using HttpResponseMessage status = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Get,
|
||||
"/v1/operator/status",
|
||||
host.ReadOnlyOperatorCredential);
|
||||
Assert.Equal(HttpStatusCode.OK, status.StatusCode);
|
||||
string statusJson = await status.Content.ReadAsStringAsync();
|
||||
Assert.Contains("not-ready", statusJson, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(host.OwnerCanary, statusJson, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain("203.0.113.25", statusJson, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain("metadata", statusJson, StringComparison.OrdinalIgnoreCase);
|
||||
OperatorStatusResponse? operatorStatus = JsonSerializer.Deserialize<OperatorStatusResponse>(
|
||||
statusJson,
|
||||
ContractJson.Options);
|
||||
Assert.Contains(operatorStatus!.Tenants, static tenant =>
|
||||
tenant.GameId == "space-game"
|
||||
&& tenant.EnvironmentId == "production"
|
||||
&& tenant.Status == "enabled");
|
||||
Assert.Contains(operatorStatus.SigningKeys, static key =>
|
||||
key.KeyId == OperatorTestHost.OperatorKeyId
|
||||
&& key.Status == "signing"
|
||||
&& key.CredentialKinds.SequenceEqual(["Operator"]));
|
||||
|
||||
await host.StartUdpAsync();
|
||||
using HttpResponseMessage readyAfterUdp = await host.Client.GetAsync("/health/ready");
|
||||
Assert.Equal(HttpStatusCode.OK, readyAfterUdp.StatusCode);
|
||||
|
||||
using HttpResponseMessage exception = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/test/exception",
|
||||
host.FullOperatorCredential);
|
||||
Assert.Equal(HttpStatusCode.InternalServerError, exception.StatusCode);
|
||||
using HttpResponseMessage saturatedPublic = await host.Client.GetAsync("/test/public");
|
||||
Assert.Equal(HttpStatusCode.TooManyRequests, saturatedPublic.StatusCode);
|
||||
|
||||
using HttpResponseMessage forbidden = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/drain",
|
||||
host.ReadOnlyOperatorCredential,
|
||||
new BeginDrainRequest { Confirmation = "DRAIN" });
|
||||
Assert.Equal(HttpStatusCode.Forbidden, forbidden.StatusCode);
|
||||
|
||||
SessionListingId listingId = host.CreateListing(host.OwnerCanary);
|
||||
using HttpResponseMessage unconfirmedListing = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/listings/revoke",
|
||||
host.FullOperatorCredential,
|
||||
new RevokeListingRequest
|
||||
{
|
||||
ListingId = listingId.ToString(),
|
||||
ConfirmListingId = Guid.NewGuid().ToString("D"),
|
||||
});
|
||||
Assert.Equal(HttpStatusCode.BadRequest, unconfirmedListing.StatusCode);
|
||||
Assert.True(host.Store.GetListing(listingId, false).Succeeded);
|
||||
|
||||
using HttpResponseMessage revokedListing = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/listings/revoke",
|
||||
host.FullOperatorCredential,
|
||||
new RevokeListingRequest
|
||||
{
|
||||
ListingId = listingId.ToString(),
|
||||
ConfirmListingId = listingId.ToString(),
|
||||
});
|
||||
Assert.Equal(HttpStatusCode.OK, revokedListing.StatusCode);
|
||||
Assert.Equal(StoreResultCode.NotFound, host.Store.GetListing(listingId, false).Code);
|
||||
|
||||
const string principalCanary = "publisher-player-canary";
|
||||
SessionListingId principalListing = host.CreateListing(principalCanary);
|
||||
using HttpResponseMessage unconfirmedPrincipal = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/principals/revoke",
|
||||
host.FullOperatorCredential,
|
||||
new RevokePrincipalRequest
|
||||
{
|
||||
Subject = principalCanary,
|
||||
ConfirmSubject = "different-subject",
|
||||
LifetimeSeconds = 60,
|
||||
});
|
||||
Assert.Equal(HttpStatusCode.BadRequest, unconfirmedPrincipal.StatusCode);
|
||||
Assert.True(host.Store.GetListing(principalListing, false).Succeeded);
|
||||
|
||||
using HttpResponseMessage revokedPrincipal = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/principals/revoke",
|
||||
host.FullOperatorCredential,
|
||||
new RevokePrincipalRequest
|
||||
{
|
||||
Subject = principalCanary,
|
||||
ConfirmSubject = principalCanary,
|
||||
LifetimeSeconds = 60,
|
||||
});
|
||||
Assert.Equal(HttpStatusCode.OK, revokedPrincipal.StatusCode);
|
||||
OperatorActionResponse? principalResult = await revokedPrincipal.Content
|
||||
.ReadFromJsonAsync<OperatorActionResponse>(ContractJson.Options);
|
||||
Assert.Equal(1, principalResult!.AffectedResources);
|
||||
Assert.Equal(StoreResultCode.NotFound, host.Store.GetListing(principalListing, false).Code);
|
||||
StoreResult<StoredListing> blockedPublisher = host.CreateListingResult(
|
||||
principalCanary,
|
||||
out _);
|
||||
Assert.Equal(StoreResultCode.Revoked, blockedPublisher.Code);
|
||||
|
||||
using HttpResponseMessage drain = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/drain",
|
||||
host.FullOperatorCredential,
|
||||
new BeginDrainRequest { Confirmation = "DRAIN" });
|
||||
Assert.Equal(HttpStatusCode.OK, drain.StatusCode);
|
||||
Assert.True(host.Store.IsDraining);
|
||||
using HttpResponseMessage liveDuringDrain = await host.Client.GetAsync("/health/live");
|
||||
Assert.Equal(HttpStatusCode.OK, liveDuringDrain.StatusCode);
|
||||
using HttpResponseMessage readyDuringDrain = await host.Client.GetAsync("/health/ready");
|
||||
Assert.Equal(HttpStatusCode.ServiceUnavailable, readyDuringDrain.StatusCode);
|
||||
|
||||
using HttpResponseMessage firstPublisherKeyRevocation = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/keys/revoke",
|
||||
host.FullOperatorCredential,
|
||||
new RevokeSigningKeyRequest { KeyId = "key-1", ConfirmKeyId = "key-1" });
|
||||
Assert.Equal(HttpStatusCode.OK, firstPublisherKeyRevocation.StatusCode);
|
||||
using HttpResponseMessage repeatedPublisherKeyRevocation = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/keys/revoke",
|
||||
host.FullOperatorCredential,
|
||||
new RevokeSigningKeyRequest { KeyId = "key-1", ConfirmKeyId = "key-1" });
|
||||
Assert.Equal(HttpStatusCode.OK, repeatedPublisherKeyRevocation.StatusCode);
|
||||
|
||||
using HttpResponseMessage unconfirmedKey = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/keys/revoke",
|
||||
host.FullOperatorCredential,
|
||||
new RevokeSigningKeyRequest
|
||||
{
|
||||
KeyId = OperatorTestHost.OperatorKeyId,
|
||||
ConfirmKeyId = "different-key",
|
||||
});
|
||||
Assert.Equal(HttpStatusCode.BadRequest, unconfirmedKey.StatusCode);
|
||||
|
||||
using HttpResponseMessage revokedKey = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Post,
|
||||
"/v1/operator/keys/revoke",
|
||||
host.FullOperatorCredential,
|
||||
new RevokeSigningKeyRequest
|
||||
{
|
||||
KeyId = OperatorTestHost.OperatorKeyId,
|
||||
ConfirmKeyId = OperatorTestHost.OperatorKeyId,
|
||||
});
|
||||
Assert.Equal(HttpStatusCode.OK, revokedKey.StatusCode);
|
||||
using HttpResponseMessage afterKeyRevocation = await SendAsync(
|
||||
host,
|
||||
HttpMethod.Get,
|
||||
"/v1/operator/status",
|
||||
host.FullOperatorCredential);
|
||||
Assert.Equal(HttpStatusCode.Unauthorized, afterKeyRevocation.StatusCode);
|
||||
|
||||
string auditText = string.Join('|', host.Audit.GetEntriesForTests());
|
||||
string logText = string.Join('|', host.AuditLogger.Messages.Concat(host.AllLogs.Messages));
|
||||
Assert.DoesNotContain(host.OwnerCanary, auditText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(principalCanary, auditText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(listingId.ToString(), auditText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(host.OwnerCanary, logText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(principalCanary, logText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain("exception-secret-canary", logText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(host.FullOperatorCredential, logText, StringComparison.Ordinal);
|
||||
meterListener.RecordObservableInstruments();
|
||||
string telemetryText = string.Join('|', telemetryData);
|
||||
Assert.DoesNotContain(host.OwnerCanary, telemetryText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(principalCanary, telemetryText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain("exception-secret-canary", telemetryText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(host.FullOperatorCredential, telemetryText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain(listingId.ToString(), telemetryText, StringComparison.Ordinal);
|
||||
Assert.DoesNotContain("203.0.113.25", telemetryText, StringComparison.Ordinal);
|
||||
Assert.Contains(host.Audit.GetEntriesForTests(), static entry =>
|
||||
entry.Action == "begin-drain" && entry.Result == "succeeded");
|
||||
Assert.Contains(host.Audit.GetEntriesForTests(), static entry =>
|
||||
entry.Action == "revoke-listing" && entry.Result == "rejected");
|
||||
Assert.Contains(host.Audit.GetEntriesForTests(), static entry =>
|
||||
entry.Action == "begin-drain" && entry.Result == "forbidden");
|
||||
}
|
||||
|
||||
private static async Task<HttpResponseMessage> SendAsync(
|
||||
OperatorTestHost host,
|
||||
HttpMethod method,
|
||||
string path,
|
||||
string bearer,
|
||||
object? body = null)
|
||||
{
|
||||
using HttpRequestMessage request = new(method, path);
|
||||
request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", bearer);
|
||||
if (body is not null)
|
||||
{
|
||||
request.Content = JsonContent.Create(body, options: ContractJson.Options);
|
||||
}
|
||||
|
||||
return await host.Client.SendAsync(request);
|
||||
}
|
||||
|
||||
private static void CaptureMeasurement<T>(
|
||||
List<string> destination,
|
||||
Instrument instrument,
|
||||
T value,
|
||||
ReadOnlySpan<KeyValuePair<string, object?>> tags)
|
||||
where T : struct
|
||||
{
|
||||
destination.Add($"{instrument.Name}={value}");
|
||||
destination.AddRange(tags.ToArray().Select(static tag => $"{tag.Key}={tag.Value}"));
|
||||
}
|
||||
|
||||
private sealed class OperatorTestHost : IAsyncDisposable
|
||||
{
|
||||
private readonly WebApplication _application;
|
||||
private int _listingSequence;
|
||||
private bool _udpStarted;
|
||||
|
||||
private OperatorTestHost(
|
||||
WebApplication application,
|
||||
HttpClient client,
|
||||
ManualRendezvousClock clock,
|
||||
InMemoryEphemeralRendezvousStore store,
|
||||
AuditTrail audit,
|
||||
CapturingLogger<AuditTrail> auditLogger,
|
||||
CapturingLoggerProvider allLogs,
|
||||
string publisherCredential,
|
||||
string readOnlyOperatorCredential,
|
||||
string fullOperatorCredential)
|
||||
{
|
||||
_application = application;
|
||||
Client = client;
|
||||
Clock = clock;
|
||||
Store = store;
|
||||
Audit = audit;
|
||||
AuditLogger = auditLogger;
|
||||
AllLogs = allLogs;
|
||||
PublisherCredential = publisherCredential;
|
||||
ReadOnlyOperatorCredential = readOnlyOperatorCredential;
|
||||
FullOperatorCredential = fullOperatorCredential;
|
||||
}
|
||||
|
||||
internal const string OperatorKeyId = "operator-key";
|
||||
internal string OwnerCanary { get; } = "publisher-owner-canary";
|
||||
internal HttpClient Client { get; }
|
||||
internal ManualRendezvousClock Clock { get; }
|
||||
internal InMemoryEphemeralRendezvousStore Store { get; }
|
||||
internal AuditTrail Audit { get; }
|
||||
internal CapturingLogger<AuditTrail> AuditLogger { get; }
|
||||
internal CapturingLoggerProvider AllLogs { get; }
|
||||
internal string PublisherCredential { get; }
|
||||
internal string ReadOnlyOperatorCredential { get; }
|
||||
internal string FullOperatorCredential { get; }
|
||||
|
||||
internal static async Task<OperatorTestHost> StartAsync()
|
||||
{
|
||||
ManualRendezvousClock clock = new(ProvisioningTestData.Now);
|
||||
EphemeralStoreOptions stateOptions = new();
|
||||
InMemoryEphemeralRendezvousStore store = new(stateOptions, clock, clock);
|
||||
SigningKeyOptions publisherKey = ProvisioningTestData.CreateKey();
|
||||
SigningKeyOptions operatorKey = ProvisioningTestData.CreateKey(
|
||||
OperatorKeyId,
|
||||
"operator-secret",
|
||||
credentialKinds: [PrincipalCredentialKind.Operator],
|
||||
gameId: null,
|
||||
environmentId: null);
|
||||
ProvisioningOptions options = ProvisioningTestData.CreateOptions();
|
||||
options.SigningKeys = [publisherKey, operatorKey];
|
||||
ProvisioningRuntime provisioning = ProvisioningRuntime.Create(
|
||||
options,
|
||||
ProvisioningTestData.CreateSecrets("secret-1", "operator-secret"),
|
||||
clock.UtcNow);
|
||||
string publisherCredential = provisioning.Credentials.Issue(
|
||||
ProvisioningTestData.CreateDedicatedPublisher(),
|
||||
clock.UtcNow);
|
||||
string readOnlyCredential = provisioning.Credentials.Issue(
|
||||
new OperatorPrincipal(
|
||||
"operator-readonly",
|
||||
clock.UtcNow.AddMinutes(10),
|
||||
[OperatorPermission.ReadPolicy]),
|
||||
clock.UtcNow);
|
||||
string fullCredential = provisioning.Credentials.Issue(
|
||||
new OperatorPrincipal(
|
||||
"operator-full",
|
||||
clock.UtcNow.AddMinutes(10),
|
||||
Enum.GetValues<OperatorPermission>()),
|
||||
clock.UtcNow);
|
||||
CapturingLogger<AuditTrail> auditLogger = new();
|
||||
CapturingLoggerProvider allLogs = new();
|
||||
EphemeralCapabilityIssuer capabilities = new();
|
||||
|
||||
WebApplicationBuilder builder = WebApplication.CreateBuilder();
|
||||
builder.WebHost.UseUrls("http://127.0.0.1:0");
|
||||
builder.Logging.ClearProviders();
|
||||
builder.Logging.SetMinimumLevel(LogLevel.Debug);
|
||||
builder.Logging.AddProvider(allLogs);
|
||||
builder.Logging.AddFilter(
|
||||
"Microsoft.AspNetCore.Diagnostics.ExceptionHandlerMiddleware",
|
||||
LogLevel.None);
|
||||
builder.Services.ConfigureHttpJsonOptions(static json =>
|
||||
ContractJson.Configure(json.SerializerOptions));
|
||||
builder.Services.Configure<RouteHandlerOptions>(static route =>
|
||||
route.ThrowOnBadRequest = true);
|
||||
builder.Services.AddProblemDetails();
|
||||
builder.Services.AddExceptionHandler<RendezvousExceptionHandler>();
|
||||
builder.Services.AddOptions<AbuseProtectionOptions>().Configure(static abuse =>
|
||||
{
|
||||
abuse.OperatorAllowedAddresses = ["127.0.0.1"];
|
||||
abuse.HttpGlobalRequestsPerWindow = 1;
|
||||
abuse.HttpOptionalRequestsPerWindow = 1;
|
||||
abuse.HttpIpPrefixRequestsPerWindow = 1;
|
||||
abuse.HttpOptionalIpPrefixRequestsPerWindow = 1;
|
||||
});
|
||||
builder.Services.AddOptions<AuditOptions>();
|
||||
builder.Services.AddOptions<UdpMediatorOptions>().Configure(static udp =>
|
||||
{
|
||||
udp.ListenAddress = "127.0.0.1";
|
||||
udp.Port = 0;
|
||||
});
|
||||
builder.Services.AddSingleton(provisioning);
|
||||
builder.Services.AddSingleton(provisioning.Policies);
|
||||
builder.Services.AddSingleton(provisioning.Credentials);
|
||||
builder.Services.AddSingleton(provisioning.PublisherAuthorization);
|
||||
builder.Services.AddSingleton(store);
|
||||
builder.Services.AddSingleton<IEphemeralRendezvousStore>(store);
|
||||
builder.Services.AddSingleton<IWallClock>(clock);
|
||||
builder.Services.AddSingleton<IMonotonicClock>(clock);
|
||||
builder.Services.AddSingleton(capabilities);
|
||||
builder.Services.AddSingleton<ISessionCapabilityService>(capabilities);
|
||||
builder.Services.AddSingleton<JoinAttemptCursorCodec>();
|
||||
builder.Services.AddSingleton<JoinAttemptService>();
|
||||
builder.Services.AddSingleton<RendezvousTelemetry>();
|
||||
builder.Services.AddSingleton<AbuseProtectionService>();
|
||||
builder.Services.AddSingleton<NatMediationProcessor>();
|
||||
builder.Services.AddSingleton<UdpMediatorService>();
|
||||
builder.Services.AddSingleton(new ProvisioningReadiness(true));
|
||||
builder.Services.AddSingleton<RendezvousReadiness>();
|
||||
builder.Services.AddSingleton<ILogger<AuditTrail>>(auditLogger);
|
||||
builder.Services.AddSingleton<AuditTrail>();
|
||||
builder.Services.AddSingleton<OperatorService>();
|
||||
|
||||
WebApplication app = builder.Build();
|
||||
app.UseMiddleware<TelemetryMiddleware>();
|
||||
app.UseExceptionHandler();
|
||||
app.UseMiddleware<HttpAbuseProtectionMiddleware>();
|
||||
app.MapOperatorEndpoints();
|
||||
app.MapRendezvousHealthEndpoints();
|
||||
app.MapGet("/test/public", static () => Results.Ok()).WithName("TestPublic");
|
||||
app.MapPost(
|
||||
"/test/exception",
|
||||
static IResult () => throw new InvalidOperationException("exception-secret-canary"))
|
||||
.WithName("TestSecretException");
|
||||
await app.StartAsync();
|
||||
IServer server = app.Services.GetRequiredService<IServer>();
|
||||
string address = Assert.Single(server.Features.Get<IServerAddressesFeature>()!.Addresses);
|
||||
return new OperatorTestHost(
|
||||
app,
|
||||
new HttpClient { BaseAddress = new Uri(address) },
|
||||
clock,
|
||||
store,
|
||||
app.Services.GetRequiredService<AuditTrail>(),
|
||||
auditLogger,
|
||||
allLogs,
|
||||
publisherCredential,
|
||||
readOnlyCredential,
|
||||
fullCredential);
|
||||
}
|
||||
|
||||
internal SessionListingId CreateListing(string owner)
|
||||
{
|
||||
StoreResult<StoredListing> result = CreateListingResult(owner, out SessionListingId listingId);
|
||||
Assert.True(result.Succeeded);
|
||||
return listingId;
|
||||
}
|
||||
|
||||
internal StoreResult<StoredListing> CreateListingResult(
|
||||
string owner,
|
||||
out SessionListingId listingId)
|
||||
{
|
||||
int sequence = Interlocked.Increment(ref _listingSequence);
|
||||
listingId = new(Guid.NewGuid());
|
||||
StoreResult<StoredListing> result = Store.CreateListing(new(
|
||||
$"operator-listing-{sequence}",
|
||||
$"operator-request-{sequence}",
|
||||
new ListingDefinition
|
||||
{
|
||||
ListingId = listingId,
|
||||
LeaseId = new(Guid.NewGuid()),
|
||||
Scope = new(new GameId("space-game"), new EnvironmentId("production")),
|
||||
OwnerSubject = owner,
|
||||
RegionId = new("eu-central"),
|
||||
ProtocolVersion = 7,
|
||||
BuildVersion = "1.0.0",
|
||||
DisplayName = "Operator test listing",
|
||||
Visibility = ListingVisibility.Public,
|
||||
TrustMode = PublisherTrustMode.ManagedDedicated,
|
||||
CurrentPlayers = 1,
|
||||
MaximumPlayers = 4,
|
||||
Metadata = new Dictionary<string, string> { ["mode"] = "online-coop" },
|
||||
LeaseFingerprint = new("lease-fingerprint"),
|
||||
HostPresenceHandle = new(Guid.NewGuid()),
|
||||
HostPresenceFingerprint = new("presence-fingerprint"),
|
||||
CapabilityDerivationSalt = new string('A', 43),
|
||||
}));
|
||||
return result;
|
||||
}
|
||||
|
||||
internal async Task StartUdpAsync()
|
||||
{
|
||||
await _application.Services.GetRequiredService<UdpMediatorService>()
|
||||
.StartAsync(CancellationToken.None);
|
||||
_udpStarted = true;
|
||||
}
|
||||
|
||||
public async ValueTask DisposeAsync()
|
||||
{
|
||||
Client.Dispose();
|
||||
if (_udpStarted)
|
||||
{
|
||||
await _application.Services.GetRequiredService<UdpMediatorService>()
|
||||
.StopAsync(CancellationToken.None);
|
||||
}
|
||||
await _application.StopAsync();
|
||||
await _application.DisposeAsync();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -94,6 +94,7 @@ public sealed class PrincipalCredentialTests
|
||||
CredentialValidationError.SignatureInvalid,
|
||||
service.Validate(tampered, ProvisioningTestData.Now).Error);
|
||||
Assert.True(keys.Revoke("key-1"));
|
||||
Assert.True(keys.Revoke("key-1"));
|
||||
Assert.Equal(
|
||||
CredentialValidationError.KeyRevoked,
|
||||
service.Validate(token, ProvisioningTestData.Now).Error);
|
||||
|
||||
@@ -6,6 +6,7 @@ using FinalFactory.Rendezvous.Server.Http;
|
||||
using Microsoft.AspNetCore.Builder;
|
||||
using Microsoft.AspNetCore.Http;
|
||||
using Microsoft.AspNetCore.HttpOverrides;
|
||||
using Microsoft.AspNetCore.Routing;
|
||||
using Microsoft.Extensions.Logging.Abstractions;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
@@ -86,6 +87,9 @@ public sealed class AbuseProtectionTests
|
||||
AbuseProtectionService protection = new(Options.Create(options));
|
||||
IPAddress source = IPAddress.Parse("198.51.100.10");
|
||||
|
||||
Assert.True(protection.IsOperatorSourceAllowed(source));
|
||||
Assert.True(protection.IsOperatorSourceAllowed(IPAddress.Parse("::ffff:198.51.100.10")));
|
||||
Assert.False(protection.IsOperatorSourceAllowed(IPAddress.Parse("198.51.100.11")));
|
||||
AssertAccepted(protection, source, "BrowseSessions");
|
||||
AssertAccepted(protection, source, "BrowseSessions");
|
||||
AssertRejected(protection, source, "BrowseSessions");
|
||||
@@ -93,6 +97,57 @@ public sealed class AbuseProtectionTests
|
||||
AssertRejected(protection, source, "RenewSessionLease");
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void PublicSaturationCannotConsumeTheOperatorPartition()
|
||||
{
|
||||
AbuseProtectionOptions options = PermissiveOptions();
|
||||
options.HttpGlobalRequestsPerWindow = 1;
|
||||
options.HttpOptionalRequestsPerWindow = 1;
|
||||
options.OperatorGlobalRequestsPerWindow = 1;
|
||||
AbuseProtectionService protection = new(Options.Create(options));
|
||||
IPAddress source = IPAddress.Parse("198.51.100.10");
|
||||
|
||||
AssertAccepted(protection, source, "BrowseSessions");
|
||||
AssertRejected(protection, source, "BrowseSessions");
|
||||
Assert.True(protection.TryAcquireOperatorIngress(source, out var lease, out _));
|
||||
lease!.Dispose();
|
||||
Assert.False(protection.TryAcquireOperatorIngress(source, out _, out _));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task DeniedOperatorSourcesConsumeTheBoundedPublicPartition()
|
||||
{
|
||||
AbuseProtectionOptions options = PermissiveOptions();
|
||||
options.OperatorAllowedAddresses = ["192.0.2.10"];
|
||||
options.HttpGlobalRequestsPerWindow = 1;
|
||||
options.HttpOptionalRequestsPerWindow = 1;
|
||||
options.HttpIpPrefixRequestsPerWindow = 1;
|
||||
options.HttpOptionalIpPrefixRequestsPerWindow = 1;
|
||||
AbuseProtectionService protection = new(Options.Create(options));
|
||||
bool dispatched = false;
|
||||
HttpAbuseProtectionMiddleware middleware = new(
|
||||
_ =>
|
||||
{
|
||||
dispatched = true;
|
||||
return Task.CompletedTask;
|
||||
},
|
||||
protection);
|
||||
|
||||
DefaultHttpContext first = Context("198.51.100.10");
|
||||
first.SetEndpoint(new Endpoint(
|
||||
_ => Task.CompletedTask,
|
||||
new EndpointMetadataCollection(new EndpointNameMetadata("GetOperatorStatus")),
|
||||
"operator-status"));
|
||||
await middleware.InvokeAsync(first);
|
||||
Assert.Equal(StatusCodes.Status404NotFound, first.Response.StatusCode);
|
||||
|
||||
DefaultHttpContext repeated = Context("198.51.100.10");
|
||||
repeated.SetEndpoint(first.GetEndpoint());
|
||||
await middleware.InvokeAsync(repeated);
|
||||
Assert.Equal(StatusCodes.Status429TooManyRequests, repeated.Response.StatusCode);
|
||||
Assert.False(dispatched);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void ResourceBudgetsRemainIsolatedAcrossTenantAndPrincipalScopes()
|
||||
{
|
||||
@@ -354,7 +409,8 @@ public sealed class AbuseProtectionTests
|
||||
{
|
||||
const string canary = "credential-canary <script> endpoint=203.0.113.8:9000";
|
||||
DefaultHttpContext context = Context("198.51.100.10");
|
||||
RendezvousExceptionHandler handler = new();
|
||||
RendezvousExceptionHandler handler = new(
|
||||
NullLogger<RendezvousExceptionHandler>.Instance);
|
||||
|
||||
Assert.True(await handler.TryHandleAsync(
|
||||
context,
|
||||
@@ -438,6 +494,11 @@ public sealed class AbuseProtectionTests
|
||||
HealthGlobalConcurrency = 10_000,
|
||||
HealthIpPrefixRequestsPerWindow = 10_000,
|
||||
HealthIpPrefixConcurrency = 1_000,
|
||||
OperatorAllowedAddresses = ["198.51.100.10"],
|
||||
OperatorGlobalRequestsPerWindow = 10_000,
|
||||
OperatorGlobalConcurrency = 10_000,
|
||||
OperatorIpPrefixRequestsPerWindow = 10_000,
|
||||
OperatorIpPrefixConcurrency = 1_000,
|
||||
HttpGlobalRequestsPerWindow = 10_000,
|
||||
HttpOptionalRequestsPerWindow = 9_000,
|
||||
HttpIpPrefixRequestsPerWindow = 10_000,
|
||||
|
||||
@@ -436,7 +436,7 @@ public sealed class InMemoryEphemeralRendezvousStoreTests
|
||||
StoreResult<int> revoked = fixture.Store.RevokePrincipal(command.Listing.OwnerSubject, TimeSpan.FromMinutes(1));
|
||||
|
||||
Assert.True(revoked.Succeeded);
|
||||
Assert.Equal(2, revoked.Value);
|
||||
Assert.Equal(4, revoked.Value);
|
||||
Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(listing.Definition.ListingId, false).Code);
|
||||
Assert.Equal(StoreResultCode.NotFound, fixture.Store.BindAttemptEndpoint(new(
|
||||
attempt.MediationHandle,
|
||||
|
||||
Reference in New Issue
Block a user