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).
|
defined in [game provisioning and signing-key lifecycle](docs/security/provisioning.md).
|
||||||
Layered HTTP/UDP budgets, overload behavior, and safe operational tuning are
|
Layered HTTP/UDP budgets, overload behavior, and safe operational tuning are
|
||||||
defined in [hostile-input and overload protection](docs/security/abuse-protection.md).
|
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
|
The scriptable host/browser/join diagnostic and its stable automation contract are
|
||||||
documented in the [TestClient integration guide](docs/integration/test-client.md).
|
documented in the [TestClient integration guide](docs/integration/test-client.md).
|
||||||
The always-on three-party scenarios, optional Linux namespace topology, and
|
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
|
Run the bootstrap server with
|
||||||
`dotnet run --project src/FinalFactory.Rendezvous.Server`. It serves HTTP health endpoints and binds
|
`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 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
|
startup fails closed until externally supplied game policies and `env:` signing
|
||||||
key references resolve to valid key material; no reusable game secret is stored
|
key references resolve to valid key material; no reusable game secret is stored
|
||||||
in this repository or the public Client package.
|
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
|
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
|
make a healthy instance fail its orchestrator probes, while health traffic is
|
||||||
still bounded.
|
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
|
3. Once an endpoint has safely derived identities, it also acquires applicable
|
||||||
tenant, principal or capability, and listing/attempt budgets. Secret
|
tenant, principal or capability, and listing/attempt budgets. Secret
|
||||||
capabilities are represented only by bounded SHA-256 fingerprints.
|
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
|
`Rendezvous:AbuseProtection:MaxTrackedKeys` is a hard combined ceiling for rate
|
||||||
and active-concurrency keys. General HTTP and UDP traffic cannot consume the
|
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
|
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
|
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
|
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[] TrustedProxyAddresses { get; set; } = [];
|
||||||
|
|
||||||
|
public string[] OperatorAllowedAddresses { get; set; } = [];
|
||||||
|
|
||||||
[Range(1, 100_000)]
|
[Range(1, 100_000)]
|
||||||
public int HealthGlobalRequestsPerWindow { get; set; } = 1_000;
|
public int HealthGlobalRequestsPerWindow { get; set; } = 1_000;
|
||||||
|
|
||||||
@@ -32,6 +34,18 @@ internal sealed class AbuseProtectionOptions
|
|||||||
[Range(1, 1_000)]
|
[Range(1, 1_000)]
|
||||||
public int HealthIpPrefixConcurrency { get; set; } = 8;
|
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)]
|
[Range(1, 1_000_000)]
|
||||||
public int HttpGlobalRequestsPerWindow { get; set; } = 20_000;
|
public int HttpGlobalRequestsPerWindow { get; set; } = 20_000;
|
||||||
|
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ using System.Buffers;
|
|||||||
using System.Net;
|
using System.Net;
|
||||||
using System.Security.Cryptography;
|
using System.Security.Cryptography;
|
||||||
using System.Text;
|
using System.Text;
|
||||||
|
using FinalFactory.Rendezvous.Server.Observability;
|
||||||
using Microsoft.Extensions.Options;
|
using Microsoft.Extensions.Options;
|
||||||
|
|
||||||
namespace FinalFactory.Rendezvous.Server.Abuse;
|
namespace FinalFactory.Rendezvous.Server.Abuse;
|
||||||
@@ -12,13 +13,23 @@ internal sealed class AbuseProtectionService
|
|||||||
private readonly TimeProvider _timeProvider;
|
private readonly TimeProvider _timeProvider;
|
||||||
private readonly TrackerState _httpTracker;
|
private readonly TrackerState _httpTracker;
|
||||||
private readonly TrackerState _udpTracker;
|
private readonly TrackerState _udpTracker;
|
||||||
|
private readonly RendezvousTelemetry? _telemetry;
|
||||||
|
private readonly HashSet<string> _operatorAllowedAddresses;
|
||||||
|
|
||||||
public AbuseProtectionService(
|
public AbuseProtectionService(
|
||||||
IOptions<AbuseProtectionOptions> options,
|
IOptions<AbuseProtectionOptions> options,
|
||||||
TimeProvider? timeProvider = null)
|
TimeProvider? timeProvider = null,
|
||||||
|
RendezvousTelemetry? telemetry = null)
|
||||||
{
|
{
|
||||||
_options = options.Value;
|
_options = options.Value;
|
||||||
_timeProvider = timeProvider ?? TimeProvider.System;
|
_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();
|
DateTimeOffset now = _timeProvider.GetUtcNow();
|
||||||
_httpTracker = new(now);
|
_httpTracker = new(now);
|
||||||
_udpTracker = new(now);
|
_udpTracker = new(now);
|
||||||
@@ -87,6 +98,35 @@ internal sealed class AbuseProtectionService
|
|||||||
out retryAfterSeconds);
|
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(
|
public bool TryAcquireHttpIdentity(
|
||||||
string operation,
|
string operation,
|
||||||
string? tenant,
|
string? tenant,
|
||||||
@@ -266,7 +306,37 @@ internal sealed class AbuseProtectionService
|
|||||||
out int retryAfterSeconds)
|
out int retryAfterSeconds)
|
||||||
{
|
{
|
||||||
TrackerState tracker = domain == TrackerDomain.Udp ? _udpTracker : _httpTracker;
|
TrackerState tracker = domain == TrackerDomain.Udp ? _udpTracker : _httpTracker;
|
||||||
|
bool accepted;
|
||||||
lock (tracker.Gate)
|
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();
|
DateTimeOffset now = _timeProvider.GetUtcNow();
|
||||||
TimeSpan window = TimeSpan.FromSeconds(_options.WindowSeconds);
|
TimeSpan window = TimeSpan.FromSeconds(_options.WindowSeconds);
|
||||||
@@ -327,7 +397,6 @@ internal sealed class AbuseProtectionService
|
|||||||
lease = new AbuseLease(this, tracker, acquiredConcurrency);
|
lease = new AbuseLease(this, tracker, acquiredConcurrency);
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
private static bool CanAcquireAll(
|
private static bool CanAcquireAll(
|
||||||
TrackerState tracker,
|
TrackerState tracker,
|
||||||
@@ -415,6 +484,9 @@ internal sealed class AbuseProtectionService
|
|||||||
return "unknown";
|
return "unknown";
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static IPAddress NormalizeAddress(IPAddress address) =>
|
||||||
|
address.IsIPv4MappedToIPv6 ? address.MapToIPv4() : address;
|
||||||
|
|
||||||
private readonly record struct RateDimension(string Key, int Limit);
|
private readonly record struct RateDimension(string Key, int Limit);
|
||||||
|
|
||||||
private enum TrackerDomain
|
private enum TrackerDomain
|
||||||
|
|||||||
@@ -19,16 +19,71 @@ internal sealed class HttpAbuseProtectionMiddleware(
|
|||||||
string operation = context.GetEndpoint()?.Metadata.GetMetadata<IEndpointNameMetadata>()
|
string operation = context.GetEndpoint()?.Metadata.GetMetadata<IEndpointNameMetadata>()
|
||||||
?.EndpointName ?? "Unmatched";
|
?.EndpointName ?? "Unmatched";
|
||||||
bool healthEndpoint = operation is "GetLiveness" or "GetReadiness";
|
bool healthEndpoint = operation is "GetLiveness" or "GetReadiness";
|
||||||
bool acquired = healthEndpoint
|
bool operatorEndpoint = operation is
|
||||||
? protection.TryAcquireHealthIngress(
|
"GetOperatorStatus"
|
||||||
|
or "RevokeOperatorListing"
|
||||||
|
or "RevokeOperatorPrincipal"
|
||||||
|
or "RevokeOperatorSigningKey"
|
||||||
|
or "BeginOperatorDrain";
|
||||||
|
if (operatorEndpoint
|
||||||
|
&& !protection.IsOperatorSourceAllowed(context.Connection.RemoteIpAddress))
|
||||||
|
{
|
||||||
|
bool deniedSourceAdmitted = protection.TryAcquireHttpIngress(
|
||||||
context.Connection.RemoteIpAddress,
|
context.Connection.RemoteIpAddress,
|
||||||
out AbuseProtectionService.AbuseLease? lease,
|
"Unmatched",
|
||||||
out int retryAfterSeconds)
|
out AbuseProtectionService.AbuseLease? deniedSourceLease,
|
||||||
: protection.TryAcquireHttpIngress(
|
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,
|
context.Connection.RemoteIpAddress,
|
||||||
operation,
|
operation,
|
||||||
out lease,
|
out lease,
|
||||||
out retryAfterSeconds);
|
out retryAfterSeconds);
|
||||||
|
}
|
||||||
|
|
||||||
if (!acquired)
|
if (!acquired)
|
||||||
{
|
{
|
||||||
context.Response.Headers.RetryAfter = retryAfterSeconds.ToString(
|
context.Response.Headers.RetryAfter = retryAfterSeconds.ToString(
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
using FinalFactory.Rendezvous.Contracts;
|
using FinalFactory.Rendezvous.Contracts;
|
||||||
|
using FinalFactory.Rendezvous.Server.Observability;
|
||||||
using FinalFactory.Rendezvous.Server.Sessions;
|
using FinalFactory.Rendezvous.Server.Sessions;
|
||||||
using FinalFactory.Rendezvous.Server.State;
|
using FinalFactory.Rendezvous.Server.State;
|
||||||
|
|
||||||
@@ -15,6 +16,10 @@ internal sealed class ConnectionOutcomeMetrics
|
|||||||
{
|
{
|
||||||
private readonly object _gate = new();
|
private readonly object _gate = new();
|
||||||
private readonly Dictionary<(ConnectionOutcomeKind, ConnectionElapsedBucket), long> _counts = [];
|
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)
|
internal void Record(ConnectionOutcomeKind outcome, ConnectionElapsedBucket elapsedBucket)
|
||||||
{
|
{
|
||||||
@@ -24,6 +29,8 @@ internal sealed class ConnectionOutcomeMetrics
|
|||||||
_counts.TryGetValue(key, out long count);
|
_counts.TryGetValue(key, out long count);
|
||||||
_counts[key] = count + 1;
|
_counts[key] = count + 1;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
_telemetry?.RecordConnectionOutcome(outcome.ToString(), elapsedBucket.ToString());
|
||||||
}
|
}
|
||||||
|
|
||||||
internal long GetCount(ConnectionOutcomeKind outcome, ConnectionElapsedBucket elapsedBucket)
|
internal long GetCount(ConnectionOutcomeKind outcome, ConnectionElapsedBucket elapsedBucket)
|
||||||
|
|||||||
@@ -4,7 +4,8 @@ using Microsoft.AspNetCore.Diagnostics;
|
|||||||
|
|
||||||
namespace FinalFactory.Rendezvous.Server.Http;
|
namespace FinalFactory.Rendezvous.Server.Http;
|
||||||
|
|
||||||
internal sealed class RendezvousExceptionHandler : IExceptionHandler
|
internal sealed partial class RendezvousExceptionHandler(
|
||||||
|
ILogger<RendezvousExceptionHandler> logger) : IExceptionHandler
|
||||||
{
|
{
|
||||||
public async ValueTask<bool> TryHandleAsync(
|
public async ValueTask<bool> TryHandleAsync(
|
||||||
HttpContext httpContext,
|
HttpContext httpContext,
|
||||||
@@ -26,6 +27,13 @@ internal sealed class RendezvousExceptionHandler : IExceptionHandler
|
|||||||
: invalidRequest
|
: invalidRequest
|
||||||
? StatusCodes.Status400BadRequest
|
? StatusCodes.Status400BadRequest
|
||||||
: StatusCodes.Status500InternalServerError;
|
: 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(
|
await httpContext.Response.WriteAsJsonAsync(
|
||||||
new ApiError
|
new ApiError
|
||||||
{
|
{
|
||||||
@@ -42,4 +50,14 @@ internal sealed class RendezvousExceptionHandler : IExceptionHandler
|
|||||||
cancellationToken).ConfigureAwait(false);
|
cancellationToken).ConfigureAwait(false);
|
||||||
return true;
|
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.ConnectionOutcomes;
|
||||||
using FinalFactory.Rendezvous.Server.Http;
|
using FinalFactory.Rendezvous.Server.Http;
|
||||||
using FinalFactory.Rendezvous.Server.JoinAttempts;
|
using FinalFactory.Rendezvous.Server.JoinAttempts;
|
||||||
|
using FinalFactory.Rendezvous.Server.Observability;
|
||||||
|
using FinalFactory.Rendezvous.Server.Operations;
|
||||||
using FinalFactory.Rendezvous.Server.Provisioning;
|
using FinalFactory.Rendezvous.Server.Provisioning;
|
||||||
using FinalFactory.Rendezvous.Server.Sessions;
|
using FinalFactory.Rendezvous.Server.Sessions;
|
||||||
using FinalFactory.Rendezvous.Server.State;
|
using FinalFactory.Rendezvous.Server.State;
|
||||||
@@ -13,6 +15,9 @@ using Microsoft.AspNetCore.HttpOverrides;
|
|||||||
using Microsoft.OpenApi;
|
using Microsoft.OpenApi;
|
||||||
|
|
||||||
WebApplicationBuilder builder = WebApplication.CreateBuilder(args);
|
WebApplicationBuilder builder = WebApplication.CreateBuilder(args);
|
||||||
|
builder.Logging.AddFilter(
|
||||||
|
"Microsoft.AspNetCore.Diagnostics.ExceptionHandlerMiddleware",
|
||||||
|
LogLevel.None);
|
||||||
bool isOpenApiGeneration = string.Equals(
|
bool isOpenApiGeneration = string.Equals(
|
||||||
System.Reflection.Assembly.GetEntryAssembly()?.GetName().Name,
|
System.Reflection.Assembly.GetEntryAssembly()?.GetName().Name,
|
||||||
"GetDocument.Insider",
|
"GetDocument.Insider",
|
||||||
@@ -61,6 +66,14 @@ builder.Services.AddOpenApi("v1", static options =>
|
|||||||
In = ParameterLocation.Header,
|
In = ParameterLocation.Header,
|
||||||
Description = "Attempt-scoped client capability returned only to the joining caller.",
|
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)
|
HashSet<string> securedOperations = new(StringComparer.Ordinal)
|
||||||
{
|
{
|
||||||
@@ -69,8 +82,17 @@ builder.Services.AddOpenApi("v1", static options =>
|
|||||||
"UpdateSession",
|
"UpdateSession",
|
||||||
"DeleteSession",
|
"DeleteSession",
|
||||||
};
|
};
|
||||||
|
HashSet<string> operatorOperations = new(StringComparer.Ordinal)
|
||||||
|
{
|
||||||
|
"GetOperatorStatus",
|
||||||
|
"RevokeOperatorListing",
|
||||||
|
"RevokeOperatorPrincipal",
|
||||||
|
"RevokeOperatorSigningKey",
|
||||||
|
"BeginOperatorDrain",
|
||||||
|
};
|
||||||
OpenApiSecuritySchemeReference reference = new(schemeName, document, null);
|
OpenApiSecuritySchemeReference reference = new(schemeName, document, null);
|
||||||
OpenApiSecuritySchemeReference attemptReference = new(attemptSchemeName, document, null);
|
OpenApiSecuritySchemeReference attemptReference = new(attemptSchemeName, document, null);
|
||||||
|
OpenApiSecuritySchemeReference operatorReference = new(operatorSchemeName, document, null);
|
||||||
foreach (OpenApiPathItem path in document.Paths.Values)
|
foreach (OpenApiPathItem path in document.Paths.Values)
|
||||||
{
|
{
|
||||||
if (path.Operations is null)
|
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)
|
foreach (OpenApiOperation operation in path.Operations.Values)
|
||||||
{
|
{
|
||||||
if (operation.Responses is null
|
if (operation.Responses is null)
|
||||||
|| !operation.Responses.TryGetValue(
|
{
|
||||||
StatusCodes.Status429TooManyRequests.ToString(
|
continue;
|
||||||
System.Globalization.CultureInfo.InvariantCulture),
|
}
|
||||||
out IOpenApiResponse? response)
|
|
||||||
|| response is not OpenApiResponse concreteResponse)
|
foreach ((string status, IOpenApiResponse response) in operation.Responses)
|
||||||
|
{
|
||||||
|
if (response is not OpenApiResponse concreteResponse)
|
||||||
{
|
{
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
concreteResponse.Headers ??=
|
concreteResponse.Headers ??=
|
||||||
new Dictionary<string, IOpenApiHeader>(StringComparer.OrdinalIgnoreCase);
|
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
|
concreteResponse.Headers["Retry-After"] = new OpenApiHeader
|
||||||
{
|
{
|
||||||
Description = "Whole seconds before the caller should retry (1-60).",
|
Description = "Whole seconds before the caller should retry (1-60).",
|
||||||
@@ -124,6 +173,8 @@ builder.Services.AddOpenApi("v1", static options =>
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
return Task.CompletedTask;
|
return Task.CompletedTask;
|
||||||
});
|
});
|
||||||
@@ -166,8 +217,18 @@ builder.Services
|
|||||||
&& addresses.All(
|
&& addresses.All(
|
||||||
static value => IPAddress.TryParse(value, out _)),
|
static value => IPAddress.TryParse(value, out _)),
|
||||||
"Trusted proxy addresses must contain at most 32 literal IP addresses.")
|
"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();
|
.ValidateOnStart();
|
||||||
builder.Services.AddSingleton<AbuseProtectionService>();
|
builder.Services.AddSingleton<AbuseProtectionService>();
|
||||||
|
builder.Services
|
||||||
|
.AddOptions<AuditOptions>()
|
||||||
|
.BindConfiguration(AuditOptions.SectionName)
|
||||||
|
.ValidateDataAnnotations()
|
||||||
|
.ValidateOnStart();
|
||||||
AbuseProtectionOptions configuredAbuseProtection = builder.Configuration
|
AbuseProtectionOptions configuredAbuseProtection = builder.Configuration
|
||||||
.GetSection(AbuseProtectionOptions.SectionName)
|
.GetSection(AbuseProtectionOptions.SectionName)
|
||||||
.Get<AbuseProtectionOptions>() ?? new AbuseProtectionOptions();
|
.Get<AbuseProtectionOptions>() ?? new AbuseProtectionOptions();
|
||||||
@@ -180,8 +241,13 @@ InMemoryEphemeralRendezvousStore stateStore = new(
|
|||||||
stateOptions,
|
stateOptions,
|
||||||
rendezvousClock,
|
rendezvousClock,
|
||||||
rendezvousClock);
|
rendezvousClock);
|
||||||
|
builder.Services.AddSingleton(stateStore);
|
||||||
builder.Services.AddSingleton<IEphemeralRendezvousStore>(stateStore);
|
builder.Services.AddSingleton<IEphemeralRendezvousStore>(stateStore);
|
||||||
builder.Services.AddSingleton<IWallClock>(rendezvousClock);
|
builder.Services.AddSingleton<IWallClock>(rendezvousClock);
|
||||||
|
builder.Services.AddSingleton<IMonotonicClock>(rendezvousClock);
|
||||||
|
builder.Services.AddSingleton<RendezvousTelemetry>();
|
||||||
|
builder.Services.AddSingleton<AuditTrail>();
|
||||||
|
builder.Services.AddSingleton<RendezvousReadiness>();
|
||||||
|
|
||||||
if (isOpenApiGeneration)
|
if (isOpenApiGeneration)
|
||||||
{
|
{
|
||||||
@@ -214,6 +280,7 @@ else
|
|||||||
builder.Services.AddSingleton<JoinAttemptService>();
|
builder.Services.AddSingleton<JoinAttemptService>();
|
||||||
builder.Services.AddSingleton<ConnectionOutcomeMetrics>();
|
builder.Services.AddSingleton<ConnectionOutcomeMetrics>();
|
||||||
builder.Services.AddSingleton<ConnectionOutcomeService>();
|
builder.Services.AddSingleton<ConnectionOutcomeService>();
|
||||||
|
builder.Services.AddSingleton<OperatorService>();
|
||||||
builder.Services.AddSingleton(new ProvisioningReadiness(true));
|
builder.Services.AddSingleton(new ProvisioningReadiness(true));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -246,34 +313,13 @@ if (TrustedProxyForwarding.IsEnabled(configuredAbuseProtection))
|
|||||||
{
|
{
|
||||||
app.UseForwardedHeaders();
|
app.UseForwardedHeaders();
|
||||||
}
|
}
|
||||||
|
app.UseMiddleware<TelemetryMiddleware>();
|
||||||
app.UseExceptionHandler();
|
app.UseExceptionHandler();
|
||||||
app.UseMiddleware<HttpAbuseProtectionMiddleware>();
|
app.UseMiddleware<HttpAbuseProtectionMiddleware>();
|
||||||
app.MapOpenApi();
|
app.MapOpenApi();
|
||||||
app.MapRendezvousContractEndpoints();
|
app.MapRendezvousContractEndpoints();
|
||||||
app.MapGet(
|
app.MapOperatorEndpoints();
|
||||||
"/health/live",
|
app.MapRendezvousHealthEndpoints();
|
||||||
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");
|
|
||||||
|
|
||||||
await app.RunAsync();
|
await app.RunAsync();
|
||||||
|
|
||||||
|
|||||||
@@ -142,8 +142,36 @@ internal sealed class SigningKeyRing : IDisposable
|
|||||||
return VerificationKeyLookup.Available;
|
return VerificationKeyLookup.Available;
|
||||||
}
|
}
|
||||||
|
|
||||||
public bool Revoke(string keyId) =>
|
public bool Revoke(string keyId)
|
||||||
_keys.ContainsKey(keyId) && _runtimeRevocations.TryAdd(keyId, 0);
|
{
|
||||||
|
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()
|
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
|
internal sealed class SigningKey : IDisposable
|
||||||
{
|
{
|
||||||
private byte[]? _material;
|
private byte[]? _material;
|
||||||
|
|||||||
@@ -297,6 +297,7 @@ internal sealed record StoredJoinAttempt
|
|||||||
public required SecretFingerprint ConnectionTicketFingerprint { get; init; }
|
public required SecretFingerprint ConnectionTicketFingerprint { get; init; }
|
||||||
public NetworkEndpoint? DedicatedFallback { get; init; }
|
public NetworkEndpoint? DedicatedFallback { get; init; }
|
||||||
public required DateTimeOffset ExpiresAt { get; init; }
|
public required DateTimeOffset ExpiresAt { get; init; }
|
||||||
|
public required TimeSpan CreatedAtMonotonic { get; init; }
|
||||||
public required DateTimeOffset ConnectionTicketExpiresAt { get; init; }
|
public required DateTimeOffset ConnectionTicketExpiresAt { get; init; }
|
||||||
public AttemptEndpointBinding? HostEndpoint { get; init; }
|
public AttemptEndpointBinding? HostEndpoint { get; init; }
|
||||||
public AttemptEndpointBinding? ClientEndpoint { get; init; }
|
public AttemptEndpointBinding? ClientEndpoint { get; init; }
|
||||||
|
|||||||
@@ -23,7 +23,10 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
private TimeSpan? _drainDeadline;
|
private TimeSpan? _drainDeadline;
|
||||||
private TimeSpan _nextUdpMaintenance;
|
private TimeSpan _nextUdpMaintenance;
|
||||||
private long _maintenanceSweepCount;
|
private long _maintenanceSweepCount;
|
||||||
|
private long _expiryChurn;
|
||||||
private bool _available = true;
|
private bool _available = true;
|
||||||
|
private EphemeralStoreSnapshot? _metricsSnapshot;
|
||||||
|
private TimeSpan _metricsSnapshotAt = TimeSpan.MinValue;
|
||||||
|
|
||||||
public InMemoryEphemeralRendezvousStore(
|
public InMemoryEphemeralRendezvousStore(
|
||||||
EphemeralStoreOptions options,
|
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(
|
public StoreResult<StoredListing> CreateListing(
|
||||||
CreateListingCommand command,
|
CreateListingCommand command,
|
||||||
CancellationToken cancellationToken = default) => Atomic<StoredListing>(now =>
|
CancellationToken cancellationToken = default) => Atomic<StoredListing>(now =>
|
||||||
@@ -400,6 +447,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
|
|
||||||
AttemptEntry attempt = new(
|
AttemptEntry attempt = new(
|
||||||
command,
|
command,
|
||||||
|
now,
|
||||||
now + _options.JoinAttemptLifetime,
|
now + _options.JoinAttemptLifetime,
|
||||||
WallDeadline(now, _options.JoinAttemptLifetime));
|
WallDeadline(now, _options.JoinAttemptLifetime));
|
||||||
_attempts.Add(command.AttemptId, attempt);
|
_attempts.Add(command.AttemptId, attempt);
|
||||||
@@ -729,6 +777,10 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
return new(StoreResultCode.CapacityExceeded);
|
return new(StoreResultCode.CapacityExceeded);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
int activeResourcesBefore = _listings.Count
|
||||||
|
+ _presence.Count
|
||||||
|
+ _attempts.Count
|
||||||
|
+ _outcomeReports.Count;
|
||||||
_revocations[subject] = now + lifetime;
|
_revocations[subject] = now + lifetime;
|
||||||
SessionListingId[] listings = _listings
|
SessionListingId[] listings = _listings
|
||||||
.Where(item => string.Equals(item.Value.Definition.OwnerSubject, subject, StringComparison.Ordinal))
|
.Where(item => string.Equals(item.Value.Definition.OwnerSubject, subject, StringComparison.Ordinal))
|
||||||
@@ -756,7 +808,11 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
_outcomeReports.Remove(attemptId);
|
_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);
|
}, cancellationToken);
|
||||||
|
|
||||||
public void BeginDrain(CancellationToken cancellationToken = default)
|
public void BeginDrain(CancellationToken cancellationToken = default)
|
||||||
@@ -768,6 +824,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
if (!_drainDeadline.HasValue)
|
if (!_drainDeadline.HasValue)
|
||||||
{
|
{
|
||||||
_drainDeadline = _monotonicClock.Elapsed + _options.GracefulDrainLifetime;
|
_drainDeadline = _monotonicClock.Elapsed + _options.GracefulDrainLifetime;
|
||||||
|
_metricsSnapshot = null;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -778,6 +835,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
{
|
{
|
||||||
_available = false;
|
_available = false;
|
||||||
ClearActiveState();
|
ClearActiveState();
|
||||||
|
_metricsSnapshot = null;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -801,7 +859,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
_nextUdpMaintenance = now + UdpMaintenanceInterval;
|
_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();
|
ClearActiveState();
|
||||||
}
|
}
|
||||||
|
|
||||||
RemoveExpired(_revocations, now);
|
_expiryChurn += RemoveExpired(_revocations, now);
|
||||||
RemoveExpired(_replay, now);
|
_expiryChurn += RemoveExpired(_replay, now);
|
||||||
foreach (string key in _idempotency
|
string[] expiredIdempotency = _idempotency
|
||||||
.Where(item => item.Value.Deadline <= now)
|
.Where(item => item.Value.Deadline <= now)
|
||||||
.Select(static item => item.Key)
|
.Select(static item => item.Key)
|
||||||
.ToArray())
|
.ToArray();
|
||||||
|
_expiryChurn += expiredIdempotency.Length;
|
||||||
|
foreach (string key in expiredIdempotency)
|
||||||
{
|
{
|
||||||
_idempotency.Remove(key);
|
_idempotency.Remove(key);
|
||||||
}
|
}
|
||||||
|
|
||||||
foreach (MediationHandle handle in _presence
|
MediationHandle[] expiredPresence = _presence
|
||||||
.Where(item => item.Value.Deadline <= now)
|
.Where(item => item.Value.Deadline <= now)
|
||||||
.Select(static item => item.Key)
|
.Select(static item => item.Key)
|
||||||
.ToArray())
|
.ToArray();
|
||||||
|
_expiryChurn += expiredPresence.Length;
|
||||||
|
foreach (MediationHandle handle in expiredPresence)
|
||||||
{
|
{
|
||||||
_presence.Remove(handle);
|
_presence.Remove(handle);
|
||||||
}
|
}
|
||||||
|
|
||||||
foreach (JoinAttemptId attemptId in _attempts
|
JoinAttemptId[] expiredAttempts = _attempts
|
||||||
.Where(item => item.Value.Deadline <= now)
|
.Where(item => item.Value.Deadline <= now)
|
||||||
.Select(static item => item.Key)
|
.Select(static item => item.Key)
|
||||||
.ToArray())
|
.ToArray();
|
||||||
|
_expiryChurn += expiredAttempts.Length;
|
||||||
|
foreach (JoinAttemptId attemptId in expiredAttempts)
|
||||||
{
|
{
|
||||||
RemoveAttempt(attemptId);
|
RemoveAttempt(attemptId);
|
||||||
}
|
}
|
||||||
|
|
||||||
foreach (JoinAttemptId attemptId in _outcomeReports
|
JoinAttemptId[] expiredOutcomes = _outcomeReports
|
||||||
.Where(item => item.Value.Deadline <= now)
|
.Where(item => item.Value.Deadline <= now)
|
||||||
.Select(static item => item.Key)
|
.Select(static item => item.Key)
|
||||||
.ToArray())
|
.ToArray();
|
||||||
|
_expiryChurn += expiredOutcomes.Length;
|
||||||
|
foreach (JoinAttemptId attemptId in expiredOutcomes)
|
||||||
{
|
{
|
||||||
_outcomeReports.Remove(attemptId);
|
_outcomeReports.Remove(attemptId);
|
||||||
}
|
}
|
||||||
|
|
||||||
foreach (SessionListingId listingId in _listings
|
SessionListingId[] expiredListings = _listings
|
||||||
.Where(item => item.Value.LeaseDeadline <= now)
|
.Where(item => item.Value.LeaseDeadline <= now)
|
||||||
.Select(static item => item.Key)
|
.Select(static item => item.Key)
|
||||||
.ToArray())
|
.ToArray();
|
||||||
|
_expiryChurn += expiredListings.Length;
|
||||||
|
foreach (SessionListingId listingId in expiredListings)
|
||||||
{
|
{
|
||||||
RemoveListing(listingId);
|
RemoveListing(listingId);
|
||||||
}
|
}
|
||||||
@@ -959,6 +1029,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
ClientCapabilityFingerprint = entry.Command.ClientCapabilityFingerprint,
|
ClientCapabilityFingerprint = entry.Command.ClientCapabilityFingerprint,
|
||||||
ConnectionTicketFingerprint = entry.Command.ConnectionTicketFingerprint,
|
ConnectionTicketFingerprint = entry.Command.ConnectionTicketFingerprint,
|
||||||
DedicatedFallback = StoredListing.CopyEndpoint(entry.Command.DedicatedFallback),
|
DedicatedFallback = StoredListing.CopyEndpoint(entry.Command.DedicatedFallback),
|
||||||
|
CreatedAtMonotonic = entry.CreatedAtMonotonic,
|
||||||
ExpiresAt = entry.WallExpiresAt,
|
ExpiresAt = entry.WallExpiresAt,
|
||||||
ConnectionTicketExpiresAt = entry.TicketWallExpiresAt ?? default,
|
ConnectionTicketExpiresAt = entry.TicketWallExpiresAt ?? default,
|
||||||
HostEndpoint = entry.HostEndpoint,
|
HostEndpoint = entry.HostEndpoint,
|
||||||
@@ -968,15 +1039,18 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
IsCancelled = entry.IsCancelled,
|
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)
|
.Where(item => item.Value <= now)
|
||||||
.Select(static item => item.Key)
|
.Select(static item => item.Key)
|
||||||
.ToArray())
|
.ToArray();
|
||||||
|
foreach (string key in expired)
|
||||||
{
|
{
|
||||||
entries.Remove(key);
|
entries.Remove(key);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
return expired.Length;
|
||||||
}
|
}
|
||||||
|
|
||||||
private static void ValidateListing(ListingDefinition listing)
|
private static void ValidateListing(ListingDefinition listing)
|
||||||
@@ -1101,10 +1175,12 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
|
|
||||||
private sealed class AttemptEntry(
|
private sealed class AttemptEntry(
|
||||||
CreateJoinAttemptCommand command,
|
CreateJoinAttemptCommand command,
|
||||||
|
TimeSpan createdAtMonotonic,
|
||||||
TimeSpan deadline,
|
TimeSpan deadline,
|
||||||
DateTimeOffset wallExpiresAt)
|
DateTimeOffset wallExpiresAt)
|
||||||
{
|
{
|
||||||
public CreateJoinAttemptCommand Command { get; } = command;
|
public CreateJoinAttemptCommand Command { get; } = command;
|
||||||
|
public TimeSpan CreatedAtMonotonic { get; } = createdAtMonotonic;
|
||||||
public SecretFingerprint HostCapabilityFingerprint { get; } = command.HostCapabilityFingerprint;
|
public SecretFingerprint HostCapabilityFingerprint { get; } = command.HostCapabilityFingerprint;
|
||||||
public SecretFingerprint ClientCapabilityFingerprint { get; } = command.ClientCapabilityFingerprint;
|
public SecretFingerprint ClientCapabilityFingerprint { get; } = command.ClientCapabilityFingerprint;
|
||||||
public SecretFingerprint ConnectionTicketFingerprint { get; } = command.ConnectionTicketFingerprint;
|
public SecretFingerprint ConnectionTicketFingerprint { get; } = command.ConnectionTicketFingerprint;
|
||||||
@@ -1137,3 +1213,16 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
|
|||||||
object ResourceId,
|
object ResourceId,
|
||||||
TimeSpan Deadline);
|
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;
|
||||||
using System.Net.Sockets;
|
using System.Net.Sockets;
|
||||||
using FinalFactory.Rendezvous.Contracts;
|
using FinalFactory.Rendezvous.Contracts;
|
||||||
using FinalFactory.Rendezvous.Server.Abuse;
|
using FinalFactory.Rendezvous.Server.Abuse;
|
||||||
using FinalFactory.Rendezvous.Server.JoinAttempts;
|
using FinalFactory.Rendezvous.Server.JoinAttempts;
|
||||||
|
using FinalFactory.Rendezvous.Server.Observability;
|
||||||
using FinalFactory.Rendezvous.Server.Sessions;
|
using FinalFactory.Rendezvous.Server.Sessions;
|
||||||
using FinalFactory.Rendezvous.Server.State;
|
using FinalFactory.Rendezvous.Server.State;
|
||||||
|
|
||||||
@@ -38,7 +40,9 @@ internal sealed class NatMediationProcessor(
|
|||||||
IEphemeralRendezvousStore store,
|
IEphemeralRendezvousStore store,
|
||||||
ISessionCapabilityService capabilities,
|
ISessionCapabilityService capabilities,
|
||||||
JoinAttemptService joinAttempts,
|
JoinAttemptService joinAttempts,
|
||||||
AbuseProtectionService? abuseProtection = null)
|
AbuseProtectionService? abuseProtection = null,
|
||||||
|
RendezvousTelemetry? telemetry = null,
|
||||||
|
IMonotonicClock? monotonicClock = null)
|
||||||
{
|
{
|
||||||
public NatMediationResult ProcessDatagram(
|
public NatMediationResult ProcessDatagram(
|
||||||
ReadOnlySpan<byte> encoded,
|
ReadOnlySpan<byte> encoded,
|
||||||
@@ -68,6 +72,26 @@ internal sealed class NatMediationProcessor(
|
|||||||
INatIntroductionSink introductionSink,
|
INatIntroductionSink introductionSink,
|
||||||
CancellationToken cancellationToken = default)
|
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 _)
|
if (!RendezvousUdpCodec.TryDecode(encoded, out PresenceDatagram? datagram, out _)
|
||||||
|| datagram is null
|
|| datagram is null
|
||||||
@@ -141,12 +165,22 @@ internal sealed class NatMediationProcessor(
|
|||||||
IPEndPoint observedPublicEndpoint,
|
IPEndPoint observedPublicEndpoint,
|
||||||
string token,
|
string token,
|
||||||
INatIntroductionSink introductionSink,
|
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,
|
claimedLocalEndpoint,
|
||||||
observedPublicEndpoint,
|
observedPublicEndpoint,
|
||||||
token,
|
token,
|
||||||
introductionSink,
|
introductionSink,
|
||||||
cancellationToken);
|
cancellationToken);
|
||||||
|
telemetry?.RecordUdp(
|
||||||
|
"litenet",
|
||||||
|
result.ToString(),
|
||||||
|
Stopwatch.GetElapsedTime(started).TotalMilliseconds);
|
||||||
|
return result;
|
||||||
|
}
|
||||||
|
|
||||||
private NatMediationResult ProcessRequestCore(
|
private NatMediationResult ProcessRequestCore(
|
||||||
IPEndPoint claimedLocalEndpoint,
|
IPEndPoint claimedLocalEndpoint,
|
||||||
@@ -255,6 +289,14 @@ internal sealed class NatMediationProcessor(
|
|||||||
try
|
try
|
||||||
{
|
{
|
||||||
introductionSink.Introduce(CreatePlan(consumed.Value, ticket.Value));
|
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;
|
return NatMediationResult.Introduced;
|
||||||
}
|
}
|
||||||
catch (Exception exception) when (exception is SocketException
|
catch (Exception exception) when (exception is SocketException
|
||||||
|
|||||||
@@ -1,5 +1,8 @@
|
|||||||
{
|
{
|
||||||
"Rendezvous": {
|
"Rendezvous": {
|
||||||
|
"AbuseProtection": {
|
||||||
|
"OperatorAllowedAddresses": ["127.0.0.1", "::1"]
|
||||||
|
},
|
||||||
"Provisioning": {
|
"Provisioning": {
|
||||||
"Issuer": "final-factory-rendezvous-development",
|
"Issuer": "final-factory-rendezvous-development",
|
||||||
"Audience": "final-factory-rendezvous",
|
"Audience": "final-factory-rendezvous",
|
||||||
@@ -14,6 +17,14 @@
|
|||||||
"NotBefore": "2025-01-01T00:00:00Z",
|
"NotBefore": "2025-01-01T00:00:00Z",
|
||||||
"SignUntil": "2035-01-01T00:00:00Z",
|
"SignUntil": "2035-01-01T00:00:00Z",
|
||||||
"VerifyUntil": "2035-01-02T00: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": [
|
"Games": [
|
||||||
|
|||||||
@@ -6,16 +6,25 @@
|
|||||||
"MaxDatagramsPerPoll": 256,
|
"MaxDatagramsPerPoll": 256,
|
||||||
"PollIntervalMilliseconds": 2
|
"PollIntervalMilliseconds": 2
|
||||||
},
|
},
|
||||||
|
"Audit": {
|
||||||
|
"MaxEntries": 10000,
|
||||||
|
"RetentionDays": 30
|
||||||
|
},
|
||||||
"AbuseProtection": {
|
"AbuseProtection": {
|
||||||
"WindowSeconds": 1,
|
"WindowSeconds": 1,
|
||||||
"MaxTrackedKeys": 100000,
|
"MaxTrackedKeys": 100000,
|
||||||
"CriticalTrackedKeyReserve": 2048,
|
"CriticalTrackedKeyReserve": 2048,
|
||||||
"UdpTrackedKeyLimit": 70000,
|
"UdpTrackedKeyLimit": 70000,
|
||||||
"TrustedProxyAddresses": [],
|
"TrustedProxyAddresses": [],
|
||||||
|
"OperatorAllowedAddresses": [],
|
||||||
"HealthGlobalRequestsPerWindow": 1000,
|
"HealthGlobalRequestsPerWindow": 1000,
|
||||||
"HealthGlobalConcurrency": 32,
|
"HealthGlobalConcurrency": 32,
|
||||||
"HealthIpPrefixRequestsPerWindow": 120,
|
"HealthIpPrefixRequestsPerWindow": 120,
|
||||||
"HealthIpPrefixConcurrency": 8,
|
"HealthIpPrefixConcurrency": 8,
|
||||||
|
"OperatorGlobalRequestsPerWindow": 1000,
|
||||||
|
"OperatorGlobalConcurrency": 32,
|
||||||
|
"OperatorIpPrefixRequestsPerWindow": 120,
|
||||||
|
"OperatorIpPrefixConcurrency": 8,
|
||||||
"HttpGlobalRequestsPerWindow": 20000,
|
"HttpGlobalRequestsPerWindow": 20000,
|
||||||
"HttpOptionalRequestsPerWindow": 18000,
|
"HttpOptionalRequestsPerWindow": 18000,
|
||||||
"HttpIpPrefixRequestsPerWindow": 500,
|
"HttpIpPrefixRequestsPerWindow": 500,
|
||||||
|
|||||||
@@ -11,6 +11,11 @@ public sealed class OpenApiCompatibilityTests
|
|||||||
"/v1/join-attempts",
|
"/v1/join-attempts",
|
||||||
"/v1/join-attempts/{attemptId}",
|
"/v1/join-attempts/{attemptId}",
|
||||||
"/v1/join-attempts/{attemptId}/outcome",
|
"/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",
|
||||||
"/v1/sessions/{listingId}",
|
"/v1/sessions/{listingId}",
|
||||||
"/v1/sessions/{listingId}/join-attempts",
|
"/v1/sessions/{listingId}/join-attempts",
|
||||||
@@ -104,6 +109,11 @@ public sealed class OpenApiCompatibilityTests
|
|||||||
Assert.Equal(
|
Assert.Equal(
|
||||||
"X-Rendezvous-Client-Punch-Capability",
|
"X-Rendezvous-Client-Punch-Capability",
|
||||||
attemptCapability.GetProperty("name").GetString());
|
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 =
|
(string Path, string Method)[] publisherOperations =
|
||||||
[
|
[
|
||||||
("/v1/sessions", "post"),
|
("/v1/sessions", "post"),
|
||||||
@@ -120,6 +130,23 @@ public sealed class OpenApiCompatibilityTests
|
|||||||
Assert.True(security[0].TryGetProperty("PublisherBearer", out _));
|
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")
|
JsonElement cancelParameters = root.GetProperty("paths")
|
||||||
.GetProperty("/v1/join-attempts/{attemptId}")
|
.GetProperty("/v1/join-attempts/{attemptId}")
|
||||||
.GetProperty("delete")
|
.GetProperty("delete")
|
||||||
@@ -157,6 +184,15 @@ public sealed class OpenApiCompatibilityTests
|
|||||||
static item => item.Name is "get" or "post" or "put" or "delete"))
|
static item => item.Name is "get" or "post" or "put" or "delete"))
|
||||||
{
|
{
|
||||||
JsonElement responses = operation.Value.GetProperty("responses");
|
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))
|
if (!responses.TryGetProperty("429", out JsonElement overloaded))
|
||||||
{
|
{
|
||||||
continue;
|
continue;
|
||||||
@@ -171,7 +207,7 @@ public sealed class OpenApiCompatibilityTests
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
Assert.Equal(12, overloadContracts);
|
Assert.Equal(17, overloadContracts);
|
||||||
(string Path, string Method)[] bodyOperations =
|
(string Path, string Method)[] bodyOperations =
|
||||||
[
|
[
|
||||||
("/v1/sessions", "post"),
|
("/v1/sessions", "post"),
|
||||||
@@ -180,6 +216,10 @@ public sealed class OpenApiCompatibilityTests
|
|||||||
("/v1/sessions/{listingId}", "delete"),
|
("/v1/sessions/{listingId}", "delete"),
|
||||||
("/v1/join-attempts", "post"),
|
("/v1/join-attempts", "post"),
|
||||||
("/v1/join-attempts/{attemptId}/outcome", "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)
|
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,
|
CredentialValidationError.SignatureInvalid,
|
||||||
service.Validate(tampered, ProvisioningTestData.Now).Error);
|
service.Validate(tampered, ProvisioningTestData.Now).Error);
|
||||||
Assert.True(keys.Revoke("key-1"));
|
Assert.True(keys.Revoke("key-1"));
|
||||||
|
Assert.True(keys.Revoke("key-1"));
|
||||||
Assert.Equal(
|
Assert.Equal(
|
||||||
CredentialValidationError.KeyRevoked,
|
CredentialValidationError.KeyRevoked,
|
||||||
service.Validate(token, ProvisioningTestData.Now).Error);
|
service.Validate(token, ProvisioningTestData.Now).Error);
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ using FinalFactory.Rendezvous.Server.Http;
|
|||||||
using Microsoft.AspNetCore.Builder;
|
using Microsoft.AspNetCore.Builder;
|
||||||
using Microsoft.AspNetCore.Http;
|
using Microsoft.AspNetCore.Http;
|
||||||
using Microsoft.AspNetCore.HttpOverrides;
|
using Microsoft.AspNetCore.HttpOverrides;
|
||||||
|
using Microsoft.AspNetCore.Routing;
|
||||||
using Microsoft.Extensions.Logging.Abstractions;
|
using Microsoft.Extensions.Logging.Abstractions;
|
||||||
using Microsoft.Extensions.Options;
|
using Microsoft.Extensions.Options;
|
||||||
|
|
||||||
@@ -86,6 +87,9 @@ public sealed class AbuseProtectionTests
|
|||||||
AbuseProtectionService protection = new(Options.Create(options));
|
AbuseProtectionService protection = new(Options.Create(options));
|
||||||
IPAddress source = IPAddress.Parse("198.51.100.10");
|
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");
|
||||||
AssertAccepted(protection, source, "BrowseSessions");
|
AssertAccepted(protection, source, "BrowseSessions");
|
||||||
AssertRejected(protection, source, "BrowseSessions");
|
AssertRejected(protection, source, "BrowseSessions");
|
||||||
@@ -93,6 +97,57 @@ public sealed class AbuseProtectionTests
|
|||||||
AssertRejected(protection, source, "RenewSessionLease");
|
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]
|
[Fact]
|
||||||
public void ResourceBudgetsRemainIsolatedAcrossTenantAndPrincipalScopes()
|
public void ResourceBudgetsRemainIsolatedAcrossTenantAndPrincipalScopes()
|
||||||
{
|
{
|
||||||
@@ -354,7 +409,8 @@ public sealed class AbuseProtectionTests
|
|||||||
{
|
{
|
||||||
const string canary = "credential-canary <script> endpoint=203.0.113.8:9000";
|
const string canary = "credential-canary <script> endpoint=203.0.113.8:9000";
|
||||||
DefaultHttpContext context = Context("198.51.100.10");
|
DefaultHttpContext context = Context("198.51.100.10");
|
||||||
RendezvousExceptionHandler handler = new();
|
RendezvousExceptionHandler handler = new(
|
||||||
|
NullLogger<RendezvousExceptionHandler>.Instance);
|
||||||
|
|
||||||
Assert.True(await handler.TryHandleAsync(
|
Assert.True(await handler.TryHandleAsync(
|
||||||
context,
|
context,
|
||||||
@@ -438,6 +494,11 @@ public sealed class AbuseProtectionTests
|
|||||||
HealthGlobalConcurrency = 10_000,
|
HealthGlobalConcurrency = 10_000,
|
||||||
HealthIpPrefixRequestsPerWindow = 10_000,
|
HealthIpPrefixRequestsPerWindow = 10_000,
|
||||||
HealthIpPrefixConcurrency = 1_000,
|
HealthIpPrefixConcurrency = 1_000,
|
||||||
|
OperatorAllowedAddresses = ["198.51.100.10"],
|
||||||
|
OperatorGlobalRequestsPerWindow = 10_000,
|
||||||
|
OperatorGlobalConcurrency = 10_000,
|
||||||
|
OperatorIpPrefixRequestsPerWindow = 10_000,
|
||||||
|
OperatorIpPrefixConcurrency = 1_000,
|
||||||
HttpGlobalRequestsPerWindow = 10_000,
|
HttpGlobalRequestsPerWindow = 10_000,
|
||||||
HttpOptionalRequestsPerWindow = 9_000,
|
HttpOptionalRequestsPerWindow = 9_000,
|
||||||
HttpIpPrefixRequestsPerWindow = 10_000,
|
HttpIpPrefixRequestsPerWindow = 10_000,
|
||||||
|
|||||||
@@ -436,7 +436,7 @@ public sealed class InMemoryEphemeralRendezvousStoreTests
|
|||||||
StoreResult<int> revoked = fixture.Store.RevokePrincipal(command.Listing.OwnerSubject, TimeSpan.FromMinutes(1));
|
StoreResult<int> revoked = fixture.Store.RevokePrincipal(command.Listing.OwnerSubject, TimeSpan.FromMinutes(1));
|
||||||
|
|
||||||
Assert.True(revoked.Succeeded);
|
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.GetListing(listing.Definition.ListingId, false).Code);
|
||||||
Assert.Equal(StoreResultCode.NotFound, fixture.Store.BindAttemptEndpoint(new(
|
Assert.Equal(StoreResultCode.NotFound, fixture.Store.BindAttemptEndpoint(new(
|
||||||
attempt.MediationHandle,
|
attempt.MediationHandle,
|
||||||
|
|||||||
Reference in New Issue
Block a user