Compare commits

...

5 Commits

Author SHA1 Message Date
KyuubiYoru 07004cd75f docs(evidence): record clean capacity candidate (#18)
quality-gate / quality (push) Failing after 1m28s
quality-gate / container (push) Has been skipped
2026-07-16 16:11:39 +02:00
KyuubiYoru cf14836d48 fix(operations): require clean candidate provenance (#18) 2026-07-16 16:04:25 +02:00
KyuubiYoru 609dad7cf1 feat(operations): add capacity and resilience gates (#18) 2026-07-16 15:57:01 +02:00
KyuubiYoru 08729ae25c feat(deployment): add secure Linux runtime (#17)
quality-gate / quality (push) Failing after 1m9s
quality-gate / container (push) Has been skipped
2026-07-16 15:03:04 +02:00
KyuubiYoru be732de7c9 feat(server): add observability and operator controls (#16)
quality-gate / quality (push) Failing after 1m1s
2026-07-16 13:22:16 +02:00
63 changed files with 7943 additions and 218 deletions
+13
View File
@@ -0,0 +1,13 @@
.git
.gitea
.idea
.vs
.codex
.agents
**/bin
**/obj
TestResults
deploy/compose/secrets
deploy/compose/.smoke.env
docs
tests
+60
View File
@@ -36,6 +36,9 @@ jobs:
- name: Test
run: dotnet test Rendezvous.slnx --configuration Release --no-build
- name: Run quick capacity and resilience gate
run: ./scripts/run-capacity-gate.sh
- name: Test privileged Linux namespace topology when available
shell: bash
run: |
@@ -88,3 +91,60 @@ jobs:
else
echo "Network namespaces/NAT tooling unavailable; deterministic loopback topology remains the required gate."
fi
container:
needs: quality
runs-on: ubuntu-latest
timeout-minutes: 15
steps:
- name: Check out repository
uses: actions/checkout@v4
- name: Install .NET SDK
uses: actions/setup-dotnet@v4
with:
dotnet-version: 10.0.301
- name: Build deployment diagnostic
run: |
dotnet restore src/FinalFactory.Rendezvous.TestClient/FinalFactory.Rendezvous.TestClient.csproj --locked-mode
dotnet build src/FinalFactory.Rendezvous.TestClient/FinalFactory.Rendezvous.TestClient.csproj --configuration Release --no-restore
- name: Build and exercise hardened container
shell: bash
run: |
set -euo pipefail
compose_file="deploy/compose/compose.yaml"
secret="deploy/compose/secrets/signing-key"
cleanup() {
RENDEZVOUS_UID=1654 RENDEZVOUS_GID=1654 \
docker compose -f "$compose_file" down --volumes >/dev/null 2>&1 || true
rm -f "$secret"
}
trap cleanup EXIT
install -d -m 0700 deploy/compose/secrets
openssl rand -out "$secret" 32
chmod 0444 "$secret"
export RENDEZVOUS_UID=1654
export RENDEZVOUS_GID=1654
docker compose -f "$compose_file" up --build --detach
container_id="$(docker compose -f "$compose_file" ps -q rendezvous)"
test -n "$container_id"
test "$(docker inspect --format '{{.Config.User}}' "$container_id")" = "1654:1654"
test "$(docker inspect --format '{{.HostConfig.ReadonlyRootfs}}' "$container_id")" = "true"
test "$(docker inspect --format '{{range .Mounts}}{{if eq .Destination \"/app/appsettings.Production.json\"}}{{.RW}}{{end}}{{end}}' "$container_id")" = "false"
test "$(docker inspect --format '{{range .Mounts}}{{if eq .Destination \"/run/secrets/rendezvous-signing-key\"}}{{.RW}}{{end}}{{end}}' "$container_id")" = "false"
for attempt in {1..100}; do
if curl --fail --silent http://127.0.0.1:8080/health/ready >/dev/null 2>&1; then
break
fi
if (( attempt == 100 )); then
docker compose -f "$compose_file" logs rendezvous
exit 1
fi
sleep 0.1
done
./scripts/smoke-deployment.sh
docker compose -f "$compose_file" stop --timeout 40 rendezvous
test "$(docker inspect --format '{{.State.Running}}' "$container_id")" = "false"
test "$(docker inspect --format '{{.State.ExitCode}}' "$container_id")" = "0"
+4
View File
@@ -6,3 +6,7 @@ TestResults/
*.suo
*.user
*.userosscache
deploy/compose/.smoke.env
artifacts/
deploy/compose/secrets/*
!deploy/compose/secrets/.gitignore
+30
View File
@@ -0,0 +1,30 @@
# syntax=docker/dockerfile:1.7@sha256:a57df69d0ea827fb7266491f2813635de6f17269be881f696fbfdf2d83dda33e
FROM mcr.microsoft.com/dotnet/sdk:10.0.301-noble@sha256:ea8bde36c11b6e7eec2656d0e59101d4462f6bd630730f2c8201ed0572b295d5 AS build
WORKDIR /source
COPY Directory.Build.props Directory.Packages.props NuGet.config global.json Rendezvous.slnx ./
COPY src/FinalFactory.Rendezvous.Contracts/FinalFactory.Rendezvous.Contracts.csproj src/FinalFactory.Rendezvous.Contracts/packages.lock.json src/FinalFactory.Rendezvous.Contracts/
COPY src/FinalFactory.Rendezvous.Server/FinalFactory.Rendezvous.Server.csproj src/FinalFactory.Rendezvous.Server/packages.lock.json src/FinalFactory.Rendezvous.Server/
RUN dotnet restore src/FinalFactory.Rendezvous.Server/FinalFactory.Rendezvous.Server.csproj --locked-mode
COPY src/FinalFactory.Rendezvous.Contracts/ src/FinalFactory.Rendezvous.Contracts/
COPY src/FinalFactory.Rendezvous.Server/ src/FinalFactory.Rendezvous.Server/
RUN dotnet publish src/FinalFactory.Rendezvous.Server/FinalFactory.Rendezvous.Server.csproj \
--configuration Release \
--no-restore \
--output /out \
/p:UseAppHost=false \
/p:OpenApiGenerateDocuments=false
FROM mcr.microsoft.com/dotnet/aspnet:10.0.9-noble-chiseled@sha256:f820c4fbfb8bb204c3bbe05c69d48cd039cd0e67aa8f13ac1cec168819b90643 AS runtime
ENV ASPNETCORE_HTTP_PORTS=8080 \
DOTNET_EnableDiagnostics=0 \
DOTNET_CLI_TELEMETRY_OPTOUT=1 \
TMPDIR=/tmp
WORKDIR /app
COPY --from=build --chown=1654:1654 /out/ ./
USER 1654:1654
EXPOSE 8080/tcp
EXPOSE 9050/udp
ENTRYPOINT ["dotnet", "FinalFactory.Rendezvous.Server.dll"]
+24 -7
View File
@@ -38,7 +38,11 @@ UDP hole punching cannot guarantee a direct connection through every network. Sy
- `FinalFactory.Rendezvous.TestClient` — thin interactive and scriptable host/browser/join diagnostic built only on the public SDK.
- `FinalFactory.Rendezvous.Tests` — unit, integration, security, and connection-lifecycle tests.
The server directory and NAT mediator begin as separate modules in one deployable service because they share session, lease, authorization, and endpoint state. Their internal boundary should allow independent deployment later if scale, availability, or security requirements diverge.
The server directory and NAT mediator are separate modules in one single-active
deployable service because they share ephemeral session, lease, authorization,
replay, and endpoint state. Their internal boundary can support a future
explicitly designed shared-state architecture; operators must not create
multiple active v1 replicas.
## Service boundaries
@@ -77,9 +81,11 @@ The initial service does not provide:
Rendezvous is under active roadmap development. The versioned contracts,
directory leases, authenticated join attempts, LiteNetLib mediator, caller-owned
SDK coordination, typed connection outcomes, and thin public-SDK diagnostic client
are implemented. Deployment hardening, the broader NAT-topology harness, and
the production-readiness roadmap remain in progress;
SDK coordination, typed connection outcomes, thin public-SDK diagnostic client,
deterministic NAT topology harness, hostile-input controls,
observability/operator surface, secure single-active Linux deployment, and
numeric capacity/resilience gates are implemented. Packaging, consumer pilots,
and final production-readiness gates remain in progress;
participating games must not treat the current repository as a finished production
service until those gates land.
@@ -91,6 +97,15 @@ Tenant policy, publisher/operator principals, and production key custody are
defined in [game provisioning and signing-key lifecycle](docs/security/provisioning.md).
Layered HTTP/UDP budgets, overload behavior, and safe operational tuning are
defined in [hostile-input and overload protection](docs/security/abuse-protection.md).
Health semantics, bounded telemetry, alerting, audit privacy, and the authenticated
operator controls are defined in the
[observability and operator runbook](docs/operations/observability-and-operator-runbook.md).
The pinned non-root container, production topology, graceful drain, Linux
hardening, smoke procedure, and recovery lifecycle are documented in
[secure single-active Linux deployment](docs/deployment/linux.md).
The numeric core-state candidate profile, public launch objectives, accelerated
soak, resilience matrix, and single-active scaling decision are recorded in
[capacity and resilience gates](docs/operations/capacity-and-resilience.md).
The scriptable host/browser/join diagnostic and its stable automation contract are
documented in the [TestClient integration guide](docs/integration/test-client.md).
The always-on three-party scenarios, optional Linux namespace topology, and
@@ -112,9 +127,11 @@ dotnet test Rendezvous.slnx --configuration Release --no-build
Run the bootstrap server with
`dotnet run --project src/FinalFactory.Rendezvous.Server`. It serves HTTP health endpoints and binds
the configured UDP mediator port; both stop through normal host cancellation.
The launch profile uses an ephemeral development-only signing key. Production
startup fails closed until externally supplied game policies and `env:` signing
key references resolve to valid key material; no reusable game secret is stored
The launch profile uses separate ephemeral development-only publisher and operator
signing keys. Production
startup fails closed until its advertised endpoints, proxy trust boundary,
externally supplied game policies, and `env:` (base64) or `file:` (raw,
absolute, non-symlink) signing-key references resolve safely; no reusable game secret is stored
in this repository or the public Client package.
The project dependency rules and supported runtime choices are documented in
[project and dependency boundaries](docs/architecture/project-boundaries.md).
+1
View File
@@ -6,6 +6,7 @@
<Project Path="src/FinalFactory.Rendezvous.TestClient/FinalFactory.Rendezvous.TestClient.csproj" />
</Folder>
<Folder Name="/tests/">
<Project Path="tests/FinalFactory.Rendezvous.Capacity/FinalFactory.Rendezvous.Capacity.csproj" />
<Project Path="tests/FinalFactory.Rendezvous.Tests/FinalFactory.Rendezvous.Tests.csproj" />
</Folder>
</Solution>
@@ -0,0 +1,58 @@
{
"AllowedHosts": "localhost;127.0.0.1",
"Rendezvous": {
"Deployment": {
"PublicHttpBaseUrl": "https://localhost/",
"PublicUdpHost": "127.0.0.1",
"PublicUdpPort": 9050,
"DrainDeadlineSeconds": 30,
"MinimumDrainSeconds": 1,
"SingleActiveInstance": true,
"AllowPrivatePublicEndpoints": true
},
"Udp": {
"ListenAddress": "0.0.0.0",
"Port": 9050
},
"AbuseProtection": {
"TrustedProxyAddresses": ["127.0.0.1"],
"OperatorAllowedAddresses": ["127.0.0.1"]
},
"Provisioning": {
"Issuer": "final-factory-rendezvous-smoke",
"Audience": "rendezvous-service",
"ClockSkewSeconds": 30,
"SigningKeys": [
{
"KeyId": "local-smoke-1",
"SecretReference": "file:/run/secrets/rendezvous-signing-key",
"CredentialKinds": ["DedicatedPublisher"],
"GameId": "space-game",
"EnvironmentId": "smoke",
"NotBefore": "2026-01-01T00:00:00Z",
"SignUntil": "2100-01-01T00:00:00Z",
"VerifyUntil": "2100-01-02T00:00:00Z"
}
],
"Games": [
{
"GameId": "space-game",
"EnvironmentId": "smoke",
"Enabled": true,
"ProtocolVersions": [1],
"Regions": ["local"],
"VisibilityModes": ["Public"],
"PublisherTrustModes": ["ManagedDedicated"],
"MetadataValueMaxBytes": {},
"RequiredMetadataKeys": [],
"MetadataMaxBytes": 512,
"MetadataMaxKeys": 0,
"MaxListingsPerPrincipal": 10,
"MaxAnonymousListingsPerAddress": 0,
"MaxActiveJoinAttempts": 100,
"FallbackPolicy": "Disabled"
}
]
}
}
}
+35
View File
@@ -0,0 +1,35 @@
name: rendezvous-local
services:
rendezvous:
image: finalfactory/rendezvous:local
build:
context: ../..
dockerfile: Dockerfile
init: true
user: "${RENDEZVOUS_UID:?set RENDEZVOUS_UID to a non-root host UID}:${RENDEZVOUS_GID:?set RENDEZVOUS_GID to its GID}"
read_only: true
tmpfs:
- /tmp:rw,noexec,nosuid,nodev,size=16m,uid=${RENDEZVOUS_UID},gid=${RENDEZVOUS_GID},mode=0700
cap_drop:
- ALL
security_opt:
- no-new-privileges:true
pids_limit: 128
mem_limit: 512m
cpus: 1.0
ulimits:
nofile:
soft: 4096
hard: 4096
stop_grace_period: 40s
restart: unless-stopped
environment:
ASPNETCORE_ENVIRONMENT: Production
ASPNETCORE_HTTP_PORTS: "8080"
volumes:
- ./appsettings.Production.json:/app/appsettings.Production.json:ro
- ./secrets/signing-key:/run/secrets/rendezvous-signing-key:ro
ports:
- "127.0.0.1:8080:8080/tcp"
- "9050:9050/udp"
+2
View File
@@ -0,0 +1,2 @@
*
!.gitignore
+48
View File
@@ -0,0 +1,48 @@
[Unit]
Description=Final Factory Rendezvous service
Documentation=https://git.finalfactory.de/HeiKyu/Rendezvous
After=network-online.target time-sync.target
Wants=network-online.target time-sync.target
[Service]
Type=simple
User=rendezvous
Group=rendezvous
WorkingDirectory=/opt/rendezvous
ExecStart=/usr/bin/dotnet /opt/rendezvous/FinalFactory.Rendezvous.Server.dll
Environment=ASPNETCORE_ENVIRONMENT=Production
Environment=ASPNETCORE_HTTP_PORTS=8080
Environment=DOTNET_EnableDiagnostics=0
EnvironmentFile=-/etc/rendezvous/rendezvous.env
Restart=on-failure
RestartSec=5s
KillSignal=SIGTERM
KillMode=mixed
TimeoutStopSec=40s
NoNewPrivileges=true
PrivateDevices=true
PrivateTmp=true
ProtectClock=true
ProtectControlGroups=true
ProtectHome=true
ProtectHostname=true
ProtectKernelLogs=true
ProtectKernelModules=true
ProtectKernelTunables=true
ProtectSystem=strict
RestrictAddressFamilies=AF_UNIX AF_INET AF_INET6
RestrictNamespaces=true
RestrictRealtime=true
RestrictSUIDSGID=true
CapabilityBoundingSet=
AmbientCapabilities=
LockPersonality=true
SystemCallArchitectures=native
UMask=0077
LimitNOFILE=4096
CPUQuota=200%
MemoryMax=2G
TasksMax=128
[Install]
WantedBy=multi-user.target
File diff suppressed because it is too large Load Diff
@@ -129,9 +129,18 @@ until the owner records:
- which games may enable anonymous unlisted player hosting;
- deployment regions, data-processing jurisdiction, and approval of the stated
30-day audit/13-month aggregate retention periods;
- the per-game dedicated fallback endpoint policy;
- the measured supported profile and whether the 99.5% single-active objective
is sufficient or shared-state/high-availability work must be brought forward.
- the per-game dedicated fallback endpoint policy.
Issue #18 measured and ratified the original 2-vCPU/2-GiB, 25,000-listing,
10,000-attempt core-state candidate profile and retained the 99.5% single-active
topology. It does not claim that core measurements prove public HTTP/UDP SLOs.
The versioned evidence, RTO, failure domains, and explicit signals that trigger
shared-state/high-availability work are recorded in the
[capacity and resilience gate](../operations/capacity-and-resilience.md). The
real-network canary in #23 must confirm that the proposed regional launch load
fits this profile and validate the public SLOs; it may lower the launch cap but
may not silently enable a
second active instance.
These are configuration and launch decisions, not permission to weaken the
tenant, replay, endpoint-verification, or secret-handling controls.
+228
View File
@@ -0,0 +1,228 @@
# Secure single-active Linux deployment
Tracking: #17
Rendezvous v1 stores listings, observed endpoints, join attempts, replay markers,
and runtime revocations only in the process that accepted them. Deploy exactly
one active instance. A second live replica would have a different directory and
replay boundary; `SingleActiveInstance=false` is therefore rejected rather than
presented as high availability.
## Pinned container
The root `Dockerfile` uses a multi-stage .NET 10 build and pins both Microsoft
base images by multi-architecture manifest digest. The runtime is the chiseled
ASP.NET image, contains only the published server, runs as UID/GID 1654, exposes
TCP 8080 and UDP 9050 explicitly, and does not require a writable application
directory. Supply a small writable `/tmp` tmpfs because runtime libraries can
legitimately need temporary space; keep the root filesystem read-only.
From a clean checkout:
```bash
docker build --pull=false --tag finalfactory/rendezvous:local .
docker inspect --format '{{.Config.User}}' finalfactory/rendezvous:local
```
The reported user must be `1654:1654`. Digest pins make a rebuild reproducible;
updating .NET is an explicit reviewed change to the tag, digest, SDK pin, and
lock files together. Do not replace the digest with `latest` in production.
The local Compose example applies a read-only root, non-root user, no Linux
capabilities, `no-new-privileges`, bounded PIDs/files/memory/CPU, and a shutdown
grace period longer than the service drain deadline:
```bash
install -d -m 0700 deploy/compose/secrets
umask 077
openssl rand -out deploy/compose/secrets/signing-key 32
export RENDEZVOUS_UID="$(id -u)"
export RENDEZVOUS_GID="$(id -g)"
test "$RENDEZVOUS_UID" -ne 0
docker compose -f deploy/compose/compose.yaml up --build --detach
```
`deploy/compose/appsettings.Production.json` is an isolated loopback smoke
profile, not an Internet template: it deliberately opts into private advertised
endpoints and has no TLS proxy. Its random key is ignored by Git and must be
deleted after use. Its deliberately long key window only keeps this disposable
local fixture usable; production keys require short, reviewed rotation windows.
Production configuration must use its real public names and must leave
`AllowPrivatePublicEndpoints` false.
## Production topology
Use one active service behind a source-preserving edge:
```text
clients -- HTTPS/443 --> TLS reverse proxy -- HTTP/8080 --> Rendezvous
clients -- UDP/9050 -------------------------------------> Rendezvous
```
- Give the HTTPS origin and UDP endpoint stable DNS names. Set
`PublicHttpBaseUrl` to the exact external HTTPS origin and `PublicUdpHost` /
`PublicUdpPort` to the endpoint given to game clients.
- Terminate TLS 1.2 or newer at a maintained reverse proxy. Bind internal HTTP
only to the private proxy network. Restrict `AllowedHosts` to the public HTTP
host; wildcard host filtering is rejected.
- Put only the proxy's exact literal addresses in
`Rendezvous:AbuseProtection:TrustedProxyAddresses`. Rendezvous ignores
forwarded headers from every other source. Keep the last proxy from replacing
the original client address and prevent direct access to TCP 8080.
- Forward UDP as UDP, without an HTTP proxy. NAT, load balancer, firewall, and
return routing must preserve the client's source IP/port and must send replies
from the same advertised IP/port. Many HTTP load balancers, Kubernetes ingress
controllers, rootless container port proxies, anycast products, and generic
L7 services cannot guarantee this. Do not deploy through one unless an actual
host/join smoke proves both observed source and reply path. A load balancer
must have exactly one healthy Rendezvous target.
- Permit inbound TCP 443 to the TLS proxy and UDP 9050 to Rendezvous. Permit the
proxy to reach TCP 8080. Permit DNS, time synchronization, image/telemetry
destinations as required by local policy, and UDP replies to client endpoints.
Deny public TCP 8080 and every unused inbound port.
Readiness is the load-balancer gate; liveness is only a process-health signal.
Remove a draining instance from new traffic when `/health/ready` becomes 503.
Do not use liveness failure to start a second active process while the old one
still owns the public UDP address.
## Required production configuration
Production startup validates all of these before binding listeners:
- an absolute path-free HTTPS `PublicHttpBaseUrl`;
- an unambiguous public `PublicUdpHost` and port;
- `SingleActiveInstance=true`, an explicit non-wildcard `AllowedHosts`, and at
least one exact trusted TLS-proxy address;
- a 1-30 second drain deadline whose minimum observation interval is shorter;
- at least one enabled game policy and an active scoped signing key.
Missing values produce an actionable startup error. The checked-in base file is
intentionally unsafe for Production so an accidental bare launch fails closed.
Signing keys support two external references:
- `env:NAME` reads 1-4096 bytes encoded as base64 from `NAME`;
- `file:/absolute/path` reads 1-4096 raw bytes from a non-symlink file.
Prefer a read-only container secret owned by the configured container identity.
For systemd, use a root-owned, `rendezvous`-group-owned `0440` file (or an
equivalent narrow ACL) so the non-root process can read but not replace it. A
signing key must contain at least 32 random bytes. Never put the key, publisher/operator
credential, or secret value in JSON, a command argument, an image layer, Compose
environment, logs, metrics, or source control. Configuration contains only the
reference and non-secret lifecycle metadata. A vault/KMS adapter can replace the
provider where local policy requires it.
Keep the host clock synchronized with authenticated NTP. Credential and key
windows use wall time; lease, timeout, drain, and rate-limit deadlines use a
monotonic clock. Alert on clock synchronization loss before rotating keys.
The checked-in Compose limits (one CPU and 512 MiB) are for its isolated smoke
profile, not a production capacity claim. The measured core-state candidate
uses 2 vCPU and 2 GiB with the same 128-PID/4096-descriptor ceilings; see
the [capacity and resilience gate](../operations/capacity-and-resilience.md).
Measure real traffic, then change resource limits and server budgets together.
Memory pressure or CPU throttling must not extend orchestrator termination past
`DrainDeadlineSeconds` plus five seconds.
## Graceful shutdown
SIGTERM and the authenticated operator drain both stop new registrations and
join attempts immediately. On process shutdown, HTTP and UDP remain available
long enough for existing join attempts to finish. The service exits as soon as
the minimum drain interval has elapsed and no attempts remain, or forcibly
clears all ephemeral state at the configured deadline. It then stops UDP and
HTTP listeners and exits. Configure Docker/systemd/Kubernetes termination grace
strictly longer than the service deadline; the examples use 40 seconds for a
30-second drain.
Never use SIGKILL for a normal rollout. After stopping, verify the process is
gone and neither `8080/tcp` nor `9050/udp` is bound before starting its
replacement on the same host. A crashed or force-killed process cannot drain;
clients recover through bounded retries and hosts re-register.
## systemd alternative
Publish the server for Linux, install the immutable output at `/opt/rendezvous`,
place production configuration beside the application read-only, place key
files below `/etc/rendezvous`, and install `deploy/systemd/rendezvous.service`:
```bash
dotnet publish src/FinalFactory.Rendezvous.Server \
--configuration Release --runtime linux-x64 --self-contained false \
--output publish/rendezvous
systemd-analyze verify deploy/systemd/rendezvous.service
sudo systemctl daemon-reload
sudo systemctl enable --now rendezvous.service
```
Create the dedicated `rendezvous` user without a login shell. Keep
`/opt/rendezvous` and `/etc/rendezvous` root-owned and non-writable by that user;
install each required key with `root:rendezvous` ownership and mode `0440`. The
unit applies the measured 2-vCPU/2-GiB core-state candidate profile plus the
same filesystem, privilege, network-family, and shutdown hardening as Compose.
## HTTP and UDP smoke
Build the diagnostic once, then exercise the actual published HTTP and UDP
paths. The test creates a public listing, sends authenticated presence and punch
traffic through UDP 9050, establishes peer-to-peer traffic, reports the outcome,
and deregisters cleanly:
```bash
dotnet build src/FinalFactory.Rendezvous.TestClient --configuration Release
./scripts/smoke-deployment.sh
```
For the local Compose profile, the script derives a ten-minute diagnostic
publisher credential from the ignored local key without printing either secret.
For production, do not copy the signing key to the smoke host. Instead inject a
short-lived, region-scoped credential through
`RENDEZVOUS_PUBLISHER_CREDENTIAL`, and set the external endpoints:
```bash
export RENDEZVOUS_PUBLISHER_CREDENTIAL='<short-lived deployment credential>'
export RENDEZVOUS_SMOKE_HTTP_URL='https://rendezvous.your-company.tld/'
export RENDEZVOUS_SMOKE_UDP_ENDPOINT='rendezvous-udp.your-company.tld:9050'
export RENDEZVOUS_SMOKE_GAME_ID='<credential game ID>'
export RENDEZVOUS_SMOKE_ENVIRONMENT_ID='<credential environment ID>'
export RENDEZVOUS_SMOKE_REGION='<credential region>'
export RENDEZVOUS_SMOKE_PROTOCOL_VERSION='<enabled protocol version>'
./scripts/smoke-deployment.sh
```
Those four scope values must match both the short-lived credential and an
enabled server policy. The defaults (`space-game`, `smoke`, `local`, protocol
`1`) are only for the checked-in local Compose profile.
The smoke fails unless both health endpoints and the complete authenticated UDP
mediation/direct-traffic flow succeed. It does not prove every consumer NAT;
run the topology harness and representative external-network tests as well.
## Restart, upgrade, rollback, and backup
Rendezvous has no durable runtime database. Restarting intentionally loses all
listings, observed endpoints, attempts, replay markers, and runtime-only
revocations. Hosts must treat registration as a renewable lease and re-register
after service recovery. Clients must re-browse and start a new bounded attempt.
Back up only reviewed configuration, policy, secret references, key material and
its custody/lifecycle records, deployment manifests, and image digest. Never
claim a backup contains live sessions or endpoints. Restore keys only through the
secret system, not into the image or repository.
For an upgrade:
1. Build and test the new pinned digest; validate configuration without starting
a second active instance.
2. Drain and stop the current process, verify both sockets are released, then
start the replacement on the same public endpoints.
3. Require live/readiness and HTTP+UDP smoke success; monitor host
re-registration, error rate, and direct-connect outcomes.
For rollback, repeat the same stop-before-start sequence with the previously
recorded image digest and compatible configuration/key set. Never run old and
new versions concurrently to avoid split ephemeral state. If a wire-incompatible
change ever becomes necessary, use a new API/protocol version rather than a
rolling two-version replica set.
@@ -0,0 +1,148 @@
{
"schemaVersion": 2,
"evidenceVersion": "v2",
"generatedAt": "2026-07-16T14:10:53.6981858+00:00",
"profile": "candidate",
"runtime": {
"framework": ".NET 10.0.9",
"operatingSystem": "CachyOS",
"kernel": "Unix 7.1.3.2",
"architecture": "X64",
"cpuModel": "AMD Ryzen 7 9800X3D 8-Core Processor",
"processorCount": 2,
"cpuAffinity": "0,1",
"cpuQuota": "not-enforced",
"memoryLimit": "not-enforced",
"garbageCollector": "workstation",
"commitSha": "cf14836d48b0b4aaa67f99433f4fba3585bcd2bb",
"treeState": "clean",
"command": "RENDEZVOUS_CAPACITY_PROFILE=candidate RENDEZVOUS_CAPACITY_CPUSET=0,1 ./scripts/run-capacity-gate.sh",
"imageDigest": "not-containerized",
"workloadSeed": "fixed-sequences-random-identifiers",
"capacityPhaseAverageCpuPercent": 56.37724115383554,
"peakWorkingSetBytes": 169705472,
"managedBytesAfterCleanup": 35615200
},
"targets": {
"visibleListings": 25000,
"activeJoinAttempts": 10000,
"coreControlOperationsPerSecond": 200,
"coreMediationOperationsPerSecond": 2000,
"coreControlP95Milliseconds": 200,
"coreMediationP95Milliseconds": 100,
"maximumAverageCpuPercent": 70,
"maximumWorkingSetBytes": 1610612736,
"soakCycles": 1000,
"soakDurationSeconds": 300
},
"measurements": [
{
"operation": "registration-and-presence",
"samples": 1000,
"p50Milliseconds": 0.003,
"p95Milliseconds": 0.0046,
"p99Milliseconds": 0.0054,
"operationsPerSecond": 282453.96000451926,
"minimumOperationsPerSecond": 200,
"budgetMilliseconds": 200,
"passed": true
},
{
"operation": "lease-renewal",
"samples": 1000,
"p50Milliseconds": 0.0004,
"p95Milliseconds": 0.0009,
"p99Milliseconds": 0.0021,
"operationsPerSecond": 968992.2480620155,
"minimumOperationsPerSecond": 200,
"budgetMilliseconds": 200,
"passed": true
},
{
"operation": "visible-session-browse",
"samples": 250,
"p50Milliseconds": 0.9046,
"p95Milliseconds": 3.3704,
"p99Milliseconds": 3.9471,
"operationsPerSecond": 695.5799787597697,
"minimumOperationsPerSecond": 200,
"budgetMilliseconds": 200,
"passed": true
},
{
"operation": "join-attempt-issuance",
"samples": 1000,
"p50Milliseconds": 0.0029,
"p95Milliseconds": 0.0045,
"p99Milliseconds": 0.0055,
"operationsPerSecond": 296428.042092782,
"minimumOperationsPerSecond": 200,
"budgetMilliseconds": 200,
"passed": true
},
{
"operation": "simultaneous-punch-pairing",
"samples": 1000,
"p50Milliseconds": 0.0043,
"p95Milliseconds": 0.0073,
"p99Milliseconds": 0.0115,
"operationsPerSecond": 109212.03516627532,
"minimumOperationsPerSecond": 2000,
"budgetMilliseconds": 100,
"passed": true
},
{
"operation": "principal-revocation",
"samples": 50,
"p50Milliseconds": 0.518,
"p95Milliseconds": 0.7049,
"p99Milliseconds": 11.8557,
"operationsPerSecond": 1320.1773262184577,
"minimumOperationsPerSecond": 50,
"budgetMilliseconds": 200,
"passed": true
},
{
"operation": "telemetry-recording",
"samples": 1000,
"p50Milliseconds": 0.0001,
"p95Milliseconds": 0.0001,
"p99Milliseconds": 0.0001,
"operationsPerSecond": 1438641.9220256077,
"minimumOperationsPerSecond": 10000,
"budgetMilliseconds": 1,
"passed": true
},
{
"operation": "coincident-listing-attempt-expiry",
"samples": 1,
"p50Milliseconds": 29.6882,
"p95Milliseconds": 29.6882,
"p99Milliseconds": 29.6882,
"operationsPerSecond": 33.682962483916384,
"minimumOperationsPerSecond": 0,
"budgetMilliseconds": 200,
"passed": true
}
],
"state": {
"peakListings": 25000,
"peakAttempts": 10000,
"peakReplayMarkers": 0,
"finalListings": 0,
"finalAttempts": 0,
"finalReplayMarkers": 0,
"expiryChurn": 94906,
"maintenanceSweeps": 36307,
"soakCyclesCompleted": 75126848,
"soakDurationSeconds": 300.0000015,
"soakPeakScheduledExpiryEntries": 7,
"soakManagedGrowthBytes": -257288,
"soakHandleGrowth": 2,
"restartStartedEmpty": true,
"overloadWasTyped": true,
"recoverySucceeded": true
},
"failures": [],
"passed": true
}
+189
View File
@@ -0,0 +1,189 @@
# Capacity, resilience, and availability gate
Tracking: #18
This gate turns the v1 budgets in ADR 0003 into a repeatable release decision.
It does not turn Rendezvous into a horizontally scalable service: v1 remains one
active process with bounded in-memory state. A second process may be a cold
standby, but it must not accept traffic until the first process has stopped and
released the public HTTP and UDP endpoints.
## Launch envelope and approved core-state profile
The approved core-state profile is one Linux process limited to 2 vCPU and
2 GiB RAM. Public HTTP/UDP numbers are launch objectives that require the #23
real-network canary before they become a supported service claim:
| Dimension | Value | Evidence status |
| --- | --- | --- |
| Visible listings | 25,000 | Enforced and measured here |
| Active join attempts | 10,000 | Enforced and measured here |
| Core control path | 200 operations/second; p95 at most 200 ms | Measured here |
| Core mediation path | 2,000 pairings/second; p95 at most 100 ms | Measured here |
| Sustained HTTP demand | 200 requests/second | #23 launch objective; not yet a supported claim |
| Sustained UDP demand | 2,000 datagrams/second | #23 launch objective; not yet a supported claim |
| Public HTTP/UDP latency | p95 at most 200 ms / 100 ms | #23 launch objective; not yet a supported claim |
| Capacity-phase average CPU / peak working memory | below 70% / below 1.5 GiB | Measured for the core candidate |
| Valid in-profile monthly availability | 99.5%, excluding announced maintenance | Operational objective |
| Process-ready RTO / host-visible recovery | 15 seconds / 90 seconds | 15 seconds automated; 90-second deployment drill required |
The proposed public-network mix is 20% registration/update, 30% lease-critical
renew/delete, 30% browse, and 20% join authorization for HTTP. The UDP mix is
60% authenticated host-presence refresh, 30% attempt contributions, and 10%
invalid or duplicate traffic that must be dropped early. A deployment may use a
lower per-game profile, but must not claim a higher one without new versioned
evidence.
The capacity harness fills the complete state ceilings, then measures
registration plus presence, renewal, a 100-item compatible browse, join
issuance, simultaneous two-peer pairing, principal revocation, and telemetry.
It applies 200/100 ms guardrails and minimum 200 control / 2,000 mediation
operations per second to the core hot path. Those measurements deliberately
exclude Kestrel, LiteNetLib, TLS, JSON, socket scheduling, and the documented
mixed traffic shape. The #23 real-network canary must exercise those layers,
rate-shape the mix, record errors and shedding, and meet the public objectives
before launch; a core result is not a public-network latency or throughput claim.
## Reproduce the evidence
Every push runs the quick profile and the selected fault matrix:
```bash
./scripts/run-capacity-gate.sh
```
Run the production candidate on an otherwise idle Linux host and restrict the
runtime to two logical CPUs. The default candidate includes a five-minute,
high-intensity expiry soak; use 3,600 seconds for a release-candidate endurance
run:
```bash
export RENDEZVOUS_CAPACITY_PROFILE=candidate
export RENDEZVOUS_CAPACITY_CPUSET=0,1
export RENDEZVOUS_CAPACITY_OUTPUT="$PWD/artifacts/capacity/candidate.json"
./scripts/run-capacity-gate.sh
# Release-candidate endurance override:
dotnet run --project tests/FinalFactory.Rendezvous.Capacity \
--configuration Release --no-build -- \
--profile candidate --soak-seconds 3600 \
--output artifacts/capacity/candidate-endurance.json
```
The machine must have at least 2 GiB available to the process. For formal
deployment evidence, run inside the same cgroup/container shape as production.
The v2 JSON embeds the commit and tree state, command, image context, CPU model,
kernel, affinity, cgroup quota/limit, collector mode, and workload seed. Supply
`RENDEZVOUS_EVIDENCE_IMAGE_DIGEST` when running a release image. Do not compare
results collected under a debugger,
concurrent build, thermal throttling, or oversubscribed CI host.
The checked-in baseline is
[`candidate-2cpu.json`](../evidence/capacity/v2/candidate-2cpu.json). It was
produced on .NET 10.0.9/Linux x64 with CPU affinity restricted to two logical
CPUs. It filled 25,000 listings and 10,000 attempts, peaked at about 162 MiB,
and cleared all active/retained state. The five-minute baseline supersedes any
earlier local probe when its timestamp and target duration differ.
## Soak and bounded-state interpretation
Each soak cycle creates a listing, repeatedly renews its lease and refreshes
presence, creates a join attempt, replay marker, and retained outcome, checks
that scheduled expiry entries remain proportional to live keys, then advances
the injected monotonic clock beyond all
deadlines, and verifies that listings, attempts, replay, idempotency, and outcome
state return to zero. The candidate also measures managed-memory and process
handle deltas after full collection. Failure is any retained state, more than
64 MiB retained managed memory, more than eight retained handles, a working set
above 1.5 GiB, an untyped capacity result, or failure to admit work after expiry.
This accelerated soak intentionally executes far more state lifecycle/cleanup
events than wall-clock traffic would permit. It catches stale deadline-queue
entries, cache growth, replay/idempotency retention, and cleanup cost. Because
it does not open Kestrel/LiteNetLib connections, its process-handle delta is only
a harness guard and is not evidence of transport stability by itself. The
selected production-process gate adds a ten-second real HTTP/UDP transport soak,
samples child-process handles and RSS, asserts bounded growth, then verifies a
clean SIGTERM and socket release. #23 must extend that into the full rate-shaped
multi-client canary while sampling queues, managed memory, and state
cardinalities. A one-hour core override remains required before tagging a
production release.
## Fault and recovery matrix
`run-capacity-gate.sh` runs these deterministic production paths before the
numeric profile:
| Fault | Required result |
| --- | --- |
| HTTP/UDP overload and tracker exhaustion | Typed HTTP `429`/`CapacityExceeded`, silent UDP drop, bounded tracker keys, recovery after the window |
| Optional traffic saturation | Lease-critical renew/update/delete capacity remains available |
| Store/dependency unavailable | Readiness fails; new authorization returns typed `ServiceUnavailable`; liveness remains independent |
| Graceful drain/SIGTERM | New work returns `Draining`; existing pairing may finish; process exits 0 and releases TCP/UDP before the deadline |
| Hard restart | In-flight state is lost; SDK reports typed `ServiceUnavailable`; a host re-registers, rebinds presence, and becomes the only browser-visible replacement |
| UDP listener bind/restart | Readiness stays false without the required listener; rebinding the advertised port restores native LiteNetLib pairing |
| Wall-clock jump/skew | Monotonic lease/attempt authority is neither shortened nor extended; credential skew remains capped at 30 seconds |
| Signing-secret rotation | New key signs, overlap verifies, retired/revoked key rejects, missing material fails startup |
| Principal revocation | Listing, presence, attempts, and outcome paths are removed atomically within the latency budget |
No external database exists in v1, so “dependency/store failure” means the
process-local atomic store is marked unavailable or a required listener/key is
unready. The service fails closed rather than pretending a degraded writable
mode exists.
## Bandwidth and amplification
- Accepted application datagrams are at most 1,200 bytes.
- Malformed, oversized, unauthenticated, stale, replayed, wrong-role, and
rate-limited traffic receives zero response bytes.
- A completing authenticated contribution produces at most one introduction to
each observed peer, and the combined response is at most 2.0 times that
contribution's bytes.
- The frozen-envelope and native LiteNetLib socket tests measure this on the real
UDP listener; the hostile corpus and allocation gate exercise 10,000+ inputs
without input-sized logs, tasks, or queues.
Bandwidth planning must therefore reserve ingress for the configured 2,000
datagrams/second plus edge overhead and egress for a worst-case verified 2.0
amplification. Actual successful pairs normally use two contributions and two
introductions; normal gameplay leaves Rendezvous entirely.
## Availability decision
Single-active remains the v1 topology. The measured core profile proves bounded
state and substantial core-path headroom, while public launch capacity remains
conditional on #23. The service has a bounded stop-before-start restart path.
Its failure domain is deliberately
one process/node/public UDP endpoint: node, kernel, host network, DNS/TLS edge,
secret configuration, or operator error can remove all readiness until the cold
replacement owns the same source-preserving endpoint.
The 99.5% objective permits about 216 minutes of unannounced downtime in a
30-day month. Operations must target process readiness within 15 seconds and
host-visible re-registration within 90 seconds, page when no ready instance
exists, and include detection plus recovery in the monthly budget. The current
in-process test validates typed downtime, same-port HTTP restart, fresh
registration, presence rebinding, and browser visibility in under five seconds;
the production-process test separately validates graceful termination, TCP/UDP
release, replacement startup on the same endpoints, UDP readiness, and the
15-second process-ready RTO. Cold-standby activation policy and the 90-second
operator-to-host recovery objective still require a deployment drill before
release. Rollout and rollback use the
deployment runbook's drain, stop, socket-release, start, smoke sequence; never
overlap old and new active processes.
Bring shared TTL/CAS state and deterministic mediator routing forward before
enabling two active instances if any of these occurs:
- one node cannot sustain 150% of the measured 30-day peak while meeting SLOs;
- CPU stays above 70%, memory above 75%, attempt depth above 70%, or limiter
drops/latency remain elevated after abusive traffic is excluded;
- the availability target rises above 99.5% or planned maintenance must preserve
listings; or
- one region requires multiple simultaneously active mediator endpoints.
Rendezvous makes no multi-instance claim today, so a two-node atomic-pairing
test is intentionally not applicable. It becomes a hard release gate with the
shared-state/routing implementation; until then `SingleActiveInstance=false`
fails production startup. Multi-region and relay remain separate evidence-driven
decisions.
@@ -0,0 +1,144 @@
# 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. A
repeat action may extend an active revocation but never shortens it; wait for its
original deadline rather than treating a shorter repeat as an un-revoke. 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.
+6 -1
View File
@@ -20,6 +20,10 @@ traffic is also isolated by its authenticated scope.
Health probes use their own source-prefix budget so public API overload cannot
make a healthy instance fail its orchestrator probes, while health traffic is
still bounded.
Operator endpoints likewise use a separate bounded rate/concurrency partition
backed by the critical tracker reserve. They first require an exact source IP
from the default-deny `OperatorAllowedAddresses` policy, so public traffic
cannot spend the incident-response budget.
3. Once an endpoint has safely derived identities, it also acquires applicable
tenant, principal or capability, and listing/attempt budgets. Secret
capabilities are represented only by bounded SHA-256 fingerprints.
@@ -41,7 +45,8 @@ load shedding.
`Rendezvous:AbuseProtection:MaxTrackedKeys` is a hard combined ceiling for rate
and active-concurrency keys. General HTTP and UDP traffic cannot consume the
configured `CriticalTrackedKeyReserve`; lease operations and health probes may
configured `CriticalTrackedKeyReserve`; lease operations, health probes, and
allowlisted operator controls may
use that reserve but never exceed the hard ceiling. A request that would exceed
its applicable ceiling fails closed without adding state. Fixed-window rate keys
are cleared at the next window boundary; concurrency keys are removed as their
+7 -4
View File
@@ -63,8 +63,9 @@ only its public key ID/lifecycle metadata and does not require retired secret
material to remain available.
Key IDs are non-secret base64url identifiers. Secret references are resolved
through `ISecretProvider`; production supports `env:<VARIABLE>` references and
the interface is replaceable by a deployment-specific vault/KMS adapter. The
through `ISecretProvider`; production supports base64 `env:<VARIABLE>` and raw
`file:/absolute/path` references to bounded non-symlink files. The interface is
replaceable by a deployment-specific vault/KMS adapter. The
committed development profile uses an in-memory random key identified by a
`development:ephemeral/...` reference. It never writes key material to disk and
all credentials become invalid when the process exits.
@@ -74,8 +75,10 @@ all credentials become invalid when the process exits.
`Rendezvous:Provisioning` supplies issuer, audience, clock skew, signing-key
descriptors, and game policies. A production key reference such as
`env:RENDEZVOUS_SIGNING_KEY_2026_01` expects that environment variable to hold at
least 32 random bytes encoded as base64. Missing, malformed, short, inactive, or
duplicate keys stop startup with a key-ID-only diagnostic. No game-wide secret
least 32 random bytes encoded as base64. `file:/run/secrets/rendezvous-signing`
expects the raw bytes in a read-only, absolute, non-symlink file. Missing,
malformed, short, inactive, or duplicate keys stop startup with a key-ID-only
diagnostic. No game-wide secret
belongs in `appsettings`, source control, examples, the Client package, URLs,
responses, logs, metrics, exceptions, or diagnostic dumps.
+58
View File
@@ -0,0 +1,58 @@
#!/usr/bin/env bash
set -euo pipefail
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
PROFILE="${RENDEZVOUS_CAPACITY_PROFILE:-quick}"
OUTPUT="${RENDEZVOUS_CAPACITY_OUTPUT:-$ROOT/artifacts/capacity/rendezvous-capacity-v2.json}"
CPUSET="${RENDEZVOUS_CAPACITY_CPUSET:-}"
PROJECT="$ROOT/tests/FinalFactory.Rendezvous.Capacity/FinalFactory.Rendezvous.Capacity.csproj"
TESTS="$ROOT/tests/FinalFactory.Rendezvous.Tests/FinalFactory.Rendezvous.Tests.csproj"
if [[ "$PROFILE" != quick && "$PROFILE" != candidate ]]; then
printf 'RENDEZVOUS_CAPACITY_PROFILE must be quick or candidate.\n' >&2
exit 2
fi
command -v dotnet >/dev/null || {
printf 'Missing required command: dotnet\n' >&2
exit 2
}
if [[ -n "$CPUSET" ]]; then
command -v taskset >/dev/null || {
printf 'taskset is required when RENDEZVOUS_CAPACITY_CPUSET is set.\n' >&2
exit 2
}
fi
cd "$ROOT"
export RENDEZVOUS_EVIDENCE_COMMIT="$(git rev-parse HEAD)"
if [[ -n "$(git status --porcelain)" ]]; then
export RENDEZVOUS_EVIDENCE_TREE_STATE=dirty
else
export RENDEZVOUS_EVIDENCE_TREE_STATE=clean
fi
if [[ "$PROFILE" == candidate && "$RENDEZVOUS_EVIDENCE_TREE_STATE" != clean ]]; then
printf 'Candidate evidence requires a clean source tree.\n' >&2
exit 2
fi
export RENDEZVOUS_EVIDENCE_CPUSET="${CPUSET:-unrestricted}"
export RENDEZVOUS_EVIDENCE_COMMAND="RENDEZVOUS_CAPACITY_PROFILE=$PROFILE RENDEZVOUS_CAPACITY_CPUSET=${CPUSET:-unrestricted} ./scripts/run-capacity-gate.sh"
dotnet restore "$ROOT/Rendezvous.slnx" --locked-mode
dotnet build "$ROOT/Rendezvous.slnx" --configuration Release --no-restore
filter='FullyQualifiedName~TrackerCapacityFailsClosedWithoutGrowingAndAWindowResetRecovers|FullyQualifiedName~OptionalTrafficCannotConsumeTheLeaseOperationReserve|FullyQualifiedName~ConcurrentAbusiveBurstStaysBoundedAndCannotBlockCriticalHttp|FullyQualifiedName~HttpOverloadIsTypedAndOversizedBodiesAreRejectedBeforeDispatch|FullyQualifiedName~WallClockMovementDoesNotExpireOrExtendLease|FullyQualifiedName~RepeatedMutableDeadlineRefreshesKeepOneScheduledEntryPerKey|FullyQualifiedName~RepeatedPrincipalRevocationCanExtendButCannotShortenProtection|FullyQualifiedName~RestartHasNewGenerationAndNoEphemeralState|FullyQualifiedName~RestartReturnsTypedUnavailabilityThenAllowsHostReregistration|FullyQualifiedName~DrainRejectsNewWorkAllowsInflightCompletionThenClearsState|FullyQualifiedName~UnavailableStoreFailsNewAuthorizationClosedAndErasesActiveState|FullyQualifiedName~KeyRotationHonorsOverlapAndRejectsRetiredKeys|FullyQualifiedName~OperatorSurfaceSeparatesAuthenticationConfirmsActionsAndRedactsInspection|FullyQualifiedName~SigtermDrainsThenReleasesHttpAndUdpSockets|FullyQualifiedName~ProductionTransportSoakKeepsHandlesMemoryAndSocketsBounded|FullyQualifiedName~NativeLiteNetLibRequestsIntroduceTheAuthorizedPair'
dotnet test "$TESTS" --configuration Release --no-build --filter "$filter" \
--logger 'console;verbosity=minimal'
mkdir -p "$(dirname "$OUTPUT")"
arguments=(
dotnet run --project "$PROJECT" --configuration Release --no-build --
--profile "$PROFILE" --output "$OUTPUT"
)
if [[ -n "$CPUSET" ]]; then
taskset -c "$CPUSET" "${arguments[@]}"
else
"${arguments[@]}"
fi
printf 'Capacity and resilience gate passed; evidence: %s\n' "$OUTPUT"
+148
View File
@@ -0,0 +1,148 @@
#!/usr/bin/env bash
set -euo pipefail
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
SERVICE_URL="${RENDEZVOUS_SMOKE_HTTP_URL:-http://127.0.0.1:8080/}"
MEDIATOR="${RENDEZVOUS_SMOKE_UDP_ENDPOINT:-127.0.0.1:9050}"
TIMEOUT_SECONDS="${RENDEZVOUS_SMOKE_TIMEOUT_SECONDS:-30}"
LOCAL_KEY="${RENDEZVOUS_SMOKE_LOCAL_KEY:-$ROOT/deploy/compose/secrets/signing-key}"
PROJECT="$ROOT/src/FinalFactory.Rendezvous.TestClient/FinalFactory.Rendezvous.TestClient.csproj"
BUILD_CONFIGURATION="${RENDEZVOUS_SMOKE_CONFIGURATION:-Release}"
GAME_ID="${RENDEZVOUS_SMOKE_GAME_ID:-space-game}"
ENVIRONMENT_ID="${RENDEZVOUS_SMOKE_ENVIRONMENT_ID:-smoke}"
REGION="${RENDEZVOUS_SMOKE_REGION:-local}"
PROTOCOL_VERSION="${RENDEZVOUS_SMOKE_PROTOCOL_VERSION:-1}"
for command in curl date dotnet jq mktemp od openssl tail tr wc; do
command -v "$command" >/dev/null || {
printf 'Missing required command: %s\n' "$command" >&2
exit 2
}
done
if [[ ! "$TIMEOUT_SECONDS" =~ ^[0-9]+$ ]] || (( TIMEOUT_SECONDS < 1 || TIMEOUT_SECONDS > 300 )); then
printf 'RENDEZVOUS_SMOKE_TIMEOUT_SECONDS must be an integer from 1 through 300.\n' >&2
exit 2
fi
if [[ ! "$PROTOCOL_VERSION" =~ ^[0-9]+$ ]] || (( PROTOCOL_VERSION < 1 )); then
printf 'RENDEZVOUS_SMOKE_PROTOCOL_VERSION must be a positive integer.\n' >&2
exit 2
fi
for scoped_value in "$GAME_ID" "$ENVIRONMENT_ID" "$REGION"; do
if [[ -z "$scoped_value" ]]; then
printf 'Smoke game, environment, and region values must not be empty.\n' >&2
exit 2
fi
done
base64url() {
openssl base64 -A | tr '+/' '-_' | tr -d '='
}
local_credential() {
if [[ ! -f "$LOCAL_KEY" ]] || [[ "$(wc -c < "$LOCAL_KEY")" -ne 32 ]]; then
printf 'Local Compose smoke key must be exactly 32 bytes: %s\n' "$LOCAL_KEY" >&2
exit 2
fi
local now expires nonce payload encoded signed hex signature
now="$(date +%s)"
expires="$((now + 600))"
nonce="$(openssl rand -hex 16)"
payload="$(jq -cn \
--arg issuer final-factory-rendezvous-smoke \
--arg audience rendezvous-service \
--arg subject local-smoke-host \
--arg kind dedicatedPublisher \
--arg gameId "$GAME_ID" \
--arg environmentId "$ENVIRONMENT_ID" \
--arg region "$REGION" \
--arg nonce "$nonce" \
--argjson now "$now" \
--argjson expires "$expires" \
'{version:1,issuer:$issuer,audience:$audience,subject:$subject,kind:$kind,gameId:$gameId,environmentId:$environmentId,regions:[$region],permissions:[],issuedAtUnixSeconds:$now,notBeforeUnixSeconds:$now,expiresAtUnixSeconds:$expires,nonce:$nonce}')"
encoded="$(printf '%s' "$payload" | base64url)"
signed="rv1.local-smoke-1.$encoded"
hex="$(od -An -v -tx1 "$LOCAL_KEY" | tr -d ' \n')"
signature="$(printf '%s' "$signed" \
| openssl dgst -sha256 -mac HMAC -macopt "hexkey:$hex" -binary \
| base64url)"
printf '%s.%s' "$signed" "$signature"
}
credential="${RENDEZVOUS_PUBLISHER_CREDENTIAL:-}"
if [[ -z "$credential" ]]; then
credential="$(local_credential)"
fi
export RENDEZVOUS_PUBLISHER_CREDENTIAL="$credential"
curl --fail --silent --show-error --max-time 5 "${SERVICE_URL%/}/health/live" >/dev/null
curl --fail --silent --show-error --max-time 5 "${SERVICE_URL%/}/health/ready" >/dev/null
temp_dir="$(mktemp -d)"
host_log="$temp_dir/host.jsonl"
join_log="$temp_dir/join.jsonl"
host_pid=''
cleanup() {
local status="$?"
if [[ -n "$host_pid" ]] && kill -0 "$host_pid" 2>/dev/null; then
kill -TERM "$host_pid" 2>/dev/null || true
wait "$host_pid" 2>/dev/null || true
fi
if [[ "$status" -ne 0 ]]; then
printf 'Deployment smoke failed; sanitized diagnostic events follow.\n' >&2
[[ -f "$host_log" ]] && jq -c . "$host_log" >&2 || true
[[ -f "$join_log" ]] && jq -c . "$join_log" >&2 || true
fi
rm -rf "$temp_dir"
return "$status"
}
trap cleanup EXIT
trap 'exit 130' INT
trap 'exit 143' TERM
dotnet run --project "$PROJECT" --configuration "$BUILD_CONFIGURATION" --no-build -- \
host --service "$SERVICE_URL" --mediator "$MEDIATOR" \
--game "$GAME_ID" --environment "$ENVIRONMENT_ID" --region "$REGION" --protocol "$PROTOCOL_VERSION" \
--script --json --exit-after-echo --timeout-seconds "$TIMEOUT_SECONDS" \
>"$host_log" 2>&1 &
host_pid="$!"
ready=false
for ((iteration = 0; iteration < TIMEOUT_SECONDS * 4; iteration++)); do
if jq -e 'select(.event == "host.ready")' "$host_log" >/dev/null 2>&1; then
ready=true
break
fi
if ! kill -0 "$host_pid" 2>/dev/null; then
printf 'Host diagnostic stopped before it became ready.\n' >&2
jq -c . "$host_log" >&2 || true
exit 1
fi
sleep 0.25
done
if [[ "$ready" != true ]]; then
printf 'Host diagnostic did not become ready within %s seconds.\n' "$TIMEOUT_SECONDS" >&2
exit 1
fi
listing_id="$(jq -r 'select(.event == "host.registered") | .listingId' "$host_log" | tail -n 1)"
if [[ -z "$listing_id" || "$listing_id" == null ]]; then
printf 'Host diagnostic did not report a listing ID.\n' >&2
exit 1
fi
dotnet run --project "$PROJECT" --configuration "$BUILD_CONFIGURATION" --no-build -- \
join --service "$SERVICE_URL" --mediator "$MEDIATOR" \
--game "$GAME_ID" --environment "$ENVIRONMENT_ID" --region "$REGION" --protocol "$PROTOCOL_VERSION" \
--listing "$listing_id" --script --json --timeout-seconds "$TIMEOUT_SECONDS" \
>"$join_log" 2>&1
wait "$host_pid"
host_pid=''
jq -e 'select(.event == "host.direct-traffic" and .status == "verified")' "$host_log" >/dev/null
jq -e 'select(.event == "host.deregistered" and .status == "complete")' "$host_log" >/dev/null
jq -e 'select(.event == "join.direct-traffic" and .status == "verified")' "$join_log" >/dev/null
jq -e 'select(.event == "join.outcome-report" and .status == "accepted")' "$join_log" >/dev/null
printf 'Rendezvous deployment smoke passed: HTTP live/ready and authenticated UDP mediation/direct traffic.\n'
@@ -20,6 +20,8 @@ internal sealed class AbuseProtectionOptions
public string[] TrustedProxyAddresses { get; set; } = [];
public string[] OperatorAllowedAddresses { get; set; } = [];
[Range(1, 100_000)]
public int HealthGlobalRequestsPerWindow { get; set; } = 1_000;
@@ -32,6 +34,18 @@ internal sealed class AbuseProtectionOptions
[Range(1, 1_000)]
public int HealthIpPrefixConcurrency { get; set; } = 8;
[Range(1, 100_000)]
public int OperatorGlobalRequestsPerWindow { get; set; } = 1_000;
[Range(1, 10_000)]
public int OperatorGlobalConcurrency { get; set; } = 32;
[Range(1, 100_000)]
public int OperatorIpPrefixRequestsPerWindow { get; set; } = 120;
[Range(1, 1_000)]
public int OperatorIpPrefixConcurrency { get; set; } = 8;
[Range(1, 1_000_000)]
public int HttpGlobalRequestsPerWindow { get; set; } = 20_000;
@@ -2,6 +2,7 @@ using System.Buffers;
using System.Net;
using System.Security.Cryptography;
using System.Text;
using FinalFactory.Rendezvous.Server.Observability;
using Microsoft.Extensions.Options;
namespace FinalFactory.Rendezvous.Server.Abuse;
@@ -12,13 +13,23 @@ internal sealed class AbuseProtectionService
private readonly TimeProvider _timeProvider;
private readonly TrackerState _httpTracker;
private readonly TrackerState _udpTracker;
private readonly RendezvousTelemetry? _telemetry;
private readonly HashSet<string> _operatorAllowedAddresses;
public AbuseProtectionService(
IOptions<AbuseProtectionOptions> options,
TimeProvider? timeProvider = null)
TimeProvider? timeProvider = null,
RendezvousTelemetry? telemetry = null)
{
_options = options.Value;
_timeProvider = timeProvider ?? TimeProvider.System;
_telemetry = telemetry;
_operatorAllowedAddresses = options.Value.OperatorAllowedAddresses
.Select(static value => IPAddress.TryParse(value, out IPAddress? address)
? NormalizeAddress(address).ToString()
: string.Empty)
.Where(static value => value.Length > 0)
.ToHashSet(StringComparer.Ordinal);
DateTimeOffset now = _timeProvider.GetUtcNow();
_httpTracker = new(now);
_udpTracker = new(now);
@@ -87,6 +98,35 @@ internal sealed class AbuseProtectionService
out retryAfterSeconds);
}
public bool IsOperatorSourceAllowed(IPAddress? remoteAddress) =>
remoteAddress is not null
&& _operatorAllowedAddresses.Contains(NormalizeAddress(remoteAddress).ToString());
public bool TryAcquireOperatorIngress(
IPAddress? remoteAddress,
out AbuseLease? lease,
out int retryAfterSeconds)
{
string prefix = GetNetworkPrefix(remoteAddress);
RateDimension[] rates =
[
new("operator:rate:global", _options.OperatorGlobalRequestsPerWindow),
new($"operator:rate:ip:{prefix}", _options.OperatorIpPrefixRequestsPerWindow),
];
RateDimension[] concurrency =
[
new("operator:concurrency:global", _options.OperatorGlobalConcurrency),
new($"operator:concurrency:ip:{prefix}", _options.OperatorIpPrefixConcurrency),
];
return TryAcquire(
rates,
concurrency,
TrackerDomain.Http,
true,
out lease,
out retryAfterSeconds);
}
public bool TryAcquireHttpIdentity(
string operation,
string? tenant,
@@ -266,67 +306,96 @@ internal sealed class AbuseProtectionService
out int retryAfterSeconds)
{
TrackerState tracker = domain == TrackerDomain.Udp ? _udpTracker : _httpTracker;
bool accepted;
lock (tracker.Gate)
{
DateTimeOffset now = _timeProvider.GetUtcNow();
TimeSpan window = TimeSpan.FromSeconds(_options.WindowSeconds);
if (now - tracker.WindowStartedAt >= window || now < tracker.WindowStartedAt)
{
tracker.WindowCounts.Clear();
tracker.WindowStartedAt = now;
}
accepted = TryAcquireLocked(
tracker,
rates,
concurrency,
domain,
canUseCriticalReserve,
out lease,
out retryAfterSeconds);
}
retryAfterSeconds = Math.Max(
1,
(int)Math.Ceiling((window - (now - tracker.WindowStartedAt)).TotalSeconds));
int stagedNewKeys = 0;
int partitionLimit = domain == TrackerDomain.Udp
? _options.UdpTrackedKeyLimit
: _options.MaxTrackedKeys - _options.UdpTrackedKeyLimit;
int maxTrackedKeys = domain == TrackerDomain.Udp || canUseCriticalReserve
? partitionLimit
: partitionLimit - _options.CriticalTrackedKeyReserve;
if (!CanAcquireAll(
tracker,
tracker.WindowCounts,
rates,
maxTrackedKeys,
ref stagedNewKeys)
|| !CanAcquireAll(
tracker,
tracker.ConcurrencyCounts,
concurrency,
maxTrackedKeys,
ref stagedNewKeys))
{
lease = null;
return false;
}
if (!accepted)
{
_telemetry?.RecordLimiterDrop(
domain == TrackerDomain.Udp ? "udp" : "http",
"rate-or-concurrency");
}
foreach (RateDimension dimension in rates)
{
tracker.WindowCounts[dimension.Key] =
tracker.WindowCounts.GetValueOrDefault(dimension.Key) + 1;
}
return accepted;
}
if (concurrency.IsEmpty)
{
lease = null;
return true;
}
private bool TryAcquireLocked(
TrackerState tracker,
ReadOnlySpan<RateDimension> rates,
ReadOnlySpan<RateDimension> concurrency,
TrackerDomain domain,
bool canUseCriticalReserve,
out AbuseLease? lease,
out int retryAfterSeconds)
{
DateTimeOffset now = _timeProvider.GetUtcNow();
TimeSpan window = TimeSpan.FromSeconds(_options.WindowSeconds);
if (now - tracker.WindowStartedAt >= window || now < tracker.WindowStartedAt)
{
tracker.WindowCounts.Clear();
tracker.WindowStartedAt = now;
}
string[] acquiredConcurrency = new string[concurrency.Length];
for (int index = 0; index < concurrency.Length; index++)
{
RateDimension dimension = concurrency[index];
tracker.ConcurrencyCounts[dimension.Key] =
tracker.ConcurrencyCounts.GetValueOrDefault(dimension.Key) + 1;
acquiredConcurrency[index] = dimension.Key;
}
retryAfterSeconds = Math.Max(
1,
(int)Math.Ceiling((window - (now - tracker.WindowStartedAt)).TotalSeconds));
int stagedNewKeys = 0;
int partitionLimit = domain == TrackerDomain.Udp
? _options.UdpTrackedKeyLimit
: _options.MaxTrackedKeys - _options.UdpTrackedKeyLimit;
int maxTrackedKeys = domain == TrackerDomain.Udp || canUseCriticalReserve
? partitionLimit
: partitionLimit - _options.CriticalTrackedKeyReserve;
if (!CanAcquireAll(
tracker,
tracker.WindowCounts,
rates,
maxTrackedKeys,
ref stagedNewKeys)
|| !CanAcquireAll(
tracker,
tracker.ConcurrencyCounts,
concurrency,
maxTrackedKeys,
ref stagedNewKeys))
{
lease = null;
return false;
}
lease = new AbuseLease(this, tracker, acquiredConcurrency);
foreach (RateDimension dimension in rates)
{
tracker.WindowCounts[dimension.Key] =
tracker.WindowCounts.GetValueOrDefault(dimension.Key) + 1;
}
if (concurrency.IsEmpty)
{
lease = null;
return true;
}
string[] acquiredConcurrency = new string[concurrency.Length];
for (int index = 0; index < concurrency.Length; index++)
{
RateDimension dimension = concurrency[index];
tracker.ConcurrencyCounts[dimension.Key] =
tracker.ConcurrencyCounts.GetValueOrDefault(dimension.Key) + 1;
acquiredConcurrency[index] = dimension.Key;
}
lease = new AbuseLease(this, tracker, acquiredConcurrency);
return true;
}
private static bool CanAcquireAll(
@@ -415,6 +484,9 @@ internal sealed class AbuseProtectionService
return "unknown";
}
private static IPAddress NormalizeAddress(IPAddress address) =>
address.IsIPv4MappedToIPv6 ? address.MapToIPv4() : address;
private readonly record struct RateDimension(string Key, int Limit);
private enum TrackerDomain
@@ -19,16 +19,71 @@ internal sealed class HttpAbuseProtectionMiddleware(
string operation = context.GetEndpoint()?.Metadata.GetMetadata<IEndpointNameMetadata>()
?.EndpointName ?? "Unmatched";
bool healthEndpoint = operation is "GetLiveness" or "GetReadiness";
bool acquired = healthEndpoint
? protection.TryAcquireHealthIngress(
bool operatorEndpoint = operation is
"GetOperatorStatus"
or "RevokeOperatorListing"
or "RevokeOperatorPrincipal"
or "RevokeOperatorSigningKey"
or "BeginOperatorDrain";
if (operatorEndpoint
&& !protection.IsOperatorSourceAllowed(context.Connection.RemoteIpAddress))
{
bool deniedSourceAdmitted = protection.TryAcquireHttpIngress(
context.Connection.RemoteIpAddress,
out AbuseProtectionService.AbuseLease? lease,
out int retryAfterSeconds)
: protection.TryAcquireHttpIngress(
"Unmatched",
out AbuseProtectionService.AbuseLease? deniedSourceLease,
out int deniedRetryAfterSeconds);
using (deniedSourceLease)
{
if (!deniedSourceAdmitted)
{
context.Response.Headers.RetryAfter = deniedRetryAfterSeconds.ToString(
System.Globalization.CultureInfo.InvariantCulture);
await WriteErrorAsync(
context,
StatusCodes.Status429TooManyRequests,
RendezvousErrorCode.RateLimited,
"The request rate limit was exceeded.",
deniedRetryAfterSeconds).ConfigureAwait(false);
return;
}
await WriteErrorAsync(
context,
StatusCodes.Status404NotFound,
RendezvousErrorCode.NotFound,
"The requested resource was not found.").ConfigureAwait(false);
}
return;
}
AbuseProtectionService.AbuseLease? lease;
int retryAfterSeconds;
bool acquired;
if (healthEndpoint)
{
acquired = protection.TryAcquireHealthIngress(
context.Connection.RemoteIpAddress,
out lease,
out retryAfterSeconds);
}
else if (operatorEndpoint)
{
acquired = protection.TryAcquireOperatorIngress(
context.Connection.RemoteIpAddress,
out lease,
out retryAfterSeconds);
}
else
{
acquired = protection.TryAcquireHttpIngress(
context.Connection.RemoteIpAddress,
operation,
out lease,
out retryAfterSeconds);
}
if (!acquired)
{
context.Response.Headers.RetryAfter = retryAfterSeconds.ToString(
@@ -1,4 +1,5 @@
using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Observability;
using FinalFactory.Rendezvous.Server.Sessions;
using FinalFactory.Rendezvous.Server.State;
@@ -15,6 +16,10 @@ internal sealed class ConnectionOutcomeMetrics
{
private readonly object _gate = new();
private readonly Dictionary<(ConnectionOutcomeKind, ConnectionElapsedBucket), long> _counts = [];
private readonly RendezvousTelemetry? _telemetry;
public ConnectionOutcomeMetrics(RendezvousTelemetry? telemetry = null) =>
_telemetry = telemetry;
internal void Record(ConnectionOutcomeKind outcome, ConnectionElapsedBucket elapsedBucket)
{
@@ -24,6 +29,8 @@ internal sealed class ConnectionOutcomeMetrics
_counts.TryGetValue(key, out long count);
_counts[key] = count + 1;
}
_telemetry?.RecordConnectionOutcome(outcome.ToString(), elapsedBucket.ToString());
}
internal long GetCount(ConnectionOutcomeKind outcome, ConnectionElapsedBucket elapsedBucket)
@@ -0,0 +1,235 @@
using System.ComponentModel.DataAnnotations;
using System.Net;
using FinalFactory.Rendezvous.Server.Abuse;
namespace FinalFactory.Rendezvous.Server.Deployment;
internal sealed record DeploymentOptions
{
public const string SectionName = "Rendezvous:Deployment";
[Required]
public string PublicHttpBaseUrl { get; init; } = string.Empty;
[Required]
public string PublicUdpHost { get; init; } = string.Empty;
[Range(1, 65_535)]
public int PublicUdpPort { get; init; } = 9050;
[Range(1, 30)]
public int DrainDeadlineSeconds { get; init; } = 30;
[Range(0, 5)]
public int MinimumDrainSeconds { get; init; } = 1;
public bool SingleActiveInstance { get; init; } = true;
public bool AllowPrivatePublicEndpoints { get; init; }
public IReadOnlyList<string> ValidateProduction(
AbuseProtectionOptions abuseProtection,
string? allowedHosts)
{
List<string> errors = [];
if (!SingleActiveInstance)
{
errors.Add("Rendezvous:Deployment:SingleActiveInstance must be true because ephemeral state is not shared between replicas.");
}
if (MinimumDrainSeconds >= DrainDeadlineSeconds)
{
errors.Add("Rendezvous:Deployment:MinimumDrainSeconds must be less than DrainDeadlineSeconds.");
}
if (DrainDeadlineSeconds is < 1 or > 30)
{
errors.Add("Rendezvous:Deployment:DrainDeadlineSeconds must be between 1 and 30.");
}
if (MinimumDrainSeconds is < 0 or > 5)
{
errors.Add("Rendezvous:Deployment:MinimumDrainSeconds must be between 0 and 5.");
}
if (PublicUdpPort is < 1 or > 65_535)
{
errors.Add("Rendezvous:Deployment:PublicUdpPort must be between 1 and 65535.");
}
ValidateHttpEndpoint(errors);
ValidateUdpEndpoint(errors);
ValidateAllowedHosts(errors, allowedHosts, PublicHttpBaseUrl);
if (abuseProtection.TrustedProxyAddresses is not { Length: > 0 })
{
errors.Add(
"Rendezvous:AbuseProtection:TrustedProxyAddresses must list the exact TLS proxy addresses; forwarded headers are rejected without this trust boundary.");
}
return errors;
}
private void ValidateHttpEndpoint(List<string> errors)
{
if (!Uri.TryCreate(PublicHttpBaseUrl, UriKind.Absolute, out Uri? endpoint)
|| !string.Equals(endpoint.Scheme, Uri.UriSchemeHttps, StringComparison.Ordinal)
|| !string.IsNullOrEmpty(endpoint.UserInfo)
|| !string.IsNullOrEmpty(endpoint.Query)
|| !string.IsNullOrEmpty(endpoint.Fragment)
|| endpoint.AbsolutePath != "/")
{
errors.Add(
"Rendezvous:Deployment:PublicHttpBaseUrl must be an absolute HTTPS origin with no credentials, path, query, or fragment (for example, https://rendezvous.your-company.tld/).");
return;
}
if (!AllowPrivatePublicEndpoints && !IsPublicHost(endpoint.Host))
{
errors.Add(
"Rendezvous:Deployment:PublicHttpBaseUrl must use a public DNS name or address; set AllowPrivatePublicEndpoints=true only for an isolated deployment smoke test.");
}
}
private void ValidateUdpEndpoint(List<string> errors)
{
if (string.IsNullOrWhiteSpace(PublicUdpHost)
|| PublicUdpHost.Contains("//", StringComparison.Ordinal)
|| PublicUdpHost.Contains(':', StringComparison.Ordinal) && !IPAddress.TryParse(PublicUdpHost, out _)
|| Uri.CheckHostName(PublicUdpHost) == UriHostNameType.Unknown)
{
errors.Add(
"Rendezvous:Deployment:PublicUdpHost must contain only the advertised DNS name or IP address; configure the port separately.");
return;
}
if (!AllowPrivatePublicEndpoints && !IsPublicHost(PublicUdpHost))
{
errors.Add(
"Rendezvous:Deployment:PublicUdpHost must use a public DNS name or address; set AllowPrivatePublicEndpoints=true only for an isolated deployment smoke test.");
}
}
private static void ValidateAllowedHosts(
List<string> errors,
string? allowedHosts,
string publicHttpBaseUrl)
{
string[] hosts = (allowedHosts ?? string.Empty).Split(
';',
StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
if (hosts.Length == 0 || hosts.Any(static host => host is "*" or "+"))
{
errors.Add(
"AllowedHosts must explicitly list the public HTTP host in production; wildcard or empty host filtering is unsafe.");
return;
}
if (Uri.TryCreate(publicHttpBaseUrl, UriKind.Absolute, out Uri? endpoint)
&& !hosts.Contains(endpoint.Host, StringComparer.OrdinalIgnoreCase))
{
errors.Add(
"AllowedHosts must contain the exact host advertised by Rendezvous:Deployment:PublicHttpBaseUrl.");
}
}
private static bool IsPublicHost(string host)
{
if (!IPAddress.TryParse(host, out IPAddress? address))
{
return Uri.CheckHostName(host) == UriHostNameType.Dns
&& host.Contains('.', StringComparison.Ordinal)
&& !IsReservedDnsName(host);
}
return IsGloballyRoutableUnicast(address);
}
private static bool IsReservedDnsName(string host)
{
string normalized = host.TrimEnd('.');
string[] reservedSuffixes =
[
"localhost",
"local",
"invalid",
"test",
"example",
"example.com",
"example.net",
"example.org",
"home.arpa",
"alt",
"onion",
];
return reservedSuffixes.Any(suffix =>
string.Equals(normalized, suffix, StringComparison.OrdinalIgnoreCase)
|| normalized.EndsWith($".{suffix}", StringComparison.OrdinalIgnoreCase));
}
private static bool IsGloballyRoutableUnicast(IPAddress address)
{
if (address.IsIPv4MappedToIPv6)
{
address = address.MapToIPv4();
}
byte[] bytes = address.GetAddressBytes();
if (address.AddressFamily == System.Net.Sockets.AddressFamily.InterNetwork)
{
return !(bytes[0] is 0 or 10 or 127
|| bytes[0] == 100 && bytes[1] is >= 64 and <= 127
|| bytes[0] == 169 && bytes[1] == 254
|| bytes[0] == 172 && bytes[1] is >= 16 and <= 31
|| bytes[0] == 192
&& (bytes[1] == 0 && bytes[2] is 0 or 2
|| bytes[1] == 88 && bytes[2] == 99
|| bytes[1] == 168)
|| bytes[0] == 198
&& (bytes[1] is 18 or 19
|| bytes[1] == 51 && bytes[2] == 100)
|| bytes[0] == 203 && bytes[1] == 0 && bytes[2] == 113
|| bytes[0] >= 224);
}
return address.AddressFamily == System.Net.Sockets.AddressFamily.InterNetworkV6
&& !IPAddress.IsLoopback(address)
&& !address.Equals(IPAddress.IPv6Any)
&& !address.IsIPv6LinkLocal
&& !address.IsIPv6SiteLocal
&& !address.IsIPv6Multicast
&& (bytes[0] & 0xe0) == 0x20
&& !HasPrefix(bytes, [0x20, 0x01, 0x00], 23)
&& !HasPrefix(bytes, [0x20, 0x01, 0x0d, 0xb8], 32)
&& !HasPrefix(bytes, [0x20, 0x02], 16)
&& !HasPrefix(bytes, [0x3f, 0xff, 0x00], 20);
}
private static bool HasPrefix(byte[] address, byte[] prefix, int bitCount)
{
int fullBytes = bitCount / 8;
for (int index = 0; index < fullBytes; index++)
{
if (address[index] != prefix[index])
{
return false;
}
}
int remainingBits = bitCount % 8;
if (remainingBits == 0)
{
return true;
}
int mask = 0xff << (8 - remainingBits);
return (address[fullBytes] & mask) == (prefix[fullBytes] & mask);
}
}
internal sealed class DeploymentConfigurationException(IReadOnlyList<string> errors)
: InvalidOperationException(
"Production deployment configuration is invalid:" + Environment.NewLine
+ string.Join(Environment.NewLine, errors.Select(static error => $"- {error}")))
{
}
@@ -0,0 +1,98 @@
using System.Diagnostics;
using FinalFactory.Rendezvous.Server.State;
using Microsoft.Extensions.Options;
namespace FinalFactory.Rendezvous.Server.Deployment;
internal sealed partial class GracefulDrainService : IHostedService, IDisposable
{
private static readonly TimeSpan PollInterval = TimeSpan.FromMilliseconds(50);
private readonly InMemoryEphemeralRendezvousStore _store;
private readonly IHostApplicationLifetime _lifetime;
private readonly DeploymentOptions _options;
private readonly ILogger<GracefulDrainService> _logger;
private readonly object _gate = new();
private CancellationTokenRegistration _stoppingRegistration;
private Task? _drainTask;
public GracefulDrainService(
InMemoryEphemeralRendezvousStore store,
IHostApplicationLifetime lifetime,
IOptions<DeploymentOptions> options,
ILogger<GracefulDrainService> logger)
{
_store = store;
_lifetime = lifetime;
_options = options.Value;
_logger = logger;
}
public Task StartAsync(CancellationToken cancellationToken)
{
cancellationToken.ThrowIfCancellationRequested();
_stoppingRegistration = _lifetime.ApplicationStopping.Register(
() => EnsureDrainAsync().GetAwaiter().GetResult());
return Task.CompletedTask;
}
public Task StopAsync(CancellationToken cancellationToken)
{
// ApplicationStopping callbacks run before hosted services and listeners
// stop. StopAsync is the idempotent fallback for directly driven hosts.
_ = cancellationToken;
return EnsureDrainAsync();
}
public void Dispose() => _stoppingRegistration.Dispose();
private Task EnsureDrainAsync()
{
lock (_gate)
{
return _drainTask ??= DrainAsync();
}
}
private async Task DrainAsync()
{
_store.BeginDrain(CancellationToken.None);
TimeSpan deadline = TimeSpan.FromSeconds(_options.DrainDeadlineSeconds);
TimeSpan minimum = TimeSpan.FromSeconds(_options.MinimumDrainSeconds);
long startedAt = Stopwatch.GetTimestamp();
LogDrainStarted(_logger, _options.DrainDeadlineSeconds);
try
{
while (Stopwatch.GetElapsedTime(startedAt) < deadline)
{
TimeSpan elapsed = Stopwatch.GetElapsedTime(startedAt);
if (elapsed >= minimum && _store.GetActiveJoinAttemptCountForDrain() == 0)
{
break;
}
TimeSpan remaining = deadline - elapsed;
await Task.Delay(
remaining < PollInterval ? remaining : PollInterval,
CancellationToken.None).ConfigureAwait(false);
}
}
finally
{
_store.MarkUnavailable();
double elapsedMilliseconds = Stopwatch.GetElapsedTime(startedAt).TotalMilliseconds;
LogDrainFinished(_logger, elapsedMilliseconds);
}
}
[LoggerMessage(
EventId = 1,
Level = LogLevel.Information,
Message = "Graceful drain started with a {DrainDeadlineSeconds}-second deadline")]
private static partial void LogDrainStarted(ILogger logger, int drainDeadlineSeconds);
[LoggerMessage(
EventId = 2,
Level = LogLevel.Information,
Message = "Graceful drain finished after {ElapsedMilliseconds:F0} ms; ephemeral state was cleared")]
private static partial void LogDrainFinished(ILogger logger, double elapsedMilliseconds);
}
@@ -4,7 +4,8 @@ using Microsoft.AspNetCore.Diagnostics;
namespace FinalFactory.Rendezvous.Server.Http;
internal sealed class RendezvousExceptionHandler : IExceptionHandler
internal sealed partial class RendezvousExceptionHandler(
ILogger<RendezvousExceptionHandler> logger) : IExceptionHandler
{
public async ValueTask<bool> TryHandleAsync(
HttpContext httpContext,
@@ -26,6 +27,13 @@ internal sealed class RendezvousExceptionHandler : IExceptionHandler
: invalidRequest
? StatusCodes.Status400BadRequest
: StatusCodes.Status500InternalServerError;
LogRequestFailure(
logger,
payloadTooLarge ? "payload-too-large" : invalidRequest ? "invalid-request" : "internal-error",
httpContext.Response.StatusCode,
httpContext.Response.Headers["X-Rendezvous-Correlation-ID"].ToString() is { Length: > 0 } value
? value
: "unavailable");
await httpContext.Response.WriteAsJsonAsync(
new ApiError
{
@@ -42,4 +50,14 @@ internal sealed class RendezvousExceptionHandler : IExceptionHandler
cancellationToken).ConfigureAwait(false);
return true;
}
[LoggerMessage(
EventId = 200,
Level = LogLevel.Warning,
Message = "Request failed with {FailureKind} and HTTP status {StatusCode}; correlation {CorrelationId}")]
private static partial void LogRequestFailure(
ILogger logger,
string failureKind,
int statusCode,
string correlationId);
}
@@ -0,0 +1,14 @@
using System.ComponentModel.DataAnnotations;
namespace FinalFactory.Rendezvous.Server.Observability;
internal sealed class AuditOptions
{
public const string SectionName = "Rendezvous:Audit";
[Range(100, 100_000)]
public int MaxEntries { get; set; } = 10_000;
[Range(1, 30)]
public int RetentionDays { get; set; } = 30;
}
@@ -0,0 +1,134 @@
using System.Security.Cryptography;
using System.Text;
using Microsoft.Extensions.Options;
namespace FinalFactory.Rendezvous.Server.Observability;
internal sealed partial class AuditTrail
{
private readonly object _gate = new();
private readonly LinkedList<AuditEntry> _entries = [];
private readonly AuditOptions _options;
private readonly TimeProvider _timeProvider;
private readonly ILogger<AuditTrail> _logger;
private readonly RendezvousTelemetry _telemetry;
public AuditTrail(
IOptions<AuditOptions> options,
ILogger<AuditTrail> logger,
RendezvousTelemetry telemetry,
TimeProvider? timeProvider = null)
{
_options = options.Value;
_logger = logger;
_telemetry = telemetry;
_timeProvider = timeProvider ?? TimeProvider.System;
}
public void Record(
string actorSubject,
string action,
string result,
string targetKind,
string targetIdentifier,
string correlationId)
{
DateTimeOffset now = _timeProvider.GetUtcNow();
AuditEntry entry = new(
now,
Fingerprint(actorSubject),
action,
result,
targetKind,
Fingerprint(targetIdentifier),
correlationId);
lock (_gate)
{
PurgeExpired(now);
while (_entries.Count >= _options.MaxEntries)
{
_entries.RemoveFirst();
}
_entries.AddLast(entry);
}
_telemetry.RecordAudit(action, result);
LogOperatorAction(
_logger,
entry.Timestamp,
entry.ActorFingerprint,
action,
result,
targetKind,
entry.TargetFingerprint,
correlationId);
}
public IReadOnlyDictionary<string, long> GetAggregateCounts()
{
lock (_gate)
{
PurgeExpired(_timeProvider.GetUtcNow());
return _entries
.GroupBy(static entry => $"{entry.Action}:{entry.Result}", StringComparer.Ordinal)
.ToDictionary(
static group => group.Key,
static group => (long)group.Count(),
StringComparer.Ordinal);
}
}
internal IReadOnlyList<AuditEntry> GetEntriesForTests()
{
lock (_gate)
{
PurgeExpired(_timeProvider.GetUtcNow());
return _entries.ToArray();
}
}
private void PurgeExpired(DateTimeOffset now)
{
DateTimeOffset oldest = now.AddDays(-_options.RetentionDays);
while (_entries.First is { Value.Timestamp: var timestamp }
&& timestamp < oldest)
{
_entries.RemoveFirst();
}
}
private static string Fingerprint(string value)
{
byte[] digest = SHA256.HashData(Encoding.UTF8.GetBytes(value));
return Convert.ToHexString(digest.AsSpan(0, 12));
}
[LoggerMessage(
EventId = 100,
Level = LogLevel.Information,
Message = "Operator audit at {Timestamp}: actor {ActorFingerprint} action {Action} completed with {Result} for {TargetKind} target {TargetFingerprint}; correlation {CorrelationId}")]
private static partial void LogOperatorAction(
ILogger logger,
DateTimeOffset timestamp,
string actorFingerprint,
string action,
string result,
string targetKind,
string targetFingerprint,
string correlationId);
}
internal sealed record AuditEntry(
DateTimeOffset Timestamp,
string ActorFingerprint,
string Action,
string Result,
string TargetKind,
string TargetFingerprint,
string CorrelationId)
{
public override string ToString() =>
$"[AuditEntry {Action}/{Result}; actor and target fingerprinted]";
}
@@ -0,0 +1,30 @@
using FinalFactory.Rendezvous.Contracts;
namespace FinalFactory.Rendezvous.Server.Observability;
internal static class HealthEndpoints
{
public static IEndpointRouteBuilder MapRendezvousHealthEndpoints(
this IEndpointRouteBuilder endpoints)
{
endpoints.MapGet(
"/health/live",
static () => Results.Ok(new HealthResponse { Status = "live" }))
.Produces<HealthResponse>()
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.WithName("GetLiveness")
.WithTags("Health");
endpoints.MapGet(
"/health/ready",
static (RendezvousReadiness readiness) =>
!readiness.GetSnapshot().IsReady
? Results.StatusCode(StatusCodes.Status503ServiceUnavailable)
: Results.Ok(new HealthResponse { Status = "ready" }))
.Produces<HealthResponse>()
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.Produces(StatusCodes.Status503ServiceUnavailable)
.WithName("GetReadiness")
.WithTags("Health");
return endpoints;
}
}
@@ -0,0 +1,41 @@
using FinalFactory.Rendezvous.Server.Provisioning;
using FinalFactory.Rendezvous.Server.State;
using FinalFactory.Rendezvous.Server.Transport;
using Microsoft.Extensions.Options;
namespace FinalFactory.Rendezvous.Server.Observability;
internal sealed class RendezvousReadiness(
UdpMediatorService mediator,
ProvisioningReadiness provisioning,
IEphemeralRendezvousStore state,
IOptions<UdpMediatorOptions> udpOptions)
{
public ReadinessSnapshot GetSnapshot()
{
bool ipv6Required = !string.IsNullOrWhiteSpace(udpOptions.Value.Ipv6ListenAddress);
return new ReadinessSnapshot(
HttpListenerReady: true,
UdpIpv4ListenerReady: mediator.LocalEndpoint is not null,
UdpIpv6ListenerReady: !ipv6Required || mediator.LocalIpv6Endpoint is not null,
ProvisioningReady: provisioning.IsReady,
StoreAvailable: state.IsAvailable,
Draining: state.IsDraining);
}
}
internal sealed record ReadinessSnapshot(
bool HttpListenerReady,
bool UdpIpv4ListenerReady,
bool UdpIpv6ListenerReady,
bool ProvisioningReady,
bool StoreAvailable,
bool Draining)
{
public bool IsReady => HttpListenerReady
&& UdpIpv4ListenerReady
&& UdpIpv6ListenerReady
&& ProvisioningReady
&& StoreAvailable
&& !Draining;
}
@@ -0,0 +1,126 @@
using System.Diagnostics;
using System.Diagnostics.Metrics;
using FinalFactory.Rendezvous.Server.State;
namespace FinalFactory.Rendezvous.Server.Observability;
internal sealed class RendezvousTelemetry : IDisposable
{
public const string MeterName = "FinalFactory.Rendezvous";
public const string ActivitySourceName = "FinalFactory.Rendezvous.Server";
private readonly InMemoryEphemeralRendezvousStore _store;
private readonly Meter _meter = new(MeterName, "1.0.0");
private readonly ActivitySource _activities = new(ActivitySourceName, "1.0.0");
private readonly Counter<long> _httpRequests;
private readonly Histogram<double> _httpDuration;
private readonly Counter<long> _udpResults;
private readonly Histogram<double> _udpDuration;
private readonly Counter<long> _limiterDrops;
private readonly Counter<long> _auditEvents;
private readonly Counter<long> _connectionOutcomes;
private readonly Counter<long> _operatorAuthentication;
private readonly Histogram<double> _pairingLatency;
public RendezvousTelemetry(InMemoryEphemeralRendezvousStore store)
{
_store = store;
_httpRequests = _meter.CreateCounter<long>("rendezvous.http.requests");
_httpDuration = _meter.CreateHistogram<double>(
"rendezvous.http.duration",
"ms");
_udpResults = _meter.CreateCounter<long>("rendezvous.udp.results");
_udpDuration = _meter.CreateHistogram<double>(
"rendezvous.udp.duration",
"ms");
_limiterDrops = _meter.CreateCounter<long>("rendezvous.limiter.drops");
_auditEvents = _meter.CreateCounter<long>("rendezvous.audit.events");
_connectionOutcomes = _meter.CreateCounter<long>("rendezvous.connection.outcomes");
_operatorAuthentication = _meter.CreateCounter<long>("rendezvous.operator.authentication");
_pairingLatency = _meter.CreateHistogram<double>(
"rendezvous.pairing.latency",
"ms");
_meter.CreateObservableGauge(
"rendezvous.store.active_listings",
() => _store.GetMetricsSnapshot().ActiveListings);
_meter.CreateObservableGauge(
"rendezvous.store.active_leases",
() => _store.GetMetricsSnapshot().ActiveListings);
_meter.CreateObservableGauge(
"rendezvous.store.active_attempts",
() => _store.GetMetricsSnapshot().ActiveJoinAttempts);
_meter.CreateObservableGauge(
"rendezvous.queue.depth",
() => _store.GetMetricsSnapshot().ActiveJoinAttempts);
_meter.CreateObservableGauge(
"rendezvous.store.replay_markers",
() => _store.GetMetricsSnapshot().ReplayMarkers);
_meter.CreateObservableGauge(
"rendezvous.store.available",
() => _store.GetMetricsSnapshot().IsAvailable ? 1 : 0);
_meter.CreateObservableCounter(
"rendezvous.store.expiry_churn",
() => _store.GetMetricsSnapshot().ExpiryChurn);
}
public Activity? StartActivity(string name, ActivityKind kind = ActivityKind.Internal) =>
_activities.StartActivity(name, kind);
public void RecordHttp(string operation, int statusCode, double elapsedMilliseconds)
{
TagList tags = new()
{
{ "operation", operation },
{ "status_code", statusCode },
};
_httpRequests.Add(1, tags);
_httpDuration.Record(elapsedMilliseconds, tags);
}
public void RecordUdp(string operation, string result, double elapsedMilliseconds)
{
TagList tags = new()
{
{ "operation", operation },
{ "result", result },
};
_udpResults.Add(1, tags);
_udpDuration.Record(elapsedMilliseconds, tags);
}
public void RecordLimiterDrop(string transport, string partition) =>
_limiterDrops.Add(1, new TagList
{
{ "transport", transport },
{ "partition", partition },
});
public void RecordAudit(string action, string result) =>
_auditEvents.Add(1, new TagList
{
{ "action", action },
{ "result", result },
});
public void RecordConnectionOutcome(string outcome, string elapsedBucket) =>
_connectionOutcomes.Add(1, new TagList
{
{ "outcome", outcome },
{ "elapsed_bucket", elapsedBucket },
});
public void RecordOperatorAuthentication(string result) =>
_operatorAuthentication.Add(1, new TagList
{
{ "result", result },
});
public void RecordPairingLatency(double elapsedMilliseconds) =>
_pairingLatency.Record(elapsedMilliseconds);
public void Dispose()
{
_activities.Dispose();
_meter.Dispose();
}
}
@@ -0,0 +1,32 @@
using System.Diagnostics;
namespace FinalFactory.Rendezvous.Server.Observability;
internal sealed class TelemetryMiddleware(
RequestDelegate next,
RendezvousTelemetry telemetry)
{
public async Task InvokeAsync(HttpContext context)
{
string operation = context.GetEndpoint()?.Metadata.GetMetadata<IEndpointNameMetadata>()
?.EndpointName ?? "Unmatched";
long started = Stopwatch.GetTimestamp();
using Activity? activity = telemetry.StartActivity(
$"HTTP {operation}",
ActivityKind.Server);
string correlationId = activity?.TraceId.ToString() ?? Guid.NewGuid().ToString("N");
context.Response.Headers["X-Rendezvous-Correlation-ID"] = correlationId;
activity?.SetTag("rendezvous.operation", operation);
try
{
await next(context).ConfigureAwait(false);
}
finally
{
telemetry.RecordHttp(
operation,
context.Response.StatusCode,
Stopwatch.GetElapsedTime(started).TotalMilliseconds);
}
}
}
@@ -0,0 +1,428 @@
using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Observability;
using FinalFactory.Rendezvous.Server.Provisioning;
using FinalFactory.Rendezvous.Server.State;
using Microsoft.AspNetCore.Mvc;
namespace FinalFactory.Rendezvous.Server.Operations;
internal static class OperatorEndpoints
{
private const string CorrelationHeader = "X-Rendezvous-Correlation-ID";
public static IEndpointRouteBuilder MapOperatorEndpoints(this IEndpointRouteBuilder endpoints)
{
RouteGroupBuilder group = endpoints.MapGroup("/v1/operator").WithTags("Operator");
group.MapGet("/status", GetStatus)
.Produces<OperatorStatusResponse>()
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
.Produces<ApiError>(StatusCodes.Status403Forbidden)
.Produces<ApiError>(StatusCodes.Status404NotFound)
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.WithName("GetOperatorStatus");
group.MapPost("/listings/revoke", RevokeListing)
.Accepts<RevokeListingRequest>("application/json")
.Produces<OperatorActionResponse>()
.Produces<ApiError>(StatusCodes.Status400BadRequest)
.Produces<ApiError>(StatusCodes.Status413PayloadTooLarge)
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
.Produces<ApiError>(StatusCodes.Status403Forbidden)
.Produces<ApiError>(StatusCodes.Status404NotFound)
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.Produces<ApiError>(StatusCodes.Status503ServiceUnavailable)
.WithName("RevokeOperatorListing");
group.MapPost("/principals/revoke", RevokePrincipal)
.Accepts<RevokePrincipalRequest>("application/json")
.Produces<OperatorActionResponse>()
.Produces<ApiError>(StatusCodes.Status400BadRequest)
.Produces<ApiError>(StatusCodes.Status413PayloadTooLarge)
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
.Produces<ApiError>(StatusCodes.Status403Forbidden)
.Produces<ApiError>(StatusCodes.Status404NotFound)
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.Produces<ApiError>(StatusCodes.Status503ServiceUnavailable)
.WithName("RevokeOperatorPrincipal");
group.MapPost("/keys/revoke", RevokeSigningKey)
.Accepts<RevokeSigningKeyRequest>("application/json")
.Produces<OperatorActionResponse>()
.Produces<ApiError>(StatusCodes.Status400BadRequest)
.Produces<ApiError>(StatusCodes.Status413PayloadTooLarge)
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
.Produces<ApiError>(StatusCodes.Status403Forbidden)
.Produces<ApiError>(StatusCodes.Status404NotFound)
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.WithName("RevokeOperatorSigningKey");
group.MapPost("/drain", BeginDrain)
.Accepts<BeginDrainRequest>("application/json")
.Produces<OperatorActionResponse>()
.Produces<ApiError>(StatusCodes.Status400BadRequest)
.Produces<ApiError>(StatusCodes.Status413PayloadTooLarge)
.Produces<ApiError>(StatusCodes.Status401Unauthorized)
.Produces<ApiError>(StatusCodes.Status403Forbidden)
.Produces<ApiError>(StatusCodes.Status404NotFound)
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.WithName("BeginOperatorDrain");
return endpoints;
}
private static IResult GetStatus(
[FromHeader(Name = "Authorization")] string? authorization,
[FromServices] PrincipalCredentialService credentials,
[FromServices] IWallClock clock,
[FromServices] OperatorService service,
[FromServices] AuditTrail audit,
[FromServices] RendezvousTelemetry telemetry,
HttpContext context)
{
if (!TryAuthorize(
authorization,
OperatorPermission.ReadPolicy,
"inspect-status",
credentials,
clock,
audit,
telemetry,
context,
out OperatorPrincipal? principal,
out IResult? failure))
{
return failure!;
}
OperatorStatusResponse response = service.GetStatus();
audit.Record(
principal!.Subject,
"inspect-status",
"succeeded",
"service",
"rendezvous",
Correlation(context));
return Results.Ok(response);
}
private static IResult RevokeListing(
[FromBody] RevokeListingRequest request,
[FromHeader(Name = "Authorization")] string? authorization,
[FromServices] PrincipalCredentialService credentials,
[FromServices] IWallClock clock,
[FromServices] OperatorService service,
[FromServices] AuditTrail audit,
[FromServices] RendezvousTelemetry telemetry,
HttpContext context,
CancellationToken cancellationToken)
{
if (!TryAuthorize(
authorization,
OperatorPermission.RevokePublisher,
"revoke-listing",
credentials,
clock,
audit,
telemetry,
context,
out OperatorPrincipal? principal,
out IResult? failure))
{
return failure!;
}
bool valid = SessionListingId.TryParse(request.ListingId, out SessionListingId listingId);
if (!valid || !string.Equals(request.ListingId, request.ConfirmListingId, StringComparison.Ordinal))
{
AuditRejected(audit, principal!, "revoke-listing", "listing", request.ListingId, context);
return BadRequest("A valid listing ID and an exact repeated confirmation are required.");
}
StoreResult<bool> result = service.RevokeListing(listingId, cancellationToken);
return StoreActionResult(
result.Code,
audit,
principal!,
"revoke-listing",
"listing",
request.ListingId,
context,
affectedResources: result.Succeeded ? 1 : null);
}
private static IResult RevokePrincipal(
[FromBody] RevokePrincipalRequest request,
[FromHeader(Name = "Authorization")] string? authorization,
[FromServices] PrincipalCredentialService credentials,
[FromServices] IWallClock clock,
[FromServices] OperatorService service,
[FromServices] AuditTrail audit,
[FromServices] RendezvousTelemetry telemetry,
HttpContext context,
CancellationToken cancellationToken)
{
if (!TryAuthorize(
authorization,
OperatorPermission.RevokePublisher,
"revoke-principal",
credentials,
clock,
audit,
telemetry,
context,
out OperatorPrincipal? principal,
out IResult? failure))
{
return failure!;
}
bool safeSubject = request.Subject is { Length: > 0 and <= 128 }
&& request.Subject.All(static character => character is >= '!' and <= '~');
if (!safeSubject
|| !string.Equals(request.Subject, request.ConfirmSubject, StringComparison.Ordinal)
|| request.LifetimeSeconds is < 1 or > 600)
{
AuditRejected(audit, principal!, "revoke-principal", "principal", request.Subject, context);
return BadRequest("A valid subject, exact repeated confirmation, and 1-600 second lifetime are required.");
}
StoreResult<int> result = service.RevokePrincipal(
request.Subject,
TimeSpan.FromSeconds(request.LifetimeSeconds),
cancellationToken);
return StoreActionResult(
result.Code,
audit,
principal!,
"revoke-principal",
"principal",
request.Subject,
context,
result.Value);
}
private static IResult RevokeSigningKey(
[FromBody] RevokeSigningKeyRequest request,
[FromHeader(Name = "Authorization")] string? authorization,
[FromServices] PrincipalCredentialService credentials,
[FromServices] IWallClock clock,
[FromServices] OperatorService service,
[FromServices] AuditTrail audit,
[FromServices] RendezvousTelemetry telemetry,
HttpContext context)
{
if (!TryAuthorize(
authorization,
OperatorPermission.RotateKeys,
"revoke-signing-key",
credentials,
clock,
audit,
telemetry,
context,
out OperatorPrincipal? principal,
out IResult? failure))
{
return failure!;
}
bool safeKeyId = request.KeyId is { Length: > 0 and <= 64 }
&& request.KeyId.All(static character => character is
>= 'A' and <= 'Z'
or >= 'a' and <= 'z'
or >= '0' and <= '9'
or '-'
or '_');
if (!safeKeyId || !string.Equals(request.KeyId, request.ConfirmKeyId, StringComparison.Ordinal))
{
AuditRejected(audit, principal!, "revoke-signing-key", "signing-key", request.KeyId, context);
return BadRequest("A valid key ID and an exact repeated confirmation are required.");
}
bool revoked = service.RevokeSigningKey(request.KeyId);
string result = revoked ? "succeeded" : "not-found";
audit.Record(
principal!.Subject,
"revoke-signing-key",
result,
"signing-key",
request.KeyId,
Correlation(context));
return revoked
? Results.Ok(new OperatorActionResponse { Status = "completed" })
: Error(RendezvousErrorCode.NotFound, "The requested resource was not found.");
}
private static IResult BeginDrain(
[FromBody] BeginDrainRequest request,
[FromHeader(Name = "Authorization")] string? authorization,
[FromServices] PrincipalCredentialService credentials,
[FromServices] IWallClock clock,
[FromServices] OperatorService service,
[FromServices] AuditTrail audit,
[FromServices] RendezvousTelemetry telemetry,
HttpContext context,
CancellationToken cancellationToken)
{
if (!TryAuthorize(
authorization,
OperatorPermission.ManagePolicy,
"begin-drain",
credentials,
clock,
audit,
telemetry,
context,
out OperatorPrincipal? principal,
out IResult? failure))
{
return failure!;
}
if (!string.Equals(request.Confirmation, "DRAIN", StringComparison.Ordinal))
{
AuditRejected(audit, principal!, "begin-drain", "service", "rendezvous", context);
return BadRequest("The confirmation value must be exactly 'DRAIN'.");
}
service.BeginDrain(cancellationToken);
audit.Record(
principal!.Subject,
"begin-drain",
"succeeded",
"service",
"rendezvous",
Correlation(context));
return Results.Ok(new OperatorActionResponse { Status = "draining" });
}
private static bool TryAuthorize(
string? authorization,
OperatorPermission requiredPermission,
string operation,
PrincipalCredentialService credentials,
IWallClock clock,
AuditTrail audit,
RendezvousTelemetry telemetry,
HttpContext context,
out OperatorPrincipal? principal,
out IResult? failure)
{
principal = null;
failure = null;
const string prefix = "Bearer ";
if (authorization is null
|| !authorization.StartsWith(prefix, StringComparison.OrdinalIgnoreCase))
{
telemetry.RecordOperatorAuthentication("rejected");
failure = AuthenticationRequired(context);
return false;
}
CredentialValidationResult validation = credentials.Validate(
authorization[prefix.Length..],
clock.UtcNow);
if (!validation.IsValid || validation.Principal is not OperatorPrincipal candidate)
{
telemetry.RecordOperatorAuthentication("rejected");
failure = AuthenticationRequired(context);
return false;
}
if (!candidate.Permissions.Contains(requiredPermission))
{
telemetry.RecordOperatorAuthentication("forbidden");
audit.Record(
candidate.Subject,
operation,
"forbidden",
"operator-operation",
operation,
Correlation(context));
failure = Error(RendezvousErrorCode.Forbidden, "The operator is not authorized for this operation.");
return false;
}
telemetry.RecordOperatorAuthentication("accepted");
principal = candidate;
return true;
}
private static IResult StoreActionResult(
StoreResultCode code,
AuditTrail audit,
OperatorPrincipal principal,
string action,
string targetKind,
string targetIdentifier,
HttpContext context,
int? affectedResources)
{
string auditResult = code == StoreResultCode.Success
? "succeeded"
: code.ToString().ToLowerInvariant();
audit.Record(
principal.Subject,
action,
auditResult,
targetKind,
targetIdentifier,
Correlation(context));
return code switch
{
StoreResultCode.Success => Results.Ok(new OperatorActionResponse
{
Status = "completed",
AffectedResources = affectedResources,
}),
StoreResultCode.NotFound => Error(
RendezvousErrorCode.NotFound,
"The requested resource was not found."),
StoreResultCode.CapacityExceeded => Error(
RendezvousErrorCode.CapacityExceeded,
"The operation could not be retained within the configured capacity."),
StoreResultCode.ServiceUnavailable or StoreResultCode.Draining => Error(
RendezvousErrorCode.ServiceUnavailable,
"The service is not available for this operation."),
_ => Error(RendezvousErrorCode.Conflict, "The operation could not be completed."),
};
}
private static void AuditRejected(
AuditTrail audit,
OperatorPrincipal principal,
string action,
string targetKind,
string? targetIdentifier,
HttpContext context) => audit.Record(
principal.Subject,
action,
"rejected",
targetKind,
targetIdentifier ?? string.Empty,
Correlation(context));
private static string Correlation(HttpContext context) =>
context.Response.Headers[CorrelationHeader].ToString() is { Length: > 0 } value
? value
: "unavailable";
private static IResult AuthenticationRequired(HttpContext context)
{
context.Response.Headers.WWWAuthenticate = "Bearer realm=\"operator\"";
return Error(
RendezvousErrorCode.AuthenticationRequired,
"A valid operator bearer credential is required.");
}
private static IResult BadRequest(string message) => Error(RendezvousErrorCode.InvalidRequest, message);
private static IResult Error(RendezvousErrorCode code, string message) => Results.Json(
new ApiError { Code = code, Message = message },
ContractJson.Options,
statusCode: code switch
{
RendezvousErrorCode.AuthenticationRequired => StatusCodes.Status401Unauthorized,
RendezvousErrorCode.Forbidden => StatusCodes.Status403Forbidden,
RendezvousErrorCode.NotFound => StatusCodes.Status404NotFound,
RendezvousErrorCode.CapacityExceeded => StatusCodes.Status429TooManyRequests,
RendezvousErrorCode.ServiceUnavailable => StatusCodes.Status503ServiceUnavailable,
RendezvousErrorCode.Conflict => StatusCodes.Status409Conflict,
_ => StatusCodes.Status400BadRequest,
});
}
@@ -0,0 +1,82 @@
namespace FinalFactory.Rendezvous.Server.Operations;
internal sealed record OperatorStatusResponse
{
public required string Status { get; init; }
public required OperatorReadinessResponse Readiness { get; init; }
public required OperatorStoreResponse Store { get; init; }
public required IReadOnlyList<OperatorTenantResponse> Tenants { get; init; }
public required IReadOnlyList<OperatorSigningKeyResponse> SigningKeys { get; init; }
public required IReadOnlyDictionary<string, long> AuditCounts { get; init; }
}
internal sealed record OperatorReadinessResponse
{
public required bool HttpListener { get; init; }
public required bool UdpIpv4Listener { get; init; }
public required bool UdpIpv6Listener { get; init; }
public required bool Provisioning { get; init; }
public required bool Store { get; init; }
public required bool Draining { get; init; }
}
internal sealed record OperatorStoreResponse
{
public required int ActiveListings { get; init; }
public required int FreshPresenceBindings { get; init; }
public required int ActiveJoinAttempts { get; init; }
public required int RetainedOutcomeReports { get; init; }
public required int ReplayMarkers { get; init; }
public required int PrincipalRevocations { get; init; }
public required int IdempotencyEntries { get; init; }
public required long MaintenanceSweeps { get; init; }
public required long ExpiryChurn { get; init; }
}
internal sealed record OperatorTenantResponse
{
public required string GameId { get; init; }
public required string EnvironmentId { get; init; }
public required string Status { get; init; }
}
internal sealed record OperatorSigningKeyResponse
{
public required string KeyId { get; init; }
public required string Status { get; init; }
public required DateTimeOffset SignUntil { get; init; }
public required DateTimeOffset VerifyUntil { get; init; }
public string? GameId { get; init; }
public string? EnvironmentId { get; init; }
public required IReadOnlyList<string> CredentialKinds { get; init; }
}
internal sealed record OperatorActionResponse
{
public required string Status { get; init; }
public int? AffectedResources { get; init; }
}
internal sealed record RevokeListingRequest
{
public required string ListingId { get; init; }
public required string ConfirmListingId { get; init; }
}
internal sealed record RevokePrincipalRequest
{
public required string Subject { get; init; }
public required string ConfirmSubject { get; init; }
public required int LifetimeSeconds { get; init; }
}
internal sealed record RevokeSigningKeyRequest
{
public required string KeyId { get; init; }
public required string ConfirmKeyId { get; init; }
}
internal sealed record BeginDrainRequest
{
public required string Confirmation { get; init; }
}
@@ -0,0 +1,81 @@
using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Observability;
using FinalFactory.Rendezvous.Server.Provisioning;
using FinalFactory.Rendezvous.Server.State;
namespace FinalFactory.Rendezvous.Server.Operations;
internal sealed class OperatorService(
InMemoryEphemeralRendezvousStore store,
ProvisioningRuntime provisioning,
RendezvousReadiness readiness,
AuditTrail audit,
IWallClock clock)
{
public OperatorStatusResponse GetStatus()
{
ReadinessSnapshot readinessSnapshot = readiness.GetSnapshot();
EphemeralStoreSnapshot storeSnapshot = store.GetSnapshot();
return new OperatorStatusResponse
{
Status = readinessSnapshot.IsReady ? "ready" : "not-ready",
Readiness = new OperatorReadinessResponse
{
HttpListener = readinessSnapshot.HttpListenerReady,
UdpIpv4Listener = readinessSnapshot.UdpIpv4ListenerReady,
UdpIpv6Listener = readinessSnapshot.UdpIpv6ListenerReady,
Provisioning = readinessSnapshot.ProvisioningReady,
Store = readinessSnapshot.StoreAvailable,
Draining = readinessSnapshot.Draining,
},
Store = new OperatorStoreResponse
{
ActiveListings = storeSnapshot.ActiveListings,
FreshPresenceBindings = storeSnapshot.FreshPresenceBindings,
ActiveJoinAttempts = storeSnapshot.ActiveJoinAttempts,
RetainedOutcomeReports = storeSnapshot.RetainedOutcomeReports,
ReplayMarkers = storeSnapshot.ReplayMarkers,
PrincipalRevocations = storeSnapshot.PrincipalRevocations,
IdempotencyEntries = storeSnapshot.IdempotencyEntries,
MaintenanceSweeps = storeSnapshot.MaintenanceSweeps,
ExpiryChurn = storeSnapshot.ExpiryChurn,
},
Tenants = provisioning.Policies.EnabledPolicies
.OrderBy(static policy => policy.GameId.Value, StringComparer.Ordinal)
.ThenBy(static policy => policy.EnvironmentId.Value, StringComparer.Ordinal)
.Select(static policy => new OperatorTenantResponse
{
GameId = policy.GameId.Value,
EnvironmentId = policy.EnvironmentId.Value,
Status = "enabled",
})
.ToArray(),
SigningKeys = provisioning.SigningKeys.GetStatuses(clock.UtcNow)
.Select(static key => new OperatorSigningKeyResponse
{
KeyId = key.KeyId,
Status = key.Status,
SignUntil = key.SignUntil,
VerifyUntil = key.VerifyUntil,
GameId = key.GameId,
EnvironmentId = key.EnvironmentId,
CredentialKinds = key.CredentialKinds,
})
.ToArray(),
AuditCounts = audit.GetAggregateCounts(),
};
}
public StoreResult<bool> RevokeListing(
SessionListingId listingId,
CancellationToken cancellationToken) => store.RevokeListing(listingId, cancellationToken);
public StoreResult<int> RevokePrincipal(
string subject,
TimeSpan lifetime,
CancellationToken cancellationToken) => store.RevokePrincipal(subject, lifetime, cancellationToken);
public bool RevokeSigningKey(string keyId) => provisioning.SigningKeys.Revoke(keyId);
public void BeginDrain(CancellationToken cancellationToken) => store.BeginDrain(cancellationToken);
}
+111 -41
View File
@@ -3,8 +3,11 @@ using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Abuse;
using FinalFactory.Rendezvous.Server.Browser;
using FinalFactory.Rendezvous.Server.ConnectionOutcomes;
using FinalFactory.Rendezvous.Server.Deployment;
using FinalFactory.Rendezvous.Server.Http;
using FinalFactory.Rendezvous.Server.JoinAttempts;
using FinalFactory.Rendezvous.Server.Observability;
using FinalFactory.Rendezvous.Server.Operations;
using FinalFactory.Rendezvous.Server.Provisioning;
using FinalFactory.Rendezvous.Server.Sessions;
using FinalFactory.Rendezvous.Server.State;
@@ -13,6 +16,9 @@ using Microsoft.AspNetCore.HttpOverrides;
using Microsoft.OpenApi;
WebApplicationBuilder builder = WebApplication.CreateBuilder(args);
builder.Logging.AddFilter(
"Microsoft.AspNetCore.Diagnostics.ExceptionHandlerMiddleware",
LogLevel.None);
bool isOpenApiGeneration = string.Equals(
System.Reflection.Assembly.GetEntryAssembly()?.GetName().Name,
"GetDocument.Insider",
@@ -61,6 +67,14 @@ builder.Services.AddOpenApi("v1", static options =>
In = ParameterLocation.Header,
Description = "Attempt-scoped client capability returned only to the joining caller.",
};
const string operatorSchemeName = "OperatorBearer";
document.Components.SecuritySchemes[operatorSchemeName] = new OpenApiSecurityScheme
{
Type = SecuritySchemeType.Http,
Scheme = "bearer",
BearerFormat = "rv1 operator credential",
Description = "Operator-only credential with an explicit permission set.",
};
HashSet<string> securedOperations = new(StringComparer.Ordinal)
{
@@ -69,8 +83,17 @@ builder.Services.AddOpenApi("v1", static options =>
"UpdateSession",
"DeleteSession",
};
HashSet<string> operatorOperations = new(StringComparer.Ordinal)
{
"GetOperatorStatus",
"RevokeOperatorListing",
"RevokeOperatorPrincipal",
"RevokeOperatorSigningKey",
"BeginOperatorDrain",
};
OpenApiSecuritySchemeReference reference = new(schemeName, document, null);
OpenApiSecuritySchemeReference attemptReference = new(attemptSchemeName, document, null);
OpenApiSecuritySchemeReference operatorReference = new(operatorSchemeName, document, null);
foreach (OpenApiPathItem path in document.Paths.Values)
{
if (path.Operations is null)
@@ -99,29 +122,58 @@ builder.Services.AddOpenApi("v1", static options =>
});
}
foreach (OpenApiOperation operation in path.Operations.Values.Where(
operation => operatorOperations.Contains(
operation.OperationId ?? string.Empty)))
{
operation.Security ??= [];
operation.Security.Add(new OpenApiSecurityRequirement
{
[operatorReference] = [],
});
}
foreach (OpenApiOperation operation in path.Operations.Values)
{
if (operation.Responses is null
|| !operation.Responses.TryGetValue(
StatusCodes.Status429TooManyRequests.ToString(
System.Globalization.CultureInfo.InvariantCulture),
out IOpenApiResponse? response)
|| response is not OpenApiResponse concreteResponse)
if (operation.Responses is null)
{
continue;
}
concreteResponse.Headers ??=
new Dictionary<string, IOpenApiHeader>(StringComparer.OrdinalIgnoreCase);
concreteResponse.Headers["Retry-After"] = new OpenApiHeader
foreach ((string status, IOpenApiResponse response) in operation.Responses)
{
Description = "Whole seconds before the caller should retry (1-60).",
Schema = new OpenApiSchema
if (response is not OpenApiResponse concreteResponse)
{
Type = JsonSchemaType.Integer,
Format = "int32",
},
};
continue;
}
concreteResponse.Headers ??=
new Dictionary<string, IOpenApiHeader>(StringComparer.OrdinalIgnoreCase);
concreteResponse.Headers["X-Rendezvous-Correlation-ID"] = new OpenApiHeader
{
Description = "Safe request correlation identifier generated by the service.",
Schema = new OpenApiSchema
{
Type = JsonSchemaType.String,
},
};
if (string.Equals(
status,
StatusCodes.Status429TooManyRequests.ToString(
System.Globalization.CultureInfo.InvariantCulture),
StringComparison.Ordinal))
{
concreteResponse.Headers["Retry-After"] = new OpenApiHeader
{
Description = "Whole seconds before the caller should retry (1-60).",
Schema = new OpenApiSchema
{
Type = JsonSchemaType.Integer,
Format = "int32",
},
};
}
}
}
}
@@ -166,22 +218,58 @@ builder.Services
&& addresses.All(
static value => IPAddress.TryParse(value, out _)),
"Trusted proxy addresses must contain at most 32 literal IP addresses.")
.Validate(
options => options.OperatorAllowedAddresses is { Length: <= 32 } addresses
&& addresses.All(
static value => IPAddress.TryParse(value, out _)),
"Operator allowed addresses must contain at most 32 literal IP addresses.")
.ValidateOnStart();
builder.Services.AddSingleton<AbuseProtectionService>();
builder.Services
.AddOptions<AuditOptions>()
.BindConfiguration(AuditOptions.SectionName)
.ValidateDataAnnotations()
.ValidateOnStart();
AbuseProtectionOptions configuredAbuseProtection = builder.Configuration
.GetSection(AbuseProtectionOptions.SectionName)
.Get<AbuseProtectionOptions>() ?? new AbuseProtectionOptions();
builder.Services.Configure<ForwardedHeadersOptions>(options =>
TrustedProxyForwarding.Configure(options, configuredAbuseProtection));
DeploymentOptions deploymentOptions = builder.Configuration
.GetSection(DeploymentOptions.SectionName)
.Get<DeploymentOptions>() ?? new DeploymentOptions();
if (!builder.Environment.IsDevelopment() && !isOpenApiGeneration)
{
IReadOnlyList<string> deploymentErrors = deploymentOptions.ValidateProduction(
configuredAbuseProtection,
builder.Configuration["AllowedHosts"]);
if (deploymentErrors.Count > 0)
{
throw new DeploymentConfigurationException(deploymentErrors);
}
}
builder.Services.AddSingleton(Microsoft.Extensions.Options.Options.Create(deploymentOptions));
builder.Services.Configure<HostOptions>(options =>
options.ShutdownTimeout = TimeSpan.FromSeconds(deploymentOptions.DrainDeadlineSeconds + 10));
SystemRendezvousClock rendezvousClock = new();
EphemeralStoreOptions stateOptions = new();
EphemeralStoreOptions stateOptions = new()
{
GracefulDrainLifetime = TimeSpan.FromSeconds(deploymentOptions.DrainDeadlineSeconds),
};
InMemoryEphemeralRendezvousStore stateStore = new(
stateOptions,
rendezvousClock,
rendezvousClock);
builder.Services.AddSingleton(stateStore);
builder.Services.AddSingleton<IEphemeralRendezvousStore>(stateStore);
builder.Services.AddSingleton<IWallClock>(rendezvousClock);
builder.Services.AddSingleton<IMonotonicClock>(rendezvousClock);
builder.Services.AddSingleton<RendezvousTelemetry>();
builder.Services.AddSingleton<AuditTrail>();
builder.Services.AddSingleton<RendezvousReadiness>();
if (isOpenApiGeneration)
{
@@ -214,6 +302,7 @@ else
builder.Services.AddSingleton<JoinAttemptService>();
builder.Services.AddSingleton<ConnectionOutcomeMetrics>();
builder.Services.AddSingleton<ConnectionOutcomeService>();
builder.Services.AddSingleton<OperatorService>();
builder.Services.AddSingleton(new ProvisioningReadiness(true));
}
@@ -237,43 +326,24 @@ if (!isOpenApiGeneration)
builder.Services.AddSingleton<NatMediationProcessor>();
builder.Services.AddHostedService(static services =>
services.GetRequiredService<UdpMediatorService>());
// Hosted services stop in reverse registration order. Drain must complete while
// Kestrel and the UDP mediator are still able to finish bounded in-flight work.
builder.Services.AddHostedService<GracefulDrainService>();
}
WebApplication app = builder.Build();
app.Lifetime.ApplicationStopping.Register(() => stateStore.BeginDrain());
if (TrustedProxyForwarding.IsEnabled(configuredAbuseProtection))
{
app.UseForwardedHeaders();
}
app.UseMiddleware<TelemetryMiddleware>();
app.UseExceptionHandler();
app.UseMiddleware<HttpAbuseProtectionMiddleware>();
app.MapOpenApi();
app.MapRendezvousContractEndpoints();
app.MapGet(
"/health/live",
static () => Results.Ok(new HealthResponse { Status = "live" }))
.Produces<HealthResponse>()
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.WithName("GetLiveness")
.WithTags("Health");
app.MapGet(
"/health/ready",
static (
UdpMediatorService mediator,
ProvisioningReadiness provisioning,
IEphemeralRendezvousStore state) =>
mediator.LocalEndpoint is null
|| !provisioning.IsReady
|| !state.IsAvailable
|| state.IsDraining
? Results.StatusCode(StatusCodes.Status503ServiceUnavailable)
: Results.Ok(new HealthResponse { Status = "ready" }))
.Produces<HealthResponse>()
.Produces<ApiError>(StatusCodes.Status429TooManyRequests)
.Produces(StatusCodes.Status503ServiceUnavailable)
.WithName("GetReadiness")
.WithTags("Health");
app.MapOperatorEndpoints();
app.MapRendezvousHealthEndpoints();
await app.RunAsync();
@@ -1,3 +1,4 @@
using System.Runtime.CompilerServices;
[assembly: InternalsVisibleTo("FinalFactory.Rendezvous.Tests")]
[assembly: InternalsVisibleTo("FinalFactory.Rendezvous.Capacity")]
@@ -40,18 +40,32 @@ internal sealed class SecretMaterial : IDisposable
internal sealed class EnvironmentSecretProvider : ISecretProvider
{
private const string Prefix = "env:";
private const string EnvironmentPrefix = "env:";
private const string FilePrefix = "file:";
private const int MaximumSecretBytes = 4096;
public bool TryGetSecret(string reference, out SecretMaterial? secret)
{
secret = null;
if (!reference.StartsWith(Prefix, StringComparison.Ordinal)
|| reference.Length == Prefix.Length)
if (reference.StartsWith(EnvironmentPrefix, StringComparison.Ordinal)
&& reference.Length > EnvironmentPrefix.Length)
{
return false;
return TryGetEnvironmentSecret(reference[EnvironmentPrefix.Length..], out secret);
}
string? encoded = Environment.GetEnvironmentVariable(reference[Prefix.Length..]);
if (reference.StartsWith(FilePrefix, StringComparison.Ordinal)
&& reference.Length > FilePrefix.Length)
{
return TryGetFileSecret(reference[FilePrefix.Length..], out secret);
}
return false;
}
private static bool TryGetEnvironmentSecret(string variableName, out SecretMaterial? secret)
{
secret = null;
string? encoded = Environment.GetEnvironmentVariable(variableName);
if (string.IsNullOrEmpty(encoded))
{
return false;
@@ -60,6 +74,12 @@ internal sealed class EnvironmentSecretProvider : ISecretProvider
try
{
byte[] bytes = Convert.FromBase64String(encoded);
if (bytes.Length is 0 or > MaximumSecretBytes)
{
CryptographicOperations.ZeroMemory(bytes);
return false;
}
secret = new SecretMaterial(bytes);
CryptographicOperations.ZeroMemory(bytes);
return true;
@@ -69,6 +89,47 @@ internal sealed class EnvironmentSecretProvider : ISecretProvider
return false;
}
}
private static bool TryGetFileSecret(string path, out SecretMaterial? secret)
{
secret = null;
byte[]? bytes = null;
try
{
FileInfo file = new(path);
if (!file.Exists
|| !Path.IsPathFullyQualified(path)
|| file.LinkTarget is not null
|| file.Length is <= 0 or > MaximumSecretBytes)
{
return false;
}
bytes = File.ReadAllBytes(path);
if (bytes.Length is 0 or > MaximumSecretBytes)
{
return false;
}
secret = new SecretMaterial(bytes);
return true;
}
catch (Exception exception) when (exception is IOException
or UnauthorizedAccessException
or ArgumentException
or NotSupportedException
or System.Security.SecurityException)
{
return false;
}
finally
{
if (bytes is not null)
{
CryptographicOperations.ZeroMemory(bytes);
}
}
}
}
internal sealed class EphemeralDevelopmentSecretProvider : ISecretProvider, IDisposable
@@ -142,8 +142,36 @@ internal sealed class SigningKeyRing : IDisposable
return VerificationKeyLookup.Available;
}
public bool Revoke(string keyId) =>
_keys.ContainsKey(keyId) && _runtimeRevocations.TryAdd(keyId, 0);
public bool Revoke(string keyId)
{
if (!_keys.ContainsKey(keyId))
{
return false;
}
_runtimeRevocations.TryAdd(keyId, 0);
return true;
}
public IReadOnlyList<SigningKeyStatus> GetStatuses(DateTimeOffset now) => _keys.Values
.OrderBy(static key => key.KeyId, StringComparer.Ordinal)
.Select(key => new SigningKeyStatus(
key.KeyId,
IsRevoked(key)
? "revoked"
: now < key.NotBefore
? "not-yet-valid"
: now < key.SignUntil
? "signing"
: now < key.VerifyUntil
? "verify-only"
: "retired",
key.SignUntil,
key.VerifyUntil,
key.GameId,
key.EnvironmentId,
key.CredentialKinds.Select(static kind => kind.ToString()).Order().ToArray()))
.ToArray();
public void Dispose()
{
@@ -210,6 +238,15 @@ internal sealed class SigningKeyRing : IDisposable
}
}
internal sealed record SigningKeyStatus(
string KeyId,
string Status,
DateTimeOffset SignUntil,
DateTimeOffset VerifyUntil,
string? GameId,
string? EnvironmentId,
IReadOnlyList<string> CredentialKinds);
internal sealed class SigningKey : IDisposable
{
private byte[]? _material;
@@ -297,6 +297,7 @@ internal sealed record StoredJoinAttempt
public required SecretFingerprint ConnectionTicketFingerprint { get; init; }
public NetworkEndpoint? DedicatedFallback { get; init; }
public required DateTimeOffset ExpiresAt { get; init; }
public required TimeSpan CreatedAtMonotonic { get; init; }
public required DateTimeOffset ConnectionTicketExpiresAt { get; init; }
public AttemptEndpointBinding? HostEndpoint { get; init; }
public AttemptEndpointBinding? ClientEndpoint { get; init; }
@@ -11,19 +11,36 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
private readonly DateTimeOffset _wallOrigin;
private readonly TimeSpan _monotonicOrigin;
private readonly Dictionary<SessionListingId, ListingEntry> _listings = [];
private readonly Dictionary<string, int> _listingCountsByOwner = new(StringComparer.Ordinal);
private readonly PriorityQueue<DeadlineEntry<SessionListingId>, long> _listingExpiries = new();
private readonly HashSet<SessionListingId> _scheduledListingExpiries = [];
private readonly Dictionary<LeaseId, SessionListingId> _leases = [];
private readonly Dictionary<MediationHandle, SessionListingId> _presenceHandles = [];
private readonly Dictionary<MediationHandle, PresenceEntry> _presence = [];
private readonly PriorityQueue<DeadlineEntry<MediationHandle>, long> _presenceExpiries = new();
private readonly HashSet<MediationHandle> _scheduledPresenceExpiries = [];
private readonly Dictionary<JoinAttemptId, AttemptEntry> _attempts = [];
private readonly Dictionary<TenantScope, int> _attemptCountsByScope = [];
private readonly Dictionary<SessionListingId, HashSet<JoinAttemptId>> _attemptsByListing = [];
private readonly PriorityQueue<DeadlineEntry<JoinAttemptId>, long> _attemptExpiries = new();
private readonly Dictionary<JoinAttemptId, OutcomeReportEntry> _outcomeReports = [];
private readonly Dictionary<SessionListingId, HashSet<JoinAttemptId>> _outcomesByListing = [];
private readonly PriorityQueue<DeadlineEntry<JoinAttemptId>, long> _outcomeExpiries = new();
private readonly Dictionary<MediationHandle, JoinAttemptId> _attemptHandles = [];
private readonly Dictionary<string, IdempotencyEntry> _idempotency = new(StringComparer.Ordinal);
private readonly PriorityQueue<DeadlineEntry<string>, long> _idempotencyExpiries = new();
private readonly Dictionary<string, TimeSpan> _replay = new(StringComparer.Ordinal);
private readonly PriorityQueue<DeadlineEntry<string>, long> _replayExpiries = new();
private readonly Dictionary<string, TimeSpan> _revocations = new(StringComparer.Ordinal);
private readonly PriorityQueue<DeadlineEntry<string>, long> _revocationExpiries = new();
private readonly HashSet<string> _scheduledRevocationExpiries = new(StringComparer.Ordinal);
private TimeSpan? _drainDeadline;
private TimeSpan _nextUdpMaintenance;
private long _maintenanceSweepCount;
private long _expiryChurn;
private bool _available = true;
private EphemeralStoreSnapshot? _metricsSnapshot;
private TimeSpan _metricsSnapshotAt = TimeSpan.MinValue;
public InMemoryEphemeralRendezvousStore(
EphemeralStoreOptions options,
@@ -43,6 +60,22 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
public Guid InstanceId { get; }
internal long MaintenanceSweepCount => Interlocked.Read(ref _maintenanceSweepCount);
internal int ScheduledExpiryEntryCount
{
get
{
lock (_gate)
{
return _listingExpiries.Count
+ _presenceExpiries.Count
+ _attemptExpiries.Count
+ _outcomeExpiries.Count
+ _idempotencyExpiries.Count
+ _replayExpiries.Count
+ _revocationExpiries.Count;
}
}
}
public bool IsAvailable
{
@@ -66,6 +99,59 @@ 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 int GetActiveJoinAttemptCountForDrain()
{
lock (_gate)
{
_expiryChurn += RemoveExpiredAttempts(_monotonicClock.Elapsed);
return _attempts.Count;
}
}
internal EphemeralStoreSnapshot GetMetricsSnapshot()
{
lock (_gate)
{
TimeSpan now = _monotonicClock.Elapsed;
if (_metricsSnapshot is null
|| now < _metricsSnapshotAt
|| now - _metricsSnapshotAt >= TimeSpan.FromMilliseconds(100))
{
Cleanup(now);
_metricsSnapshot = CreateSnapshot();
_metricsSnapshotAt = now;
}
return _metricsSnapshot;
}
}
private EphemeralStoreSnapshot CreateSnapshot() => new(
_listings.Count,
_presence.Count,
_attempts.Count,
_outcomeReports.Count,
_replay.Count,
_revocations.Count,
_idempotency.Count,
_maintenanceSweepCount,
_expiryChurn,
_available,
_drainDeadline.HasValue);
public StoreResult<StoredListing> CreateListing(
CreateListingCommand command,
CancellationToken cancellationToken = default) => Atomic<StoredListing>(now =>
@@ -104,10 +190,8 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
if (_listings.Count >= _options.MaxListings
|| _idempotency.Count >= _options.MaxIdempotencyEntries
|| _listings.Values.Count(entry => string.Equals(
entry.Definition.OwnerSubject,
command.Listing.OwnerSubject,
StringComparison.Ordinal)) >= command.OwnerListingLimit)
|| _listingCountsByOwner.GetValueOrDefault(command.Listing.OwnerSubject)
>= command.OwnerListingLimit)
{
return new(StoreResultCode.CapacityExceeded);
}
@@ -126,12 +210,21 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
WallDeadline(now, _options.LeaseLifetime),
version: 1);
_listings.Add(frozen.ListingId, entry);
ScheduleMutableDeadline(
_listingExpiries,
_scheduledListingExpiries,
frozen.ListingId,
entry.LeaseDeadline);
_listingCountsByOwner[frozen.OwnerSubject] =
_listingCountsByOwner.GetValueOrDefault(frozen.OwnerSubject) + 1;
_leases.Add(frozen.LeaseId, frozen.ListingId);
_presenceHandles.Add(frozen.HostPresenceHandle, frozen.ListingId);
_idempotency.Add(idempotencyKey, new(
IdempotencyEntry idempotency = new(
command.RequestFingerprint,
frozen.ListingId,
now + _options.IdempotencyLifetime));
now + _options.IdempotencyLifetime);
_idempotency.Add(idempotencyKey, idempotency);
EnqueueDeadline(_idempotencyExpiries, idempotencyKey, idempotency.Deadline);
return new(StoreResultCode.Success, Snapshot(entry));
}, cancellationToken);
@@ -289,10 +382,20 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
return new(StoreResultCode.CapacityExceeded);
}
_presence[command.Handle] = new(
bool isNewPresence = !_presence.ContainsKey(command.Handle);
PresenceEntry presence = new(
command.PublicEndpoint,
command.LocalEndpoint,
now + _options.PresenceLifetime);
_presence[command.Handle] = presence;
if (isNewPresence)
{
ScheduleMutableDeadline(
_presenceExpiries,
_scheduledPresenceExpiries,
command.Handle,
presence.Deadline);
}
return new(StoreResultCode.Success, Snapshot(entry));
}, cancellationToken, eagerCleanup: false);
@@ -387,8 +490,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
if (_attempts.Count >= _options.MaxJoinAttempts
|| _outcomeReports.Count >= _options.MaxOutcomeReports
|| _idempotency.Count >= _options.MaxIdempotencyEntries
|| _attempts.Values.Count(entry => entry.Command.Scope == command.Scope)
>= command.ScopeAttemptLimit)
|| _attemptCountsByScope.GetValueOrDefault(command.Scope) >= command.ScopeAttemptLimit)
{
return new(StoreResultCode.CapacityExceeded);
}
@@ -400,19 +502,29 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
AttemptEntry attempt = new(
command,
now,
now + _options.JoinAttemptLifetime,
WallDeadline(now, _options.JoinAttemptLifetime));
_attempts.Add(command.AttemptId, attempt);
_outcomeReports.Add(command.AttemptId, new(
AddToIndex(_attemptsByListing, command.ListingId, command.AttemptId);
_attemptCountsByScope[command.Scope] =
_attemptCountsByScope.GetValueOrDefault(command.Scope) + 1;
EnqueueDeadline(_attemptExpiries, command.AttemptId, attempt.Deadline);
OutcomeReportEntry outcome = new(
command.ListingId,
command.ClientSubject,
command.ClientCapabilityFingerprint,
now + _options.JoinAttemptLifetime + _options.IdempotencyLifetime));
now + _options.JoinAttemptLifetime + _options.IdempotencyLifetime);
_outcomeReports.Add(command.AttemptId, outcome);
AddToIndex(_outcomesByListing, command.ListingId, command.AttemptId);
EnqueueDeadline(_outcomeExpiries, command.AttemptId, outcome.Deadline);
_attemptHandles.Add(command.MediationHandle, command.AttemptId);
_idempotency.Add(idempotencyKey, new(
IdempotencyEntry idempotency = new(
command.RequestFingerprint,
command.AttemptId,
now + _options.IdempotencyLifetime));
now + _options.IdempotencyLifetime);
_idempotency.Add(idempotencyKey, idempotency);
EnqueueDeadline(_idempotencyExpiries, idempotencyKey, idempotency.Deadline);
return new(StoreResultCode.Success, Snapshot(attempt));
}, cancellationToken);
@@ -696,7 +808,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
throw new ArgumentOutOfRangeException(nameof(consumption), "Replay lifetime exceeds the configured ceiling.");
}
_replay.Add(key, now + lifetime);
TimeSpan deadline = now + lifetime;
_replay.Add(key, deadline);
EnqueueDeadline(_replayExpiries, key, deadline);
return new(StoreResultCode.Success, true);
}, cancellationToken);
@@ -729,7 +843,28 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
return new(StoreResultCode.CapacityExceeded);
}
_revocations[subject] = now + lifetime;
int activeResourcesBefore = _listings.Count
+ _presence.Count
+ _attempts.Count
+ _outcomeReports.Count;
TimeSpan deadline = now + lifetime;
bool isNewRevocation = !_revocations.TryGetValue(subject, out TimeSpan existingDeadline);
if (!isNewRevocation && existingDeadline > deadline)
{
// A repeated operator action may extend protection but cannot silently
// shorten an already-authoritative security revocation.
deadline = existingDeadline;
}
_revocations[subject] = deadline;
if (isNewRevocation)
{
ScheduleMutableDeadline(
_revocationExpiries,
_scheduledRevocationExpiries,
subject,
deadline);
}
SessionListingId[] listings = _listings
.Where(item => string.Equals(item.Value.Definition.OwnerSubject, subject, StringComparison.Ordinal))
.Select(static item => item.Key)
@@ -753,10 +888,14 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
}
foreach (JoinAttemptId attemptId in outcomeReports)
{
_outcomeReports.Remove(attemptId);
RemoveOutcome(attemptId);
}
return new(StoreResultCode.Success, listings.Length + attempts.Length);
int activeResourcesAfter = _listings.Count
+ _presence.Count
+ _attempts.Count
+ _outcomeReports.Count;
return new(StoreResultCode.Success, activeResourcesBefore - activeResourcesAfter);
}, cancellationToken);
public void BeginDrain(CancellationToken cancellationToken = default)
@@ -768,6 +907,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
if (!_drainDeadline.HasValue)
{
_drainDeadline = _monotonicClock.Elapsed + _options.GracefulDrainLifetime;
_metricsSnapshot = null;
}
}
}
@@ -778,6 +918,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
{
_available = false;
ClearActiveState();
_metricsSnapshot = null;
}
}
@@ -801,7 +942,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
_nextUdpMaintenance = now + UdpMaintenanceInterval;
}
return operation(now);
StoreResult<T> result = operation(now);
_metricsSnapshot = null;
return result;
}
}
@@ -838,60 +981,42 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
ClearActiveState();
}
RemoveExpired(_revocations, now);
RemoveExpired(_replay, now);
foreach (string key in _idempotency
.Where(item => item.Value.Deadline <= now)
.Select(static item => item.Key)
.ToArray())
{
_idempotency.Remove(key);
}
foreach (MediationHandle handle in _presence
.Where(item => item.Value.Deadline <= now)
.Select(static item => item.Key)
.ToArray())
{
_presence.Remove(handle);
}
foreach (JoinAttemptId attemptId in _attempts
.Where(item => item.Value.Deadline <= now)
.Select(static item => item.Key)
.ToArray())
{
RemoveAttempt(attemptId);
}
foreach (JoinAttemptId attemptId in _outcomeReports
.Where(item => item.Value.Deadline <= now)
.Select(static item => item.Key)
.ToArray())
{
_outcomeReports.Remove(attemptId);
}
foreach (SessionListingId listingId in _listings
.Where(item => item.Value.LeaseDeadline <= now)
.Select(static item => item.Key)
.ToArray())
{
RemoveListing(listingId);
}
_expiryChurn += RemoveExpiredMutableDeadlines(
_revocations,
_revocationExpiries,
_scheduledRevocationExpiries,
now);
_expiryChurn += RemoveExpiredDeadlines(_replay, _replayExpiries, now);
_expiryChurn += RemoveExpiredIdempotency(now);
_expiryChurn += RemoveExpiredPresence(now);
_expiryChurn += RemoveExpiredAttempts(now);
_expiryChurn += RemoveExpiredOutcomes(now);
_expiryChurn += RemoveExpiredListings(now);
}
private void ClearActiveState()
{
_listings.Clear();
_listingCountsByOwner.Clear();
_listingExpiries.Clear();
_scheduledListingExpiries.Clear();
_leases.Clear();
_presenceHandles.Clear();
_presence.Clear();
_presenceExpiries.Clear();
_scheduledPresenceExpiries.Clear();
_attempts.Clear();
_attemptCountsByScope.Clear();
_attemptsByListing.Clear();
_attemptExpiries.Clear();
_outcomeReports.Clear();
_outcomesByListing.Clear();
_outcomeExpiries.Clear();
_attemptHandles.Clear();
_idempotency.Clear();
_idempotencyExpiries.Clear();
_replay.Clear();
_replayExpiries.Clear();
}
private void RemoveListing(SessionListingId listingId)
@@ -902,23 +1027,23 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
}
_leases.Remove(listing.Definition.LeaseId);
DecrementCount(_listingCountsByOwner, listing.Definition.OwnerSubject);
_presenceHandles.Remove(listing.Definition.HostPresenceHandle);
_presence.Remove(listing.Definition.HostPresenceHandle);
foreach (JoinAttemptId attemptId in _attempts
.Where(item => item.Value.Command.ListingId == listingId)
.Select(static item => item.Key)
.ToArray())
if (_attemptsByListing.TryGetValue(listingId, out HashSet<JoinAttemptId>? attempts))
{
RemoveAttempt(attemptId);
foreach (JoinAttemptId attemptId in attempts.ToArray())
{
RemoveAttempt(attemptId);
}
}
foreach (JoinAttemptId attemptId in _outcomeReports
.Where(item => item.Value.ListingId == listingId)
.Select(static item => item.Key)
.ToArray())
if (_outcomesByListing.TryGetValue(listingId, out HashSet<JoinAttemptId>? outcomes))
{
_outcomeReports.Remove(attemptId);
foreach (JoinAttemptId attemptId in outcomes.ToArray())
{
RemoveOutcome(attemptId);
}
}
}
@@ -927,9 +1052,260 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
if (_attempts.Remove(attemptId, out AttemptEntry? attempt))
{
_attemptHandles.Remove(attempt.Command.MediationHandle);
DecrementCount(_attemptCountsByScope, attempt.Command.Scope);
RemoveFromIndex(_attemptsByListing, attempt.Command.ListingId, attemptId);
}
}
private void RemoveOutcome(JoinAttemptId attemptId)
{
if (_outcomeReports.Remove(attemptId, out OutcomeReportEntry? outcome))
{
RemoveFromIndex(_outcomesByListing, outcome.ListingId, attemptId);
}
}
private static void AddToIndex<TKey>(
Dictionary<TKey, HashSet<JoinAttemptId>> index,
TKey key,
JoinAttemptId attemptId)
where TKey : notnull
{
if (!index.TryGetValue(key, out HashSet<JoinAttemptId>? values))
{
values = [];
index.Add(key, values);
}
values.Add(attemptId);
}
private static void RemoveFromIndex<TKey>(
Dictionary<TKey, HashSet<JoinAttemptId>> index,
TKey key,
JoinAttemptId attemptId)
where TKey : notnull
{
if (index.TryGetValue(key, out HashSet<JoinAttemptId>? values)
&& values.Remove(attemptId)
&& values.Count == 0)
{
index.Remove(key);
}
}
private static void DecrementCount<TKey>(Dictionary<TKey, int> counts, TKey key)
where TKey : notnull
{
int remaining = counts[key] - 1;
if (remaining == 0)
{
counts.Remove(key);
}
else
{
counts[key] = remaining;
}
}
private int RemoveExpiredAttempts(TimeSpan now)
{
int removed = 0;
while (_attemptExpiries.TryPeek(
out DeadlineEntry<JoinAttemptId> candidate,
out long deadlineTicks)
&& deadlineTicks <= now.Ticks)
{
_attemptExpiries.Dequeue();
if (_attempts.TryGetValue(candidate.Key, out AttemptEntry? current)
&& current.Deadline == candidate.Deadline)
{
RemoveAttempt(candidate.Key);
removed++;
}
}
return removed;
}
private int RemoveExpiredListings(TimeSpan now)
{
int removed = 0;
while (_listingExpiries.TryPeek(
out DeadlineEntry<SessionListingId> candidate,
out long deadlineTicks)
&& deadlineTicks <= now.Ticks)
{
_listingExpiries.Dequeue();
_scheduledListingExpiries.Remove(candidate.Key);
if (!_listings.TryGetValue(candidate.Key, out ListingEntry? current))
{
continue;
}
if (current.LeaseDeadline > now)
{
ScheduleMutableDeadline(
_listingExpiries,
_scheduledListingExpiries,
candidate.Key,
current.LeaseDeadline);
}
else
{
RemoveListing(candidate.Key);
removed++;
}
}
return removed;
}
private int RemoveExpiredPresence(TimeSpan now)
{
int removed = 0;
while (_presenceExpiries.TryPeek(
out DeadlineEntry<MediationHandle> candidate,
out long deadlineTicks)
&& deadlineTicks <= now.Ticks)
{
_presenceExpiries.Dequeue();
_scheduledPresenceExpiries.Remove(candidate.Key);
if (!_presence.TryGetValue(candidate.Key, out PresenceEntry? current))
{
continue;
}
if (current.Deadline > now)
{
ScheduleMutableDeadline(
_presenceExpiries,
_scheduledPresenceExpiries,
candidate.Key,
current.Deadline);
}
else if (_presence.Remove(candidate.Key))
{
removed++;
}
}
return removed;
}
private int RemoveExpiredOutcomes(TimeSpan now)
{
int removed = 0;
while (_outcomeExpiries.TryPeek(
out DeadlineEntry<JoinAttemptId> candidate,
out long deadlineTicks)
&& deadlineTicks <= now.Ticks)
{
_outcomeExpiries.Dequeue();
if (_outcomeReports.TryGetValue(candidate.Key, out OutcomeReportEntry? current)
&& current.Deadline == candidate.Deadline)
{
RemoveOutcome(candidate.Key);
removed++;
}
}
return removed;
}
private int RemoveExpiredIdempotency(TimeSpan now)
{
int removed = 0;
while (_idempotencyExpiries.TryPeek(
out DeadlineEntry<string> candidate,
out long deadlineTicks)
&& deadlineTicks <= now.Ticks)
{
_idempotencyExpiries.Dequeue();
if (_idempotency.TryGetValue(candidate.Key, out IdempotencyEntry? current)
&& current.Deadline == candidate.Deadline
&& _idempotency.Remove(candidate.Key))
{
removed++;
}
}
return removed;
}
private static int RemoveExpiredDeadlines<TKey>(
Dictionary<TKey, TimeSpan> entries,
PriorityQueue<DeadlineEntry<TKey>, long> expiries,
TimeSpan now)
where TKey : notnull
{
int removed = 0;
while (expiries.TryPeek(out DeadlineEntry<TKey> candidate, out long deadlineTicks)
&& deadlineTicks <= now.Ticks)
{
expiries.Dequeue();
if (entries.TryGetValue(candidate.Key, out TimeSpan current)
&& current == candidate.Deadline
&& entries.Remove(candidate.Key))
{
removed++;
}
}
return removed;
}
private static int RemoveExpiredMutableDeadlines<TKey>(
Dictionary<TKey, TimeSpan> entries,
PriorityQueue<DeadlineEntry<TKey>, long> expiries,
HashSet<TKey> scheduled,
TimeSpan now)
where TKey : notnull
{
int removed = 0;
while (expiries.TryPeek(out DeadlineEntry<TKey> candidate, out long deadlineTicks)
&& deadlineTicks <= now.Ticks)
{
expiries.Dequeue();
scheduled.Remove(candidate.Key);
if (!entries.TryGetValue(candidate.Key, out TimeSpan current))
{
continue;
}
if (current > now)
{
ScheduleMutableDeadline(expiries, scheduled, candidate.Key, current);
}
else if (entries.Remove(candidate.Key))
{
removed++;
}
}
return removed;
}
private static void ScheduleMutableDeadline<TKey>(
PriorityQueue<DeadlineEntry<TKey>, long> expiries,
HashSet<TKey> scheduled,
TKey key,
TimeSpan deadline)
where TKey : notnull
{
if (scheduled.Add(key))
{
EnqueueDeadline(expiries, key, deadline);
}
}
private static void EnqueueDeadline<TKey>(
PriorityQueue<DeadlineEntry<TKey>, long> expiries,
TKey key,
TimeSpan deadline)
where TKey : notnull =>
expiries.Enqueue(new(key, deadline), deadline.Ticks);
private bool HandleExists(MediationHandle handle) =>
_presenceHandles.ContainsKey(handle) || _attemptHandles.ContainsKey(handle);
@@ -959,6 +1335,7 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
ClientCapabilityFingerprint = entry.Command.ClientCapabilityFingerprint,
ConnectionTicketFingerprint = entry.Command.ConnectionTicketFingerprint,
DedicatedFallback = StoredListing.CopyEndpoint(entry.Command.DedicatedFallback),
CreatedAtMonotonic = entry.CreatedAtMonotonic,
ExpiresAt = entry.WallExpiresAt,
ConnectionTicketExpiresAt = entry.TicketWallExpiresAt ?? default,
HostEndpoint = entry.HostEndpoint,
@@ -968,17 +1345,6 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
IsCancelled = entry.IsCancelled,
};
private static void RemoveExpired(Dictionary<string, TimeSpan> entries, TimeSpan now)
{
foreach (string key in entries
.Where(item => item.Value <= now)
.Select(static item => item.Key)
.ToArray())
{
entries.Remove(key);
}
}
private static void ValidateListing(ListingDefinition listing)
{
ArgumentNullException.ThrowIfNull(listing);
@@ -1101,10 +1467,12 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
private sealed class AttemptEntry(
CreateJoinAttemptCommand command,
TimeSpan createdAtMonotonic,
TimeSpan deadline,
DateTimeOffset wallExpiresAt)
{
public CreateJoinAttemptCommand Command { get; } = command;
public TimeSpan CreatedAtMonotonic { get; } = createdAtMonotonic;
public SecretFingerprint HostCapabilityFingerprint { get; } = command.HostCapabilityFingerprint;
public SecretFingerprint ClientCapabilityFingerprint { get; } = command.ClientCapabilityFingerprint;
public SecretFingerprint ConnectionTicketFingerprint { get; } = command.ConnectionTicketFingerprint;
@@ -1119,6 +1487,9 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
public bool IsCancelled { get; set; }
}
private readonly record struct DeadlineEntry<TKey>(TKey Key, TimeSpan Deadline)
where TKey : notnull;
private sealed class OutcomeReportEntry(
SessionListingId listingId,
string clientSubject,
@@ -1137,3 +1508,16 @@ internal sealed class InMemoryEphemeralRendezvousStore : IEphemeralRendezvousSto
object ResourceId,
TimeSpan Deadline);
}
internal sealed record EphemeralStoreSnapshot(
int ActiveListings,
int FreshPresenceBindings,
int ActiveJoinAttempts,
int RetainedOutcomeReports,
int ReplayMarkers,
int PrincipalRevocations,
int IdempotencyEntries,
long MaintenanceSweeps,
long ExpiryChurn,
bool IsAvailable,
bool IsDraining);
@@ -1,8 +1,10 @@
using System.Diagnostics;
using System.Net;
using System.Net.Sockets;
using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Abuse;
using FinalFactory.Rendezvous.Server.JoinAttempts;
using FinalFactory.Rendezvous.Server.Observability;
using FinalFactory.Rendezvous.Server.Sessions;
using FinalFactory.Rendezvous.Server.State;
@@ -38,7 +40,9 @@ internal sealed class NatMediationProcessor(
IEphemeralRendezvousStore store,
ISessionCapabilityService capabilities,
JoinAttemptService joinAttempts,
AbuseProtectionService? abuseProtection = null)
AbuseProtectionService? abuseProtection = null,
RendezvousTelemetry? telemetry = null,
IMonotonicClock? monotonicClock = null)
{
public NatMediationResult ProcessDatagram(
ReadOnlySpan<byte> encoded,
@@ -68,6 +72,26 @@ internal sealed class NatMediationProcessor(
INatIntroductionSink introductionSink,
CancellationToken cancellationToken = default)
{
long started = Stopwatch.GetTimestamp();
using Activity? activity = telemetry?.StartActivity("UDP frozen", ActivityKind.Server);
NatMediationResult result = ProcessDatagramCore(
encoded,
observedPublicEndpoint,
introductionSink,
cancellationToken);
telemetry?.RecordUdp(
"frozen",
result.ToString(),
Stopwatch.GetElapsedTime(started).TotalMilliseconds);
return result;
}
private NatMediationResult ProcessDatagramCore(
ReadOnlySpan<byte> encoded,
IPEndPoint observedPublicEndpoint,
INatIntroductionSink introductionSink,
CancellationToken cancellationToken)
{
if (!RendezvousUdpCodec.TryDecode(encoded, out PresenceDatagram? datagram, out _)
|| datagram is null
@@ -141,12 +165,22 @@ internal sealed class NatMediationProcessor(
IPEndPoint observedPublicEndpoint,
string token,
INatIntroductionSink introductionSink,
CancellationToken cancellationToken = default) => ProcessRequestCore(
CancellationToken cancellationToken = default)
{
long started = Stopwatch.GetTimestamp();
using Activity? activity = telemetry?.StartActivity("UDP litenet", ActivityKind.Server);
NatMediationResult result = ProcessRequestCore(
claimedLocalEndpoint,
observedPublicEndpoint,
token,
introductionSink,
cancellationToken);
telemetry?.RecordUdp(
"litenet",
result.ToString(),
Stopwatch.GetElapsedTime(started).TotalMilliseconds);
return result;
}
private NatMediationResult ProcessRequestCore(
IPEndPoint claimedLocalEndpoint,
@@ -255,6 +289,14 @@ internal sealed class NatMediationProcessor(
try
{
introductionSink.Introduce(CreatePlan(consumed.Value, ticket.Value));
if (telemetry is not null && monotonicClock is not null)
{
telemetry.RecordPairingLatency(Math.Max(
0,
(monotonicClock.Elapsed - consumed.Value.Attempt.CreatedAtMonotonic)
.TotalMilliseconds));
}
return NatMediationResult.Introduced;
}
catch (Exception exception) when (exception is SocketException
@@ -1,5 +1,8 @@
{
"Rendezvous": {
"AbuseProtection": {
"OperatorAllowedAddresses": ["127.0.0.1", "::1"]
},
"Provisioning": {
"Issuer": "final-factory-rendezvous-development",
"Audience": "final-factory-rendezvous",
@@ -14,6 +17,14 @@
"NotBefore": "2025-01-01T00:00:00Z",
"SignUntil": "2035-01-01T00:00:00Z",
"VerifyUntil": "2035-01-02T00:00:00Z"
},
{
"KeyId": "development-operator-1",
"SecretReference": "development:ephemeral/rendezvous-operator-signing",
"CredentialKinds": ["Operator"],
"NotBefore": "2025-01-01T00:00:00Z",
"SignUntil": "2035-01-01T00:00:00Z",
"VerifyUntil": "2035-01-02T00:00:00Z"
}
],
"Games": [
@@ -1,21 +1,39 @@
{
"Rendezvous": {
"Deployment": {
"PublicHttpBaseUrl": "",
"PublicUdpHost": "",
"PublicUdpPort": 9050,
"DrainDeadlineSeconds": 30,
"MinimumDrainSeconds": 1,
"SingleActiveInstance": true,
"AllowPrivatePublicEndpoints": false
},
"Udp": {
"ListenAddress": "0.0.0.0",
"Port": 9050,
"MaxDatagramsPerPoll": 256,
"PollIntervalMilliseconds": 2
},
"Audit": {
"MaxEntries": 10000,
"RetentionDays": 30
},
"AbuseProtection": {
"WindowSeconds": 1,
"MaxTrackedKeys": 100000,
"CriticalTrackedKeyReserve": 2048,
"UdpTrackedKeyLimit": 70000,
"TrustedProxyAddresses": [],
"OperatorAllowedAddresses": [],
"HealthGlobalRequestsPerWindow": 1000,
"HealthGlobalConcurrency": 32,
"HealthIpPrefixRequestsPerWindow": 120,
"HealthIpPrefixConcurrency": 8,
"OperatorGlobalRequestsPerWindow": 1000,
"OperatorGlobalConcurrency": 32,
"OperatorIpPrefixRequestsPerWindow": 120,
"OperatorIpPrefixConcurrency": 8,
"HttpGlobalRequestsPerWindow": 20000,
"HttpOptionalRequestsPerWindow": 18000,
"HttpIpPrefixRequestsPerWindow": 500,
@@ -0,0 +1,91 @@
namespace FinalFactory.Rendezvous.Capacity;
internal sealed record CapacityOptions
{
public required string Profile { get; init; }
public required int Listings { get; init; }
public required int Attempts { get; init; }
public required int Samples { get; init; }
public required int SoakCycles { get; init; }
public required int SoakSeconds { get; init; }
public string? OutputPath { get; init; }
public static CapacityOptions Parse(string[] args)
{
Dictionary<string, string> values = ParseArguments(args);
string profile = values.GetValueOrDefault("--profile") ?? "quick";
(int listings, int attempts, int samples, int soakCycles, int soakSeconds) = profile switch
{
"quick" => (1_000, 500, 100, 20, 0),
"candidate" => (25_000, 10_000, 1_000, 1_000, 300),
_ => throw new ArgumentException("--profile must be 'quick' or 'candidate'."),
};
return new()
{
Profile = profile,
Listings = ParsePositive(values.GetValueOrDefault("--listings"), listings, "--listings"),
Attempts = ParsePositive(values.GetValueOrDefault("--attempts"), attempts, "--attempts"),
Samples = ParsePositive(values.GetValueOrDefault("--samples"), samples, "--samples"),
SoakCycles = ParsePositive(
values.GetValueOrDefault("--soak-cycles"),
soakCycles,
"--soak-cycles"),
SoakSeconds = ParseNonNegative(
values.GetValueOrDefault("--soak-seconds"),
soakSeconds,
"--soak-seconds"),
OutputPath = values.GetValueOrDefault("--output"),
};
}
private static Dictionary<string, string> ParseArguments(string[] args)
{
HashSet<string> allowed =
[
"--profile",
"--listings",
"--attempts",
"--samples",
"--soak-cycles",
"--soak-seconds",
"--output",
];
Dictionary<string, string> values = new(StringComparer.Ordinal);
for (int index = 0; index < args.Length; index += 2)
{
string option = args[index];
if (!allowed.Contains(option))
{
throw new ArgumentException($"Unknown option: {option}.");
}
if (index == args.Length - 1
|| args[index + 1].StartsWith("--", StringComparison.Ordinal))
{
throw new ArgumentException($"{option} requires a value.");
}
if (!values.TryAdd(option, args[index + 1]))
{
throw new ArgumentException($"{option} may be supplied only once.");
}
}
return values;
}
private static int ParsePositive(string? value, int fallback, string option) =>
value is null
? fallback
: int.TryParse(value, out int parsed) && parsed > 0
? parsed
: throw new ArgumentException($"{option} must be a positive integer.");
private static int ParseNonNegative(string? value, int fallback, string option) =>
value is null
? fallback
: int.TryParse(value, out int parsed) && parsed >= 0
? parsed
: throw new ArgumentException($"{option} must be a non-negative integer.");
}
@@ -0,0 +1,76 @@
namespace FinalFactory.Rendezvous.Capacity;
internal sealed record CapacityReport
{
public required int SchemaVersion { get; init; }
public required string EvidenceVersion { get; init; }
public required DateTimeOffset GeneratedAt { get; init; }
public required string Profile { get; init; }
public required RuntimeEvidence Runtime { get; init; }
public required CapacityTargets Targets { get; init; }
public required IReadOnlyList<CapacityMeasurement> Measurements { get; init; }
public required StateEvidence State { get; init; }
public required IReadOnlyList<string> Failures { get; init; }
public required bool Passed { get; init; }
}
internal sealed record RuntimeEvidence(
string Framework,
string OperatingSystem,
string Kernel,
string Architecture,
string CpuModel,
int ProcessorCount,
string CpuAffinity,
string CpuQuota,
string MemoryLimit,
string GarbageCollector,
string CommitSha,
string TreeState,
string Command,
string ImageDigest,
string WorkloadSeed,
double CapacityPhaseAverageCpuPercent,
long PeakWorkingSetBytes,
long ManagedBytesAfterCleanup);
internal sealed record CapacityTargets(
int VisibleListings,
int ActiveJoinAttempts,
int CoreControlOperationsPerSecond,
int CoreMediationOperationsPerSecond,
double CoreControlP95Milliseconds,
double CoreMediationP95Milliseconds,
double MaximumAverageCpuPercent,
long MaximumWorkingSetBytes,
int SoakCycles,
int SoakDurationSeconds);
internal sealed record CapacityMeasurement(
string Operation,
int Samples,
double P50Milliseconds,
double P95Milliseconds,
double P99Milliseconds,
double OperationsPerSecond,
double MinimumOperationsPerSecond,
double BudgetMilliseconds,
bool Passed);
internal sealed record StateEvidence(
int PeakListings,
int PeakAttempts,
int PeakReplayMarkers,
int FinalListings,
int FinalAttempts,
int FinalReplayMarkers,
long ExpiryChurn,
long MaintenanceSweeps,
int SoakCyclesCompleted,
double SoakDurationSeconds,
int SoakPeakScheduledExpiryEntries,
long SoakManagedGrowthBytes,
int SoakHandleGrowth,
bool RestartStartedEmpty,
bool OverloadWasTyped,
bool RecoverySucceeded);
@@ -0,0 +1,580 @@
using System.Collections.Concurrent;
using System.Diagnostics;
using System.Runtime;
using System.Runtime.InteropServices;
using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Observability;
using FinalFactory.Rendezvous.Server.State;
namespace FinalFactory.Rendezvous.Capacity;
internal static class CapacityRunner
{
private static readonly TenantScope Scope = new(new("space-game"), new("production"));
private const uint ProtocolVersion = 1;
public static Task<CapacityReport> RunAsync(CapacityOptions options)
{
ArgumentNullException.ThrowIfNull(options);
Process process = Process.GetCurrentProcess();
TimeSpan cpuBefore = process.TotalProcessorTime;
Stopwatch capacityPhaseTime = Stopwatch.StartNew();
List<string> failures = [];
List<CapacityMeasurement> measurements = [];
ManualClock clock = new();
InMemoryEphemeralRendezvousStore store = CreateStore(options, clock);
List<StoredListing> listings = new(options.Listings);
List<CreateJoinAttemptCommand> attempts = new(options.Attempts);
int registrationSamples = Math.Min(options.Samples, options.Listings);
int attemptSamples = Math.Min(options.Samples, options.Attempts);
for (int index = 0; index < options.Listings - registrationSamples; index++)
{
listings.Add(CreateVisibleListing(store, index));
}
measurements.Add(Measure(
"registration-and-presence",
registrationSamples,
budgetMilliseconds: 200,
minimumOperationsPerSecond: 200,
index => listings.Add(CreateVisibleListing(
store,
options.Listings - registrationSamples + index))));
measurements.Add(Measure(
"lease-renewal",
registrationSamples,
budgetMilliseconds: 200,
minimumOperationsPerSecond: 200,
index =>
{
StoredListing listing = listings[index];
StoreResult<StoredListing> renewed = store.RenewLease(new(
listing.Definition.ListingId,
listing.Definition.LeaseId,
listing.Definition.LeaseFingerprint,
listing.Definition.OwnerSubject,
listing.Version));
RequireSuccess(renewed, "renewal");
listings[index] = renewed.Value!;
}));
int browseSamples = Math.Min(options.Samples, 250);
measurements.Add(Measure(
"visible-session-browse",
browseSamples,
budgetMilliseconds: 200,
minimumOperationsPerSecond: 200,
_ => RequireSuccess(
store.BrowseVisibleListings(new(
Scope,
ProtocolVersion,
new RegionId("eu-central"),
ContractLimits.BrowserPageMaxItems,
ExcludeFull: true)),
"browse")));
for (int index = 0; index < options.Attempts - attemptSamples; index++)
{
CreateJoinAttemptCommand command = CreateAttempt(index, listings[index % listings.Count]);
RequireSuccess(store.CreateJoinAttempt(command), "join issuance");
attempts.Add(command);
}
measurements.Add(Measure(
"join-attempt-issuance",
attemptSamples,
budgetMilliseconds: 200,
minimumOperationsPerSecond: 200,
index =>
{
int sequence = options.Attempts - attemptSamples + index;
CreateJoinAttemptCommand command = CreateAttempt(
sequence,
listings[sequence % listings.Count]);
RequireSuccess(store.CreateJoinAttempt(command), "join issuance");
attempts.Add(command);
}));
int punchSamples = Math.Min(options.Samples, attempts.Count);
measurements.Add(MeasureConcurrentPunch(store, attempts, punchSamples));
EphemeralStoreSnapshot peak = store.GetSnapshot();
StoreResult<StoredJoinAttempt> overloaded = store.CreateJoinAttempt(
CreateAttempt(options.Attempts + 1, listings[^1]));
bool overloadWasTyped = peak.ActiveJoinAttempts == options.Attempts
&& overloaded.Code == StoreResultCode.CapacityExceeded;
if (!overloadWasTyped)
{
failures.Add(
$"Expected typed CapacityExceeded at {options.Attempts} active attempts, "
+ $"observed count={peak.ActiveJoinAttempts}, result={overloaded.Code}.");
}
int revocationSamples = Math.Min(Math.Max(1, options.Samples / 20), listings.Count / 2);
measurements.Add(Measure(
"principal-revocation",
revocationSamples,
budgetMilliseconds: 200,
minimumOperationsPerSecond: 50,
index => RequireSuccess(
store.RevokePrincipal(
listings[index].Definition.OwnerSubject,
TimeSpan.FromMinutes(1)),
"principal revocation")));
using (RendezvousTelemetry telemetry = new(store))
{
measurements.Add(Measure(
"telemetry-recording",
Math.Max(100, options.Samples),
budgetMilliseconds: 1,
minimumOperationsPerSecond: 10_000,
_ =>
{
telemetry.RecordHttp("browse", 200, 1);
telemetry.RecordUdp("contribution", "accepted", 1);
telemetry.RecordPairingLatency(2);
}));
}
clock.Advance(TimeSpan.FromSeconds(61));
EphemeralStoreSnapshot? afterCoincidentExpiry = null;
measurements.Add(Measure(
"coincident-listing-attempt-expiry",
samples: 1,
budgetMilliseconds: 200,
minimumOperationsPerSecond: 0,
_ => afterCoincidentExpiry = store.GetSnapshot()));
if (afterCoincidentExpiry!.ActiveJoinAttempts != 0
|| afterCoincidentExpiry.ActiveListings != 0)
{
failures.Add("Coincident 60-second cleanup retained expired listings or join attempts.");
}
clock.Advance(TimeSpan.FromSeconds(90));
_ = store.GetSnapshot();
StoredListing recoveryListing = CreateVisibleListing(store, options.Listings + 1);
StoreResult<StoredJoinAttempt> recovered = store.CreateJoinAttempt(
CreateAttempt(options.Attempts + 2, recoveryListing));
bool recoverySucceeded = recovered.Succeeded;
if (!recoverySucceeded)
{
failures.Add($"Store did not recover after attempt expiry: {recovered.Code}.");
}
clock.Advance(TimeSpan.FromSeconds(151));
EphemeralStoreSnapshot final = store.GetSnapshot();
if (final.ActiveListings != 0
|| final.ActiveJoinAttempts != 0
|| final.ReplayMarkers != 0
|| final.IdempotencyEntries != 0
|| final.RetainedOutcomeReports != 0)
{
failures.Add("Expiry cleanup left active or retained state after every configured deadline.");
}
capacityPhaseTime.Stop();
process.Refresh();
double capacityPhaseCpuPercent = 100
* (process.TotalProcessorTime - cpuBefore).TotalSeconds
/ Math.Max(capacityPhaseTime.Elapsed.TotalSeconds * Environment.ProcessorCount, 0.000_001);
if (options.Profile == "candidate" && capacityPhaseCpuPercent > 70)
{
failures.Add(
$"Capacity-phase CPU {capacityPhaseCpuPercent:F1}% exceeded the 70% candidate budget.");
}
SoakEvidence soak = RunAcceleratedSoak(options, failures);
ManualClock restartClock = new();
EphemeralStoreSnapshot restarted = CreateStore(options, restartClock).GetSnapshot();
bool restartStartedEmpty = restarted.ActiveListings == 0
&& restarted.ActiveJoinAttempts == 0
&& restarted.ReplayMarkers == 0;
if (!restartStartedEmpty)
{
failures.Add("A restarted store did not begin empty.");
}
foreach (CapacityMeasurement measurement in measurements.Where(static item => !item.Passed))
{
failures.Add(
$"{measurement.Operation} missed its budget: p95={measurement.P95Milliseconds:F3} ms, "
+ $"rate={measurement.OperationsPerSecond:F1}/s.");
}
GC.Collect();
GC.WaitForPendingFinalizers();
GC.Collect();
process.Refresh();
long managedAfterCleanup = GC.GetTotalMemory(forceFullCollection: true);
long memoryBudget = 1_610_612_736;
if (process.PeakWorkingSet64 > memoryBudget)
{
failures.Add(
$"Peak working set {process.PeakWorkingSet64} exceeded the 1.5 GiB profile budget.");
}
if (options.Profile == "candidate" && Environment.ProcessorCount != 2)
{
failures.Add(
$"Candidate evidence must expose exactly two CPUs; runtime exposed "
+ $"{Environment.ProcessorCount}.");
}
CapacityReport report = new()
{
SchemaVersion = 2,
EvidenceVersion = "v2",
GeneratedAt = DateTimeOffset.UtcNow,
Profile = options.Profile,
Runtime = new(
RuntimeInformation.FrameworkDescription,
RuntimeInformation.OSDescription,
Environment.OSVersion.VersionString,
RuntimeInformation.ProcessArchitecture.ToString(),
ReadCpuModel(),
Environment.ProcessorCount,
Environment.GetEnvironmentVariable("RENDEZVOUS_EVIDENCE_CPUSET") ?? "unrestricted",
ReadCgroupValue("/sys/fs/cgroup/cpu.max"),
ReadCgroupValue("/sys/fs/cgroup/memory.max"),
GCSettings.IsServerGC ? "server" : "workstation",
Environment.GetEnvironmentVariable("RENDEZVOUS_EVIDENCE_COMMIT") ?? "unrecorded",
Environment.GetEnvironmentVariable("RENDEZVOUS_EVIDENCE_TREE_STATE") ?? "unrecorded",
Environment.GetEnvironmentVariable("RENDEZVOUS_EVIDENCE_COMMAND") ?? "unrecorded",
Environment.GetEnvironmentVariable("RENDEZVOUS_EVIDENCE_IMAGE_DIGEST") ?? "not-containerized",
"fixed-sequences-random-identifiers",
capacityPhaseCpuPercent,
process.PeakWorkingSet64,
managedAfterCleanup),
Targets = new(
options.Listings,
options.Attempts,
200,
2_000,
200,
100,
70,
memoryBudget,
options.SoakCycles,
options.SoakSeconds),
Measurements = measurements,
State = new(
peak.ActiveListings,
peak.ActiveJoinAttempts,
peak.ReplayMarkers,
final.ActiveListings,
final.ActiveJoinAttempts,
final.ReplayMarkers,
final.ExpiryChurn,
final.MaintenanceSweeps,
soak.Cycles,
soak.Duration.TotalSeconds,
soak.PeakScheduledExpiryEntries,
soak.ManagedGrowthBytes,
soak.HandleGrowth,
restartStartedEmpty,
overloadWasTyped,
recoverySucceeded),
Failures = failures,
Passed = failures.Count == 0,
};
return Task.FromResult(report);
}
private static InMemoryEphemeralRendezvousStore CreateStore(
CapacityOptions options,
ManualClock clock) => new(
new EphemeralStoreOptions
{
MaxListings = options.Listings,
MaxPresenceBindings = options.Listings,
MaxJoinAttempts = options.Attempts,
MaxOutcomeReports = options.Attempts,
MaxIdempotencyEntries = options.Listings + options.Attempts + 1,
},
clock,
clock);
private static StoredListing CreateVisibleListing(
InMemoryEphemeralRendezvousStore store,
int sequence)
{
string owner = $"publisher-{sequence}";
SecretFingerprint leaseFingerprint = new($"lease-{sequence}");
SecretFingerprint presenceFingerprint = new($"presence-{sequence}");
ListingDefinition definition = new()
{
ListingId = new(Guid.NewGuid()),
LeaseId = new(Guid.NewGuid()),
Scope = Scope,
OwnerSubject = owner,
RegionId = new("eu-central"),
ProtocolVersion = ProtocolVersion,
BuildVersion = "1.0.0",
DisplayName = $"Capacity host {sequence}",
Visibility = ListingVisibility.Public,
TrustMode = PublisherTrustMode.ManagedDedicated,
CurrentPlayers = 1,
MaximumPlayers = 8,
Metadata = new Dictionary<string, string>(StringComparer.Ordinal),
LeaseFingerprint = leaseFingerprint,
HostPresenceHandle = new(Guid.NewGuid()),
HostPresenceFingerprint = presenceFingerprint,
CapabilityDerivationSalt = "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
};
StoreResult<StoredListing> created = store.CreateListing(new(
$"register-{sequence}",
$"register-request-{sequence}",
definition));
RequireSuccess(created, "registration");
StoreResult<StoredListing> bound = store.BindHostPresence(new(
definition.HostPresenceHandle,
presenceFingerprint,
PublicEndpoint(10_000 + sequence % 50_000),
null));
RequireSuccess(bound, "host presence");
return bound.Value!;
}
private static CreateJoinAttemptCommand CreateAttempt(int sequence, StoredListing listing) => new()
{
IdempotencyOwner = $"client-{sequence}",
IdempotencyKey = $"join-{sequence}",
RequestFingerprint = $"join-request-{sequence}",
ClientSubject = $"client-{sequence}",
AttemptId = new(Guid.NewGuid()),
MediationHandle = new(Guid.NewGuid()),
Scope = Scope,
ListingId = listing.Definition.ListingId,
ProtocolVersion = ProtocolVersion,
HostCapabilityFingerprint = new($"host-capability-{sequence}"),
ClientCapabilityFingerprint = new($"client-capability-{sequence}"),
ConnectionTicketFingerprint = new($"ticket-{sequence}"),
CapabilityDerivationSalt = "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
};
private static CapacityMeasurement MeasureConcurrentPunch(
InMemoryEphemeralRendezvousStore store,
List<CreateJoinAttemptCommand> attempts,
int samples)
{
ConcurrentBag<double> latencies = [];
Stopwatch total = Stopwatch.StartNew();
Parallel.ForEach(
Enumerable.Range(0, samples),
new ParallelOptions { MaxDegreeOfParallelism = Math.Min(64, Environment.ProcessorCount * 4) },
index =>
{
CreateJoinAttemptCommand attempt = attempts[index];
Stopwatch elapsed = Stopwatch.StartNew();
RequireSuccess(store.BindAttemptEndpoint(new(
attempt.MediationHandle,
AttemptPeerRole.Host,
attempt.HostCapabilityFingerprint,
PublicEndpoint(20_000 + index % 20_000),
null)), "host punch");
RequireSuccess(store.BindAttemptEndpoint(new(
attempt.MediationHandle,
AttemptPeerRole.Client,
attempt.ClientCapabilityFingerprint,
PublicEndpoint(40_000 + index % 20_000),
null)), "client punch");
RequireSuccess(store.ConsumeIntroduction(attempt.MediationHandle), "introduction");
latencies.Add(elapsed.Elapsed.TotalMilliseconds);
});
total.Stop();
return BuildMeasurement(
"simultaneous-punch-pairing",
latencies.ToArray(),
total.Elapsed,
budgetMilliseconds: 100,
minimumOperationsPerSecond: 2_000);
}
private static CapacityMeasurement Measure(
string operation,
int samples,
double budgetMilliseconds,
double minimumOperationsPerSecond,
Action<int> action)
{
double[] latencies = new double[samples];
Stopwatch total = Stopwatch.StartNew();
for (int index = 0; index < samples; index++)
{
long started = Stopwatch.GetTimestamp();
action(index);
latencies[index] = Stopwatch.GetElapsedTime(started).TotalMilliseconds;
}
total.Stop();
return BuildMeasurement(
operation,
latencies,
total.Elapsed,
budgetMilliseconds,
minimumOperationsPerSecond);
}
private static CapacityMeasurement BuildMeasurement(
string operation,
double[] latencies,
TimeSpan elapsed,
double budgetMilliseconds,
double minimumOperationsPerSecond)
{
Array.Sort(latencies);
double operationsPerSecond = latencies.Length / Math.Max(elapsed.TotalSeconds, 0.000_001);
double p95 = Percentile(latencies, 0.95);
return new(
operation,
latencies.Length,
Percentile(latencies, 0.50),
p95,
Percentile(latencies, 0.99),
operationsPerSecond,
minimumOperationsPerSecond,
budgetMilliseconds,
p95 <= budgetMilliseconds && operationsPerSecond >= minimumOperationsPerSecond);
}
private static double Percentile(double[] sorted, double percentile)
{
int index = Math.Clamp((int)Math.Ceiling(sorted.Length * percentile) - 1, 0, sorted.Length - 1);
return sorted[index];
}
private static SoakEvidence RunAcceleratedSoak(
CapacityOptions options,
List<string> failures)
{
GC.Collect();
GC.WaitForPendingFinalizers();
GC.Collect();
long managedBefore = GC.GetTotalMemory(forceFullCollection: true);
int handlesBefore = Process.GetCurrentProcess().HandleCount;
ManualClock clock = new();
CapacityOptions soakOptions = options with { Listings = 100, Attempts = 100 };
InMemoryEphemeralRendezvousStore store = CreateStore(soakOptions, clock);
Stopwatch elapsed = Stopwatch.StartNew();
int cycle = 0;
int peakScheduledExpiryEntries = 0;
while (cycle < options.SoakCycles
|| elapsed.Elapsed < TimeSpan.FromSeconds(options.SoakSeconds))
{
StoredListing listing = CreateVisibleListing(store, cycle);
for (int refresh = 0; refresh < 10; refresh++)
{
clock.Advance(TimeSpan.FromTicks(1));
listing = RequireSuccess(store.RenewLease(new(
listing.Definition.ListingId,
listing.Definition.LeaseId,
listing.Definition.LeaseFingerprint,
listing.Definition.OwnerSubject,
listing.Version)), "soak lease refresh");
listing = RequireSuccess(store.BindHostPresence(new(
listing.Definition.HostPresenceHandle,
listing.Definition.HostPresenceFingerprint,
PublicEndpoint(10_000 + cycle % 50_000),
null)), "soak presence refresh");
}
CreateJoinAttemptCommand attempt = CreateAttempt(cycle, listing);
RequireSuccess(store.CreateJoinAttempt(attempt), "soak join issuance");
RequireSuccess(store.ConsumeReplay(new("capacity-soak", $"replay-{cycle}")), "soak replay");
peakScheduledExpiryEntries = Math.Max(
peakScheduledExpiryEntries,
store.ScheduledExpiryEntryCount);
if (store.ScheduledExpiryEntryCount > 7)
{
failures.Add(
$"Mutable deadline refresh grew the expiry queue to "
+ $"{store.ScheduledExpiryEntryCount} entries for one lifecycle.");
break;
}
clock.Advance(TimeSpan.FromSeconds(151));
EphemeralStoreSnapshot snapshot = store.GetSnapshot();
if (snapshot.ActiveListings != 0
|| snapshot.ActiveJoinAttempts != 0
|| snapshot.ReplayMarkers != 0
|| snapshot.IdempotencyEntries != 0
|| snapshot.RetainedOutcomeReports != 0)
{
failures.Add($"Accelerated soak retained state after cycle {cycle}.");
break;
}
cycle++;
}
elapsed.Stop();
GC.Collect();
GC.WaitForPendingFinalizers();
GC.Collect();
long managedGrowth = GC.GetTotalMemory(forceFullCollection: true) - managedBefore;
int handleGrowth = Process.GetCurrentProcess().HandleCount - handlesBefore;
if (managedGrowth > 67_108_864)
{
failures.Add($"Soak retained {managedGrowth} managed bytes; budget is 64 MiB.");
}
if (handleGrowth > 8)
{
failures.Add($"Soak retained {handleGrowth} process handles; budget is 8.");
}
return new(cycle, elapsed.Elapsed, peakScheduledExpiryEntries, managedGrowth, handleGrowth);
}
private static ObservedEndpoint PublicEndpoint(int port) =>
new(AddressFamilyKind.Ipv4, "203.0.113.10", port);
private static T RequireSuccess<T>(StoreResult<T> result, string operation)
{
if (!result.Succeeded)
{
throw new InvalidOperationException($"{operation} failed with {result.Code}.");
}
return result.Value!;
}
private static string ReadCpuModel()
{
const string cpuInfoPath = "/proc/cpuinfo";
if (!File.Exists(cpuInfoPath))
{
return "unavailable";
}
string? model = File.ReadLines(cpuInfoPath)
.FirstOrDefault(static line => line.StartsWith("model name", StringComparison.Ordinal));
int separator = model?.IndexOf(':') ?? -1;
return separator >= 0 ? model![(separator + 1)..].Trim() : "unavailable";
}
private static string ReadCgroupValue(string path) =>
File.Exists(path) ? File.ReadAllText(path).Trim() : "not-enforced";
private sealed class ManualClock : IWallClock, IMonotonicClock
{
public DateTimeOffset UtcNow { get; private set; } = DateTimeOffset.UtcNow;
public TimeSpan Elapsed { get; private set; }
public void Advance(TimeSpan duration)
{
UtcNow += duration;
Elapsed += duration;
}
}
private readonly record struct SoakEvidence(
int Cycles,
TimeSpan Duration,
int PeakScheduledExpiryEntries,
long ManagedGrowthBytes,
int HandleGrowth);
}
@@ -0,0 +1,13 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<OutputType>Exe</OutputType>
<TargetFramework>net10.0</TargetFramework>
<AssemblyName>FinalFactory.Rendezvous.Capacity</AssemblyName>
<RootNamespace>FinalFactory.Rendezvous.Capacity</RootNamespace>
<IsPackable>false</IsPackable>
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="../../src/FinalFactory.Rendezvous.Contracts/FinalFactory.Rendezvous.Contracts.csproj" />
<ProjectReference Include="../../src/FinalFactory.Rendezvous.Server/FinalFactory.Rendezvous.Server.csproj" />
</ItemGroup>
</Project>
@@ -0,0 +1,37 @@
using System.Text.Json;
namespace FinalFactory.Rendezvous.Capacity;
internal static class Program
{
private static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web)
{
WriteIndented = true,
};
public static async Task<int> Main(string[] args)
{
CapacityOptions options;
try
{
options = CapacityOptions.Parse(args);
}
catch (ArgumentException exception)
{
Console.Error.WriteLine(exception.Message);
return 2;
}
CapacityReport report = await CapacityRunner.RunAsync(options).ConfigureAwait(false);
string json = JsonSerializer.Serialize(report, JsonOptions);
Console.WriteLine(json);
if (options.OutputPath is not null)
{
string fullPath = Path.GetFullPath(options.OutputPath);
Directory.CreateDirectory(Path.GetDirectoryName(fullPath)!);
await File.WriteAllTextAsync(fullPath, json + Environment.NewLine).ConfigureAwait(false);
}
return report.Passed ? 0 : 1;
}
}
@@ -0,0 +1,39 @@
{
"version": 2,
"dependencies": {
"net10.0": {
"finalfactory.rendezvous.contracts": {
"type": "Project"
},
"finalfactory.rendezvous.server": {
"type": "Project",
"dependencies": {
"FinalFactory.Rendezvous.Contracts": "[1.0.0, )",
"LiteNetLib": "[2.1.4, )",
"Microsoft.AspNetCore.OpenApi": "[10.0.9, )"
}
},
"LiteNetLib": {
"type": "CentralTransitive",
"requested": "[2.1.4, )",
"resolved": "2.1.4",
"contentHash": "KWlxvMw3Urpqj9joD96LRiK+LC62pQNs/zkXRJc+rHnxgkGp+vV703xzDrxRmv+V1YhCFfIGzs5nrVWtREIlyA=="
},
"Microsoft.AspNetCore.OpenApi": {
"type": "CentralTransitive",
"requested": "[10.0.9, )",
"resolved": "10.0.9",
"contentHash": "1ihb8FO9cGgEK1/m3CTtT/SfnynwmiZib0W2pcDVj3KSWk/Sca4VOXEtaptKQc582zpFrzTFiwkGRCglt6H+WQ==",
"dependencies": {
"Microsoft.OpenApi": "2.0.0"
}
},
"Microsoft.OpenApi": {
"type": "CentralTransitive",
"requested": "[2.7.5, )",
"resolved": "2.7.5",
"contentHash": "0FA67RSnRM4tcBKqiqVu/HPdZ9+QOKbmeRjxRUGTCjPU4C0bmUhd97Dso7Yild5P7nOV6GxJ2xrK0Kv/O9xp0w=="
}
}
}
}
@@ -1,3 +1,6 @@
using System.Diagnostics;
using System.Net;
using System.Net.Sockets;
using FinalFactory.Rendezvous.Client;
using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Abuse;
@@ -19,6 +22,49 @@ namespace FinalFactory.Rendezvous.Tests.Client;
public sealed class RendezvousClientIntegrationTests
{
[Fact]
public async Task RestartReturnsTypedUnavailabilityThenAllowsHostReregistration()
{
int port = ReserveTcpPort();
string address = $"http://127.0.0.1:{port}";
using HttpClient client = new() { BaseAddress = new Uri(address) };
RendezvousClientOptions noRetry = new()
{
MaximumSafeRetries = 0,
RequestTimeout = TimeSpan.FromSeconds(1),
};
ClientTestHost first = await ClientTestHost.StartAsync(address);
RendezvousPublisherClient publisher = new(client, noRetry);
RendezvousClientResult<PublishedSession> registered = await publisher.RegisterAsync(
CreateRegistration(100),
first.PublisherCredential);
PublishedSession initialSession = AssertSuccess(registered);
BindPresence(first, initialSession, 41_100);
RendezvousSessionBrowserClient browser = new(client, noRetry);
BrowseSessionsResponse beforeRestart = AssertSuccess(await browser.BrowseAsync(BrowseRequest()));
Assert.Equal(initialSession.ListingId, Assert.Single(beforeRestart.Items).ListingId);
Stopwatch restart = Stopwatch.StartNew();
await first.DisposeAsync();
RendezvousClientResult<BrowseSessionsResponse> unavailable = await browser.BrowseAsync(BrowseRequest());
Assert.False(unavailable.IsSuccess);
Assert.Equal(RendezvousErrorCode.ServiceUnavailable, unavailable.Error);
await using ClientTestHost second = await ClientTestHost.StartAsync(address);
RendezvousClientResult<PublishedSession> reregistered = await publisher.RegisterAsync(
CreateRegistration(101),
second.PublisherCredential);
PublishedSession replacementSession = AssertSuccess(reregistered);
Assert.NotEqual(initialSession.ListingId, replacementSession.ListingId);
BindPresence(second, replacementSession, 41_101);
BrowseSessionsResponse afterRestart = AssertSuccess(await browser.BrowseAsync(BrowseRequest()));
Assert.Equal(replacementSession.ListingId, Assert.Single(afterRestart.Items).ListingId);
Assert.DoesNotContain(afterRestart.Items, item => item.ListingId == initialSession.ListingId);
Assert.True(
restart.Elapsed < TimeSpan.FromSeconds(5),
$"Local restart and host re-registration took {restart.Elapsed}.");
}
[Fact]
public async Task PublisherAndBrowserClientsCompleteTheRealSessionLifecycleAndPaging()
{
@@ -102,6 +148,27 @@ public sealed class RendezvousClientIntegrationTests
return Assert.IsAssignableFrom<T>(result.Value);
}
private static void BindPresence(ClientTestHost host, PublishedSession session, int port)
{
Assert.True(host.Capabilities.TryFingerprint(
session.HostPresenceCapability,
out SecretFingerprint fingerprint));
Assert.Equal(StoreResultCode.Success, host.Store.BindHostPresence(new(
session.HostPresenceHandle,
fingerprint,
new(AddressFamilyKind.Ipv4, "203.0.113.80", port),
null)).Code);
}
private static BrowseSessionsRequest BrowseRequest() => new()
{
GameId = new("space-game"),
EnvironmentId = new("production"),
RegionId = new("eu-central"),
ProtocolVersion = 7,
PageSize = 10,
};
private static RegisterSessionRequest CreateRegistration(int index) => new()
{
IdempotencyKey = $"sdk-integration-{index}",
@@ -139,7 +206,7 @@ public sealed class RendezvousClientIntegrationTests
internal EphemeralCapabilityIssuer Capabilities { get; }
internal string PublisherCredential { get; }
internal static async Task<ClientTestHost> StartAsync()
internal static async Task<ClientTestHost> StartAsync(string? bindAddress = null)
{
ManualRendezvousClock clock = new(ProvisioningTestData.Now);
EphemeralStoreOptions stateOptions = new();
@@ -153,7 +220,7 @@ public sealed class RendezvousClientIntegrationTests
string credential = provisioning.Credentials.Issue(principal, clock.UtcNow);
WebApplicationBuilder builder = WebApplication.CreateBuilder();
builder.WebHost.UseUrls("http://127.0.0.1:0");
builder.WebHost.UseUrls(bindAddress ?? "http://127.0.0.1:0");
builder.Services.ConfigureHttpJsonOptions(static options =>
ContractJson.Configure(options.SerializerOptions));
builder.Services.Configure<RouteHandlerOptions>(static options =>
@@ -180,10 +247,10 @@ public sealed class RendezvousClientIntegrationTests
app.MapRendezvousContractEndpoints();
await app.StartAsync();
IServer server = app.Services.GetRequiredService<IServer>();
string address = Assert.Single(server.Features.Get<IServerAddressesFeature>()!.Addresses);
string serviceAddress = Assert.Single(server.Features.Get<IServerAddressesFeature>()!.Addresses);
return new(
app,
new HttpClient { BaseAddress = new Uri(address) },
new HttpClient { BaseAddress = new Uri(serviceAddress) },
store,
capabilities,
credential);
@@ -196,4 +263,13 @@ public sealed class RendezvousClientIntegrationTests
await _application.DisposeAsync();
}
}
private static int ReserveTcpPort()
{
TcpListener listener = new(IPAddress.Loopback, 0);
listener.Start();
int port = ((IPEndPoint)listener.LocalEndpoint).Port;
listener.Stop();
return port;
}
}
@@ -11,6 +11,11 @@ public sealed class OpenApiCompatibilityTests
"/v1/join-attempts",
"/v1/join-attempts/{attemptId}",
"/v1/join-attempts/{attemptId}/outcome",
"/v1/operator/drain",
"/v1/operator/keys/revoke",
"/v1/operator/listings/revoke",
"/v1/operator/principals/revoke",
"/v1/operator/status",
"/v1/sessions",
"/v1/sessions/{listingId}",
"/v1/sessions/{listingId}/join-attempts",
@@ -104,6 +109,11 @@ public sealed class OpenApiCompatibilityTests
Assert.Equal(
"X-Rendezvous-Client-Punch-Capability",
attemptCapability.GetProperty("name").GetString());
JsonElement operatorBearer = root.GetProperty("components")
.GetProperty("securitySchemes")
.GetProperty("OperatorBearer");
Assert.Equal("http", operatorBearer.GetProperty("type").GetString());
Assert.Equal("bearer", operatorBearer.GetProperty("scheme").GetString());
(string Path, string Method)[] publisherOperations =
[
("/v1/sessions", "post"),
@@ -120,6 +130,23 @@ public sealed class OpenApiCompatibilityTests
Assert.True(security[0].TryGetProperty("PublisherBearer", out _));
}
(string Path, string Method)[] operatorOperations =
[
("/v1/operator/status", "get"),
("/v1/operator/listings/revoke", "post"),
("/v1/operator/principals/revoke", "post"),
("/v1/operator/keys/revoke", "post"),
("/v1/operator/drain", "post"),
];
foreach ((string operationPath, string method) in operatorOperations)
{
JsonElement security = root.GetProperty("paths")
.GetProperty(operationPath)
.GetProperty(method)
.GetProperty("security");
Assert.True(security[0].TryGetProperty("OperatorBearer", out _));
}
JsonElement cancelParameters = root.GetProperty("paths")
.GetProperty("/v1/join-attempts/{attemptId}")
.GetProperty("delete")
@@ -157,6 +184,15 @@ public sealed class OpenApiCompatibilityTests
static item => item.Name is "get" or "post" or "put" or "delete"))
{
JsonElement responses = operation.Value.GetProperty("responses");
foreach (JsonProperty response in responses.EnumerateObject())
{
JsonElement correlation = response.Value.GetProperty("headers")
.GetProperty("X-Rendezvous-Correlation-ID");
Assert.Equal(
"string",
correlation.GetProperty("schema").GetProperty("type").GetString());
}
if (!responses.TryGetProperty("429", out JsonElement overloaded))
{
continue;
@@ -171,7 +207,7 @@ public sealed class OpenApiCompatibilityTests
}
}
Assert.Equal(12, overloadContracts);
Assert.Equal(17, overloadContracts);
(string Path, string Method)[] bodyOperations =
[
("/v1/sessions", "post"),
@@ -180,6 +216,10 @@ public sealed class OpenApiCompatibilityTests
("/v1/sessions/{listingId}", "delete"),
("/v1/join-attempts", "post"),
("/v1/join-attempts/{attemptId}/outcome", "post"),
("/v1/operator/listings/revoke", "post"),
("/v1/operator/principals/revoke", "post"),
("/v1/operator/keys/revoke", "post"),
("/v1/operator/drain", "post"),
];
foreach ((string operationPath, string method) in bodyOperations)
{
@@ -0,0 +1,137 @@
using FinalFactory.Rendezvous.Server.Abuse;
using FinalFactory.Rendezvous.Server.Deployment;
namespace FinalFactory.Rendezvous.Tests.Deployment;
public sealed class DeploymentOptionsTests
{
[Fact]
public void ProductionConfigurationAcceptsExplicitPublicEndpointsAndTrustBoundary()
{
DeploymentOptions options = ValidOptions();
AbuseProtectionOptions abuse = new()
{
TrustedProxyAddresses = ["192.0.2.10"],
};
IReadOnlyList<string> errors = options.ValidateProduction(
abuse,
"rendezvous.finalfactory.at");
Assert.Empty(errors);
}
[Fact]
public void ProductionConfigurationRejectsUnsafeAndAmbiguousDefaults()
{
DeploymentOptions options = new()
{
SingleActiveInstance = false,
};
IReadOnlyList<string> errors = options.ValidateProduction(
new AbuseProtectionOptions(),
"*");
Assert.Contains(errors, error => error.Contains("SingleActiveInstance", StringComparison.Ordinal));
Assert.Contains(errors, error => error.Contains("PublicHttpBaseUrl", StringComparison.Ordinal));
Assert.Contains(errors, error => error.Contains("PublicUdpHost", StringComparison.Ordinal));
Assert.Contains(errors, error => error.Contains("TrustedProxyAddresses", StringComparison.Ordinal));
Assert.Contains(errors, error => error.Contains("AllowedHosts", StringComparison.Ordinal));
}
[Theory]
[InlineData("http://rendezvous.example.com/")]
[InlineData("https://user@example.com/")]
[InlineData("https://rendezvous.example.com/path")]
[InlineData("https://localhost/")]
[InlineData("https://10.0.0.1/")]
[InlineData("https://192.0.2.10/")]
[InlineData("https://[::ffff:10.0.0.1]/")]
[InlineData("https://[ff02::1]/")]
[InlineData("https://rendezvous.invalid/")]
public void ProductionConfigurationRejectsUnsafeHttpEndpoint(string endpoint)
{
DeploymentOptions options = ValidOptions() with { PublicHttpBaseUrl = endpoint };
IReadOnlyList<string> errors = options.ValidateProduction(
new AbuseProtectionOptions { TrustedProxyAddresses = ["192.0.2.10"] },
"rendezvous.finalfactory.at");
Assert.Contains(errors, error => error.Contains("PublicHttpBaseUrl", StringComparison.Ordinal));
}
[Theory]
[InlineData("0.1.2.3")]
[InlineData("100.64.0.1")]
[InlineData("192.0.2.1")]
[InlineData("198.18.0.1")]
[InlineData("198.51.100.1")]
[InlineData("203.0.113.1")]
[InlineData("224.0.0.1")]
[InlineData("255.255.255.255")]
[InlineData("::ffff:192.168.1.1")]
[InlineData("2001:db8::1")]
[InlineData("ff02::1")]
[InlineData("rendezvous.example.com")]
[InlineData("rendezvous.home.arpa")]
[InlineData("rendezvous.alt")]
[InlineData("service.test")]
public void ProductionConfigurationRejectsNonPublicUdpEndpoint(string endpoint)
{
DeploymentOptions options = ValidOptions() with { PublicUdpHost = endpoint };
IReadOnlyList<string> errors = options.ValidateProduction(
new AbuseProtectionOptions { TrustedProxyAddresses = ["192.0.2.10"] },
"rendezvous.finalfactory.at");
Assert.Contains(errors, error => error.Contains("PublicUdpHost", StringComparison.Ordinal));
}
[Fact]
public void IsolatedSmokeTestRequiresExplicitPrivateEndpointOverride()
{
DeploymentOptions options = ValidOptions() with
{
PublicHttpBaseUrl = "https://localhost/",
PublicUdpHost = "127.0.0.1",
AllowPrivatePublicEndpoints = true,
};
IReadOnlyList<string> errors = options.ValidateProduction(
new AbuseProtectionOptions { TrustedProxyAddresses = ["127.0.0.1"] },
"localhost");
Assert.Empty(errors);
}
[Fact]
public void ProductionConfigurationRejectsInvalidPortsDeadlinesAndMismatchedHostFilter()
{
DeploymentOptions options = ValidOptions() with
{
PublicUdpPort = 0,
DrainDeadlineSeconds = 31,
MinimumDrainSeconds = 6,
};
IReadOnlyList<string> errors = options.ValidateProduction(
new AbuseProtectionOptions { TrustedProxyAddresses = ["192.0.2.10"] },
"different.finalfactory.at");
Assert.Contains(errors, error => error.Contains("PublicUdpPort", StringComparison.Ordinal));
Assert.Contains(errors, error => error.Contains("DrainDeadlineSeconds", StringComparison.Ordinal));
Assert.Contains(errors, error => error.Contains("MinimumDrainSeconds", StringComparison.Ordinal));
Assert.Contains(errors, error => error.Contains("exact host", StringComparison.Ordinal));
}
private static DeploymentOptions ValidOptions() => new()
{
PublicHttpBaseUrl = "https://rendezvous.finalfactory.at/",
PublicUdpHost = "rendezvous-udp.finalfactory.at",
PublicUdpPort = 9050,
DrainDeadlineSeconds = 10,
MinimumDrainSeconds = 1,
SingleActiveInstance = true,
};
}
@@ -0,0 +1,149 @@
using System.Diagnostics;
using FinalFactory.Rendezvous.Server.Deployment;
using FinalFactory.Rendezvous.Server.State;
using FinalFactory.Rendezvous.Tests.State;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
namespace FinalFactory.Rendezvous.Tests.Deployment;
public sealed class GracefulDrainServiceTests
{
[Fact]
public async Task ShutdownRejectsNewWorkAndClearsStateAfterBoundedAttemptDeadline()
{
EphemeralStoreOptions stateOptions = new()
{
GracefulDrainLifetime = TimeSpan.FromSeconds(1),
};
EphemeralStateFixture fixture = new(stateOptions);
StoredListing listing = fixture.CreateVisibleListing(out _);
StoreResult<StoredJoinAttempt> attempt = fixture.Store.CreateJoinAttempt(
fixture.AttemptCommand(listing));
Assert.True(attempt.Succeeded);
using FakeApplicationLifetime lifetime = new();
using GracefulDrainService service = CreateService(fixture.Store, lifetime);
await service.StartAsync(CancellationToken.None);
Task stopping = Task.Run(lifetime.StopApplication);
await WaitUntilAsync(() => fixture.Store.IsDraining, TimeSpan.FromSeconds(1));
StoreResult<StoredListing> rejected = fixture.Store.CreateListing(fixture.ListingCommand());
long sweepsAfterAdmissionCheck = fixture.Store.MaintenanceSweepCount;
Stopwatch elapsed = Stopwatch.StartNew();
await stopping;
await service.StopAsync(CancellationToken.None);
Assert.Equal(StoreResultCode.Draining, rejected.Code);
Assert.Equal(sweepsAfterAdmissionCheck, fixture.Store.MaintenanceSweepCount);
Assert.InRange(elapsed.Elapsed, TimeSpan.FromMilliseconds(850), TimeSpan.FromSeconds(2));
Assert.False(fixture.Store.IsAvailable);
Assert.Equal(0, fixture.Store.GetSnapshot().ActiveJoinAttempts);
Assert.Equal(0, fixture.Store.GetSnapshot().ActiveListings);
}
[Fact]
public async Task ShutdownWithoutAttemptsStopsAfterTheConfiguredMinimumOnly()
{
EphemeralStateFixture fixture = new(new EphemeralStoreOptions
{
GracefulDrainLifetime = TimeSpan.FromSeconds(2),
});
fixture.CreateVisibleListing(out _);
using FakeApplicationLifetime lifetime = new();
DeploymentOptions options = new()
{
DrainDeadlineSeconds = 2,
MinimumDrainSeconds = 0,
};
using GracefulDrainService service = new(
fixture.Store,
lifetime,
Options.Create(options),
NullLogger<GracefulDrainService>.Instance);
await service.StartAsync(CancellationToken.None);
Stopwatch elapsed = Stopwatch.StartNew();
await service.StopAsync(CancellationToken.None);
Assert.True(elapsed.Elapsed < TimeSpan.FromMilliseconds(500));
Assert.False(fixture.Store.IsAvailable);
Assert.Equal(0, fixture.Store.GetSnapshot().ActiveListings);
}
[Fact]
public async Task ShutdownStopsWhenTheLastAttemptExpiresWithoutAFullStateSweep()
{
EphemeralStateFixture fixture = new(new EphemeralStoreOptions
{
JoinAttemptLifetime = TimeSpan.FromMilliseconds(100),
ConnectionTicketLifetime = TimeSpan.FromMilliseconds(50),
GracefulDrainLifetime = TimeSpan.FromSeconds(2),
});
StoredListing listing = fixture.CreateVisibleListing(out _);
Assert.True(fixture.Store.CreateJoinAttempt(fixture.AttemptCommand(listing)).Succeeded);
using FakeApplicationLifetime lifetime = new();
using GracefulDrainService service = new(
fixture.Store,
lifetime,
Options.Create(new DeploymentOptions
{
DrainDeadlineSeconds = 2,
MinimumDrainSeconds = 0,
}),
NullLogger<GracefulDrainService>.Instance);
await service.StartAsync(CancellationToken.None);
Task stopping = Task.Run(lifetime.StopApplication);
await WaitUntilAsync(() => fixture.Store.IsDraining, TimeSpan.FromSeconds(1));
long sweepsBeforeExpiry = fixture.Store.MaintenanceSweepCount;
fixture.Clock.Advance(TimeSpan.FromMilliseconds(150));
using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(1));
await stopping.WaitAsync(timeout.Token);
Assert.Equal(sweepsBeforeExpiry, fixture.Store.MaintenanceSweepCount);
Assert.False(fixture.Store.IsAvailable);
}
private static GracefulDrainService CreateService(
InMemoryEphemeralRendezvousStore store,
IHostApplicationLifetime lifetime) => new(
store,
lifetime,
Options.Create(new DeploymentOptions
{
DrainDeadlineSeconds = 1,
MinimumDrainSeconds = 0,
}),
NullLogger<GracefulDrainService>.Instance);
private static async Task WaitUntilAsync(Func<bool> predicate, TimeSpan timeout)
{
Stopwatch elapsed = Stopwatch.StartNew();
while (!predicate())
{
Assert.True(elapsed.Elapsed < timeout, "The service did not enter drain in time.");
await Task.Delay(10);
}
}
private sealed class FakeApplicationLifetime : IHostApplicationLifetime, IDisposable
{
private readonly CancellationTokenSource _started = new();
private readonly CancellationTokenSource _stopping = new();
private readonly CancellationTokenSource _stopped = new();
public CancellationToken ApplicationStarted => _started.Token;
public CancellationToken ApplicationStopping => _stopping.Token;
public CancellationToken ApplicationStopped => _stopped.Token;
public void StopApplication() => _stopping.Cancel();
public void Dispose()
{
_started.Dispose();
_stopping.Dispose();
_stopped.Dispose();
}
}
}
@@ -0,0 +1,575 @@
using System.Diagnostics;
using System.Net;
using System.Net.Sockets;
using System.Security.Cryptography;
namespace FinalFactory.Rendezvous.Tests.Deployment;
public sealed class ProductionProcessTests
{
[Fact]
public async Task SigtermDrainsThenReleasesHttpAndUdpSockets()
{
if (!OperatingSystem.IsLinux())
{
return;
}
int httpPort = ReserveTcpPort();
int udpPort = ReserveUdpPort();
string secretPath = Path.Combine(
Path.GetTempPath(),
$"rendezvous-process-secret-{Guid.NewGuid():N}");
await File.WriteAllBytesAsync(secretPath, RandomNumberGenerator.GetBytes(32));
Process? process = null;
try
{
ProcessStartInfo startInfo = CreateStartInfo(httpPort, udpPort, secretPath);
process = Process.Start(startInfo)
?? throw new InvalidOperationException("The production server process did not start.");
Task<string> standardOutput = process.StandardOutput.ReadToEndAsync();
Task<string> standardError = process.StandardError.ReadToEndAsync();
await WaitForReadyAsync(httpPort, process, TimeSpan.FromSeconds(10));
AssertUdpPortIsBound(udpPort);
Stopwatch shutdown = Stopwatch.StartNew();
ProcessStartInfo signalInfo = new()
{
FileName = "/bin/kill",
UseShellExecute = false,
ArgumentList =
{
"-TERM",
process.Id.ToString(System.Globalization.CultureInfo.InvariantCulture),
},
};
using (Process signal = Process.Start(signalInfo)
?? throw new InvalidOperationException("Could not send SIGTERM."))
{
await signal.WaitForExitAsync();
Assert.Equal(0, signal.ExitCode);
}
using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(6));
await process.WaitForExitAsync(timeout.Token);
string output = await standardOutput;
string error = await standardError;
Assert.True(
process.ExitCode == 0,
$"Server exited with {process.ExitCode}. stdout: {output} stderr: {error}");
Assert.InRange(shutdown.Elapsed, TimeSpan.FromMilliseconds(700), TimeSpan.FromSeconds(5));
AssertTcpPortIsReleased(httpPort);
AssertUdpPortIsReleased(udpPort);
process.Dispose();
process = null;
Stopwatch replacementReady = Stopwatch.StartNew();
ProcessStartInfo replacementInfo = CreateStartInfo(httpPort, udpPort, secretPath);
process = Process.Start(replacementInfo)
?? throw new InvalidOperationException("The replacement production process did not start.");
Task<string> replacementOutput = process.StandardOutput.ReadToEndAsync();
Task<string> replacementError = process.StandardError.ReadToEndAsync();
await WaitForReadyAsync(httpPort, process, TimeSpan.FromSeconds(15));
Assert.True(
replacementReady.Elapsed < TimeSpan.FromSeconds(15),
$"Replacement readiness took {replacementReady.Elapsed}.");
AssertUdpPortIsBound(udpPort);
await SendSigtermAsync(process);
using CancellationTokenSource replacementTimeout = new(TimeSpan.FromSeconds(6));
await process.WaitForExitAsync(replacementTimeout.Token);
string replacementFinalOutput = await replacementOutput;
string replacementFinalError = await replacementError;
Assert.True(
process.ExitCode == 0,
$"Replacement exited with {process.ExitCode}. "
+ $"stdout: {replacementFinalOutput} stderr: {replacementFinalError}");
AssertTcpPortIsReleased(httpPort);
AssertUdpPortIsReleased(udpPort);
}
finally
{
if (process is not null)
{
if (!process.HasExited)
{
process.Kill(entireProcessTree: true);
await process.WaitForExitAsync();
}
process.Dispose();
}
File.Delete(secretPath);
}
}
[Fact]
public async Task ProductionTransportSoakKeepsHandlesMemoryAndSocketsBounded()
{
if (!OperatingSystem.IsLinux())
{
return;
}
int httpPort = ReserveTcpPort();
int udpPort = ReserveUdpPort();
string secretPath = Path.Combine(
Path.GetTempPath(),
$"rendezvous-transport-soak-secret-{Guid.NewGuid():N}");
await File.WriteAllBytesAsync(secretPath, RandomNumberGenerator.GetBytes(32));
Process? process = null;
try
{
process = Process.Start(CreateStartInfo(httpPort, udpPort, secretPath))
?? throw new InvalidOperationException("The production soak process did not start.");
Task<string> standardOutput = process.StandardOutput.ReadToEndAsync();
Task<string> standardError = process.StandardError.ReadToEndAsync();
await WaitForReadyAsync(httpPort, process, TimeSpan.FromSeconds(10));
process.Refresh();
int baselineHandles = process.HandleCount;
long baselineWorkingSet = process.WorkingSet64;
int peakHandles = baselineHandles;
byte[] invalidDatagram = RandomNumberGenerator.GetBytes(64);
IPEndPoint udpEndpoint = new(IPAddress.Loopback, udpPort);
using HttpClient client = new() { Timeout = TimeSpan.FromSeconds(1) };
using UdpClient udp = new();
Stopwatch soak = Stopwatch.StartNew();
int cycles = 0;
int accepted = 0;
int shed = 0;
while (soak.Elapsed < TimeSpan.FromSeconds(10))
{
using HttpResponseMessage response = await client.GetAsync(
$"http://127.0.0.1:{httpPort}/health/live");
if (response.StatusCode == HttpStatusCode.OK)
{
accepted++;
}
else
{
Assert.Equal(HttpStatusCode.TooManyRequests, response.StatusCode);
shed++;
}
await udp.SendAsync(invalidDatagram, udpEndpoint);
cycles++;
if (cycles % 100 == 0)
{
process.Refresh();
peakHandles = Math.Max(peakHandles, process.HandleCount);
}
}
Assert.True(cycles >= 100, $"Transport soak completed only {cycles} cycles.");
Assert.True(accepted > 0, "Transport soak never admitted a health request.");
Assert.True(shed > 0, "Transport soak never exercised typed HTTP load shedding.");
await Task.Delay(TimeSpan.FromSeconds(2));
using (HttpResponseMessage recovered = await client.GetAsync(
$"http://127.0.0.1:{httpPort}/health/live"))
{
Assert.Equal(HttpStatusCode.OK, recovered.StatusCode);
}
process.Refresh();
Assert.InRange(peakHandles, 0, baselineHandles + 32);
Assert.InRange(process.HandleCount, 0, baselineHandles + 16);
Assert.InRange(process.WorkingSet64, 0, baselineWorkingSet + 67_108_864);
AssertUdpPortIsBound(udpPort);
await SendSigtermAsync(process);
using CancellationTokenSource shutdownTimeout = new(TimeSpan.FromSeconds(6));
await process.WaitForExitAsync(shutdownTimeout.Token);
string output = await standardOutput;
string error = await standardError;
Assert.True(
process.ExitCode == 0,
$"Transport soak process failed. stdout: {output} stderr: {error}");
AssertTcpPortIsReleased(httpPort);
AssertUdpPortIsReleased(udpPort);
}
finally
{
await StopProcessTreeAsync(process);
File.Delete(secretPath);
}
}
[Fact]
public async Task DocumentedSmokeScriptReachesHttpAndAuthenticatedUdpFlow()
{
if (!OperatingSystem.IsLinux())
{
return;
}
string root = RepositoryRoot();
int httpPort = ReserveTcpPort();
int udpPort = ReserveUdpPort();
string secretPath = Path.Combine(
Path.GetTempPath(),
$"rendezvous-smoke-secret-{Guid.NewGuid():N}");
await File.WriteAllBytesAsync(secretPath, RandomNumberGenerator.GetBytes(32));
Process? server = null;
Process? smoke = null;
try
{
DateTimeOffset now = DateTimeOffset.UtcNow;
string assembly = typeof(Program).Assembly.Location;
ProcessStartInfo serverInfo = new()
{
FileName = "dotnet",
WorkingDirectory = root,
RedirectStandardOutput = true,
RedirectStandardError = true,
UseShellExecute = false,
ArgumentList =
{
assembly,
"--contentRoot", Path.Combine(root, "deploy", "compose"),
"--Rendezvous:Provisioning:SigningKeys:0:SecretReference", $"file:{secretPath}",
"--Rendezvous:Provisioning:SigningKeys:0:NotBefore", now.AddHours(-1).ToString("O"),
"--Rendezvous:Provisioning:SigningKeys:0:SignUntil", now.AddHours(1).ToString("O"),
"--Rendezvous:Provisioning:SigningKeys:0:VerifyUntil", now.AddHours(2).ToString("O"),
"--Rendezvous:Udp:Port", udpPort.ToString(System.Globalization.CultureInfo.InvariantCulture),
"--Rendezvous:Deployment:PublicUdpPort", udpPort.ToString(System.Globalization.CultureInfo.InvariantCulture),
},
};
serverInfo.Environment["ASPNETCORE_ENVIRONMENT"] = "Production";
serverInfo.Environment["ASPNETCORE_URLS"] = $"http://127.0.0.1:{httpPort}";
server = Process.Start(serverInfo)
?? throw new InvalidOperationException("The smoke server process did not start.");
Task<string> serverOutput = server.StandardOutput.ReadToEndAsync();
Task<string> serverError = server.StandardError.ReadToEndAsync();
await WaitForReadyAsync(httpPort, server, TimeSpan.FromSeconds(10));
ProcessStartInfo smokeInfo = new()
{
FileName = Path.Combine(root, "scripts", "smoke-deployment.sh"),
WorkingDirectory = root,
RedirectStandardOutput = true,
RedirectStandardError = true,
UseShellExecute = false,
};
smokeInfo.Environment.Remove("RENDEZVOUS_PUBLISHER_CREDENTIAL");
smokeInfo.Environment["RENDEZVOUS_SMOKE_HTTP_URL"] = $"http://127.0.0.1:{httpPort}/";
smokeInfo.Environment["RENDEZVOUS_SMOKE_UDP_ENDPOINT"] = $"127.0.0.1:{udpPort}";
smokeInfo.Environment["RENDEZVOUS_SMOKE_LOCAL_KEY"] = secretPath;
smokeInfo.Environment["RENDEZVOUS_SMOKE_TIMEOUT_SECONDS"] = "15";
smokeInfo.Environment["RENDEZVOUS_SMOKE_CONFIGURATION"] = BuildConfiguration();
smoke = Process.Start(smokeInfo)
?? throw new InvalidOperationException("The deployment smoke process did not start.");
Task<string> smokeOutput = smoke.StandardOutput.ReadToEndAsync();
Task<string> smokeError = smoke.StandardError.ReadToEndAsync();
using (CancellationTokenSource timeout = new(TimeSpan.FromSeconds(25)))
{
await smoke.WaitForExitAsync(timeout.Token);
}
string output = await smokeOutput;
string error = await smokeError;
Assert.True(
smoke.ExitCode == 0,
$"Smoke exited with {smoke.ExitCode}. stdout: {output} stderr: {error}");
Assert.Contains("deployment smoke passed", output, StringComparison.OrdinalIgnoreCase);
Assert.DoesNotContain("rv1.", output, StringComparison.Ordinal);
Assert.DoesNotContain("rv1.", error, StringComparison.Ordinal);
await SendSigtermAsync(server);
using CancellationTokenSource shutdownTimeout = new(TimeSpan.FromSeconds(6));
await server.WaitForExitAsync(shutdownTimeout.Token);
string finalServerOutput = await serverOutput;
string finalServerError = await serverError;
Assert.True(
server.ExitCode == 0,
$"Smoke server failed. stdout: {finalServerOutput} stderr: {finalServerError}");
}
finally
{
await StopProcessTreeAsync(smoke);
if (server is not null)
{
await StopProcessTreeAsync(server);
}
File.Delete(secretPath);
}
}
[Theory]
[InlineData("missing-deployment", "PublicHttpBaseUrl")]
[InlineData("wildcard-host", "AllowedHosts")]
[InlineData("reserved-endpoint", "public DNS name or address")]
[InlineData("missing-key", "process-test-key")]
public async Task UnsafeProductionConfigurationFailsBeforeBinding(
string scenario,
string expectedDiagnostic)
{
if (!OperatingSystem.IsLinux())
{
return;
}
int httpPort = ReserveTcpPort();
int udpPort = ReserveUdpPort();
string secretPath = Path.Combine(
Path.GetTempPath(),
$"rendezvous-rejected-secret-{Guid.NewGuid():N}");
await File.WriteAllBytesAsync(secretPath, RandomNumberGenerator.GetBytes(32));
Process? process = null;
try
{
ProcessStartInfo startInfo = scenario == "missing-deployment"
? CreateBareProductionStartInfo(httpPort)
: CreateStartInfo(httpPort, udpPort, secretPath);
if (scenario == "wildcard-host")
{
startInfo.Environment["AllowedHosts"] = "*";
}
else if (scenario == "reserved-endpoint")
{
startInfo.Environment["Rendezvous__Deployment__AllowPrivatePublicEndpoints"] = "false";
startInfo.Environment["Rendezvous__Deployment__PublicHttpBaseUrl"] = "https://192.0.2.1/";
startInfo.Environment["Rendezvous__Deployment__PublicUdpHost"] = "203.0.113.1";
startInfo.Environment["AllowedHosts"] = "192.0.2.1";
}
else if (scenario == "missing-key")
{
File.Delete(secretPath);
}
process = Process.Start(startInfo)
?? throw new InvalidOperationException("The rejected production process did not start.");
Task<string> standardOutput = process.StandardOutput.ReadToEndAsync();
Task<string> standardError = process.StandardError.ReadToEndAsync();
using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(6));
await process.WaitForExitAsync(timeout.Token);
string diagnostic = $"{await standardOutput}\n{await standardError}";
Assert.NotEqual(0, process.ExitCode);
Assert.Contains(expectedDiagnostic, diagnostic, StringComparison.OrdinalIgnoreCase);
AssertTcpPortIsReleased(httpPort);
AssertUdpPortIsReleased(udpPort);
}
finally
{
await StopProcessTreeAsync(process);
File.Delete(secretPath);
}
}
private static ProcessStartInfo CreateStartInfo(int httpPort, int udpPort, string secretPath)
{
string assembly = typeof(Program).Assembly.Location;
ProcessStartInfo info = new()
{
FileName = "dotnet",
WorkingDirectory = Path.GetDirectoryName(assembly)!,
RedirectStandardOutput = true,
RedirectStandardError = true,
UseShellExecute = false,
};
info.ArgumentList.Add(assembly);
Dictionary<string, string> settings = new(StringComparer.Ordinal)
{
["ASPNETCORE_ENVIRONMENT"] = "Production",
["ASPNETCORE_URLS"] = $"http://127.0.0.1:{httpPort}",
["AllowedHosts"] = "127.0.0.1",
["Rendezvous__Deployment__PublicHttpBaseUrl"] = "https://127.0.0.1/",
["Rendezvous__Deployment__PublicUdpHost"] = "127.0.0.1",
["Rendezvous__Deployment__PublicUdpPort"] = udpPort.ToString(System.Globalization.CultureInfo.InvariantCulture),
["Rendezvous__Deployment__DrainDeadlineSeconds"] = "3",
["Rendezvous__Deployment__MinimumDrainSeconds"] = "1",
["Rendezvous__Deployment__SingleActiveInstance"] = "true",
["Rendezvous__Deployment__AllowPrivatePublicEndpoints"] = "true",
["Rendezvous__Udp__ListenAddress"] = "127.0.0.1",
["Rendezvous__Udp__Port"] = udpPort.ToString(System.Globalization.CultureInfo.InvariantCulture),
["Rendezvous__AbuseProtection__TrustedProxyAddresses__0"] = "127.0.0.1",
["Rendezvous__Provisioning__Issuer"] = "rendezvous-process-test",
["Rendezvous__Provisioning__Audience"] = "rendezvous-service",
["Rendezvous__Provisioning__ClockSkewSeconds"] = "30",
["Rendezvous__Provisioning__SigningKeys__0__KeyId"] = "process-test-key",
["Rendezvous__Provisioning__SigningKeys__0__SecretReference"] = $"file:{secretPath}",
["Rendezvous__Provisioning__SigningKeys__0__CredentialKinds__0"] = "DedicatedPublisher",
["Rendezvous__Provisioning__SigningKeys__0__GameId"] = "space-game",
["Rendezvous__Provisioning__SigningKeys__0__EnvironmentId"] = "process-test",
["Rendezvous__Provisioning__SigningKeys__0__NotBefore"] = DateTimeOffset.UtcNow.AddHours(-1).ToString("O"),
["Rendezvous__Provisioning__SigningKeys__0__SignUntil"] = DateTimeOffset.UtcNow.AddDays(1).ToString("O"),
["Rendezvous__Provisioning__SigningKeys__0__VerifyUntil"] = DateTimeOffset.UtcNow.AddDays(2).ToString("O"),
["Rendezvous__Provisioning__Games__0__GameId"] = "space-game",
["Rendezvous__Provisioning__Games__0__EnvironmentId"] = "process-test",
["Rendezvous__Provisioning__Games__0__Enabled"] = "true",
["Rendezvous__Provisioning__Games__0__ProtocolVersions__0"] = "1",
["Rendezvous__Provisioning__Games__0__Regions__0"] = "local",
["Rendezvous__Provisioning__Games__0__VisibilityModes__0"] = "Public",
["Rendezvous__Provisioning__Games__0__PublisherTrustModes__0"] = "ManagedDedicated",
["Rendezvous__Provisioning__Games__0__MetadataMaxBytes"] = "512",
["Rendezvous__Provisioning__Games__0__MetadataMaxKeys"] = "0",
["Rendezvous__Provisioning__Games__0__MaxListingsPerPrincipal"] = "10",
["Rendezvous__Provisioning__Games__0__MaxAnonymousListingsPerAddress"] = "0",
["Rendezvous__Provisioning__Games__0__MaxActiveJoinAttempts"] = "100",
["Rendezvous__Provisioning__Games__0__FallbackPolicy"] = "Disabled",
};
foreach ((string key, string value) in settings)
{
info.Environment[key] = value;
}
return info;
}
private static ProcessStartInfo CreateBareProductionStartInfo(int httpPort)
{
string assembly = typeof(Program).Assembly.Location;
ProcessStartInfo info = new()
{
FileName = "dotnet",
WorkingDirectory = Path.GetDirectoryName(assembly)!,
RedirectStandardOutput = true,
RedirectStandardError = true,
UseShellExecute = false,
};
info.ArgumentList.Add(assembly);
foreach (string key in info.Environment.Keys
.Where(static key => key.StartsWith("Rendezvous__", StringComparison.OrdinalIgnoreCase))
.ToArray())
{
info.Environment.Remove(key);
}
info.Environment["ASPNETCORE_ENVIRONMENT"] = "Production";
info.Environment["ASPNETCORE_URLS"] = $"http://127.0.0.1:{httpPort}";
info.Environment["AllowedHosts"] = "*";
return info;
}
private static async Task StopProcessTreeAsync(Process? process)
{
if (process is null)
{
return;
}
try
{
if (!process.HasExited)
{
process.Kill(entireProcessTree: true);
using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(3));
await process.WaitForExitAsync(timeout.Token);
}
}
finally
{
process.Dispose();
}
}
private static async Task SendSigtermAsync(Process process)
{
ProcessStartInfo signalInfo = new()
{
FileName = "/bin/kill",
UseShellExecute = false,
ArgumentList =
{
"-TERM",
process.Id.ToString(System.Globalization.CultureInfo.InvariantCulture),
},
};
using Process signal = Process.Start(signalInfo)
?? throw new InvalidOperationException("Could not send SIGTERM.");
await signal.WaitForExitAsync();
Assert.Equal(0, signal.ExitCode);
}
private static string BuildConfiguration()
{
string path = typeof(ProductionProcessTests).Assembly.Location;
return path.Contains(
$"{Path.DirectorySeparatorChar}Release{Path.DirectorySeparatorChar}",
StringComparison.Ordinal)
? "Release"
: "Debug";
}
private static string RepositoryRoot()
{
DirectoryInfo? directory = new(AppContext.BaseDirectory);
while (directory is not null && !File.Exists(Path.Combine(directory.FullName, "Rendezvous.slnx")))
{
directory = directory.Parent;
}
return directory?.FullName
?? throw new InvalidOperationException("Could not locate the repository root.");
}
private static async Task WaitForReadyAsync(int port, Process process, TimeSpan timeout)
{
using HttpClient client = new() { Timeout = TimeSpan.FromMilliseconds(500) };
Stopwatch elapsed = Stopwatch.StartNew();
while (elapsed.Elapsed < timeout)
{
if (process.HasExited)
{
throw new InvalidOperationException("The production server exited before readiness.");
}
try
{
using HttpResponseMessage response = await client.GetAsync(
$"http://127.0.0.1:{port}/health/ready");
if (response.StatusCode == HttpStatusCode.OK)
{
return;
}
}
catch (HttpRequestException)
{
}
catch (TaskCanceledException)
{
}
await Task.Delay(50);
}
throw new TimeoutException("The production server did not become ready.");
}
private static int ReserveTcpPort()
{
TcpListener listener = new(IPAddress.Loopback, 0);
listener.Start();
int port = ((IPEndPoint)listener.LocalEndpoint).Port;
listener.Stop();
return port;
}
private static int ReserveUdpPort()
{
using UdpClient client = new(new IPEndPoint(IPAddress.Loopback, 0));
return ((IPEndPoint)client.Client.LocalEndPoint!).Port;
}
private static void AssertUdpPortIsBound(int port)
{
using Socket socket = new(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp);
Assert.Throws<SocketException>(() => socket.Bind(new IPEndPoint(IPAddress.Loopback, port)));
}
private static void AssertTcpPortIsReleased(int port)
{
TcpListener listener = new(IPAddress.Loopback, port);
listener.Start();
listener.Stop();
}
private static void AssertUdpPortIsReleased(int port)
{
using Socket socket = new(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp);
socket.Bind(new IPEndPoint(IPAddress.Loopback, port));
}
}
@@ -0,0 +1,253 @@
using System.Collections.Concurrent;
using System.Diagnostics;
using System.Diagnostics.Metrics;
using System.Net;
using FinalFactory.Rendezvous.Server.Observability;
using FinalFactory.Rendezvous.Server.State;
using FinalFactory.Rendezvous.Server.Transport;
using FinalFactory.Rendezvous.Tests.State;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
namespace FinalFactory.Rendezvous.Tests.Observability;
[CollectionDefinition(RendezvousTelemetryIsolation.Name, DisableParallelization = true)]
public sealed class RendezvousTelemetryIsolation
{
public const string Name = "Rendezvous telemetry";
}
[Collection(RendezvousTelemetryIsolation.Name)]
public sealed class ObservabilityTests
{
private static readonly HashSet<string> AllowedTagKeys =
[
"operation",
"status_code",
"result",
"transport",
"partition",
"action",
"outcome",
"elapsed_bucket",
];
[Fact]
public void MetricsAndTracesUseBoundedDimensionsWithoutSensitiveValues()
{
EphemeralStateFixture fixture = new();
using RendezvousTelemetry telemetry = new(fixture.Store);
List<Measurement> measurements = [];
using MeterListener meterListener = new();
meterListener.InstrumentPublished = (instrument, listener) =>
{
if (instrument.Meter.Name == RendezvousTelemetry.MeterName)
{
listener.EnableMeasurementEvents(instrument);
}
};
meterListener.SetMeasurementEventCallback<long>((instrument, value, tags, _) =>
measurements.Add(new(instrument.Name, value, Tags(tags))));
meterListener.SetMeasurementEventCallback<int>((instrument, value, tags, _) =>
measurements.Add(new(instrument.Name, value, Tags(tags))));
meterListener.SetMeasurementEventCallback<double>((instrument, value, tags, _) =>
measurements.Add(new(instrument.Name, value, Tags(tags))));
meterListener.Start();
Activity? observed = null;
List<string> activityData = [];
using ActivityListener activityListener = new()
{
ShouldListenTo = source => source.Name == RendezvousTelemetry.ActivitySourceName,
Sample = static (ref ActivityCreationOptions<ActivityContext> _) =>
ActivitySamplingResult.AllData,
ActivityStopped = activity =>
{
observed = activity;
activityData.Add(activity.DisplayName);
activityData.AddRange(activity.TagObjects.Select(static tag => $"{tag.Key}={tag.Value}"));
},
};
ActivitySource.AddActivityListener(activityListener);
const string secret = "secret-player-token-canary";
using (telemetry.StartActivity("HTTP GetOperatorStatus", ActivityKind.Server))
{
telemetry.RecordHttp("GetOperatorStatus", 200, 3.5);
telemetry.RecordUdp("frozen", "Introduced", 1.25);
telemetry.RecordLimiterDrop("udp", "rate-or-concurrency");
telemetry.RecordAudit("revoke-listing", "succeeded");
telemetry.RecordConnectionOutcome("Connected", "UnderOneSecond");
telemetry.RecordOperatorAuthentication("accepted");
telemetry.RecordPairingLatency(12.5);
}
NatMediationProcessor processor = new(null!, null!, null!, telemetry: telemetry);
Assert.Equal(
NatMediationResult.Dropped,
processor.ProcessRequest(
new IPEndPoint(IPAddress.Parse("10.0.0.8"), 9000),
new IPEndPoint(IPAddress.Parse("203.0.113.8"), 50000),
secret,
NoopIntroductionSink.Instance));
meterListener.RecordObservableInstruments();
Assert.NotNull(observed);
string flattened = string.Join('|', measurements.Select(static item => item.ToString()));
Assert.DoesNotContain(secret, flattened, StringComparison.Ordinal);
Assert.DoesNotContain(secret, string.Join('|', activityData), StringComparison.Ordinal);
Assert.Contains(measurements, static item => item.Name == "rendezvous.http.requests");
Assert.Contains(measurements, static item => item.Name == "rendezvous.udp.results");
Assert.Contains(measurements, static item => item.Name == "rendezvous.limiter.drops");
Assert.Contains(measurements, static item => item.Name == "rendezvous.queue.depth");
Assert.Contains(measurements, static item => item.Name == "rendezvous.store.active_leases");
Assert.Contains(measurements, static item => item.Name == "rendezvous.store.expiry_churn");
Assert.Contains(measurements, static item => item.Name == "rendezvous.pairing.latency");
Assert.All(measurements.SelectMany(static item => item.Tags), static tag =>
Assert.Contains(tag.Key, AllowedTagKeys));
}
[Fact]
public void MetricScrapeExpiresIdleStateAndReportsExpiryChurn()
{
EphemeralStateFixture fixture = new();
Assert.True(fixture.Store.CreateListing(fixture.ListingCommand()).Succeeded);
using RendezvousTelemetry telemetry = new(fixture.Store);
List<Measurement> measurements = [];
using MeterListener listener = new();
listener.InstrumentPublished = (instrument, meterListener) =>
{
if (instrument.Meter.Name == RendezvousTelemetry.MeterName)
{
meterListener.EnableMeasurementEvents(instrument);
}
};
listener.SetMeasurementEventCallback<long>((instrument, value, tags, _) =>
measurements.Add(new(instrument.Name, value, Tags(tags))));
listener.SetMeasurementEventCallback<int>((instrument, value, tags, _) =>
measurements.Add(new(instrument.Name, value, Tags(tags))));
listener.Start();
listener.RecordObservableInstruments();
Assert.Equal(
1,
Assert.Single(measurements, static item => item.Name == "rendezvous.store.active_leases").Value);
fixture.Clock.Advance(TimeSpan.FromSeconds(61));
measurements.Clear();
listener.RecordObservableInstruments();
Assert.Equal(
0,
Assert.Single(measurements, static item => item.Name == "rendezvous.store.active_leases").Value);
Assert.True(Assert.Single(
measurements,
static item => item.Name == "rendezvous.store.expiry_churn").Value >= 1);
}
[Fact]
public void AuditTrailFingerprintsIdentifiersEnforcesRetentionAndBoundsCapacity()
{
EphemeralStateFixture fixture = new();
using RendezvousTelemetry telemetry = new(fixture.Store);
CapturingLogger<AuditTrail> logger = new();
ManualTimeProvider time = new(new DateTimeOffset(2026, 7, 16, 0, 0, 0, TimeSpan.Zero));
AuditTrail audit = new(
Options.Create(new AuditOptions { MaxEntries = 100, RetentionDays = 1 }),
logger,
telemetry,
time);
const string actor = "operator-secret-subject";
const string target = "player-secret-subject";
for (int index = 0; index < 101; index++)
{
audit.Record(actor, "revoke-principal", "succeeded", "principal", target, "safe-correlation");
}
IReadOnlyList<AuditEntry> bounded = audit.GetEntriesForTests();
Assert.Equal(100, bounded.Count);
Assert.All(bounded, entry =>
{
Assert.DoesNotContain(actor, entry.ToString(), StringComparison.Ordinal);
Assert.DoesNotContain(target, entry.ToString(), StringComparison.Ordinal);
Assert.NotEqual(actor, entry.ActorFingerprint);
Assert.NotEqual(target, entry.TargetFingerprint);
});
Assert.DoesNotContain(actor, string.Join('|', logger.Messages), StringComparison.Ordinal);
Assert.DoesNotContain(target, string.Join('|', logger.Messages), StringComparison.Ordinal);
time.Advance(TimeSpan.FromDays(2));
Assert.Empty(audit.GetEntriesForTests());
Assert.Empty(audit.GetAggregateCounts());
audit.Record(actor, "inspect-status", "succeeded", "service", "rendezvous", "safe-correlation");
Assert.Single(audit.GetEntriesForTests());
Assert.Equal(1, audit.GetAggregateCounts()["inspect-status:succeeded"]);
}
private static KeyValuePair<string, object?>[] Tags(
ReadOnlySpan<KeyValuePair<string, object?>> tags) => tags.ToArray();
private sealed record Measurement(
string Name,
double Value,
KeyValuePair<string, object?>[] Tags);
private sealed class ManualTimeProvider(DateTimeOffset now) : TimeProvider
{
private DateTimeOffset _now = now;
public override DateTimeOffset GetUtcNow() => _now;
public void Advance(TimeSpan duration) => _now += duration;
}
private sealed class NoopIntroductionSink : INatIntroductionSink
{
public static NoopIntroductionSink Instance { get; } = new();
public void Introduce(NatIntroductionPlan plan)
{
}
}
}
internal sealed class CapturingLogger<T> : ILogger<T>
{
public List<string> Messages { get; } = [];
public IDisposable? BeginScope<TState>(TState state) where TState : notnull => null;
public bool IsEnabled(LogLevel logLevel) => true;
public void Log<TState>(
LogLevel logLevel,
EventId eventId,
TState state,
Exception? exception,
Func<TState, Exception?, string> formatter) => Messages.Add(formatter(state, exception));
}
internal sealed class CapturingLoggerProvider : ILoggerProvider
{
public ConcurrentQueue<string> Messages { get; } = new();
public ILogger CreateLogger(string categoryName) => new Sink(Messages);
public void Dispose() => GC.SuppressFinalize(this);
private sealed class Sink(ConcurrentQueue<string> messages) : ILogger
{
public IDisposable BeginScope<TState>(TState state) where TState : notnull => Scope.Instance;
public bool IsEnabled(LogLevel logLevel) => true;
public void Log<TState>(
LogLevel logLevel,
EventId eventId,
TState state,
Exception? exception,
Func<TState, Exception?, string> formatter) => messages.Enqueue(formatter(state, exception));
}
private sealed class Scope : IDisposable
{
public static Scope Instance { get; } = new();
public void Dispose()
{
}
}
}
@@ -0,0 +1,519 @@
using System.Diagnostics;
using System.Diagnostics.Metrics;
using System.Net;
using System.Net.Http.Headers;
using System.Net.Http.Json;
using System.Text.Json;
using FinalFactory.Rendezvous.Contracts;
using FinalFactory.Rendezvous.Server.Abuse;
using FinalFactory.Rendezvous.Server.Http;
using FinalFactory.Rendezvous.Server.JoinAttempts;
using FinalFactory.Rendezvous.Server.Observability;
using FinalFactory.Rendezvous.Server.Operations;
using FinalFactory.Rendezvous.Server.Provisioning;
using FinalFactory.Rendezvous.Server.Sessions;
using FinalFactory.Rendezvous.Server.State;
using FinalFactory.Rendezvous.Server.Transport;
using FinalFactory.Rendezvous.Tests.Observability;
using FinalFactory.Rendezvous.Tests.Provisioning;
using FinalFactory.Rendezvous.Tests.State;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Hosting.Server;
using Microsoft.AspNetCore.Hosting.Server.Features;
using Microsoft.AspNetCore.Http;
using Microsoft.AspNetCore.Routing;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
namespace FinalFactory.Rendezvous.Tests.Operations;
[Collection(RendezvousTelemetryIsolation.Name)]
public sealed class OperatorEndpointTests
{
[Fact]
public async Task OperatorSurfaceSeparatesAuthenticationConfirmsActionsAndRedactsInspection()
{
await using OperatorTestHost host = await OperatorTestHost.StartAsync();
List<string> telemetryData = [];
using MeterListener meterListener = new();
meterListener.InstrumentPublished = (instrument, listener) =>
{
if (instrument.Meter.Name == RendezvousTelemetry.MeterName)
{
listener.EnableMeasurementEvents(instrument);
}
};
meterListener.SetMeasurementEventCallback<long>((instrument, value, tags, _) =>
CaptureMeasurement(telemetryData, instrument, value, tags));
meterListener.SetMeasurementEventCallback<int>((instrument, value, tags, _) =>
CaptureMeasurement(telemetryData, instrument, value, tags));
meterListener.SetMeasurementEventCallback<double>((instrument, value, tags, _) =>
CaptureMeasurement(telemetryData, instrument, value, tags));
meterListener.Start();
using ActivityListener activityListener = new()
{
ShouldListenTo = static source => source.Name == RendezvousTelemetry.ActivitySourceName,
Sample = static (ref ActivityCreationOptions<ActivityContext> _) =>
ActivitySamplingResult.AllData,
ActivityStopped = activity =>
{
telemetryData.Add(activity.DisplayName);
telemetryData.AddRange(activity.TagObjects.Select(static tag => $"{tag.Key}={tag.Value}"));
},
};
ActivitySource.AddActivityListener(activityListener);
using HttpResponseMessage liveBeforeDependencies = await host.Client.GetAsync("/health/live");
Assert.Equal(HttpStatusCode.OK, liveBeforeDependencies.StatusCode);
using HttpResponseMessage readyBeforeUdp = await host.Client.GetAsync("/health/ready");
Assert.Equal(HttpStatusCode.ServiceUnavailable, readyBeforeUdp.StatusCode);
using HttpResponseMessage unauthenticated = await host.Client.GetAsync("/v1/operator/status");
Assert.Equal(HttpStatusCode.Unauthorized, unauthenticated.StatusCode);
AuthenticationHeaderValue challenge = Assert.Single(
unauthenticated.Headers.WwwAuthenticate);
Assert.Equal("Bearer", challenge.Scheme);
Assert.Equal("realm=\"operator\"", challenge.Parameter);
using HttpResponseMessage publisher = await SendAsync(
host,
HttpMethod.Get,
"/v1/operator/status",
host.PublisherCredential);
Assert.Equal(HttpStatusCode.Unauthorized, publisher.StatusCode);
using HttpResponseMessage status = await SendAsync(
host,
HttpMethod.Get,
"/v1/operator/status",
host.ReadOnlyOperatorCredential);
Assert.Equal(HttpStatusCode.OK, status.StatusCode);
string statusJson = await status.Content.ReadAsStringAsync();
Assert.Contains("not-ready", statusJson, StringComparison.Ordinal);
Assert.DoesNotContain(host.OwnerCanary, statusJson, StringComparison.Ordinal);
Assert.DoesNotContain("203.0.113.25", statusJson, StringComparison.Ordinal);
Assert.DoesNotContain("metadata", statusJson, StringComparison.OrdinalIgnoreCase);
OperatorStatusResponse? operatorStatus = JsonSerializer.Deserialize<OperatorStatusResponse>(
statusJson,
ContractJson.Options);
Assert.Contains(operatorStatus!.Tenants, static tenant =>
tenant.GameId == "space-game"
&& tenant.EnvironmentId == "production"
&& tenant.Status == "enabled");
Assert.Contains(operatorStatus.SigningKeys, static key =>
key.KeyId == OperatorTestHost.OperatorKeyId
&& key.Status == "signing"
&& key.CredentialKinds.SequenceEqual(["Operator"]));
await host.StartUdpAsync();
using HttpResponseMessage readyAfterUdp = await host.Client.GetAsync("/health/ready");
Assert.Equal(HttpStatusCode.OK, readyAfterUdp.StatusCode);
using HttpResponseMessage exception = await SendAsync(
host,
HttpMethod.Post,
"/test/exception",
host.FullOperatorCredential);
Assert.Equal(HttpStatusCode.InternalServerError, exception.StatusCode);
using HttpResponseMessage saturatedPublic = await host.Client.GetAsync("/test/public");
Assert.Equal(HttpStatusCode.TooManyRequests, saturatedPublic.StatusCode);
using HttpResponseMessage forbidden = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/drain",
host.ReadOnlyOperatorCredential,
new BeginDrainRequest { Confirmation = "DRAIN" });
Assert.Equal(HttpStatusCode.Forbidden, forbidden.StatusCode);
SessionListingId listingId = host.CreateListing(host.OwnerCanary);
using HttpResponseMessage unconfirmedListing = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/listings/revoke",
host.FullOperatorCredential,
new RevokeListingRequest
{
ListingId = listingId.ToString(),
ConfirmListingId = Guid.NewGuid().ToString("D"),
});
Assert.Equal(HttpStatusCode.BadRequest, unconfirmedListing.StatusCode);
Assert.True(host.Store.GetListing(listingId, false).Succeeded);
using HttpResponseMessage revokedListing = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/listings/revoke",
host.FullOperatorCredential,
new RevokeListingRequest
{
ListingId = listingId.ToString(),
ConfirmListingId = listingId.ToString(),
});
Assert.Equal(HttpStatusCode.OK, revokedListing.StatusCode);
Assert.Equal(StoreResultCode.NotFound, host.Store.GetListing(listingId, false).Code);
const string principalCanary = "publisher-player-canary";
SessionListingId principalListing = host.CreateListing(principalCanary);
using HttpResponseMessage unconfirmedPrincipal = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/principals/revoke",
host.FullOperatorCredential,
new RevokePrincipalRequest
{
Subject = principalCanary,
ConfirmSubject = "different-subject",
LifetimeSeconds = 60,
});
Assert.Equal(HttpStatusCode.BadRequest, unconfirmedPrincipal.StatusCode);
Assert.True(host.Store.GetListing(principalListing, false).Succeeded);
using HttpResponseMessage revokedPrincipal = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/principals/revoke",
host.FullOperatorCredential,
new RevokePrincipalRequest
{
Subject = principalCanary,
ConfirmSubject = principalCanary,
LifetimeSeconds = 60,
});
Assert.Equal(HttpStatusCode.OK, revokedPrincipal.StatusCode);
OperatorActionResponse? principalResult = await revokedPrincipal.Content
.ReadFromJsonAsync<OperatorActionResponse>(ContractJson.Options);
Assert.Equal(1, principalResult!.AffectedResources);
Assert.Equal(StoreResultCode.NotFound, host.Store.GetListing(principalListing, false).Code);
StoreResult<StoredListing> blockedPublisher = host.CreateListingResult(
principalCanary,
out _);
Assert.Equal(StoreResultCode.Revoked, blockedPublisher.Code);
using HttpResponseMessage drain = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/drain",
host.FullOperatorCredential,
new BeginDrainRequest { Confirmation = "DRAIN" });
Assert.Equal(HttpStatusCode.OK, drain.StatusCode);
Assert.True(host.Store.IsDraining);
using HttpResponseMessage liveDuringDrain = await host.Client.GetAsync("/health/live");
Assert.Equal(HttpStatusCode.OK, liveDuringDrain.StatusCode);
using HttpResponseMessage readyDuringDrain = await host.Client.GetAsync("/health/ready");
Assert.Equal(HttpStatusCode.ServiceUnavailable, readyDuringDrain.StatusCode);
using HttpResponseMessage firstPublisherKeyRevocation = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/keys/revoke",
host.FullOperatorCredential,
new RevokeSigningKeyRequest { KeyId = "key-1", ConfirmKeyId = "key-1" });
Assert.Equal(HttpStatusCode.OK, firstPublisherKeyRevocation.StatusCode);
using HttpResponseMessage repeatedPublisherKeyRevocation = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/keys/revoke",
host.FullOperatorCredential,
new RevokeSigningKeyRequest { KeyId = "key-1", ConfirmKeyId = "key-1" });
Assert.Equal(HttpStatusCode.OK, repeatedPublisherKeyRevocation.StatusCode);
using HttpResponseMessage unconfirmedKey = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/keys/revoke",
host.FullOperatorCredential,
new RevokeSigningKeyRequest
{
KeyId = OperatorTestHost.OperatorKeyId,
ConfirmKeyId = "different-key",
});
Assert.Equal(HttpStatusCode.BadRequest, unconfirmedKey.StatusCode);
using HttpResponseMessage revokedKey = await SendAsync(
host,
HttpMethod.Post,
"/v1/operator/keys/revoke",
host.FullOperatorCredential,
new RevokeSigningKeyRequest
{
KeyId = OperatorTestHost.OperatorKeyId,
ConfirmKeyId = OperatorTestHost.OperatorKeyId,
});
Assert.Equal(HttpStatusCode.OK, revokedKey.StatusCode);
using HttpResponseMessage afterKeyRevocation = await SendAsync(
host,
HttpMethod.Get,
"/v1/operator/status",
host.FullOperatorCredential);
Assert.Equal(HttpStatusCode.Unauthorized, afterKeyRevocation.StatusCode);
string auditText = string.Join('|', host.Audit.GetEntriesForTests());
string logText = string.Join('|', host.AuditLogger.Messages.Concat(host.AllLogs.Messages));
Assert.DoesNotContain(host.OwnerCanary, auditText, StringComparison.Ordinal);
Assert.DoesNotContain(principalCanary, auditText, StringComparison.Ordinal);
Assert.DoesNotContain(listingId.ToString(), auditText, StringComparison.Ordinal);
Assert.DoesNotContain(host.OwnerCanary, logText, StringComparison.Ordinal);
Assert.DoesNotContain(principalCanary, logText, StringComparison.Ordinal);
Assert.DoesNotContain("exception-secret-canary", logText, StringComparison.Ordinal);
Assert.DoesNotContain(host.FullOperatorCredential, logText, StringComparison.Ordinal);
meterListener.RecordObservableInstruments();
string telemetryText = string.Join('|', telemetryData);
Assert.DoesNotContain(host.OwnerCanary, telemetryText, StringComparison.Ordinal);
Assert.DoesNotContain(principalCanary, telemetryText, StringComparison.Ordinal);
Assert.DoesNotContain("exception-secret-canary", telemetryText, StringComparison.Ordinal);
Assert.DoesNotContain(host.FullOperatorCredential, telemetryText, StringComparison.Ordinal);
Assert.DoesNotContain(listingId.ToString(), telemetryText, StringComparison.Ordinal);
Assert.DoesNotContain("203.0.113.25", telemetryText, StringComparison.Ordinal);
Assert.Contains(host.Audit.GetEntriesForTests(), static entry =>
entry.Action == "begin-drain" && entry.Result == "succeeded");
Assert.Contains(host.Audit.GetEntriesForTests(), static entry =>
entry.Action == "revoke-listing" && entry.Result == "rejected");
Assert.Contains(host.Audit.GetEntriesForTests(), static entry =>
entry.Action == "begin-drain" && entry.Result == "forbidden");
}
private static async Task<HttpResponseMessage> SendAsync(
OperatorTestHost host,
HttpMethod method,
string path,
string bearer,
object? body = null)
{
using HttpRequestMessage request = new(method, path);
request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", bearer);
if (body is not null)
{
request.Content = JsonContent.Create(body, options: ContractJson.Options);
}
return await host.Client.SendAsync(request);
}
private static void CaptureMeasurement<T>(
List<string> destination,
Instrument instrument,
T value,
ReadOnlySpan<KeyValuePair<string, object?>> tags)
where T : struct
{
destination.Add($"{instrument.Name}={value}");
destination.AddRange(tags.ToArray().Select(static tag => $"{tag.Key}={tag.Value}"));
}
private sealed class OperatorTestHost : IAsyncDisposable
{
private readonly WebApplication _application;
private int _listingSequence;
private bool _udpStarted;
private OperatorTestHost(
WebApplication application,
HttpClient client,
ManualRendezvousClock clock,
InMemoryEphemeralRendezvousStore store,
AuditTrail audit,
CapturingLogger<AuditTrail> auditLogger,
CapturingLoggerProvider allLogs,
string publisherCredential,
string readOnlyOperatorCredential,
string fullOperatorCredential)
{
_application = application;
Client = client;
Clock = clock;
Store = store;
Audit = audit;
AuditLogger = auditLogger;
AllLogs = allLogs;
PublisherCredential = publisherCredential;
ReadOnlyOperatorCredential = readOnlyOperatorCredential;
FullOperatorCredential = fullOperatorCredential;
}
internal const string OperatorKeyId = "operator-key";
internal string OwnerCanary { get; } = "publisher-owner-canary";
internal HttpClient Client { get; }
internal ManualRendezvousClock Clock { get; }
internal InMemoryEphemeralRendezvousStore Store { get; }
internal AuditTrail Audit { get; }
internal CapturingLogger<AuditTrail> AuditLogger { get; }
internal CapturingLoggerProvider AllLogs { get; }
internal string PublisherCredential { get; }
internal string ReadOnlyOperatorCredential { get; }
internal string FullOperatorCredential { get; }
internal static async Task<OperatorTestHost> StartAsync()
{
ManualRendezvousClock clock = new(ProvisioningTestData.Now);
EphemeralStoreOptions stateOptions = new();
InMemoryEphemeralRendezvousStore store = new(stateOptions, clock, clock);
SigningKeyOptions publisherKey = ProvisioningTestData.CreateKey();
SigningKeyOptions operatorKey = ProvisioningTestData.CreateKey(
OperatorKeyId,
"operator-secret",
credentialKinds: [PrincipalCredentialKind.Operator],
gameId: null,
environmentId: null);
ProvisioningOptions options = ProvisioningTestData.CreateOptions();
options.SigningKeys = [publisherKey, operatorKey];
ProvisioningRuntime provisioning = ProvisioningRuntime.Create(
options,
ProvisioningTestData.CreateSecrets("secret-1", "operator-secret"),
clock.UtcNow);
string publisherCredential = provisioning.Credentials.Issue(
ProvisioningTestData.CreateDedicatedPublisher(),
clock.UtcNow);
string readOnlyCredential = provisioning.Credentials.Issue(
new OperatorPrincipal(
"operator-readonly",
clock.UtcNow.AddMinutes(10),
[OperatorPermission.ReadPolicy]),
clock.UtcNow);
string fullCredential = provisioning.Credentials.Issue(
new OperatorPrincipal(
"operator-full",
clock.UtcNow.AddMinutes(10),
Enum.GetValues<OperatorPermission>()),
clock.UtcNow);
CapturingLogger<AuditTrail> auditLogger = new();
CapturingLoggerProvider allLogs = new();
EphemeralCapabilityIssuer capabilities = new();
WebApplicationBuilder builder = WebApplication.CreateBuilder();
builder.WebHost.UseUrls("http://127.0.0.1:0");
builder.Logging.ClearProviders();
builder.Logging.SetMinimumLevel(LogLevel.Debug);
builder.Logging.AddProvider(allLogs);
builder.Logging.AddFilter(
"Microsoft.AspNetCore.Diagnostics.ExceptionHandlerMiddleware",
LogLevel.None);
builder.Services.ConfigureHttpJsonOptions(static json =>
ContractJson.Configure(json.SerializerOptions));
builder.Services.Configure<RouteHandlerOptions>(static route =>
route.ThrowOnBadRequest = true);
builder.Services.AddProblemDetails();
builder.Services.AddExceptionHandler<RendezvousExceptionHandler>();
builder.Services.AddOptions<AbuseProtectionOptions>().Configure(static abuse =>
{
abuse.OperatorAllowedAddresses = ["127.0.0.1"];
abuse.HttpGlobalRequestsPerWindow = 1;
abuse.HttpOptionalRequestsPerWindow = 1;
abuse.HttpIpPrefixRequestsPerWindow = 1;
abuse.HttpOptionalIpPrefixRequestsPerWindow = 1;
});
builder.Services.AddOptions<AuditOptions>();
builder.Services.AddOptions<UdpMediatorOptions>().Configure(static udp =>
{
udp.ListenAddress = "127.0.0.1";
udp.Port = 0;
});
builder.Services.AddSingleton(provisioning);
builder.Services.AddSingleton(provisioning.Policies);
builder.Services.AddSingleton(provisioning.Credentials);
builder.Services.AddSingleton(provisioning.PublisherAuthorization);
builder.Services.AddSingleton(store);
builder.Services.AddSingleton<IEphemeralRendezvousStore>(store);
builder.Services.AddSingleton<IWallClock>(clock);
builder.Services.AddSingleton<IMonotonicClock>(clock);
builder.Services.AddSingleton(capabilities);
builder.Services.AddSingleton<ISessionCapabilityService>(capabilities);
builder.Services.AddSingleton<JoinAttemptCursorCodec>();
builder.Services.AddSingleton<JoinAttemptService>();
builder.Services.AddSingleton<RendezvousTelemetry>();
builder.Services.AddSingleton<AbuseProtectionService>();
builder.Services.AddSingleton<NatMediationProcessor>();
builder.Services.AddSingleton<UdpMediatorService>();
builder.Services.AddSingleton(new ProvisioningReadiness(true));
builder.Services.AddSingleton<RendezvousReadiness>();
builder.Services.AddSingleton<ILogger<AuditTrail>>(auditLogger);
builder.Services.AddSingleton<AuditTrail>();
builder.Services.AddSingleton<OperatorService>();
WebApplication app = builder.Build();
app.UseMiddleware<TelemetryMiddleware>();
app.UseExceptionHandler();
app.UseMiddleware<HttpAbuseProtectionMiddleware>();
app.MapOperatorEndpoints();
app.MapRendezvousHealthEndpoints();
app.MapGet("/test/public", static () => Results.Ok()).WithName("TestPublic");
app.MapPost(
"/test/exception",
static IResult () => throw new InvalidOperationException("exception-secret-canary"))
.WithName("TestSecretException");
await app.StartAsync();
IServer server = app.Services.GetRequiredService<IServer>();
string address = Assert.Single(server.Features.Get<IServerAddressesFeature>()!.Addresses);
return new OperatorTestHost(
app,
new HttpClient { BaseAddress = new Uri(address) },
clock,
store,
app.Services.GetRequiredService<AuditTrail>(),
auditLogger,
allLogs,
publisherCredential,
readOnlyCredential,
fullCredential);
}
internal SessionListingId CreateListing(string owner)
{
StoreResult<StoredListing> result = CreateListingResult(owner, out SessionListingId listingId);
Assert.True(result.Succeeded);
return listingId;
}
internal StoreResult<StoredListing> CreateListingResult(
string owner,
out SessionListingId listingId)
{
int sequence = Interlocked.Increment(ref _listingSequence);
listingId = new(Guid.NewGuid());
StoreResult<StoredListing> result = Store.CreateListing(new(
$"operator-listing-{sequence}",
$"operator-request-{sequence}",
new ListingDefinition
{
ListingId = listingId,
LeaseId = new(Guid.NewGuid()),
Scope = new(new GameId("space-game"), new EnvironmentId("production")),
OwnerSubject = owner,
RegionId = new("eu-central"),
ProtocolVersion = 7,
BuildVersion = "1.0.0",
DisplayName = "Operator test listing",
Visibility = ListingVisibility.Public,
TrustMode = PublisherTrustMode.ManagedDedicated,
CurrentPlayers = 1,
MaximumPlayers = 4,
Metadata = new Dictionary<string, string> { ["mode"] = "online-coop" },
LeaseFingerprint = new("lease-fingerprint"),
HostPresenceHandle = new(Guid.NewGuid()),
HostPresenceFingerprint = new("presence-fingerprint"),
CapabilityDerivationSalt = new string('A', 43),
}));
return result;
}
internal async Task StartUdpAsync()
{
await _application.Services.GetRequiredService<UdpMediatorService>()
.StartAsync(CancellationToken.None);
_udpStarted = true;
}
public async ValueTask DisposeAsync()
{
Client.Dispose();
if (_udpStarted)
{
await _application.Services.GetRequiredService<UdpMediatorService>()
.StopAsync(CancellationToken.None);
}
await _application.StopAsync();
await _application.DisposeAsync();
}
}
}
@@ -94,6 +94,7 @@ public sealed class PrincipalCredentialTests
CredentialValidationError.SignatureInvalid,
service.Validate(tampered, ProvisioningTestData.Now).Error);
Assert.True(keys.Revoke("key-1"));
Assert.True(keys.Revoke("key-1"));
Assert.Equal(
CredentialValidationError.KeyRevoked,
service.Validate(token, ProvisioningTestData.Now).Error);
@@ -0,0 +1,66 @@
using System.Security.Cryptography;
using FinalFactory.Rendezvous.Server.Provisioning;
namespace FinalFactory.Rendezvous.Tests.Provisioning;
public sealed class ProductionSecretProviderTests : IDisposable
{
private readonly List<string> _paths = [];
[Fact]
public void ReadsBoundedSecretFromAbsoluteReadOnlyFile()
{
byte[] expected = RandomNumberGenerator.GetBytes(32);
string path = CreateSecretFile(expected);
EnvironmentSecretProvider provider = new();
bool found = provider.TryGetSecret($"file:{path}", out SecretMaterial? material);
Assert.True(found);
using (material)
{
Assert.Equal(expected, material!.CopyBytes());
}
}
[Fact]
public void RejectsRelativeSymlinkEmptyAndOversizedFileReferences()
{
string empty = CreateSecretFile([]);
string oversized = CreateSecretFile(new byte[4097]);
string target = CreateSecretFile(RandomNumberGenerator.GetBytes(32));
string symlink = Path.Combine(Path.GetTempPath(), $"rendezvous-secret-link-{Guid.NewGuid():N}");
EnvironmentSecretProvider provider = new();
Assert.False(provider.TryGetSecret("file:relative-secret", out _));
Assert.False(provider.TryGetSecret($"file:{empty}", out _));
Assert.False(provider.TryGetSecret($"file:{oversized}", out _));
Assert.False(provider.TryGetSecret("file:/tmp/invalid\0path", out _));
if (!OperatingSystem.IsWindows())
{
File.CreateSymbolicLink(symlink, target);
_paths.Add(symlink);
Assert.False(provider.TryGetSecret($"file:{symlink}", out _));
}
}
public void Dispose()
{
foreach (string path in _paths)
{
File.Delete(path);
}
}
private string CreateSecretFile(byte[] bytes)
{
string path = Path.Combine(Path.GetTempPath(), $"rendezvous-secret-{Guid.NewGuid():N}");
File.WriteAllBytes(path, bytes);
if (OperatingSystem.IsLinux() || OperatingSystem.IsMacOS() || OperatingSystem.IsFreeBSD())
{
File.SetUnixFileMode(path, UnixFileMode.UserRead);
}
_paths.Add(path);
return path;
}
}
@@ -6,6 +6,7 @@ using FinalFactory.Rendezvous.Server.Http;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Http;
using Microsoft.AspNetCore.HttpOverrides;
using Microsoft.AspNetCore.Routing;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
@@ -86,6 +87,9 @@ public sealed class AbuseProtectionTests
AbuseProtectionService protection = new(Options.Create(options));
IPAddress source = IPAddress.Parse("198.51.100.10");
Assert.True(protection.IsOperatorSourceAllowed(source));
Assert.True(protection.IsOperatorSourceAllowed(IPAddress.Parse("::ffff:198.51.100.10")));
Assert.False(protection.IsOperatorSourceAllowed(IPAddress.Parse("198.51.100.11")));
AssertAccepted(protection, source, "BrowseSessions");
AssertAccepted(protection, source, "BrowseSessions");
AssertRejected(protection, source, "BrowseSessions");
@@ -93,6 +97,57 @@ public sealed class AbuseProtectionTests
AssertRejected(protection, source, "RenewSessionLease");
}
[Fact]
public void PublicSaturationCannotConsumeTheOperatorPartition()
{
AbuseProtectionOptions options = PermissiveOptions();
options.HttpGlobalRequestsPerWindow = 1;
options.HttpOptionalRequestsPerWindow = 1;
options.OperatorGlobalRequestsPerWindow = 1;
AbuseProtectionService protection = new(Options.Create(options));
IPAddress source = IPAddress.Parse("198.51.100.10");
AssertAccepted(protection, source, "BrowseSessions");
AssertRejected(protection, source, "BrowseSessions");
Assert.True(protection.TryAcquireOperatorIngress(source, out var lease, out _));
lease!.Dispose();
Assert.False(protection.TryAcquireOperatorIngress(source, out _, out _));
}
[Fact]
public async Task DeniedOperatorSourcesConsumeTheBoundedPublicPartition()
{
AbuseProtectionOptions options = PermissiveOptions();
options.OperatorAllowedAddresses = ["192.0.2.10"];
options.HttpGlobalRequestsPerWindow = 1;
options.HttpOptionalRequestsPerWindow = 1;
options.HttpIpPrefixRequestsPerWindow = 1;
options.HttpOptionalIpPrefixRequestsPerWindow = 1;
AbuseProtectionService protection = new(Options.Create(options));
bool dispatched = false;
HttpAbuseProtectionMiddleware middleware = new(
_ =>
{
dispatched = true;
return Task.CompletedTask;
},
protection);
DefaultHttpContext first = Context("198.51.100.10");
first.SetEndpoint(new Endpoint(
_ => Task.CompletedTask,
new EndpointMetadataCollection(new EndpointNameMetadata("GetOperatorStatus")),
"operator-status"));
await middleware.InvokeAsync(first);
Assert.Equal(StatusCodes.Status404NotFound, first.Response.StatusCode);
DefaultHttpContext repeated = Context("198.51.100.10");
repeated.SetEndpoint(first.GetEndpoint());
await middleware.InvokeAsync(repeated);
Assert.Equal(StatusCodes.Status429TooManyRequests, repeated.Response.StatusCode);
Assert.False(dispatched);
}
[Fact]
public void ResourceBudgetsRemainIsolatedAcrossTenantAndPrincipalScopes()
{
@@ -354,7 +409,8 @@ public sealed class AbuseProtectionTests
{
const string canary = "credential-canary <script> endpoint=203.0.113.8:9000";
DefaultHttpContext context = Context("198.51.100.10");
RendezvousExceptionHandler handler = new();
RendezvousExceptionHandler handler = new(
NullLogger<RendezvousExceptionHandler>.Instance);
Assert.True(await handler.TryHandleAsync(
context,
@@ -438,6 +494,11 @@ public sealed class AbuseProtectionTests
HealthGlobalConcurrency = 10_000,
HealthIpPrefixRequestsPerWindow = 10_000,
HealthIpPrefixConcurrency = 1_000,
OperatorAllowedAddresses = ["198.51.100.10"],
OperatorGlobalRequestsPerWindow = 10_000,
OperatorGlobalConcurrency = 10_000,
OperatorIpPrefixRequestsPerWindow = 10_000,
OperatorIpPrefixConcurrency = 1_000,
HttpGlobalRequestsPerWindow = 10_000,
HttpOptionalRequestsPerWindow = 9_000,
HttpIpPrefixRequestsPerWindow = 10_000,
@@ -411,6 +411,92 @@ public sealed class InMemoryEphemeralRendezvousStoreTests
Assert.True(fixture.Store.CreateJoinAttempt(attempt).Succeeded);
Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.CreateJoinAttempt(
fixture.AttemptCommand(listing, "client-quota-2") with { ScopeAttemptLimit = 1 }).Code);
fixture.Clock.Advance(TimeSpan.FromSeconds(30));
Assert.True(fixture.Store.BindHostPresence(new(
first.Listing.HostPresenceHandle,
first.Listing.HostPresenceFingerprint,
EphemeralStateFixture.PublicEndpoint(42_101),
null)).Succeeded);
Assert.True(fixture.Store.CreateJoinAttempt(
fixture.AttemptCommand(listing, "client-quota-3") with { ScopeAttemptLimit = 1 }).Succeeded);
}
[Fact]
public void RemovingAListingReleasesItsConstantTimeOwnerQuota()
{
EphemeralStateFixture fixture = new();
CreateListingCommand first = fixture.ListingCommand(owner: "publisher-quota") with
{
OwnerListingLimit = 1,
};
CreateListingCommand second = fixture.ListingCommand(owner: "publisher-quota") with
{
OwnerListingLimit = 1,
};
StoreResult<StoredListing> created = fixture.Store.CreateListing(first);
Assert.True(created.Succeeded);
Assert.Equal(StoreResultCode.CapacityExceeded, fixture.Store.CreateListing(second).Code);
Assert.True(fixture.Store.DeleteListing(new(
first.Listing.ListingId,
first.Listing.LeaseId,
first.Listing.LeaseFingerprint,
first.Listing.OwnerSubject)).Succeeded);
Assert.True(fixture.Store.CreateListing(second).Succeeded);
}
[Fact]
public void RepeatedMutableDeadlineRefreshesKeepOneScheduledEntryPerKey()
{
EphemeralStateFixture fixture = new();
StoredListing listing = fixture.CreateVisibleListing(out _);
for (int iteration = 0; iteration < 10_000; iteration++)
{
fixture.Clock.Advance(TimeSpan.FromTicks(1));
StoreResult<StoredListing> renewed = fixture.Store.RenewLease(new(
listing.Definition.ListingId,
listing.Definition.LeaseId,
listing.Definition.LeaseFingerprint,
listing.Definition.OwnerSubject,
listing.Version));
Assert.True(renewed.Succeeded, $"Renewal {iteration} failed with {renewed.Code}.");
listing = renewed.Value!;
StoreResult<StoredListing> presence = fixture.Store.BindHostPresence(new(
listing.Definition.HostPresenceHandle,
listing.Definition.HostPresenceFingerprint,
EphemeralStateFixture.PublicEndpoint(42_200),
null));
Assert.True(presence.Succeeded, $"Presence {iteration} failed with {presence.Code}.");
StoreResult<int> revoked = fixture.Store.RevokePrincipal(
"repeated-revocation",
TimeSpan.FromMinutes(1));
Assert.True(revoked.Succeeded, $"Revocation {iteration} failed with {revoked.Code}.");
}
Assert.Equal(4, fixture.Store.ScheduledExpiryEntryCount);
fixture.Clock.Advance(TimeSpan.FromSeconds(19));
Assert.True(fixture.Store.GetListing(listing.Definition.ListingId, true).Succeeded);
Assert.Equal(4, fixture.Store.ScheduledExpiryEntryCount);
}
[Fact]
public void RepeatedPrincipalRevocationCanExtendButCannotShortenProtection()
{
EphemeralStateFixture fixture = new();
const string subject = "protected-publisher";
Assert.True(fixture.Store.RevokePrincipal(subject, TimeSpan.FromMinutes(10)).Succeeded);
fixture.Clock.Advance(TimeSpan.FromSeconds(1));
Assert.True(fixture.Store.RevokePrincipal(subject, TimeSpan.FromSeconds(1)).Succeeded);
fixture.Clock.Advance(TimeSpan.FromSeconds(2));
Assert.Equal(
StoreResultCode.Revoked,
fixture.Store.CreateListing(fixture.ListingCommand(owner: subject)).Code);
fixture.Clock.Advance(TimeSpan.FromSeconds(597));
Assert.True(fixture.Store.CreateListing(fixture.ListingCommand(owner: subject)).Succeeded);
}
[Fact]
@@ -436,7 +522,7 @@ public sealed class InMemoryEphemeralRendezvousStoreTests
StoreResult<int> revoked = fixture.Store.RevokePrincipal(command.Listing.OwnerSubject, TimeSpan.FromMinutes(1));
Assert.True(revoked.Succeeded);
Assert.Equal(2, revoked.Value);
Assert.Equal(4, revoked.Value);
Assert.Equal(StoreResultCode.NotFound, fixture.Store.GetListing(listing.Definition.ListingId, false).Code);
Assert.Equal(StoreResultCode.NotFound, fixture.Store.BindAttemptEndpoint(new(
attempt.MediationHandle,
@@ -666,15 +666,30 @@ public sealed class TestClientProcessIntegrationTests
private static Dictionary<string, string> ServerEnvironment(
string signingKey,
DateTimeOffset now,
string listenAddress)
string listenAddress,
string? readinessAddress)
{
string advertisedAddress = readinessAddress ?? listenAddress;
string allowedHosts = string.Join(
';',
new[] { advertisedAddress, listenAddress, "127.0.0.1", "localhost" }
.Distinct(StringComparer.OrdinalIgnoreCase));
Dictionary<string, string> values = new(StringComparer.Ordinal)
{
["ASPNETCORE_ENVIRONMENT"] = "Production",
["ASPNETCORE_URLS"] = $"http://{listenAddress}:0",
["AllowedHosts"] = allowedHosts,
["Rendezvous__Deployment__PublicHttpBaseUrl"] = $"https://{advertisedAddress}/",
["Rendezvous__Deployment__PublicUdpHost"] = advertisedAddress,
["Rendezvous__Deployment__PublicUdpPort"] = "9050",
["Rendezvous__Deployment__DrainDeadlineSeconds"] = "3",
["Rendezvous__Deployment__MinimumDrainSeconds"] = "1",
["Rendezvous__Deployment__SingleActiveInstance"] = "true",
["Rendezvous__Deployment__AllowPrivatePublicEndpoints"] = "true",
["Rendezvous__Udp__ListenAddress"] = listenAddress,
["Rendezvous__Udp__Port"] = "0",
["Rendezvous__Udp__PollIntervalMilliseconds"] = "1",
["Rendezvous__AbuseProtection__TrustedProxyAddresses__0"] = "127.0.0.1",
["Rendezvous__Provisioning__Issuer"] = "rendezvous-process-test",
["Rendezvous__Provisioning__Audience"] = "rendezvous-process-test-client",
["Rendezvous__Provisioning__ClockSkewSeconds"] = "5",
@@ -1346,7 +1361,11 @@ public sealed class TestClientProcessIntegrationTests
string signingKeyText = Convert.ToBase64String(signingKey);
string publisherCredential = IssuePublisherCredential(signingKey, now);
CryptographicOperations.ZeroMemory(signingKey);
Dictionary<string, string> environment = ServerEnvironment(signingKeyText, now, listenAddress);
Dictionary<string, string> environment = ServerEnvironment(
signingKeyText,
now,
listenAddress,
readinessAddress);
ProcessCapture server = processNamespace is null
? Start(serverAssembly, [], environment)
: StartInNamespace(processNamespace, serverAssembly, [], environment);