From e890e0c365ed9df0dc283130b884bbeacb802fbc Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Mon, 20 Jul 2026 09:30:20 -0500 Subject: [PATCH] measure HTTP and getBlock serving --- README.md | 3 +- docs/grafana-dashboard.md | 11 ++ docs/semantic-parity.md | 2 +- docs/upstream-harness.md | 15 ++ justfile | 28 ++++ src/internal/metrics.zig | 254 ++++++++++++++++++++++++++++++++- src/internal/server.zig | 151 +++++++++++++++++--- tests/http_metrics_contract.py | 193 +++++++++++++++++++++++++ 8 files changed, 635 insertions(+), 22 deletions(-) create mode 100644 tests/http_metrics_contract.py diff --git a/README.md b/README.md index be87c24..eb26458 100644 --- a/README.md +++ b/README.md @@ -131,7 +131,8 @@ product/copy context for public surfaces. `docs/semantic-parity.md` is the top-level deployment gate. `docs/grafana-dashboard.md` records the checksum-pinned upstream source, deterministic Stream adaptation, and its fail-closed metric contract. The dashboard and process receipts are fully -offline: run `just dashboard-test`, `just process-metrics-contract`, then +offline: run `just dashboard-test`, `just process-metrics-contract`, +`just http-metrics-contract`, then `just dashboard-contract http://127.0.0.1:6008/metrics` against a local candidate. The final command intentionally remains red until every upstream lifecycle family has a real producer. diff --git a/docs/grafana-dashboard.md b/docs/grafana-dashboard.md index 87182c5..6042a32 100644 --- a/docs/grafana-dashboard.md +++ b/docs/grafana-dashboard.md @@ -42,6 +42,13 @@ counters. Values come from `getrusage` plus Linux procfs or Darwin `just process-metrics-contract` launches the production binary with only loopback dependencies and validates types, values, invariants, monotonicity, and the absence of `go_*` series across independent HTTP scrapes. +`just http-metrics-contract` seeds a real sealed archive and exercises the +ReleaseSafe binary through ordinary HTTP, Range/conditional getBlock requests, +and a real WebSocket upgrade/close. It proves the same public-mux middleware +boundary and `xrpc/` label as upstream, all four getBlock outcomes, exact +histogram counts, and full-200-only served-byte accounting. Upstream's debug +`/metrics` and `/healthz` endpoints and unmatched public requests remain +outside the HTTP histogram rather than manufacturing extra handler labels. ## Implemented canonical producers @@ -81,6 +88,10 @@ The following dashboard families already have non-placeholder producers: - `jetstream_ingest_dropped_events_total` - `jetstream_segment_seal_duration_seconds` - `jetstream_data_dir_free_bytes` +- `jetstream_http_request_duration_seconds` +- `jetstream_getblock_requests_total` +- `jetstream_getblock_served_bytes_total` +- `jetstream_getblock_duration_seconds` - `jetstream_verifier_failures_total` - `jetstream_backfill_completed_total` - `jetstream_backfill_failed_total` diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 8494917..7969b8c 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -20,7 +20,7 @@ this document or `bootstrap-semantic-parity.md` is open. | Sync 1.1 resync ordering | Upstream serializes per DID, drops stale async repairs, and emits a sync tombstone plus authoritative replacements | **closed at the implementation boundary.** The live coordinator has 32 real fetch workers, a 64-job queue, a five-minute fetch budget plus the 1,000 B/s-for-30-seconds slow-transfer guard, a bounded 16,384-entry 5/minute-per-DID limiter, 2,048 pending commits with drop-oldest overflow, fetch outside the DID lane, authenticated apply with older/equal-contradictory head rejection, ordered pending replay, and trigger-ticket-gated outbox delivery. Commit divergence schedules async repair; `#sync` divergence attempts inline repair and queues one retry only for transient failure, preserving the original relay cursor only on inline success. The writer emits `sync` then bounded 1,024-row `create_resync` batches and stages the fetched/pending chain checkpoints under the archive lock before any matching fsync. Loopback-only tests exercise a real HTTP getRepo mmap, real DID resolution/signature, buffering while the response is held open, slow-transfer cancellation, durable RocksDB promotion/cursor coupling, v2 tail publication, and sealed JSS row order. | | Timestamp import | Canonical dashboard and API include timestamp import lifecycle/rewrite metrics and behavior | **Closed at the implementation boundary.** The format preserves distinct `witnessed_at` and sentinel-zero `indexed_at` columns. Subscriber v1/v2 encoding applies upstream's `indexed_at != 0 ? indexed_at : witnessed_at` display rule while ranges and timestamp cursors remain anchored to immutable witness envelopes. The strict seekable RFC 4180 parser matches upstream's header, row rejection, byte-offset, CID, RFC3339, 64 KiB, and bounded-sampling contracts in a forward-only allocation-free pass. A separate Bloom-filtered RocksDB rule store bulk-loads sorted external SST chunks with CSV-order last-write-wins, reconstructs its resident collection gate on open, and implements specific-CID precedence with all-version fallback. Its central archive-lock hook stamps live, repair, bootstrap, and merge materializations before sequence assignment, and feeds the mutated rows to immediate delivery so hot and cold encodings agree. The sealed-JSS patch primitive preserves block topology/envelopes, clean frames, and the opaque bloom/collection footer while atomically changing only `indexed_at`; it is exercised against an upstream-produced segment. A concrete sealed-segment bloom catalog, bounded DID/FD bucketer, fsynced packed offsets, positioned revalidation, and collision-safe per-segment patch plan implement Phase B/C with specific-CID precedence and idempotent disk-to-disk coverage. A dedicated archive rewrite mutex serializes the real delete compactor with timestamp patching while leaving live append available. The unified metadata DB synchronously persists the current-job pointer, complete job record, and per-segment done set. The real runner activates rules before force-seal, durably crosses the bucketed handoff, checkpoints only after patch fsync/rename/dir-sync, retains partial cancellation progress, and resumes after full archive/rule/metadata reopen. The background manager canonicalizes symlinks, confines regular files, enforces one durable nonterminal job, and adopts it on startup. Bearer-gated import/status XRPC matches the pinned lexicons, fixed-width digest comparison, error names/statuses, steady-state admission, and current/by-id status behavior; a real loopback HTTP run exercised it against the local simulator. Every canonical import metric is produced at its real parse, route, patch, or durable-terminal boundary; cancellation remains a resumable pause and is not counted as terminal failure. | | Status/diagnostics | Repo attempts/error class/final PDS and per-host aggregates survive restart and are operator-queryable | **closed at the implementation boundary.** Durable v5 repo rows distinguish initial Backfill.Rev from latest Rev/UpdatedAt, retain resolved Handle/PDS and reserved record/byte fields, and atomically maintain `handle/`. Merge source cursor commits refresh latest revisions without mutating backfill watermarks. Same-batch `host/` totals, active/status counts, cumulative error classes, and five bounded recent samples survive restart; host moves and active flips are atomic with the source repo row. `/status?tab=hosts` reads those rows without a whole-network scan and matches upstream ordering and error presentation. Account queries preserve resolver-first handle semantics, the legacy `did`/`handle` aliases, and missing-identity hydration, then reconstruct records last-writer-wins from checksum-verified, DID-bloom-pruned sealed JSS plus rotation-safe active/pending and bootstrap-live snapshots. The resulting canonical MST is compared with the PDS commit root through real Sync 1.1 `getLatestCommit`/`getBlocks`; the page renders the same match fields and error state as upstream. `HEAD` performs lookup without verification; successful status responses are `no-store` and carry a generation timestamp. Expensive verification uses upstream's exact 4/source-IP/minute fixed-window limiter, 4,096-entry bound, stale pruning, oldest eviction, and explicit disable option. The offline receipt uses a real HTTP Stream route, real RocksDB row, real JSS archive, signed commit, and loopback PLC/PDS; four separate connections verify successfully, a fifth changing only its ephemeral port is blocked without another upstream request, HEAD remains side-effect-free, and disabling the limiter performs a new authoritative verification. | -| Prometheus/Grafana | Every canonical dashboard query maps to a real, same-semantics producer; runtime-specific panels use honest Zig/process metrics | exact checksum-pinned upstream dashboard is vendored with an offline fail-closed scrape contract; canonical build/live-firehose/archive-append/drop/verifier/backfill/timestamp-import producers are present. The complete canonical subscriber family set is exported from real scan, delivery, rejection, cursor, and disconnect boundaries; `none`, `deflate`, and `zstd` labels now arise only from the actual negotiated delivery path. Live frame inspection distinguishes gaps, bounded-label server errors, forward-compatible unknowns, and malformed CBOR at the raw transport boundary; a Zat reconnect callback measures attempts rather than connection failures. The only dashboard divergence is now a deterministic, checksum-pinned process-runtime row: same-semantics start/CPU/RSS/FD measurements, real VM/thread/page-fault/context-switch signals, no `go_*` aliases, and no host-network counters mislabeled as process I/O. A production-binary loopback receipt proves values and cumulative monotonicity on independent scrapes. The hot readable log now matches upstream's durability pin: rows above the archive watermark may exceed the byte budget but cannot be evicted, publication releases their pinned-byte accounting, and only then may retention advance the floor; all three canonical readable-log gauges come from that exact state. Data-directory free bytes are collected fresh from the archive directory's filesystem at scrape time, and the canonical seal histogram measures only successful flush-through-footer/header-fsync work with upstream's exact 15 slow-latency buckets. **Remaining API/manifest/store/orchestrator/backfill/compaction families open.** | +| Prometheus/Grafana | Every canonical dashboard query maps to a real, same-semantics producer; runtime-specific panels use honest Zig/process metrics | exact checksum-pinned upstream dashboard is vendored with an offline fail-closed scrape contract; canonical build/live-firehose/archive-append/drop/verifier/backfill/timestamp-import producers are present. The complete canonical subscriber family set is exported from real scan, delivery, rejection, cursor, and disconnect boundaries; `none`, `deflate`, and `zstd` labels now arise only from the actual negotiated delivery path. Live frame inspection distinguishes gaps, bounded-label server errors, forward-compatible unknowns, and malformed CBOR at the raw transport boundary; a Zat reconnect callback measures attempts rather than connection failures. The only dashboard divergence is now a deterministic, checksum-pinned process-runtime row: same-semantics start/CPU/RSS/FD measurements, real VM/thread/page-fault/context-switch signals, no `go_*` aliases, and no host-network counters mislabeled as process I/O. A production-binary loopback receipt proves values and cumulative monotonicity on independent scrapes. The hot readable log now matches upstream's durability pin: rows above the archive watermark may exceed the byte budget but cannot be evicted, publication releases their pinned-byte accounting, and only then may retention advance the floor; all three canonical readable-log gauges come from that exact state. Data-directory free bytes are collected fresh from the archive directory's filesystem at scrape time, and the canonical seal histogram measures only successful flush-through-footer/header-fsync work with upstream's exact 15 slow-latency buckets. HTTP duration uses upstream's public-mux boundary and stable labels, including inline WebSocket lifetime and the single `xrpc/` subtree; debug and unmatched routes remain unobserved. getBlock outcomes, duration, and served bytes come from actual response completion, with bytes counted only for complete 200 responses. A ReleaseSafe offline receipt proves 200/206/304/400/404/416 behavior and WebSocket close accounting. **Remaining manifest/store/orchestrator/backfill/compaction families open.** | | Oracle | Event-log equivalence, final-state convergence, crash/power-loss, verifier repair, hostile input, and anti-vacuity receipts | **open — current crash matrix is necessary but substantially smaller than upstream's oracle** | ## Experiment admission diff --git a/docs/upstream-harness.md b/docs/upstream-harness.md index 025bf3b..1a9fda3 100644 --- a/docs/upstream-harness.md +++ b/docs/upstream-harness.md @@ -74,6 +74,21 @@ Zig's checksum-pinned dependencies must already be in its global cache. A missing pin, Go toolchain/module, Python environment, row, sentinel, or cursor transition fails the command rather than weakening the assertion. +## HTTP and getBlock metrics contract + +`just http-metrics-contract` seeds the same real sealed archive and launches a +ReleaseSafe Stream process with loopback-only dependencies. Standard-library +HTTP requests exercise getBlock 200, conditional 304, Range 206, invalid +Range 416, malformed 400, and missing-block 404 responses. The receipt proves +upstream's `ok`, `bad_request`, `not_found`, and `error` partition; duration +cardinality; and that served bytes increase by exactly one completely written +200 frame, never by partial or conditional responses. It also sends a real +masked WebSocket close frame and waits for the inline subscription handler to +finish before checking its HTTP lifetime metric. Handler labels reproduce the +pinned upstream public mux: one `xrpc/` subtree, individual root/status/ +subscribe routes, and no observations for debug health/metrics or unmatched +requests. + ## status diagnostics contract `just status-contract` seeds repository transitions through the real RocksDB diff --git a/justfile b/justfile index 0a8af8c..3d94d65 100644 --- a/justfile +++ b/justfile @@ -45,6 +45,34 @@ process-metrics-contract: done STREAM_BASE_URL="http://127.0.0.1:$PORT" python3 tests/process_metrics_contract.py +# Seed a real sealed archive and prove the upstream HTTP middleware and +# getBlock metric families against actual HTTP and WebSocket connections. +http-metrics-contract: + #!/usr/bin/env bash + set -euo pipefail + DATA=$(mktemp -d) + LOG=$(mktemp) + PORT=${STREAM_HTTP_METRICS_PORT:-6021} + cleanup() { + if [[ -n "${pid:-}" ]]; then kill "$pid" 2>/dev/null || true; wait "$pid" 2>/dev/null || true; fi + rm -rf "$DATA" "$LOG" + } + trap cleanup EXIT + python3 -c 'import socket,sys; s=socket.socket(); s.bind(("127.0.0.1",int(sys.argv[1]))); s.close()' "$PORT" + zig build -Doptimize=ReleaseSafe + zig build write-sample -- --archive "$DATA" + ./zig-out/bin/stream --port="$PORT" --data-dir="$DATA" \ + --upstream=ws://127.0.0.1:17999 --plc=http://127.0.0.1:17999 \ + --relay-http=http://127.0.0.1:17999 --compaction-interval=0 \ + --retry-interval=0 --no-verify >"$LOG" 2>&1 & + pid=$! + for _ in {1..100}; do + if curl -fsS "http://127.0.0.1:$PORT/healthz" >/dev/null 2>&1; then break; fi + if ! kill -0 "$pid" 2>/dev/null; then cat "$LOG" >&2; exit 1; fi + sleep 0.1 + done + STREAM_BASE_URL="http://127.0.0.1:$PORT" python3 tests/http_metrics_contract.py + # Seed a real sealed archive, serve it, then exercise all four archive XRPCs # plus whole-segment, block-plan, bounded, and live-cutover paths through the # exact pinned upstream Go client. Go and uv are forced offline; Zig's pinned diff --git a/src/internal/metrics.zig b/src/internal/metrics.zig index e5641b3..0516e56 100644 --- a/src/internal/metrics.zig +++ b/src/internal/metrics.zig @@ -20,6 +20,19 @@ pub const Counter = std.atomic.Value(u64); pub const Compression = enum { none, deflate, zstd }; pub const OptionsErrorReason = enum { oversize, bad_envelope_json, bad_payload_json, invalid_options }; pub const CursorMode = enum { live, seq, time_us, clamped, disabled, unavailable, too_old, resolve_failed }; +pub const HttpHandler = enum { + root, + status, + subscribe, + subscribe_v2, + xrpc, +}; +pub const HttpMethod = enum { get, head, post, put, patch, delete, options, connect, trace, other }; +pub const HttpCode = enum { c200, c201, c206, c304, c400, c401, c404, c405, c409, c416, c500, c503, other }; +pub const GetBlockResult = enum { ok, not_found, bad_request, error_result }; + +const http_series_count = @typeInfo(HttpHandler).@"enum".fields.len * + @typeInfo(HttpMethod).@"enum".fields.len * @typeInfo(HttpCode).@"enum".fields.len; pub const Stats = struct { start_time_s: i64 = 0, @@ -33,6 +46,14 @@ pub const Stats = struct { segment_seal_duration_count: Counter = .init(0), segment_seal_duration_sum_us: Counter = .init(0), segment_seal_duration_buckets: [15]Counter = @splat(.init(0)), + http_request_duration_count: [http_series_count]Counter = @splat(.init(0)), + http_request_duration_sum_us: [http_series_count]Counter = @splat(.init(0)), + http_request_duration_buckets: [http_series_count][11]Counter = @splat(@splat(.init(0))), + getblock_requests: std.enums.EnumArray(GetBlockResult, Counter) = .initFill(.init(0)), + getblock_served_bytes_total: Counter = .init(0), + getblock_duration_count: Counter = .init(0), + getblock_duration_sum_us: Counter = .init(0), + getblock_duration_buckets: [12]Counter = @splat(.init(0)), upstream_seq: std.atomic.Value(i64) = .init(0), drops: std.enums.EnumArray(convert.DropReason, Counter) = .initFill(.init(0)), verify_valid: Counter = .init(0), @@ -133,8 +154,91 @@ pub const Stats = struct { _ = self.segment_seal_duration_buckets[i].fetchAdd(1, .monotonic); } } + + pub fn observeHttp(self: *Stats, handler: HttpHandler, method: HttpMethod, code: u16, duration_us: u64) void { + const index = httpIndex(handler, method, httpCode(code)); + _ = self.http_request_duration_count[index].fetchAdd(1, .monotonic); + _ = self.http_request_duration_sum_us[index].fetchAdd(duration_us, .monotonic); + const bounds_us = [_]u64{ 5_000, 10_000, 25_000, 50_000, 100_000, 250_000, 500_000, 1_000_000, 2_500_000, 5_000_000, 10_000_000 }; + inline for (bounds_us, 0..) |upper, i| { + if (duration_us <= upper) + _ = self.http_request_duration_buckets[index][i].fetchAdd(1, .monotonic); + } + } + + pub fn observeGetBlock(self: *Stats, code: u16, served_bytes: u64, duration_us: u64) void { + const result: GetBlockResult = switch (code) { + 200, 206, 304 => .ok, + 400 => .bad_request, + 404 => .not_found, + else => .error_result, + }; + _ = self.getblock_requests.getPtr(result).fetchAdd(1, .monotonic); + if (code == 200 and served_bytes > 0) + _ = self.getblock_served_bytes_total.fetchAdd(served_bytes, .monotonic); + _ = self.getblock_duration_count.fetchAdd(1, .monotonic); + _ = self.getblock_duration_sum_us.fetchAdd(duration_us, .monotonic); + inline for (0..12) |i| { + const upper_us = @as(u64, 100) << @intCast(i); + if (duration_us <= upper_us) + _ = self.getblock_duration_buckets[i].fetchAdd(1, .monotonic); + } + } }; +fn httpIndex(handler: HttpHandler, method: HttpMethod, code: HttpCode) usize { + const method_count = @typeInfo(HttpMethod).@"enum".fields.len; + const code_count = @typeInfo(HttpCode).@"enum".fields.len; + return (@intFromEnum(handler) * method_count + @intFromEnum(method)) * code_count + @intFromEnum(code); +} + +fn httpCode(code: u16) HttpCode { + return switch (code) { + 200 => .c200, + 201 => .c201, + 206 => .c206, + 304 => .c304, + 400 => .c400, + 401 => .c401, + 404 => .c404, + 405 => .c405, + 409 => .c409, + 416 => .c416, + 500 => .c500, + 503 => .c503, + else => .other, + }; +} + +/// Mirror upstream's public mux, including its middleware boundaries. The +/// debug /metrics and /healthz routes and unmatched public requests are not +/// wrapped. All XRPC methods share the stable `xrpc/` handler label because +/// upstream registers one `/xrpc/` subtree handler. +pub fn httpHandler(path: []const u8, method: []const u8) ?HttpHandler { + if (std.mem.startsWith(u8, path, "/xrpc/")) return .xrpc; + if (std.mem.eql(u8, method, "GET")) { + if (std.mem.eql(u8, path, "/")) return .root; + if (std.mem.eql(u8, path, "/status")) return .status; + if (std.mem.eql(u8, path, "/subscribe")) return .subscribe; + if (std.mem.eql(u8, path, "/subscribe-v2")) return .subscribe_v2; + } + if (std.mem.eql(u8, method, "HEAD") and std.mem.eql(u8, path, "/status")) return .status; + return null; +} + +pub fn httpMethod(method: []const u8) HttpMethod { + if (std.mem.eql(u8, method, "GET")) return .get; + if (std.mem.eql(u8, method, "HEAD")) return .head; + if (std.mem.eql(u8, method, "POST")) return .post; + if (std.mem.eql(u8, method, "PUT")) return .put; + if (std.mem.eql(u8, method, "PATCH")) return .patch; + if (std.mem.eql(u8, method, "DELETE")) return .delete; + if (std.mem.eql(u8, method, "OPTIONS")) return .options; + if (std.mem.eql(u8, method, "CONNECT")) return .connect; + if (std.mem.eql(u8, method, "TRACE")) return .trace; + return .other; +} + pub const TailGauges = struct { entries: usize, bytes: usize, @@ -312,6 +416,63 @@ pub fn format( } w.print("jetstream_segment_seal_duration_seconds_bucket{{le=\"+Inf\"}} {d}\n", .{seal_count}) catch {}; + w.print("# TYPE jetstream_http_request_duration_seconds histogram\n", .{}) catch {}; + for (comptime std.enums.values(HttpHandler)) |handler| { + for (comptime std.enums.values(HttpMethod)) |method| { + for (comptime std.enums.values(HttpCode)) |code| { + const index = httpIndex(handler, method, code); + const count = stats.http_request_duration_count[index].load(.monotonic); + if (count != 0) { + const labels = .{ httpHandlerLabel(handler), httpMethodLabel(method), httpCodeLabel(code) }; + const sum_us = stats.http_request_duration_sum_us[index].load(.monotonic); + w.print("jetstream_http_request_duration_seconds_sum{{handler=\"{s}\",method=\"{s}\",code=\"{s}\"}} {d}.{d:0>6}\n", .{ + labels[0], labels[1], labels[2], sum_us / std.time.us_per_s, sum_us % std.time.us_per_s, + }) catch {}; + w.print("jetstream_http_request_duration_seconds_count{{handler=\"{s}\",method=\"{s}\",code=\"{s}\"}} {d}\n", .{ + labels[0], labels[1], labels[2], count, + }) catch {}; + const bounds = [_][]const u8{ "0.005", "0.01", "0.025", "0.05", "0.1", "0.25", "0.5", "1", "2.5", "5", "10" }; + for (bounds, 0..) |upper, i| w.print("jetstream_http_request_duration_seconds_bucket{{handler=\"{s}\",method=\"{s}\",code=\"{s}\",le=\"{s}\"}} {d}\n", .{ + labels[0], labels[1], labels[2], upper, stats.http_request_duration_buckets[index][i].load(.monotonic), + }) catch {}; + w.print("jetstream_http_request_duration_seconds_bucket{{handler=\"{s}\",method=\"{s}\",code=\"{s}\",le=\"+Inf\"}} {d}\n", .{ + labels[0], labels[1], labels[2], count, + }) catch {}; + } + } + } + } + + w.print("# TYPE jetstream_getblock_requests_total counter\n", .{}) catch {}; + inline for (comptime std.enums.values(GetBlockResult)) |result| { + const count = stats.getblock_requests.getPtrConst(result).load(.monotonic); + if (count != 0) w.print( + "jetstream_getblock_requests_total{{result=\"{s}\"}} {d}\n", + .{ getBlockResultLabel(result), count }, + ) catch {}; + } + w.print( + \\# TYPE jetstream_getblock_served_bytes_total counter + \\jetstream_getblock_served_bytes_total {d} + \\# TYPE jetstream_getblock_duration_seconds histogram + \\jetstream_getblock_duration_seconds_sum {d}.{d:0>6} + \\jetstream_getblock_duration_seconds_count {d} + \\ + , .{ + stats.getblock_served_bytes_total.load(.monotonic), + stats.getblock_duration_sum_us.load(.monotonic) / std.time.us_per_s, + stats.getblock_duration_sum_us.load(.monotonic) % std.time.us_per_s, + stats.getblock_duration_count.load(.monotonic), + }) catch {}; + const getblock_bounds = [_][]const u8{ "0.0001", "0.0002", "0.0004", "0.0008", "0.0016", "0.0032", "0.0064", "0.0128", "0.0256", "0.0512", "0.1024", "0.2048" }; + inline for (getblock_bounds, 0..) |upper, i| w.print( + "jetstream_getblock_duration_seconds_bucket{{le=\"{s}\"}} {d}\n", + .{ upper, stats.getblock_duration_buckets[i].load(.monotonic) }, + ) catch {}; + w.print("jetstream_getblock_duration_seconds_bucket{{le=\"+Inf\"}} {d}\n", .{ + stats.getblock_duration_count.load(.monotonic), + }) catch {}; + // Canonical upstream names. These are aliases only where Stream owns the // same lifecycle boundary; absent families stay absent so the dashboard // contract cannot be satisfied by placeholder zeroes. @@ -614,8 +775,90 @@ pub fn format( return w.buffered(); } +fn httpHandlerLabel(value: HttpHandler) []const u8 { + return switch (value) { + .root => "root", + .status => "status", + .subscribe => "subscribe", + .subscribe_v2 => "subscribe-v2", + .xrpc => "xrpc/", + }; +} + +fn httpMethodLabel(value: HttpMethod) []const u8 { + return switch (value) { + .get => "GET", + .head => "HEAD", + .post => "POST", + .put => "PUT", + .patch => "PATCH", + .delete => "DELETE", + .options => "OPTIONS", + .connect => "CONNECT", + .trace => "TRACE", + .other => "OTHER", + }; +} + +fn httpCodeLabel(value: HttpCode) []const u8 { + return switch (value) { + .c200 => "200", + .c201 => "201", + .c206 => "206", + .c304 => "304", + .c400 => "400", + .c401 => "401", + .c404 => "404", + .c405 => "405", + .c409 => "409", + .c416 => "416", + .c500 => "500", + .c503 => "503", + .other => "other", + }; +} + +fn getBlockResultLabel(value: GetBlockResult) []const u8 { + return switch (value) { + .ok => "ok", + .not_found => "not_found", + .bad_request => "bad_request", + .error_result => "error", + }; +} + // === tests === +test "http handler labels match upstream public middleware boundaries" { + try std.testing.expectEqual(HttpHandler.root, httpHandler("/", "GET").?); + try std.testing.expectEqual(HttpHandler.status, httpHandler("/status", "GET").?); + try std.testing.expectEqual(HttpHandler.status, httpHandler("/status", "HEAD").?); + try std.testing.expectEqual(HttpHandler.subscribe, httpHandler("/subscribe", "GET").?); + try std.testing.expectEqual(HttpHandler.subscribe_v2, httpHandler("/subscribe-v2", "GET").?); + try std.testing.expectEqual(HttpHandler.xrpc, httpHandler("/xrpc/network.bsky.jetstream.getBlock", "GET").?); + try std.testing.expectEqual(HttpHandler.xrpc, httpHandler("/xrpc/not-implemented", "POST").?); + try std.testing.expectEqual(null, httpHandler("/metrics", "GET")); + try std.testing.expectEqual(null, httpHandler("/healthz", "GET")); + try std.testing.expectEqual(null, httpHandler("/", "POST")); + try std.testing.expectEqual(null, httpHandler("/not-found", "GET")); +} + +test "getBlock metrics preserve upstream result and served-byte semantics" { + var stats: Stats = .{}; + stats.observeGetBlock(200, 100, 99); + stats.observeGetBlock(206, 50, 100); + stats.observeGetBlock(304, 0, 101); + stats.observeGetBlock(400, 0, 102); + stats.observeGetBlock(404, 0, 103); + stats.observeGetBlock(416, 0, 104); + try std.testing.expectEqual(3, stats.getblock_requests.getPtrConst(.ok).load(.monotonic)); + try std.testing.expectEqual(1, stats.getblock_requests.getPtrConst(.bad_request).load(.monotonic)); + try std.testing.expectEqual(1, stats.getblock_requests.getPtrConst(.not_found).load(.monotonic)); + try std.testing.expectEqual(1, stats.getblock_requests.getPtrConst(.error_result).load(.monotonic)); + try std.testing.expectEqual(100, stats.getblock_served_bytes_total.load(.monotonic)); + try std.testing.expectEqual(6, stats.getblock_duration_count.load(.monotonic)); +} + test "format renders counters" { var stats: Stats = .{}; _ = stats.events_total.fetchAdd(7, .monotonic); @@ -662,8 +905,11 @@ test "format renders counters" { _ = stats.subscribe_clean_disconnects_total.fetchAdd(15, .monotonic); _ = stats.subscribe_options_updates_total.fetchAdd(16, .monotonic); stats.observeSegmentSeal(20_000); + stats.observeHttp(.xrpc, .get, 200, 12_000); + stats.observeGetBlock(200, 321, 800); + stats.observeGetBlock(404, 0, 200); - var buf: [16 * 1024]u8 = undefined; + var buf: [64 * 1024]u8 = undefined; const out = format(&buf, &stats, 2, .{ .entries = 5, .bytes = 100, @@ -706,6 +952,12 @@ test "format renders counters" { try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_segment_seal_duration_seconds_sum 0.020000") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_segment_seal_duration_seconds_bucket{le=\"0.02\"} 1") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_segment_seal_duration_seconds_count 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_http_request_duration_seconds_count{handler=\"xrpc/\",method=\"GET\",code=\"200\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_http_request_duration_seconds_bucket{handler=\"xrpc/\",method=\"GET\",code=\"200\",le=\"0.025\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_getblock_requests_total{result=\"ok\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_getblock_requests_total{result=\"not_found\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_getblock_served_bytes_total 321") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_getblock_duration_seconds_count 2") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_repos_total{result=\"complete\"} 11") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_workers_active 4") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_queue_depth 23") != null); diff --git a/src/internal/server.zig b/src/internal/server.zig index de74e3b..aaac880 100644 --- a/src/internal/server.zig +++ b/src/internal/server.zig @@ -315,6 +315,34 @@ const Subscriber = struct { } }; +const HttpObservation = struct { + hub: *Hub, + started: Io.Timestamp, + handler: ?metrics.HttpHandler, + method: metrics.HttpMethod, + get_block: bool = false, + code: u16 = 500, + served_bytes: u64 = 0, + + fn init(hub: *Hub, path: []const u8, method: []const u8) HttpObservation { + return .{ + .hub = hub, + .started = Io.Timestamp.now(hub.io, .awake), + .handler = metrics.httpHandler(path, method), + .method = metrics.httpMethod(method), + }; + } + + fn finish(self: *HttpObservation) void { + const elapsed = self.started.durationTo(Io.Timestamp.now(self.hub.io, .awake)).toMicroseconds(); + const duration_us: u64 = @intCast(@max(0, elapsed)); + if (self.handler) |handler| + self.hub.stats.observeHttp(handler, self.method, self.code, duration_us); + if (self.get_block) + self.hub.stats.observeGetBlock(self.code, self.served_bytes, duration_us); + } +}; + pub const Handler = struct { hub: *Hub, conn: *websocket.Conn, @@ -328,20 +356,33 @@ pub const Handler = struct { compress: bool, v2: bool, clean_disconnect_counted: bool = false, + http_started: Io.Timestamp, + http_handler: metrics.HttpHandler, + http_observed: bool = false, pub fn init(handshake: *const websocket.Handshake, conn: *websocket.Conn, hub: *Hub) !Handler { const url = handshake.url; const qi = std.mem.indexOfScalar(u8, url, '?'); const query = if (qi) |i| url[i + 1 ..] else ""; + const path = if (qi) |i| url[0..i] else url; + const http_started = Io.Timestamp.now(hub.io, .awake); + const http_handler = metrics.httpHandler(path, "GET"); + const finishHttp = struct { + fn call(h: *Hub, started: Io.Timestamp, handler: ?metrics.HttpHandler, code: u16) void { + const label = handler orelse return; + const elapsed = started.durationTo(Io.Timestamp.now(h.io, .awake)).toMicroseconds(); + h.stats.observeHttp(label, .get, code, @intCast(@max(0, elapsed))); + } + }.call; if (!hub.serving.load(.acquire)) { if (queryParam(query, "cursor")) |raw| { if (raw.len != 0) _ = hub.stats.subscribe_cursor_requests.getPtr(.unavailable).fetchAdd(1, .monotonic); } respond(conn, "503 Service Unavailable", "text/plain", "bootstrap in progress\n"); + finishHttp(hub, http_started, http_handler, 503); return error.Close; // response already written; Close suppresses the 400 } - const path = if (qi) |i| url[0..i] else url; const v2 = std.mem.eql(u8, path, "/subscribe-v2"); if (!v2 and !std.mem.eql(u8, path, "/subscribe")) return error.NotFound; @@ -354,21 +395,25 @@ pub const Handler = struct { if (v2) { if (legacy_compress or header_zstd) { respond(conn, "400 Bad Request", "text/plain", "compress=true / Socket-Encoding: zstd is the /subscribe (v1) opt-in; /subscribe-v2 uses zstdDictionary= with the dictionary from getZstdDictionary\n"); + finishHttp(hub, http_started, http_handler, 400); return error.Close; } if (queryParam(query, "zstdDictionary")) |raw| if (raw.len != 0) { const id = std.fmt.parseInt(u32, raw, 10) catch { respond(conn, "400 Bad Request", "text/plain", "zstdDictionary must be a positive integer zstd dictionary ID\n"); + finishHttp(hub, http_started, http_handler, 400); return error.Close; }; if (id == 0) { respond(conn, "400 Bad Request", "text/plain", "zstdDictionary must be a positive integer zstd dictionary ID\n"); + finishHttp(hub, http_started, http_handler, 400); return error.Close; } if (id != zstd.v2_dictionary_id) { var buf: [224]u8 = undefined; const body = std.fmt.bufPrint(&buf, "unknown zstd dictionary id {d}; current dictionary id is {d} (fetch it via getZstdDictionary and reconnect)\n", .{ id, zstd.v2_dictionary_id }) catch unreachable; respond(conn, "400 Bad Request", "text/plain", body); + finishHttp(hub, http_started, http_handler, 400); return error.Close; } compress = true; @@ -378,6 +423,7 @@ pub const Handler = struct { if (compress) if (handshake.headers.get("sec-websocket-extensions")) |extensions| { if (offersPerMessageDeflate(extensions)) { respond(conn, "400 Bad Request", "text/plain", "choose one compression scheme: custom zstd (compress=true / Socket-Encoding: zstd) or RFC 7692 permessage-deflate, not both\n"); + finishHttp(hub, http_started, http_handler, 400); return error.Close; } }; @@ -388,6 +434,7 @@ pub const Handler = struct { var filter = filter_mod.Filter.parseQuery(hub.allocator, query) catch |err| { log.debug("subscribe: invalid options: {s}", .{@errorName(err)}); + finishHttp(hub, http_started, http_handler, if (err == error.OutOfMemory) 500 else 400); return error.InvalidRequest; }; errdefer filter.deinit(); @@ -401,10 +448,12 @@ pub const Handler = struct { if (err == error.CursorResolveFailed) { _ = hub.stats.subscribe_cursor_requests.getPtr(.resolve_failed).fetchAdd(1, .monotonic); respond(conn, "503 Service Unavailable", "text/plain", "service not ready: cursor resolution failed\n"); + finishHttp(hub, http_started, http_handler, 503); return error.Close; } if (err == error.CursorTooOld) _ = hub.stats.subscribe_cursor_requests.getPtr(.too_old).fetchAdd(1, .monotonic); + finishHttp(hub, http_started, http_handler, 400); return error.InvalidRequest; }; const resolve_finished = Io.Timestamp.now(hub.io, .real).toMicroseconds(); @@ -417,7 +466,10 @@ pub const Handler = struct { else false; - try conn.writeTimeout(frame_write_timeout_ms); + conn.writeTimeout(frame_write_timeout_ms) catch |err| { + finishHttp(hub, http_started, http_handler, 500); + return err; + }; return .{ .hub = hub, @@ -428,6 +480,8 @@ pub const Handler = struct { .awaiting_hello = awaiting_hello, .compress = compress, .v2 = v2, + .http_started = http_started, + .http_handler = http_handler orelse unreachable, }; } @@ -612,6 +666,11 @@ pub const Handler = struct { self.subscriber = null; } if (self.pending_filter) |*f| f.deinit(); + if (!self.http_observed) { + self.http_observed = true; + const elapsed = self.http_started.durationTo(Io.Timestamp.now(self.hub.io, .awake)).toMicroseconds(); + self.hub.stats.observeHttp(self.http_handler, .get, 200, @intCast(@max(0, elapsed))); + } } fn signalStop(self: *Handler, sub: *Subscriber) void { @@ -640,16 +699,21 @@ pub const Handler = struct { ) void { const path = if (std.mem.indexOfScalar(u8, url, '?')) |i| url[0..i] else url; const query = if (std.mem.indexOfScalar(u8, url, '?')) |i| url[i + 1 ..] else ""; + var observation = HttpObservation.init(hub, path, method); + defer observation.finish(); // the homepage stays up during bootstrap too (it says so) const import_ops = std.mem.eql(u8, path, "/xrpc/network.bsky.jetstream.importTimestamps") or std.mem.eql(u8, path, "/xrpc/network.bsky.jetstream.getImportStatus"); const dictionary_path = std.mem.eql(u8, path, "/xrpc/network.bsky.jetstream.getZstdDictionary"); const ops_path = std.mem.eql(u8, path, "/healthz") or std.mem.eql(u8, path, "/metrics") or std.mem.eql(u8, path, "/status") or std.mem.eql(u8, path, "/") or import_ops or dictionary_path; if (!hub.serving.load(.acquire) and !ops_path) { + observation.code = 503; return respond(conn, "503 Service Unavailable", "text/plain", "bootstrap in progress\n"); } if (hub.xrpc) |api| { - var responder: XrpcResponder = .{ .conn = conn }; + var responder: XrpcResponder = .{ .conn = conn, .observation = &observation }; + observation.get_block = std.mem.eql(u8, method, "GET") and + std.mem.eql(u8, path, "/xrpc/network.bsky.jetstream.getBlock"); if (api.handle( &responder, method, @@ -663,9 +727,16 @@ pub const Handler = struct { )) return; } const status_head = std.mem.eql(u8, method, "HEAD") and std.mem.eql(u8, path, "/status"); - if (!std.mem.eql(u8, method, "GET") and !status_head) return respond(conn, "405 Method Not Allowed", "text/plain", "method not allowed\n"); - if (std.mem.eql(u8, path, "/healthz")) return respond(conn, "200 OK", "text/plain", "ok\n"); + if (!std.mem.eql(u8, method, "GET") and !status_head) { + observation.code = 405; + return respond(conn, "405 Method Not Allowed", "text/plain", "method not allowed\n"); + } + if (std.mem.eql(u8, path, "/healthz")) { + observation.code = 200; + return respond(conn, "200 OK", "text/plain", "ok\n"); + } if (std.mem.eql(u8, path, "/")) { + observation.code = 200; const now_s: i64 = @divTrunc(Io.Timestamp.now(hub.io, .real).toMicroseconds(), std.time.us_per_s); const args: homepage.Args = .{ .seq = hub.stats.upstream_seq.load(.monotonic), @@ -682,11 +753,16 @@ pub const Handler = struct { } if (std.mem.eql(u8, path, "/status")) { const wants_html = if (headers.get("accept")) |a| std.mem.indexOf(u8, a, "text/html") != null else false; - const account_query = status_page.accountFromQueryAlloc(hub.allocator, query) catch + const account_query = status_page.accountFromQueryAlloc(hub.allocator, query) catch { + observation.code = 400; return respond(conn, "400 Bad Request", "text/plain", "invalid account query\n"); + }; if (account_query) |raw_account| { defer hub.allocator.free(raw_account); - const store = hub.repo_status orelse return respond(conn, "500 Internal Server Error", "text/plain", "account diagnostics unavailable\n"); + const store = hub.repo_status orelse { + observation.code = 500; + return respond(conn, "500 Internal Server Error", "text/plain", "account diagnostics unavailable\n"); + }; var arena = std.heap.ArenaAllocator.init(hub.allocator); defer arena.deinit(); const account = status_page.collectAccount( @@ -696,7 +772,10 @@ pub const Handler = struct { hub.plc_url, raw_account, hub.status_identity_resolution, - ) catch return respond(conn, "500 Internal Server Error", "text/plain", "failed to read account diagnostics\n"); + ) catch { + observation.code = 500; + return respond(conn, "500 Internal Server Error", "text/plain", "failed to read account diagnostics\n"); + }; const now_us = Io.Timestamp.now(hub.io, .real).toMicroseconds(); var verification: ?repo_export.VerifyReport = null; var verification_error: []const u8 = ""; @@ -728,20 +807,31 @@ pub const Handler = struct { break :verify; }; }; - const page = status_page.renderAccountAlloc(hub.allocator, account, now_us, wants_html, verification, verification_error) catch + const page = status_page.renderAccountAlloc(hub.allocator, account, now_us, wants_html, verification, verification_error) catch { + observation.code = 500; return respond(conn, "500 Internal Server Error", "text/plain", "failed to render account diagnostics\n"); + }; defer hub.allocator.free(page); + observation.code = 200; return respondStatus(conn, if (wants_html) "text/html; charset=utf-8" else "text/plain; charset=utf-8", page, status_head, now_us); } if (status_page.hostSortFromQuery(query)) |sort| { - const store = hub.repo_status orelse return respond(conn, "500 Internal Server Error", "text/plain", "host diagnostics unavailable\n"); - const statuses = store.listHostStatuses(hub.allocator) catch + const store = hub.repo_status orelse { + observation.code = 500; + return respond(conn, "500 Internal Server Error", "text/plain", "host diagnostics unavailable\n"); + }; + const statuses = store.listHostStatuses(hub.allocator) catch { + observation.code = 500; return respond(conn, "500 Internal Server Error", "text/plain", "failed to read host diagnostics\n"); + }; defer store.freeHostStatuses(hub.allocator, statuses); const now_us = Io.Timestamp.now(hub.io, .real).toMicroseconds(); - const page = status_page.renderHostsAlloc(hub.allocator, statuses, sort, now_us, wants_html) catch + const page = status_page.renderHostsAlloc(hub.allocator, statuses, sort, now_us, wants_html) catch { + observation.code = 500; return respond(conn, "500 Internal Server Error", "text/plain", "failed to render host diagnostics\n"); + }; defer hub.allocator.free(page); + observation.code = 200; return respondStatus(conn, if (wants_html) "text/html; charset=utf-8" else "text/plain; charset=utf-8", page, status_head, now_us); } const now_s: i64 = @divTrunc(Io.Timestamp.now(hub.io, .real).toMicroseconds(), std.time.us_per_s); @@ -759,14 +849,16 @@ pub const Handler = struct { .verify_invalid = hub.stats.verify_invalid.load(.monotonic), }; var page_buf: [32 * 1024]u8 = undefined; + observation.code = 200; if (wants_html) return respondStatus(conn, "text/html; charset=utf-8", homepage.renderStatusHtml(&page_buf, args), status_head, now_s * std.time.us_per_s); return respondStatus(conn, "text/plain; charset=utf-8", homepage.renderStatus(&page_buf, args), status_head, now_s * std.time.us_per_s); } if (std.mem.eql(u8, path, "/metrics")) { + observation.code = 200; const g = hub.tail.gauges(); const now_s: i64 = @divTrunc(Io.Timestamp.now(hub.io, .real).toMicroseconds(), std.time.us_per_s); - var buf: [24 * 1024]u8 = undefined; + var buf: [256 * 1024]u8 = undefined; const out = metrics.format(&buf, hub.stats, hub.active.load(.monotonic), .{ .entries = g.entries, .bytes = g.bytes, @@ -784,10 +876,9 @@ pub const Handler = struct { "stream_pipeline_ticket{{stage=\"claimed\"}} {d}\n" ++ "stream_pipeline_ticket{{stage=\"emitted\"}} {d}\n", .{ v[0], v[1], v[2] }) catch ""; } - var joined: [17 * 1024]u8 = undefined; - const full = std.fmt.bufPrint(&joined, "{s}{s}", .{ out, extra }) catch out; - return respond(conn, "200 OK", "text/plain; version=0.0.4", full); + return respondParts(conn, "200 OK", "text/plain; version=0.0.4", out, extra); } + observation.code = 404; respond(conn, "404 Not Found", "text/plain", "not found\n"); } }; @@ -795,21 +886,24 @@ pub const Handler = struct { /// adapter the xrpc api writes through (json / raw+etag / xrpc error) const XrpcResponder = struct { conn: *websocket.Conn, + observation: *HttpObservation, pub fn json(self: *XrpcResponder, status: u16, body: []const u8) void { + self.observation.code = status; respond(self.conn, statusLine(status), "application/json", body); } pub fn raw(self: *XrpcResponder, status: u16, body: []const u8, etag_hex: []const u8) void { - _ = status; var buf: [320]u8 = undefined; const header = std.fmt.bufPrint( &buf, - "HTTP/1.1 200 OK\r\nContent-Type: application/octet-stream\r\nContent-Length: {d}\r\nAccept-Ranges: bytes\r\nETag: \"{s}\"\r\nCache-Control: public, no-cache\r\nConnection: close\r\nServer: stream\r\n\r\n", - .{ body.len, etag_hex }, + "HTTP/1.1 {s}\r\nContent-Type: application/octet-stream\r\nContent-Length: {d}\r\nAccept-Ranges: bytes\r\nETag: \"{s}\"\r\nCache-Control: public, no-cache\r\nConnection: close\r\nServer: stream\r\n\r\n", + .{ statusLine(status), body.len, etag_hex }, ) catch return; + self.observation.code = status; self.conn.writeFramed(header) catch return; if (body.len > 0) self.conn.writeFramed(body) catch return; + if (status == 200) self.observation.served_bytes = body.len; } pub fn dictionary(self: *XrpcResponder, body: []const u8, id: u32, etag: []const u8, not_modified: bool) void { @@ -826,6 +920,7 @@ const XrpcResponder = struct { "HTTP/1.1 200 OK\r\nContent-Type: application/octet-stream\r\nContent-Length: {d}\r\nETag: {s}\r\nCache-Control: public, max-age=31536000, immutable\r\nX-Zstd-Dictionary-Id: {d}\r\nConnection: close\r\nServer: stream\r\n\r\n", .{ body.len, etag, id }, ) catch return; + self.observation.code = if (not_modified) 304 else 200; self.conn.writeFramed(header) catch return; if (!not_modified) self.conn.writeFramed(body) catch return; } @@ -837,6 +932,7 @@ const XrpcResponder = struct { "HTTP/1.1 206 Partial Content\r\nContent-Type: application/octet-stream\r\nContent-Length: {d}\r\nContent-Range: bytes {d}-{d}/{d}\r\nAccept-Ranges: bytes\r\nETag: \"{s}\"\r\nCache-Control: public, no-cache\r\nConnection: close\r\nServer: stream\r\n\r\n", .{ body.len, start, end, total, etag_hex }, ) catch return; + self.observation.code = 206; self.conn.writeFramed(header) catch return; if (body.len > 0) self.conn.writeFramed(body) catch return; } @@ -862,6 +958,7 @@ const XrpcResponder = struct { "HTTP/1.1 200 OK\r\nContent-Type: application/octet-stream\r\nContent-Length: {d}\r\nAccept-Ranges: bytes\r\nETag: \"{s}\"\r\nCache-Control: public, no-cache\r\nConnection: close\r\nServer: stream\r\n\r\n", .{ len, etag_hex }, ) catch return; + self.observation.code = if (partial_response) 206 else 200; self.conn.writeFramed(header) catch return; var offset = start; @@ -875,6 +972,7 @@ const XrpcResponder = struct { offset += n; remaining -= n; } + if (!partial_response) self.observation.served_bytes = len; } pub fn rangeNotSatisfiable(self: *XrpcResponder, total: u64) void { @@ -884,6 +982,7 @@ const XrpcResponder = struct { "HTTP/1.1 416 Range Not Satisfiable\r\nContent-Range: bytes */{d}\r\nContent-Length: 0\r\nConnection: close\r\nServer: stream\r\n\r\n", .{total}, ) catch return; + self.observation.code = 416; self.conn.writeFramed(header) catch return; } @@ -894,6 +993,7 @@ const XrpcResponder = struct { "HTTP/1.1 304 Not Modified\r\nETag: \"{s}\"\r\nCache-Control: public, no-cache\r\nConnection: close\r\nServer: stream\r\n\r\n", .{etag_hex}, ) catch return; + self.observation.code = 304; self.conn.writeFramed(header) catch return; } @@ -909,6 +1009,7 @@ const XrpcResponder = struct { 503 => "503 Service Unavailable", else => "400 Bad Request", }; + self.observation.code = status; respond(self.conn, status_line, "application/json", body); } }; @@ -939,6 +1040,18 @@ fn respond(conn: *websocket.Conn, status: []const u8, content_type: []const u8, if (resp_body.len > 0) conn.writeFramed(resp_body) catch return; } +fn respondParts(conn: *websocket.Conn, status: []const u8, content_type: []const u8, first: []const u8, second: []const u8) void { + var buf: [256]u8 = undefined; + const header = std.fmt.bufPrint( + &buf, + "HTTP/1.1 {s}\r\nContent-Type: {s}\r\nContent-Length: {d}\r\nConnection: close\r\nServer: stream\r\n\r\n", + .{ status, content_type, first.len + second.len }, + ) catch return; + conn.writeFramed(header) catch return; + if (first.len > 0) conn.writeFramed(first) catch return; + if (second.len > 0) conn.writeFramed(second) catch return; +} + fn respondStatus(conn: *websocket.Conn, content_type: []const u8, body: []const u8, head: bool, generated_us: i64) void { var time_buf: [40]u8 = undefined; const generated_at = formatRfc3339(&time_buf, generated_us); diff --git a/tests/http_metrics_contract.py b/tests/http_metrics_contract.py new file mode 100644 index 0000000..3ea4b56 --- /dev/null +++ b/tests/http_metrics_contract.py @@ -0,0 +1,193 @@ +#!/usr/bin/env python3 +"""Offline production-binary contract for upstream HTTP/getBlock metrics.""" + +from __future__ import annotations + +import base64 +import hashlib +import json +import os +import re +import socket +import struct +import time +import urllib.error +import urllib.parse +import urllib.request + + +BASE = os.environ.get("STREAM_BASE_URL", "http://127.0.0.1:6021").rstrip("/") +XRPC = BASE + "/xrpc/network.bsky.jetstream." +SAMPLE_RE = re.compile(r"^([a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{([^}]*)\})?\s+([^\s]+)$") +LABEL_RE = re.compile(r'(\w+)="([^"]*)"') + + +def request(path: str, *, method: str = "GET", headers: dict[str, str] | None = None): + req = urllib.request.Request(BASE + path, method=method, headers=headers or {}) + try: + with urllib.request.urlopen(req, timeout=5) as response: + return response.status, response.read(), {key.lower(): value for key, value in response.headers.items()} + except urllib.error.HTTPError as error: + return error.code, error.read(), {key.lower(): value for key, value in error.headers.items()} + + +def scrape() -> dict[tuple[str, tuple[tuple[str, str], ...]], float]: + status, raw, _ = request("/metrics") + assert status == 200 + samples: dict[tuple[str, tuple[tuple[str, str], ...]], float] = {} + for line in raw.decode().splitlines(): + if line.startswith("#"): + continue + match = SAMPLE_RE.match(line) + if not match: + continue + name, labels_raw, value = match.groups() + labels = tuple(sorted(LABEL_RE.findall(labels_raw or ""))) + samples[(name, labels)] = float(value) + return samples + + +def value(samples, name: str, **labels: str) -> float: + return samples.get((name, tuple(sorted(labels.items()))), 0.0) + + +def delta(before, after, name: str, **labels: str) -> float: + return value(after, name, **labels) - value(before, name, **labels) + + +def websocket_round_trip(path: str) -> None: + parsed = urllib.parse.urlsplit(BASE) + host = parsed.hostname or "127.0.0.1" + port = parsed.port or 80 + key = base64.b64encode(os.urandom(16)).decode() + expected_accept = base64.b64encode( + hashlib.sha1((key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").encode()).digest() + ).decode() + conn = socket.create_connection((host, port), timeout=5) + try: + conn.sendall( + ( + f"GET {path} HTTP/1.1\r\n" + f"Host: {host}:{port}\r\n" + "Upgrade: websocket\r\n" + "Connection: Upgrade\r\n" + f"Sec-WebSocket-Key: {key}\r\n" + "Sec-WebSocket-Version: 13\r\n\r\n" + ).encode() + ) + response = b"" + while b"\r\n\r\n" not in response: + chunk = conn.recv(4096) + assert chunk, "websocket handshake closed before headers" + response += chunk + headers = response.split(b"\r\n\r\n", 1)[0].decode().split("\r\n") + assert headers[0].startswith("HTTP/1.1 101 "), headers[0] + parsed_headers = { + key.strip().lower(): value.strip() + for key, value in (line.split(":", 1) for line in headers[1:]) + } + assert parsed_headers.get("sec-websocket-accept") == expected_accept + + # A real masked client close frame; the server must finish the inline + # subscription handler before its HTTP middleware observation appears. + payload = struct.pack("!H", 1000) + mask = os.urandom(4) + masked = bytes(byte ^ mask[i % 4] for i, byte in enumerate(payload)) + conn.sendall(bytes((0x88, 0x80 | len(payload))) + mask + masked) + conn.settimeout(2) + try: + while conn.recv(4096): + pass + except (TimeoutError, socket.timeout): + pass + finally: + conn.close() + + +def main() -> None: + # The archive is real and sealed. Startup may still be completing its + # offline live-source attempt, so wait for the XRPC readiness gate. + listing = None + for _ in range(100): + status, body, _ = request("/xrpc/network.bsky.jetstream.listSegments") + if status == 200: + listing = json.loads(body) + break + time.sleep(0.05) + assert listing is not None, "archive XRPC never became ready" + segment = listing["segments"][0] + name = urllib.parse.quote(segment["name"]) + + before = scrape() + + block_path = f"/xrpc/network.bsky.jetstream.getBlock?segment={name}&blockIndex=0" + status, frame, headers = request(block_path) + assert status == 200 and frame + etag = headers["etag"] + assert request(block_path, headers={"If-None-Match": etag})[0] == 304 + range_status, partial, _ = request(block_path, headers={"Range": "bytes=0-0"}) + assert range_status == 206 and partial == frame[:1] + assert request(block_path, headers={"Range": "bytes=999999999-"})[0] == 416 + assert request("/xrpc/network.bsky.jetstream.getBlock")[0] == 400 + assert request(f"/xrpc/network.bsky.jetstream.getBlock?segment={name}&blockIndex=999999")[0] == 404 + assert request("/xrpc/network.bsky.jetstream.notImplemented")[0] == 400 + assert request(block_path, method="POST")[0] == 405 + assert request("/")[0] == 200 + assert request("/status", method="HEAD")[0] == 200 + assert request("/healthz")[0] == 200 + assert request("/not-found")[0] == 404 + websocket_round_trip("/subscribe") + + after = None + for _ in range(100): + candidate = scrape() + if delta(before, candidate, "jetstream_http_request_duration_seconds_count", handler="subscribe", method="GET", code="200") == 1: + after = candidate + break + time.sleep(0.05) + assert after is not None, "subscription HTTP lifetime was not observed" + + expected_http = { + ("xrpc/", "GET", "200"): 1, + ("xrpc/", "GET", "304"): 1, + ("xrpc/", "GET", "206"): 1, + ("xrpc/", "GET", "416"): 1, + ("xrpc/", "GET", "400"): 2, + ("xrpc/", "GET", "404"): 1, + ("xrpc/", "POST", "405"): 1, + ("root", "GET", "200"): 1, + ("status", "HEAD", "200"): 1, + ("subscribe", "GET", "200"): 1, + } + for (handler, method, code), expected in expected_http.items(): + got = delta( + before, + after, + "jetstream_http_request_duration_seconds_count", + handler=handler, + method=method, + code=code, + ) + assert got == expected, ((handler, method, code), got, expected) + + total_http_delta = sum( + value(after, name, **dict(labels)) - value(before, name, **dict(labels)) + for name, labels in set(before) | set(after) + if name == "jetstream_http_request_duration_seconds_count" + ) + assert total_http_delta == sum(expected_http.values()), total_http_delta + + expected_results = {"ok": 3, "bad_request": 1, "not_found": 1, "error": 1} + for result, expected in expected_results.items(): + assert delta(before, after, "jetstream_getblock_requests_total", result=result) == expected + assert delta(before, after, "jetstream_getblock_served_bytes_total") == len(frame) + assert delta(before, after, "jetstream_getblock_duration_seconds_count") == 6 + + print( + "HTTP/getBlock metrics PASS: " + f"http={int(total_http_delta)} getBlock=6 full_bytes={len(frame)} websocket=1" + ) + + +if __name__ == "__main__": + main() -- 2.51.2