From 68abcffd18def69d38e3d12b556be868d22e1893 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Thu, 13 Aug 2026 00:09:53 -0500 Subject: [PATCH] fix: start firehose subscriptions at live cursor --- docs/benchmarking.md | 2 +- src/atproto/sync.zig | 52 +++++++++++++++++++++++++++++++++------- src/storage/store.zig | 4 ++++ tools/websocket_churn.sh | 40 ++++++++++++++++++++++++++++++- 4 files changed, 88 insertions(+), 10 deletions(-) diff --git a/docs/benchmarking.md b/docs/benchmarking.md index 8cdae98..16a598c 100644 --- a/docs/benchmarking.md +++ b/docs/benchmarking.md @@ -140,7 +140,7 @@ first request. `just bench websocket-churn` covers the WebSocket lifecycle separately. It repeatedly opens and closes cohorts of `subscribeRepos` clients at the server's -connection limit, then checks that detached connection workers return to the +live-connection limit, then checks that detached connection workers return to the baseline thread count. A warm second round also guards against renewed, unbounded RSS growth while allowing the allocator to retain its first-round high-water capacity. diff --git a/src/atproto/sync.zig b/src/atproto/sync.zig index 1a15394..6af21eb 100644 --- a/src/atproto/sync.zig +++ b/src/atproto/sync.zig @@ -10,9 +10,11 @@ const httpz = @import("httpz"); const http = std.http; const max_subscribe_repos_connections = 32; +const max_subscribe_repos_backfills = 4; const notify_threshold_ms = 20 * 60 * 1000; var subscribe_repos_connections: usize = 0; +var subscribe_repos_backfills: usize = 0; var get_repo_connections: usize = 0; var crawler_last_notified_ms: std.atomic.Value(i64) = .init(0); var crawler_notify_in_flight: std.atomic.Value(bool) = .init(false); @@ -22,11 +24,12 @@ pub const SubscribeReposClient = struct { pub const Context = struct { cursor: u64, + backfill: bool, }; pub fn init(conn: *httpz.websocket.Conn, ctx: *const Context) !SubscribeReposClient { const state = try std.heap.smp_allocator.create(StreamState); - state.* = .{ .conn = conn, .cursor = ctx.cursor }; + state.* = .{ .conn = conn, .cursor = ctx.cursor, .backfill = ctx.backfill }; return .{ .state = state }; } @@ -52,6 +55,7 @@ pub const SubscribeReposClient = struct { const StreamState = struct { conn: *httpz.websocket.Conn, cursor: u64, + backfill: bool, closed: bool = false, counted: bool = true, refs: usize = 1, @@ -69,6 +73,9 @@ pub const SubscribeReposClient = struct { fn finish(self: *StreamState) void { if (@atomicRmw(bool, &self.counted, .Xchg, false, .acq_rel)) { _ = @atomicRmw(usize, &subscribe_repos_connections, .Sub, 1, .monotonic); + if (self.backfill) { + _ = @atomicRmw(usize, &subscribe_repos_backfills, .Sub, 1, .monotonic); + } } } @@ -292,30 +299,59 @@ pub fn listRepos(request: *http_api.Request) !void { } pub fn subscribeRepos(request: *http_api.Request) !void { + var cursor_buf: [32]u8 = undefined; + const cursor_param = http_api.queryParam(request.url.raw, "cursor", &cursor_buf); + const current_seq = store.latestSequence(); + const start = subscriptionStart(cursor_param, current_seq) catch + return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Malformed cursor"); + const cursor = start.cursor; + const backfill = start.backfill; + const active = @atomicRmw(usize, &subscribe_repos_connections, .Add, 1, .monotonic); if (active >= max_subscribe_repos_connections) { _ = @atomicRmw(usize, &subscribe_repos_connections, .Sub, 1, .monotonic); log.debug("sync subscribeRepos rejected too_many_connections active={d} max={d}\n", .{ active + 1, max_subscribe_repos_connections }); return http_api.xrpcError(request, .too_many_requests, "RateLimitExceeded", "too many subscribeRepos connections"); } + if (backfill) { + const active_backfills = @atomicRmw(usize, &subscribe_repos_backfills, .Add, 1, .monotonic); + if (active_backfills >= max_subscribe_repos_backfills) { + _ = @atomicRmw(usize, &subscribe_repos_backfills, .Sub, 1, .monotonic); + _ = @atomicRmw(usize, &subscribe_repos_connections, .Sub, 1, .monotonic); + log.debug("sync subscribeRepos rejected too_many_backfills active={d} max={d}\n", .{ active_backfills + 1, max_subscribe_repos_backfills }); + return http_api.xrpcError(request, .too_many_requests, "RateLimitExceeded", "too many subscribeRepos backfills"); + } + } var upgraded = false; defer if (!upgraded) { _ = @atomicRmw(usize, &subscribe_repos_connections, .Sub, 1, .monotonic); + if (backfill) _ = @atomicRmw(usize, &subscribe_repos_backfills, .Sub, 1, .monotonic); }; - var cursor_buf: [32]u8 = undefined; - const cursor: u64 = if (http_api.queryParam(request.url.raw, "cursor", &cursor_buf)) |raw| - std.fmt.parseInt(u64, raw, 10) catch 0 - else - 0; - - const ctx = SubscribeReposClient.Context{ .cursor = cursor }; + const ctx = SubscribeReposClient.Context{ .cursor = cursor, .backfill = backfill }; upgraded = try http_api.upgradeWebsocket(SubscribeReposClient, request, &ctx); if (!upgraded) { return http_api.xrpcError(request, .upgrade_required, "InvalidRequest", "Expected WebSocket upgrade"); } } +const SubscriptionStart = struct { + cursor: u64, + backfill: bool, +}; + +fn subscriptionStart(cursor_param: ?[]const u8, current_seq: u64) !SubscriptionStart { + const cursor = if (cursor_param) |raw| try std.fmt.parseInt(u64, raw, 10) else current_seq; + return .{ .cursor = cursor, .backfill = cursor < current_seq }; +} + +test "subscribeRepos starts live unless a historical cursor is explicit" { + try std.testing.expectEqual(SubscriptionStart{ .cursor = 42, .backfill = false }, try subscriptionStart(null, 42)); + try std.testing.expectEqual(SubscriptionStart{ .cursor = 0, .backfill = true }, try subscriptionStart("0", 42)); + try std.testing.expectEqual(SubscriptionStart{ .cursor = 42, .backfill = false }, try subscriptionStart("42", 42)); + try std.testing.expectError(error.InvalidCharacter, subscriptionStart("nope", 42)); +} + pub fn getRepoStatus(request: *http_api.Request) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); diff --git a/src/storage/store.zig b/src/storage/store.zig index 4fe75c5..84a26b2 100644 --- a/src/storage/store.zig +++ b/src/storage/store.zig @@ -604,6 +604,10 @@ pub fn currentIo() Io { return store_io; } +pub fn latestSequence() u64 { + return next_seq.load(.acquire) -| 1; +} + pub fn resolveRepo(repo: []const u8) ?auth.Account { return findAccount(std.heap.page_allocator, repo) catch null; } diff --git a/tools/websocket_churn.sh b/tools/websocket_churn.sh index be6301a..4117035 100755 --- a/tools/websocket_churn.sh +++ b/tools/websocket_churn.sh @@ -9,7 +9,7 @@ db="${TMPDIR:-/tmp}/zds-websocket-churn.sqlite3" blob_root="${TMPDIR:-/tmp}/zds-websocket-churn-blobs" log="${TMPDIR:-/tmp}/zds-websocket-churn.log" base="http://127.0.0.1:${port}" -socket="ws://127.0.0.1:${port}/xrpc/com.atproto.sync.subscribeRepos?cursor=0" +socket="ws://127.0.0.1:${port}/xrpc/com.atproto.sync.subscribeRepos" cleanup() { if [ -n "${server_pid:-}" ]; then @@ -58,6 +58,44 @@ while [ "$i" -lt 100 ]; do done curl -fsS "$base/xrpc/_health" >/dev/null +# Restart over a seeded event log so an omitted cursor proves it starts live +# instead of accidentally replaying retained history. +kill "$server_pid" +wait "$server_pid" 2>/dev/null || true +server_pid="" +sqlite3 "$db" <<'SQL' +INSERT INTO accounts (did, handle, email, password_hash, activated_at, email_confirmed_at) +VALUES ('did:plc:websocketchurn', 'websocket-churn.test', 'websocket-churn@test.com', 'password', unixepoch(), unixepoch()); +WITH RECURSIVE events(seq) AS ( + SELECT 1 + UNION ALL + SELECT seq + 1 FROM events WHERE seq < 1000 +) +INSERT INTO seq_events (seq, did, commit_cid, evt) +SELECT seq, 'did:plc:websocketchurn', 'bafyreiaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa', zeroblob(65536) +FROM events; +SQL +./zig-out/bin/zds \ + --host 127.0.0.1 \ + --port "$port" \ + --db "$db" \ + --blobstore-path "$blob_root" \ + --public-url "$base" \ + --server-did did:web:localhost \ + --handle-domains .test \ + >"$log" 2>&1 & +server_pid=$! + +i=0 +while [ "$i" -lt 100 ]; do + if curl -fsS "$base/xrpc/_health" >/dev/null 2>&1; then + break + fi + i=$((i + 1)) + sleep 0.1 +done +curl -fsS "$base/xrpc/_health" >/dev/null + baseline_threads=$(threads) baseline_rss=$(rss_kb) node tools/websocket_churn.mjs "$socket" "$rounds" "$concurrency" -- 2.51.2