From 455b34878a7c6f347fb91548ae1e9b8730d07668 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Thu, 16 Jul 2026 20:32:51 -0500 Subject: [PATCH] match canonical subscriber lifecycle semantics --- build.zig.zon | 8 +- docs/grafana-dashboard.md | 10 + docs/semantic-parity.md | 4 +- docs/upstream-subscribe-protocol.md | 5 +- src/internal/cold.zig | 158 +++++++++++-- src/internal/ingest.zig | 9 +- src/internal/metrics.zig | 107 +++++++++ src/internal/server.zig | 338 ++++++++++++++++++++++------ src/internal/tail.zig | 71 +++++- 9 files changed, 606 insertions(+), 104 deletions(-) diff --git a/build.zig.zon b/build.zig.zon index 89052dd..e05a795 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -5,12 +5,12 @@ .minimum_zig_version = "0.16.0", .dependencies = .{ .zat = .{ - .url = "git+https://tangled.org/zat.dev/zat#6fb7ac137610db71ce2b2636955fab710e34a587", - .hash = "zat-0.3.16-5PuC7o4pCwBjGxkaZdIt-qcGW-u3sM4mgsV1xRTrSgxI", + .url = "git+https://tangled.org/zat.dev/zat#03383118cd30e6e40fbb115eeb69c7a0c724a7ad", + .hash = "zat-0.3.16-5PuC7o8pCwAgpjdP8uWtJJVX5YU7P3NnoVIIRiqgKSeD", }, .websocket = .{ - .url = "https://tangled.org/zzstoatzz.io/websocket.zig/archive/v0.1.9.tar.gz", - .hash = "websocket-0.1.9-ZPISdYtnBABW10-3-CAGGR7fd6q5t8vlOj6Sxd-6Qu-C", + .url = "https://tangled.org/zzstoatzz.io/websocket.zig/archive/ac5d16e.tar.gz", + .hash = "websocket-0.1.9-ZPISdWJuBAD1NNhJiyV9_BwI9ksRWFOnzpzB4WqFM93N", }, .zstd = .{ .url = "https://github.com/facebook/zstd/releases/download/v1.5.7/zstd-1.5.7.tar.gz", diff --git a/docs/grafana-dashboard.md b/docs/grafana-dashboard.md index 55c48f0..b6d213e 100644 --- a/docs/grafana-dashboard.md +++ b/docs/grafana-dashboard.md @@ -64,6 +64,16 @@ The following dashboard families already have non-placeholder producers: - `jetstream_subscribe_bytes_encoded_total` - `jetstream_subscribe_events_filtered_total` - `jetstream_subscribe_events_oversize_total` +- `jetstream_subscribe_events_skipped_total` +- `jetstream_subscribe_encode_errors_total` +- `jetstream_subscribe_options_updates_total` +- `jetstream_subscribe_options_update_errors_total` +- `jetstream_subscribe_cursor_requests_total` +- `jetstream_subscribe_cursor_resolve_seconds` +- `jetstream_subscribe_hot_reads_total` +- `jetstream_subscribe_cold_reads_total` +- `jetstream_subscribe_adversarial_drops_total` +- `jetstream_subscribe_clean_disconnects_total` - `jetstream_import_jobs_total` - `jetstream_import_job_duration_seconds` - `jetstream_import_phase` diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 45d096e..ad85799 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -15,12 +15,12 @@ this document or `bootstrap-semantic-parity.md` is open. | Bootstrap, merge, retry | Detailed row-by-row contract and adversity receipt | See `bootstrap-semantic-parity.md`; adversity remains open | | JSS v1 storage | Upstream-produced fixture, reciprocal segment parsing, checksum/index/gloom vectors | closed for sealed format; rerun reciprocal corpus before experiment | | Archive XRPC | Official Go client plus Stream conformance replay over listSegments/getSegment/getBlock/planBackfill | implementation present; fresh full suite required | -| Subscribe v1/v2 | Wire/filter/cursor/zstd unit and local e2e coverage | implementation present; official-client replay and below-floor/adversarial suite required | +| Subscribe v1/v2 | Wire/filter/cursor/zstd unit and local e2e coverage | Cursor parsing/resolution now precedes upgrade; both endpoints use the seq/time-us magnitude split, v1 clamps below-floor seqs, and v2 rejects them. Hot/cold scan, skip, encode, options-update, oversize-parser, clean-peer-close, cursor-resolution, and sustained adversarial-rate boundaries feed the canonical metrics, with deterministic offline tests derived from upstream. **Still open:** official-client replay, 30-second ping/write deadlines, graceful-shutdown close accounting, v2 dictionary negotiation/download, and pre-upgrade timestamp translation fault classification. | | Sync 1.1 live verification | Upstream requires durable per-DID chain/hosting state, MST inversion, op-CID consistency, rev replay/future guards, default acceptance of legacy-shaped commits, transparent whole-repo repair, and Atmos delivery scheduling | **closed at the implementation boundary.** Diff-CAR/MST inversion (including upstream's narrow default lenient carve-out), post-state op-CID proof, signed inner/outer consistency, exact decimal size/future-rev/replay gates, exact legacy-shape detection with default `LegacyAccept`, two-phase durable chain/hosting state, and account/identity replay ratchets are implemented. Recoverable decode/inversion/duplicate-path/op-CID/chain failures route to repair; signature failures and outer/inner producer-integrity mismatches bypass it, and no failed event is archived. `prepareRepo` retains and authenticates the canonical signed complete head without a second multi-gigabyte parse/walk. The live path now uses Atmos's 32 worker slots, one worker-held FIFO chain per DID, a 64-pending drop-oldest boundary, completion-order batches of 50/500 ms across DIDs, and `min(inflight)-1` cursor watermarks. Dropped and silent events leave the inflight set without inventing a delivery batch; cursor advancement remains coupled to durable archive boundaries. Offline receipts use real encoded frames, real RocksDB verifier state, and actual worker blocking rather than fixture verdicts or scheduler mocks. See `live-scheduler.md`. | | 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 | durable repo fields implemented; **operator repo/host views 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. Subscriber gauges and successful-delivery bytes/events are split by actual `none`/`zstd` compression, and filter/oversize drops are counted at their decision boundaries. 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. **Remaining subscriber/storage/API/orchestrator metric families and Go-row adaptation 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; compression labels currently represent actual `none`/v1 `zstd` delivery only, while deflate remains zero until the WebSocket implementation genuinely compresses frames. 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. **Remaining storage/API/orchestrator families and Go-row adaptation 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-subscribe-protocol.md b/docs/upstream-subscribe-protocol.md index 2bc0a25..55373b1 100644 --- a/docs/upstream-subscribe-protocol.md +++ b/docs/upstream-subscribe-protocol.md @@ -20,10 +20,11 @@ endpoint: `GET /subscribe` → websocket upgrade. pre-upgrade errors are plain H ### cursor - absent/empty → live tail - not base-10 int64 or negative → HTTP 400 pre-upgrade -- v1 unit = time_us (µs since epoch). upstream v2 splits by magnitude: < 1e15 = seq, >= 1e15 = time_us +- both endpoints split by magnitude: < 1e15 = seq, >= 1e15 = time_us (µs since epoch) - future cursor → silently live mode, no error. cursor 0 → floored to first event. - too-old cursor on v1 → **silently clamped** to oldest retained (never rejected) -- replay: first event with witnessed_at >= cursor (inclusive-overlapping; clients dedupe) +- seq replay: first event with local seq >= cursor; timestamp replay: first + event with witnessed_at >= cursor (inclusive-overlapping; clients dedupe) ### compress - `compress=true` or header `Socket-Encoding: zstd` → per-message zstd, binary frames diff --git a/src/internal/cold.zig b/src/internal/cold.zig index 8b76887..dbbd443 100644 --- a/src/internal/cold.zig +++ b/src/internal/cold.zig @@ -18,11 +18,15 @@ const filter_mod = @import("filter.zig"); const segment = @import("segment.zig"); const archive_mod = @import("archive.zig"); const wire = @import("wire.zig"); +const tail = @import("tail.zig"); const Io = std.Io; const Allocator = std.mem.Allocator; const log = std.log.scoped(.stream); +pub const ReplayResult = enum { sent, filtered, oversize, skipped_sync, skipped_resync, encode_error }; +pub const WriteResult = enum { sent, oversize, encode_error }; + /// stream all events with witnessed_at in [from_us, until_us) matching /// `filter` to `sink` (fn (ctx, json: []const u8) !void). returns the /// number of frames written. @@ -31,9 +35,10 @@ pub fn replay( io: Io, archive: *archive_mod.Archive, filter: anytype, // *Subscriber-shaped: wantsPred via filter fn below - comptime wants: fn (@TypeOf(filter), wire.Kind, []const u8, []const u8) bool, + comptime wants: fn (@TypeOf(filter), wire.Kind, []const u8, []const u8, tail.SkipV1) bool, sink: anytype, - comptime write: fn (@TypeOf(sink), []const u8) anyerror!void, + comptime write: fn (@TypeOf(sink), []const u8) anyerror!WriteResult, + comptime observe: fn (@TypeOf(sink), u64, ReplayResult) anyerror!void, from_us: i64, until_us: i64, ) !u64 { @@ -71,37 +76,62 @@ pub fn replay( defer block.deinit(allocator); for (block.events) |ev| { if (ev.witnessed_at < from_us or ev.witnessed_at >= until_us) continue; + const kind: wire.Kind = switch (ev.kind) { + .create, .update, .delete, .create_resync => .commit, + .identity => .identity, + .account => .account, + .sync => .sync, + }; + const skip: tail.SkipV1 = switch (ev.kind) { + .sync => .sync, + .create_resync => .resync_replacement, + else => .none, + }; + if (skip != .none) { + try observe(sink, ev.seq, if (skip == .sync) .skipped_sync else .skipped_resync); + continue; + } + if (!wants(filter, kind, ev.did, ev.collection, .none)) { + try observe(sink, ev.seq, .filtered); + continue; + } var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const json = encodeEvent(arena.allocator(), ev) catch |err| { log.warn("cold: undecodable row seq={d}: {s}", .{ ev.seq, @errorName(err) }); + try observe(sink, ev.seq, .encode_error); + continue; + } orelse { + try observe(sink, ev.seq, .encode_error); continue; - } orelse continue; - const kind: wire.Kind = switch (ev.kind) { - .create, .update, .delete => .commit, - .identity => .identity, - .account => .account, - else => continue, }; - if (!wants(filter, kind, ev.did, ev.collection)) continue; - try write(sink, json); - frames += 1; + switch (try write(sink, json)) { + .sent => { + try observe(sink, ev.seq, .sent); + frames += 1; + }, + .oversize => try observe(sink, ev.seq, .oversize), + .encode_error => try observe(sink, ev.seq, .encode_error), + } } } } return frames; } -/// v2: stream events with seq >= from_seq up to the hot tail's seq floor, -/// encoded on the v2 wire (segment rows carry everything v2 needs). +/// stream events with seq >= from_seq up to the hot tail's seq floor on the +/// requested endpoint wire. Both v1 and v2 accept seq cursors; only v2 emits +/// sync/resync-replacement rows and the v2 envelope fields. pub fn replaySeq( allocator: Allocator, io: Io, archive: *archive_mod.Archive, filter: anytype, - comptime wants: fn (@TypeOf(filter), wire.Kind, []const u8, []const u8) bool, + comptime wants: fn (@TypeOf(filter), wire.Kind, []const u8, []const u8, tail.SkipV1) bool, sink: anytype, - comptime write: fn (@TypeOf(sink), []const u8) anyerror!void, + comptime write: fn (@TypeOf(sink), []const u8) anyerror!WriteResult, + comptime observe: fn (@TypeOf(sink), u64, ReplayResult) anyerror!void, + v2: bool, from_seq: u64, ) !u64 { var names: std.ArrayList([]u8) = .empty; @@ -130,18 +160,47 @@ pub fn replaySeq( defer block.deinit(allocator); for (block.events) |ev| { if (ev.seq < from_seq) continue; - var arena = std.heap.ArenaAllocator.init(allocator); - defer arena.deinit(); - const json = encodeEventV2(arena.allocator(), ev) catch continue orelse continue; const kind: wire.Kind = switch (ev.kind) { .create, .update, .delete, .create_resync => .commit, .identity => .identity, .account => .account, .sync => .sync, }; - if (!wants(filter, kind, ev.did, ev.collection)) continue; - try write(sink, json); - frames += 1; + if (!v2) switch (ev.kind) { + .sync => { + try observe(sink, ev.seq, .skipped_sync); + continue; + }, + .create_resync => { + try observe(sink, ev.seq, .skipped_resync); + continue; + }, + else => {}, + }; + if (!wants(filter, kind, ev.did, ev.collection, .none)) { + try observe(sink, ev.seq, .filtered); + continue; + } + var arena = std.heap.ArenaAllocator.init(allocator); + defer arena.deinit(); + const json = (if (v2) + encodeEventV2(arena.allocator(), ev) + else + encodeEvent(arena.allocator(), ev)) catch { + try observe(sink, ev.seq, .encode_error); + continue; + } orelse { + try observe(sink, ev.seq, .encode_error); + continue; + }; + switch (try write(sink, json)) { + .sent => { + try observe(sink, ev.seq, .sent); + frames += 1; + }, + .oversize => try observe(sink, ev.seq, .oversize), + .encode_error => try observe(sink, ev.seq, .encode_error), + } } } } @@ -293,3 +352,60 @@ test "indexed timestamp overrides wire display time without changing witness" { const legacy = (try encodeEvent(allocator, fallback)).?; try testing.expect(std.mem.indexOf(u8, legacy, "\"time_us\":1700000000000000") != null); } + +test "offline seq replay uses the selected endpoint wire" { + const testing = std.testing; + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tmp = testing.tmpDir(.{}); + defer tmp.cleanup(); + var path_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const path_len = try tmp.dir.realPath(io, &path_buf); + var archive = try archive_mod.Archive.init(testing.allocator, io, path_buf[0..path_len]); + defer archive.deinit(); + archive.max_events_per_block = 1; + archive.writer.max_events_per_block = 1; + _ = try archive.append(.{ + .seq = 0, + .witnessed_at = 1_700_000_000_000_000, + .indexed_at = 0, + .kind = .create, + .did = "did:plc:coldwirefixture", + .collection = "app.bsky.feed.post", + .rkey = "3k2abcdefghij", + .rev = "3k2abcdefghij", + .payload = "\xa1\x64text\x62hi", + }, 1_700_000_000_000_000); + try archive.rotate(); + + const Sink = struct { + allocator: Allocator, + frames: std.ArrayList([]u8) = .empty, + + fn deinit(self: *@This()) void { + for (self.frames.items) |frame| self.allocator.free(frame); + self.frames.deinit(self.allocator); + } + fn wants(_: *@This(), _: wire.Kind, _: []const u8, _: []const u8, _: tail.SkipV1) bool { + return true; + } + fn write(self: *@This(), json: []const u8) anyerror!WriteResult { + try self.frames.append(self.allocator, try self.allocator.dupe(u8, json)); + return .sent; + } + fn observe(_: *@This(), _: u64, _: ReplayResult) anyerror!void {} + }; + + var v1: Sink = .{ .allocator = testing.allocator }; + defer v1.deinit(); + try testing.expectEqual(@as(u64, 1), try replaySeq(testing.allocator, io, &archive, &v1, Sink.wants, &v1, Sink.write, Sink.observe, false, 1)); + try testing.expect(std.mem.indexOf(u8, v1.frames.items[0], "\"seq\":1") == null); + try testing.expect(std.mem.indexOf(u8, v1.frames.items[0], "record_cbor") == null); + + var v2: Sink = .{ .allocator = testing.allocator }; + defer v2.deinit(); + try testing.expectEqual(@as(u64, 1), try replaySeq(testing.allocator, io, &archive, &v2, Sink.wants, &v2, Sink.write, Sink.observe, true, 1)); + try testing.expect(std.mem.indexOf(u8, v2.frames.items[0], "\"seq\":1") != null); + try testing.expect(std.mem.indexOf(u8, v2.frames.items[0], "record_cbor") != null); +} diff --git a/src/internal/ingest.zig b/src/internal/ingest.zig index e6044f5..1a5944f 100644 --- a/src/internal/ingest.zig +++ b/src/internal/ingest.zig @@ -234,7 +234,12 @@ pub const Consumer = struct { .account => .account, .sync => .sync, }; - try sink.consumer.tail.append(kind, row.did, row.collection, stamped.displayTime(), stamped.seq, json, json_v2); + const skip_v1: tail_mod.SkipV1 = switch (row.kind) { + .sync => .sync, + .create_resync => .resync_replacement, + else => .none, + }; + try sink.consumer.tail.append(kind, row.did, row.collection, stamped.displayTime(), stamped.seq, json, json_v2, skip_v1); } sink.last_seq = first + sink.batch.items.len - 1; sink.batch.clearRetainingCapacity(); @@ -381,7 +386,7 @@ pub const Consumer = struct { } for (frames.items) |frame| { - try self.tail.append(frame.kind, frame.did, frame.collection, frame.time_us, frame.seq, frame.json, frame.json_v2); + try self.tail.append(frame.kind, frame.did, frame.collection, frame.time_us, frame.seq, frame.json, frame.json_v2, if (frame.kind == .sync) .sync else .none); } if (persisted.last == 0) { diff --git a/src/internal/metrics.zig b/src/internal/metrics.zig index 7a3afc6..b049670 100644 --- a/src/internal/metrics.zig +++ b/src/internal/metrics.zig @@ -17,6 +17,8 @@ const Io = std.Io; 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 Stats = struct { start_time_s: i64 = 0, @@ -88,6 +90,19 @@ pub const Stats = struct { subscribe_bytes_encoded: std.enums.EnumArray(Compression, Counter) = .initFill(.init(0)), subscribe_events_filtered_total: Counter = .init(0), subscribe_events_oversize_total: Counter = .init(0), + subscribe_events_skipped_sync_total: Counter = .init(0), + subscribe_events_skipped_resync_total: Counter = .init(0), + subscribe_encode_errors_total: Counter = .init(0), + subscribe_options_update_errors: std.enums.EnumArray(OptionsErrorReason, Counter) = .initFill(.init(0)), + subscribe_cursor_requests: std.enums.EnumArray(CursorMode, Counter) = .initFill(.init(0)), + subscribe_cursor_resolve_count: Counter = .init(0), + subscribe_cursor_resolve_sum_us: Counter = .init(0), + subscribe_cursor_resolve_buckets: [8]Counter = @splat(.init(0)), + subscribe_hot_reads_total: Counter = .init(0), + subscribe_cold_reads_total: Counter = .init(0), + subscribe_adversarial_drops_total: Counter = .init(0), + subscribe_clean_disconnects_total: Counter = .init(0), + subscribe_options_updates_total: Counter = .init(0), pub fn addDrops(self: *Stats, before: *const convert.Drops, after: *const convert.Drops) void { inline for (comptime std.enums.values(convert.DropReason)) |reason| { @@ -95,6 +110,15 @@ pub const Stats = struct { if (delta > 0) _ = self.drops.getPtr(reason).fetchAdd(delta, .monotonic); } } + + pub fn observeCursorResolve(self: *Stats, duration_us: u64) void { + _ = self.subscribe_cursor_resolve_count.fetchAdd(1, .monotonic); + _ = self.subscribe_cursor_resolve_sum_us.fetchAdd(duration_us, .monotonic); + const bounds = [_]u64{ 100, 400, 1_600, 6_400, 25_600, 102_400, 409_600, 1_638_400 }; + inline for (bounds, 0..) |upper, i| { + if (duration_us <= upper) _ = self.subscribe_cursor_resolve_buckets[i].fetchAdd(1, .monotonic); + } + } }; pub const TailGauges = struct { @@ -277,6 +301,67 @@ pub fn format( stats.subscribe_events_filtered_total.load(.monotonic), stats.subscribe_events_oversize_total.load(.monotonic), }) catch {}; + w.print( + \\# TYPE jetstream_subscribe_events_skipped_total counter + \\jetstream_subscribe_events_skipped_total{{reason="sync"}} {d} + \\jetstream_subscribe_events_skipped_total{{reason="resync_replacement"}} {d} + \\# TYPE jetstream_subscribe_encode_errors_total counter + \\jetstream_subscribe_encode_errors_total {d} + \\# TYPE jetstream_subscribe_hot_reads_total counter + \\jetstream_subscribe_hot_reads_total {d} + \\# TYPE jetstream_subscribe_cold_reads_total counter + \\jetstream_subscribe_cold_reads_total {d} + \\# TYPE jetstream_subscribe_adversarial_drops_total counter + \\jetstream_subscribe_adversarial_drops_total {d} + \\# TYPE jetstream_subscribe_clean_disconnects_total counter + \\jetstream_subscribe_clean_disconnects_total {d} + \\# TYPE jetstream_subscribe_options_updates_total counter + \\jetstream_subscribe_options_updates_total {d} + \\ + , .{ + stats.subscribe_events_skipped_sync_total.load(.monotonic), + stats.subscribe_events_skipped_resync_total.load(.monotonic), + stats.subscribe_encode_errors_total.load(.monotonic), + stats.subscribe_hot_reads_total.load(.monotonic), + stats.subscribe_cold_reads_total.load(.monotonic), + stats.subscribe_adversarial_drops_total.load(.monotonic), + stats.subscribe_clean_disconnects_total.load(.monotonic), + stats.subscribe_options_updates_total.load(.monotonic), + }) catch {}; + w.print("# TYPE jetstream_subscribe_options_update_errors_total counter\n", .{}) catch {}; + inline for (comptime std.enums.values(OptionsErrorReason)) |reason| { + w.print("jetstream_subscribe_options_update_errors_total{{reason=\"{s}\"}} {d}\n", .{ + @tagName(reason), stats.subscribe_options_update_errors.getPtrConst(reason).load(.monotonic), + }) catch {}; + } + w.print("# TYPE jetstream_subscribe_cursor_requests_total counter\n", .{}) catch {}; + inline for (comptime std.enums.values(CursorMode)) |mode| { + w.print("jetstream_subscribe_cursor_requests_total{{mode=\"{s}\"}} {d}\n", .{ + @tagName(mode), stats.subscribe_cursor_requests.getPtrConst(mode).load(.monotonic), + }) catch {}; + } + const cursor_sum_us = stats.subscribe_cursor_resolve_sum_us.load(.monotonic); + w.print( + \\# TYPE jetstream_subscribe_cursor_resolve_seconds histogram + \\jetstream_subscribe_cursor_resolve_seconds_sum {d}.{d:0>6} + \\jetstream_subscribe_cursor_resolve_seconds_count {d} + \\ + , .{ + cursor_sum_us / std.time.us_per_s, + cursor_sum_us % std.time.us_per_s, + stats.subscribe_cursor_resolve_count.load(.monotonic), + }) catch {}; + const cursor_bucket_us = [_]u64{ 100, 400, 1_600, 6_400, 25_600, 102_400, 409_600, 1_638_400 }; + inline for (cursor_bucket_us, 0..) |upper_us, i| { + w.print("jetstream_subscribe_cursor_resolve_seconds_bucket{{le=\"{d}.{d:0>6}\"}} {d}\n", .{ + upper_us / std.time.us_per_s, + upper_us % std.time.us_per_s, + stats.subscribe_cursor_resolve_buckets[i].load(.monotonic), + }) catch {}; + } + w.print("jetstream_subscribe_cursor_resolve_seconds_bucket{{le=\"+Inf\"}} {d}\n", .{ + stats.subscribe_cursor_resolve_count.load(.monotonic), + }) catch {}; w.print( \\# TYPE stream_verify_total counter @@ -457,6 +542,17 @@ test "format renders counters" { _ = stats.subscribe_bytes_encoded.getPtr(.zstd).fetchAdd(800, .monotonic); _ = stats.subscribe_events_filtered_total.fetchAdd(5, .monotonic); _ = stats.subscribe_events_oversize_total.fetchAdd(6, .monotonic); + _ = stats.subscribe_events_skipped_sync_total.fetchAdd(7, .monotonic); + _ = stats.subscribe_events_skipped_resync_total.fetchAdd(8, .monotonic); + _ = stats.subscribe_encode_errors_total.fetchAdd(9, .monotonic); + _ = stats.subscribe_options_update_errors.getPtr(.invalid_options).fetchAdd(10, .monotonic); + _ = stats.subscribe_cursor_requests.getPtr(.seq).fetchAdd(11, .monotonic); + stats.observeCursorResolve(400); + _ = stats.subscribe_hot_reads_total.fetchAdd(12, .monotonic); + _ = stats.subscribe_cold_reads_total.fetchAdd(13, .monotonic); + _ = stats.subscribe_adversarial_drops_total.fetchAdd(14, .monotonic); + _ = stats.subscribe_clean_disconnects_total.fetchAdd(15, .monotonic); + _ = stats.subscribe_options_updates_total.fetchAdd(16, .monotonic); var buf: [16 * 1024]u8 = undefined; const out = format(&buf, &stats, 2, .{ .entries = 5, .bytes = 100, .base = 1, .tip = 6 }, 60); @@ -493,6 +589,17 @@ test "format renders counters" { try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_bytes_encoded_total{compression=\"zstd\"} 800") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_events_filtered_total 5") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_events_oversize_total 6") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_events_skipped_total{reason=\"sync\"} 7") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_events_skipped_total{reason=\"resync_replacement\"} 8") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_encode_errors_total 9") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_options_update_errors_total{reason=\"invalid_options\"} 10") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_cursor_requests_total{mode=\"seq\"} 11") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_cursor_resolve_seconds_bucket{le=\"0.000400\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_hot_reads_total 12") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_cold_reads_total 13") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_adversarial_drops_total 14") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_clean_disconnects_total 15") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_subscribe_options_updates_total 16") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_livestream_replayed_account_events_dropped_total 3") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_livestream_replayed_identity_events_dropped_total 4") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_failed_repo_retry_attempts_total 5") != null); diff --git a/src/internal/server.zig b/src/internal/server.zig index 2061807..f10a09d 100644 --- a/src/internal/server.zig +++ b/src/internal/server.zig @@ -5,8 +5,8 @@ //! time; options_update swaps them atomically. contract and quirks per //! docs/upstream-subscribe-protocol.md. //! -//! not yet (tracked in README): zstd compress (needs libzstd — std can only -//! decompress), cold storage below the tail floor, /subscribe-v2. +//! v1 and v2 share the pull fanout while retaining their distinct wire, +//! cursor-floor, and skipped-event contracts. const std = @import("std"); const websocket = @import("websocket"); @@ -26,6 +26,46 @@ const log = std.log.scoped(.stream); /// 8 MB, matching the Io.Threaded stacks (docs/lessons-from-zlay.md #1) const subscriber_stack_size: usize = 8 * 1024 * 1024; +const slow_window_us: i64 = 60 * std.time.us_per_s; +const slow_lag_threshold: u64 = 100_000; +const slow_min_rate: u64 = 5; + +const SlowDetector = struct { + streak_start_us: i64 = 0, + streak_start_pos: u64 = 0, + in_streak: bool = false, + + fn observe(self: *SlowDetector, now_us: i64, pos: u64, lag: u64) bool { + if (lag <= slow_lag_threshold) { + self.in_streak = false; + return false; + } + if (!self.in_streak) { + self.in_streak = true; + self.streak_start_us = now_us; + self.streak_start_pos = pos; + return false; + } + const elapsed_us = now_us - self.streak_start_us; + if (elapsed_us < slow_window_us) return false; + const progress = pos -| self.streak_start_pos; + if (progress *| std.time.us_per_s >= @as(u64, @intCast(elapsed_us)) *| slow_min_rate) { + self.streak_start_us = now_us; + self.streak_start_pos = pos; + return false; + } + return true; + } +}; + +const CursorPlan = struct { + /// null means live: anchor at the published tip when delivery starts, + /// after requireHello if requested. Replay plans are fixed pre-upgrade. + start_idx: ?u64 = null, + cold_from_us: ?i64 = null, + cold_from_seq: ?u64 = null, + mode: metrics.CursorMode = .live, +}; pub const Hub = struct { allocator: Allocator, @@ -73,6 +113,7 @@ const Subscriber = struct { /// /subscribe-v2: superset wire (seq, record_cbor, #sync), strict cursors v2: bool = false, thread: ?std.Thread = null, + slow: SlowDetector = .{}, /// the pull loop. cleanup happens in Handler.close, which joins this /// thread — the subscriber never frees itself, so the handler can always @@ -81,9 +122,15 @@ const Subscriber = struct { if (self.cold_from_seq) |from_seq| self.coldPhaseSeq(from_seq); if (self.cold_from_us) |from_us| self.coldPhase(from_us); while (self.alive.load(.acquire)) { - const json = self.hub.tail.next(self.hub.allocator, &self.idx, &self.stopped, self.v2, self, wantsPred) catch return orelse return; + const json = self.hub.tail.next(self.hub.allocator, &self.idx, &self.stopped, self.v2, self, hotWantsPred) catch return orelse return; defer self.hub.allocator.free(json); if (!self.alive.load(.acquire)) return; + const tip = self.hub.tail.publishedTip(); + const lag = if (tip > self.idx) tip - self.idx else 0; + if (self.slow.observe(Io.Timestamp.now(self.hub.io, .real).toMicroseconds(), self.idx, lag)) { + _ = self.hub.stats.subscribe_adversarial_drops_total.fetchAdd(1, .monotonic); + return; + } // oversize events are silently skipped, cursor advances (v1 quirk: // cap compares the uncompressed JSON length) const cap = self.max_msg_size.load(.monotonic); @@ -91,17 +138,20 @@ const Subscriber = struct { _ = self.hub.stats.subscribe_events_oversize_total.fetchAdd(1, .monotonic); continue; } - self.write(json) catch return; + if ((self.write(json) catch return) == .encode_error) + _ = self.hub.stats.subscribe_encode_errors_total.fetchAdd(1, .monotonic); } } - fn write(self: *Subscriber, json: []const u8) !void { + fn write(self: *Subscriber, json: []const u8) !cold.WriteResult { self.write_mutex.lockUncancelable(self.hub.io); defer self.write_mutex.unlock(self.hub.io); if (!self.alive.load(.acquire)) return error.Closed; const scheme = self.compression(); if (self.compressor) |*comp| { - const frame = try comp.compress(self.hub.allocator, json); + const frame = comp.compress(self.hub.allocator, json) catch { + return .encode_error; + }; defer self.hub.allocator.free(frame); try self.conn.writeBin(frame); _ = self.hub.stats.subscribe_events_sent.getPtr(scheme).fetchAdd(1, .monotonic); @@ -112,16 +162,19 @@ const Subscriber = struct { _ = self.hub.stats.subscribe_bytes_sent.getPtr(scheme).fetchAdd(@intCast(json.len), .monotonic); } _ = self.hub.stats.subscribe_bytes_encoded.getPtr(scheme).fetchAdd(@intCast(json.len), .monotonic); + return .sent; } fn compression(self: *const Subscriber) metrics.Compression { - return if (self.compressor != null) .zstd else .none; + if (self.compressor != null) return .zstd; + return if (self.conn.compression != null) .deflate else .none; } fn coldPhase(self: *Subscriber, from_us: i64) void { const hub = self.hub; const archive = hub.archive orelse return; const until_us = hub.tail.oldestTime() orelse std.math.maxInt(i64); + _ = hub.stats.subscribe_cold_reads_total.fetchAdd(1, .monotonic); const n = cold.replay( hub.allocator, hub.io, @@ -130,6 +183,7 @@ const Subscriber = struct { wantsPred, self, coldWrite, + coldObserve, from_us, until_us, ) catch |err| { @@ -142,6 +196,7 @@ const Subscriber = struct { fn coldPhaseSeq(self: *Subscriber, from_seq: u64) void { const hub = self.hub; const archive = hub.archive orelse return; + _ = hub.stats.subscribe_cold_reads_total.fetchAdd(1, .monotonic); const n = cold.replaySeq( hub.allocator, hub.io, @@ -150,25 +205,53 @@ const Subscriber = struct { wantsPred, self, coldWrite, + coldObserve, + self.v2, from_seq, ) catch |err| { - log.warn("v2 cold replay failed: {s}", .{@errorName(err)}); + log.warn("seq cold replay failed: {s}", .{@errorName(err)}); return; }; - log.debug("v2 cold replay served {d} frames", .{n}); + log.debug("seq cold replay served {d} frames", .{n}); } - fn coldWrite(self: *Subscriber, json: []const u8) anyerror!void { + fn coldWrite(self: *Subscriber, json: []const u8) anyerror!cold.WriteResult { if (!self.alive.load(.acquire)) return error.Closed; const cap = self.max_msg_size.load(.monotonic); if (cap != 0 and json.len > cap) { _ = self.hub.stats.subscribe_events_oversize_total.fetchAdd(1, .monotonic); - return; + return .oversize; } - try self.write(json); + return self.write(json); } - fn wantsPred(self: *Subscriber, kind: wire.Kind, did: []const u8, collection: []const u8) bool { + fn coldObserve(self: *Subscriber, seq: u64, result: cold.ReplayResult) anyerror!void { + switch (result) { + .skipped_sync => _ = self.hub.stats.subscribe_events_skipped_sync_total.fetchAdd(1, .monotonic), + .skipped_resync => _ = self.hub.stats.subscribe_events_skipped_resync_total.fetchAdd(1, .monotonic), + .encode_error => _ = self.hub.stats.subscribe_encode_errors_total.fetchAdd(1, .monotonic), + .sent, .filtered, .oversize => {}, + } + const tip = if (self.hub.archive) |archive| archive.committed_seq.load(.acquire) else seq; + const lag = if (tip > seq) tip - seq else 0; + if (self.slow.observe(Io.Timestamp.now(self.hub.io, .real).toMicroseconds(), seq, lag)) { + _ = self.hub.stats.subscribe_adversarial_drops_total.fetchAdd(1, .monotonic); + return error.AdversarialSlow; + } + } + + fn wantsPred(self: *Subscriber, kind: wire.Kind, did: []const u8, collection: []const u8, skip: tail_mod.SkipV1) bool { + switch (skip) { + .sync => { + _ = self.hub.stats.subscribe_events_skipped_sync_total.fetchAdd(1, .monotonic); + return false; + }, + .resync_replacement => { + _ = self.hub.stats.subscribe_events_skipped_resync_total.fetchAdd(1, .monotonic); + return false; + }, + .none => {}, + } self.opts_mutex.lockUncancelable(self.hub.io); defer self.opts_mutex.unlock(self.hub.io); const wanted = self.filter.wants(kind, did, collection); @@ -176,6 +259,11 @@ const Subscriber = struct { return wanted; } + fn hotWantsPred(self: *Subscriber, kind: wire.Kind, did: []const u8, collection: []const u8, skip: tail_mod.SkipV1) bool { + _ = self.hub.stats.subscribe_hot_reads_total.fetchAdd(1, .monotonic); + return self.wantsPred(kind, did, collection, skip); + } + /// atomically replace the whole filter (options_update contract) fn updateOptions(self: *Subscriber, new_filter: filter_mod.Filter, max_msg_size: u64) void { self.opts_mutex.lockUncancelable(self.hub.io); @@ -195,22 +283,26 @@ pub const Handler = struct { /// starts (immediately, or on hello when requireHello) pending_filter: ?filter_mod.Filter, pending_max_msg_size: u64, - cursor: ?i64, + cursor_plan: CursorPlan, awaiting_hello: bool, compress: bool, v2: bool, 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 ""; 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"); return error.Close; // response already written; Close suppresses the 400 } - const url = handshake.url; - const qi = std.mem.indexOfScalar(u8, url, '?'); 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; - const query = if (qi) |i| url[i + 1 ..] else ""; var filter = filter_mod.Filter.parseQuery(hub.allocator, query) catch |err| { log.debug("subscribe: invalid options: {s}", .{@errorName(err)}); @@ -218,16 +310,19 @@ pub const Handler = struct { }; errdefer filter.deinit(); - // cursor: absent → live tail; not an int64 or negative → reject - // pre-upgrade; too-old → clamped to the floor; future → live tail - var cursor: ?i64 = null; - if (queryParam(query, "cursor")) |raw| { - if (raw.len > 0) { - const parsed = std.fmt.parseInt(i64, raw, 10) catch return error.InvalidRequest; - if (parsed < 0) return error.InvalidRequest; - cursor = parsed; - } - } + // Cursor resolution includes parsing and happens before upgrade, so + // malformed and too-old v2 cursors remain ordinary HTTP failures. + const resolve_started = Io.Timestamp.now(hub.io, .real).toMicroseconds(); + const cursor_plan = resolveCursor(hub, v2, queryParam(query, "cursor")) catch |err| { + const finished = Io.Timestamp.now(hub.io, .real).toMicroseconds(); + hub.stats.observeCursorResolve(@intCast(@max(finished - resolve_started, 0))); + if (err == error.CursorTooOld) + _ = hub.stats.subscribe_cursor_requests.getPtr(.too_old).fetchAdd(1, .monotonic); + return error.InvalidRequest; + }; + const resolve_finished = Io.Timestamp.now(hub.io, .real).toMicroseconds(); + hub.stats.observeCursorResolve(@intCast(@max(resolve_finished - resolve_started, 0))); + _ = hub.stats.subscribe_cursor_requests.getPtr(cursor_plan.mode).fetchAdd(1, .monotonic); // quirk: requireHello is true iff the value is exactly "true" const awaiting_hello = if (queryParam(query, "requireHello")) |v| @@ -251,7 +346,7 @@ pub const Handler = struct { .conn = conn, .pending_filter = filter, .pending_max_msg_size = parseMaxMsgSize(queryParam(query, "maxMessageSizeBytes")), - .cursor = cursor, + .cursor_plan = cursor_plan, .awaiting_hello = awaiting_hello, .compress = compress, .v2 = v2, @@ -270,31 +365,14 @@ pub const Handler = struct { fn startSubscriber(self: *Handler) !void { const hub = self.hub; - // v2 cursors below 1e15 are seqs; v1 (and v2 timestamp) cursors are - // time_us. below-floor: v1 clamps silently; v2 seq cursors reject - // (the client must know about the gap — upstream §5.1) - var start_idx: u64 = hub.tail.publishedTip(); - var cold_from: ?i64 = null; - var cold_from_seq: ?u64 = null; - if (self.cursor) |c| { - if (self.v2 and c < 1_000_000_000_000_000) { - const seq_cursor: u64 = @intCast(c); - if (hub.tail.indexForSeq(seq_cursor)) |idx| { - start_idx = idx; - } else if (hub.archive != null) { - start_idx = hub.tail.base; - cold_from_seq = seq_cursor; - } else { - return error.InvalidRequest; // below floor, no archive - } - } else { - start_idx = hub.tail.indexForCursor(c); - if (hub.archive != null) { - const floor = hub.tail.oldestTime(); - if (floor == null or c < floor.?) cold_from = c; - } - } - } + const plan = self.cursor_plan; + const start_idx = plan.start_idx orelse hub.tail.publishedTip(); + + var compressor: ?zstd.DictCompressor = if (self.compress) + try zstd.DictCompressor.init(zstd.v1_dictionary) + else + null; + errdefer if (compressor) |*comp| comp.deinit(); const sub = try hub.allocator.create(Subscriber); sub.* = .{ @@ -303,13 +381,10 @@ pub const Handler = struct { .filter = self.pending_filter.?, .max_msg_size = .init(self.pending_max_msg_size), .idx = start_idx, - .cold_from_us = cold_from, - .cold_from_seq = cold_from_seq, + .cold_from_us = plan.cold_from_us, + .cold_from_seq = plan.cold_from_seq, .v2 = self.v2, - .compressor = if (self.compress) - zstd.DictCompressor.init(zstd.v1_dictionary) catch null - else - null, + .compressor = compressor, }; self.pending_filter = null; // ownership moved to the subscriber self.subscriber = sub; @@ -334,24 +409,38 @@ pub const Handler = struct { if (tpe != .text) return; const parsed = std.json.parseFromSlice(std.json.Value, allocator, data, .{}) catch { + self.noteOptionsError(.bad_envelope_json); return self.closeWith(1007, "bad SubscriberSourcedMessage envelope"); }; defer parsed.deinit(); const root = parsed.value; - if (root != .object) return self.closeWith(1007, "bad SubscriberSourcedMessage envelope"); + if (root != .object) { + self.noteOptionsError(.bad_envelope_json); + return self.closeWith(1007, "bad SubscriberSourcedMessage envelope"); + } - const type_val = root.object.get("type") orelse + const type_val = root.object.get("type") orelse { + self.noteOptionsError(.bad_envelope_json); return self.closeWith(1007, "bad SubscriberSourcedMessage envelope"); - if (type_val != .string) return self.closeWith(1007, "bad SubscriberSourcedMessage envelope"); + }; + if (type_val != .string) { + self.noteOptionsError(.bad_envelope_json); + return self.closeWith(1007, "bad SubscriberSourcedMessage envelope"); + } if (!std.mem.eql(u8, type_val.string, "options_update")) { // quirk: unknown types are logged and ignored, never fatal log.debug("ignoring client message type: {s}", .{type_val.string}); return; } - const payload = root.object.get("payload") orelse + const payload = root.object.get("payload") orelse { + self.noteOptionsError(.bad_payload_json); + return self.closeWith(1007, "bad options_update payload"); + }; + if (payload != .object) { + self.noteOptionsError(.bad_payload_json); return self.closeWith(1007, "bad options_update payload"); - if (payload != .object) return self.closeWith(1007, "bad options_update payload"); + } var new_filter = filter_mod.Filter.init(self.hub.allocator); var ok = true; @@ -364,6 +453,7 @@ pub const Handler = struct { }; if (!ok) { new_filter.deinit(); + self.noteOptionsError(.invalid_options); return self.closeWith(1008, reason); } @@ -386,6 +476,17 @@ pub const Handler = struct { return self.closeWith(1011, "internal error"); }; } + _ = self.hub.stats.subscribe_options_updates_total.fetchAdd(1, .monotonic); + } + + fn noteOptionsError(self: *Handler, reason: metrics.OptionsErrorReason) void { + _ = self.hub.stats.subscribe_options_update_errors.getPtr(reason).fetchAdd(1, .monotonic); + } + + /// Parser-level failures happen before clientMessage. The websocket + /// dependency reports its actual size-limit error through this hook. + pub fn clientError(self: *Handler, err: anyerror) void { + if (err == error.TooLarge) self.noteOptionsError(.oversize); } fn closeWith(self: *Handler, code: u16, reason: []const u8) !void { @@ -395,6 +496,7 @@ pub const Handler = struct { pub fn clientClose(self: *Handler, _: []const u8) !void { if (self.subscriber) |sub| self.signalStop(sub); + _ = self.hub.stats.subscribe_clean_disconnects_total.fetchAdd(1, .monotonic); self.conn.close(.{}) catch {}; } @@ -706,6 +808,59 @@ fn queryParam(query: []const u8, name: []const u8) ?[]const u8 { return null; } +fn resolveCursor(hub: *Hub, v2: bool, raw_param: ?[]const u8) !CursorPlan { + const raw = raw_param orelse return .{}; + if (raw.len == 0) return .{}; + const parsed = std.fmt.parseInt(i64, raw, 10) catch return error.InvalidCursor; + if (parsed < 0) return error.InvalidCursor; + + // Both endpoints use the magnitude split. V2 differs only by rejecting + // a seq cursor below the retained floor; v1 preserves its clamp. + if (parsed < 1_000_000_000_000_000) { + var seq: u64 = @intCast(parsed); + var next_seq = hub.tail.publishedNextSeq(); + if (hub.archive) |archive| { + const durable = archive.committed_seq.load(.acquire); + if (durable != 0) next_seq = @max(next_seq, durable +| 1); + } + if (next_seq == 0 or seq >= next_seq) + return .{ .mode = .clamped }; + var clamped = false; + if (seq == 0) { + seq = 1; + clamped = true; + } + const gauges = hub.tail.gauges(); + if (gauges.entries != 0) if (hub.tail.indexForSeq(seq)) |idx| + return .{ .start_idx = idx, .mode = if (clamped) .clamped else .seq }; + if (hub.archive != null) + return .{ + .start_idx = gauges.base, + .cold_from_seq = seq, + .mode = if (clamped) .clamped else .seq, + }; + if (v2) return error.CursorTooOld; + return .{ .start_idx = gauges.base, .mode = .clamped }; + } + + const now_us = Io.Timestamp.now(hub.io, .real).toMicroseconds(); + if (parsed > now_us) return .{ .mode = .clamped }; + var mode: metrics.CursorMode = .time_us; + if (hub.archive == null) if (hub.tail.oldestTime()) |floor| { + if (parsed < floor) mode = .clamped; + }; + const start_idx = hub.tail.indexForCursor(parsed); + if (hub.archive != null) { + const floor = hub.tail.oldestTime(); + if (floor == null or parsed < floor.?) return .{ + .start_idx = hub.tail.gauges().base, + .cold_from_us = parsed, + .mode = mode, + }; + } + return .{ .start_idx = start_idx, .mode = mode }; +} + // === tests === test "queryParam" { @@ -727,3 +882,58 @@ test "permessage-deflate offer detection" { try std.testing.expect(offersPerMessageDeflate("foo, permessage-deflate; client_max_window_bits")); try std.testing.expect(!offersPerMessageDeflate("x-webkit-deflate-frame")); } + +test "slow detector matches upstream sustained-rate semantics" { + var stalled: SlowDetector = .{}; + var progressing: SlowDetector = .{}; + var caught_up: SlowDetector = .{}; + var stalled_pos: u64 = 0; + var progressing_pos: u64 = 0; + var dropped = false; + for (1..121) |second| { + const now_us: i64 = @intCast(second * std.time.us_per_s); + stalled_pos += 1; + progressing_pos += 100; + if (stalled.observe(now_us, stalled_pos, 1_000_000)) dropped = true; + try std.testing.expect(!progressing.observe(now_us, progressing_pos, 1_000_000)); + try std.testing.expect(!caught_up.observe(now_us, 0, 0)); + } + try std.testing.expect(dropped); +} + +test "slow detector recovery resets its full window" { + var detector: SlowDetector = .{}; + var pos: u64 = 0; + for (1..31) |second| { + pos += 1; + try std.testing.expect(!detector.observe(@intCast(second * std.time.us_per_s), pos, 1_000_000)); + } + pos += 1_000_000; + try std.testing.expect(!detector.observe(31 * std.time.us_per_s, pos, 100)); + for (32..62) |second| { + pos += 1; + try std.testing.expect(!detector.observe(@intCast(second * std.time.us_per_s), pos, 1_000_000)); + } +} + +test "cursor resolution is pre-upgrade and preserves v1-v2 floor semantics" { + var threaded: Io.Threaded = .init(std.testing.allocator, .{}); + defer threaded.deinit(); + var tail = tail_mod.Tail.init(std.testing.allocator, threaded.io(), 1 << 20); + defer tail.deinit(); + try tail.append(.commit, "did:plc:test", "app.bsky.feed.post", 100, 100, "{}", "{}", .none); + var stats: metrics.Stats = .{}; + var hub: Hub = .{ + .allocator = std.testing.allocator, + .io = threaded.io(), + .tail = &tail, + .stats = &stats, + }; + + try std.testing.expectEqual(metrics.CursorMode.live, (try resolveCursor(&hub, false, null)).mode); + try std.testing.expectEqual(metrics.CursorMode.clamped, (try resolveCursor(&hub, false, "50")).mode); + try std.testing.expectError(error.CursorTooOld, resolveCursor(&hub, true, "50")); + try std.testing.expectEqual(metrics.CursorMode.clamped, (try resolveCursor(&hub, true, "101")).mode); + try std.testing.expectError(error.InvalidCursor, resolveCursor(&hub, true, "nope")); + try std.testing.expectError(error.InvalidCursor, resolveCursor(&hub, true, "-1")); +} diff --git a/src/internal/tail.zig b/src/internal/tail.zig index 9c9c940..f0e87a3 100644 --- a/src/internal/tail.zig +++ b/src/internal/tail.zig @@ -14,6 +14,8 @@ const wire = @import("wire.zig"); const Io = std.Io; const Allocator = std.mem.Allocator; +pub const SkipV1 = enum { none, sync, resync_replacement }; + pub const Entry = struct { kind: wire.Kind, did: []u8, @@ -25,6 +27,7 @@ pub const Entry = struct { json: []u8, /// v2 superset wire json_v2: []u8, + skip_v1: SkipV1, fn deinit(self: Entry, allocator: Allocator) void { allocator.free(self.did); @@ -75,7 +78,7 @@ pub const Tail = struct { } /// append a frame; copies all slices. wakes blocked readers. - pub fn append(self: *Tail, kind: wire.Kind, did: []const u8, collection: []const u8, time_us: i64, seq: u64, json: []const u8, json_v2: []const u8) !void { + pub fn append(self: *Tail, kind: wire.Kind, did: []const u8, collection: []const u8, time_us: i64, seq: u64, json: []const u8, json_v2: []const u8, skip_v1: SkipV1) !void { const entry: Entry = .{ .kind = kind, .did = try self.allocator.dupe(u8, did), @@ -84,6 +87,7 @@ pub const Tail = struct { .seq = seq, .json = try self.allocator.dupe(u8, json), .json_v2 = try self.allocator.dupe(u8, json_v2), + .skip_v1 = skip_v1, }; self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); @@ -138,6 +142,20 @@ pub const Tail = struct { return self.base + visible; } + /// Next local archive seq after the visible hot prefix. Zero means no + /// seq-bearing row is visible yet, matching upstream's unwarmed writer + /// sentinel used by cursor resolution. + pub fn publishedNextSeq(self: *Tail) u64 { + self.mutex.lockUncancelable(self.io); + defer self.mutex.unlock(self.io); + var next_seq: u64 = 0; + for (self.entries.items) |entry| { + if (self.durability_gated and entry.seq != 0 and entry.seq > self.published_seq) break; + if (entry.seq != 0) next_seq = entry.seq +| 1; + } + return next_seq; + } + /// time of the oldest retained entry, or null when empty. pub fn oldestTime(self: *Tail) ?i64 { self.mutex.lockUncancelable(self.io); @@ -168,7 +186,7 @@ pub const Tail = struct { stop: *const std.atomic.Value(bool), v2: bool, ctx: anytype, - comptime pred: fn (@TypeOf(ctx), wire.Kind, []const u8, []const u8) bool, + comptime pred: fn (@TypeOf(ctx), wire.Kind, []const u8, []const u8, SkipV1) bool, ) !?[]u8 { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); @@ -186,9 +204,13 @@ pub const Tail = struct { continue; } idx.* += 1; + if (!v2 and e.skip_v1 != .none) { + _ = pred(ctx, e.kind, e.did, e.collection, e.skip_v1); + continue; + } const json = if (v2) e.json_v2 else e.json; if (json.len == 0) continue; // kind not on this wire (#sync on v1) - if (!pred(ctx, e.kind, e.did, e.collection)) continue; + if (!pred(ctx, e.kind, e.did, e.collection, .none)) continue; return try allocator.dupe(u8, json); } } @@ -237,7 +259,7 @@ pub const Tail = struct { const testing = std.testing; -fn acceptAll(_: void, _: wire.Kind, _: []const u8, _: []const u8) bool { +fn acceptAll(_: void, _: wire.Kind, _: []const u8, _: []const u8, _: SkipV1) bool { return true; } @@ -247,8 +269,8 @@ test "append, read, evict, cursor clamp" { var t = Tail.init(testing.allocator, threaded.io(), 300); defer t.deinit(); - try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 100, 1, "{\"n\":1}", "{\"v2\":1}"); - try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 200, 2, "{\"n\":2}", "{\"v2\":2}"); + try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 100, 1, "{\"n\":1}", "{\"v2\":1}", .none); + try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 200, 2, "{\"n\":2}", "{\"v2\":2}", .none); var idx: u64 = 0; var stop: std.atomic.Value(bool) = .init(false); @@ -257,8 +279,8 @@ test "append, read, evict, cursor clamp" { try testing.expectEqualStrings("{\"n\":1}", j1); // small budget: appending more evicts the oldest - try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 300, 3, "{\"n\":3}", "{\"v2\":3}"); - try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 400, 4, "{\"n\":4}", "{\"v2\":4}"); + try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 300, 3, "{\"n\":3}", "{\"v2\":3}", .none); + try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 400, 4, "{\"n\":4}", "{\"v2\":4}", .none); try testing.expect(t.base > 0); try testing.expectEqual(t.base, t.indexForCursor(0)); @@ -272,7 +294,7 @@ test "durability gate hides appended frames until publication" { defer t.deinit(); t.enableDurabilityGate(0); - try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 100, 1, "{\"n\":1}", "{\"v2\":1}"); + try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 100, 1, "{\"n\":1}", "{\"v2\":1}", .none); try testing.expectEqual(@as(u64, 1), t.tip()); try testing.expectEqual(@as(u64, 0), t.publishedTip()); try testing.expectEqual(@as(?u64, 0), t.indexForSeq(1)); @@ -286,6 +308,37 @@ test "durability gate hides appended frames until publication" { try testing.expectEqualStrings("{\"n\":1}", json); } +test "v1 observes skipped rows while v2 receives them" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + var t = Tail.init(testing.allocator, threaded.io(), 1 << 20); + defer t.deinit(); + try t.append(.sync, "did:plc:a", "", 100, 1, "", "{\"kind\":\"sync\"}", .sync); + try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 200, 2, "{\"kind\":\"commit\"}", "{\"kind\":\"commit\",\"seq\":2}", .none); + + const Observer = struct { + skipped: usize = 0, + fn wants(self: *@This(), _: wire.Kind, _: []const u8, _: []const u8, skip: SkipV1) bool { + if (skip != .none) self.skipped += 1; + return true; + } + }; + var stop: std.atomic.Value(bool) = .init(false); + var v1_observer: Observer = .{}; + var v1_idx: u64 = 0; + const v1 = (try t.next(testing.allocator, &v1_idx, &stop, false, &v1_observer, Observer.wants)).?; + defer testing.allocator.free(v1); + try testing.expectEqual(@as(usize, 1), v1_observer.skipped); + try testing.expectEqualStrings("{\"kind\":\"commit\"}", v1); + + var v2_observer: Observer = .{}; + var v2_idx: u64 = 0; + const v2 = (try t.next(testing.allocator, &v2_idx, &stop, true, &v2_observer, Observer.wants)).?; + defer testing.allocator.free(v2); + try testing.expectEqual(@as(usize, 0), v2_observer.skipped); + try testing.expectEqualStrings("{\"kind\":\"sync\"}", v2); +} + test "close unblocks readers" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); -- 2.51.2