From 55babfacea25734bdc09343e923f38aea9466f23 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Thu, 30 Jul 2026 16:17:21 -0500 Subject: [PATCH] Fix firehose event read contention --- bench/README.md | 5 ++ bench/justfile | 4 ++ bench/main.zig | 55 +++++++++++++++++++- src/storage/store.zig | 118 +++++++++++++++++++++++------------------- 4 files changed, 129 insertions(+), 53 deletions(-) diff --git a/bench/README.md b/bench/README.md index 2ff1c89..1491598 100644 --- a/bench/README.md +++ b/bench/README.md @@ -25,6 +25,7 @@ just bench get-record 10 1000 just bench list-records 10 1000 just bench repo-size 5000 just bench get-repo 5000 4 5 +just bench firehose 1000 32 100 just bench repo-large 100000 just bench space 1000 just bench http @@ -60,6 +61,10 @@ just bench run --scenario write --records 10000 additional commit, and concurrent full exports against one seeded repo. Its isolation lane runs account lookups and writes during those exports so a throughput improvement cannot hide resident-facing database stalls. +- `firehose`: runs the maximum 32 synthetic `subscribeRepos` pollers while + measuring resident account-read latency. Event-history reads must remain + ordinary indexed range reads; migration or repair work in this path causes + global database-lock contention that stalls login, writes, and blob reads. - `space`: experimental permissioned-data storage probes for space discovery, private record writes, private record reads/lists, repo oplog catch-up, and blob readback. diff --git a/bench/justfile b/bench/justfile index c2dfd93..dcdfd08 100644 --- a/bench/justfile +++ b/bench/justfile @@ -65,6 +65,10 @@ repo-size max_records="5000": get-repo records="5000" callers="4" ops_per_caller="5": {{zig}} build bench -Doptimize=ReleaseFast -- --scenario get-repo --records {{records}} --callers {{callers}} --ops-per-caller {{ops_per_caller}} +# benchmark resident-facing reads while firehose subscribers poll event history +firehose records="1000" subscribers="32" polls_per_subscriber="100": + {{zig}} build bench -Doptimize=ReleaseFast -- --scenario firehose --records {{records}} --callers {{subscribers}} --ops-per-caller {{polls_per_subscriber}} + # deliberately slow large-repo lane for Paul/Jerry-shaped synthetic repositories repo-large records="100000": {{zig}} build bench -Doptimize=ReleaseFast -- --scenario repo-size --records {{records}} diff --git a/bench/main.zig b/bench/main.zig index db58574..e78478b 100644 --- a/bench/main.zig +++ b/bench/main.zig @@ -16,6 +16,7 @@ const Scenario = enum { list_records, repo_size, get_repo, + firehose, space, write_profile, }; @@ -125,6 +126,7 @@ pub fn main(init: std.process.Init) !void { .list_records => try benchListRecords(allocator, options), .repo_size => try benchRepoSizes(allocator, options.records), .get_repo => try benchGetRepo(allocator, options), + .firehose => try benchFirehoseContention(allocator, options), .space => try benchSpace(allocator, options), .write_profile => try benchWriteProfile(allocator, options), } @@ -569,6 +571,56 @@ fn benchGetRepo(allocator: std.mem.Allocator, options: Options) !void { try benchRepoExportIsolation(allocator, state.account, options.records + 2, level); } +fn benchFirehoseContention(allocator: std.mem.Allocator, options: Options) !void { + const subscribers = options.callers orelse 32; + const polls_per_subscriber = options.ops_per_caller orelse 100; + var state = try initBench(allocator); + defer state.deinit(); + try seedIndexedRecords(allocator, state.account, options.records); + + const threads = try allocator.alloc(std.Thread, subscribers); + defer allocator.free(threads); + const probe_latencies = try allocator.alloc(u64, 500); + defer allocator.free(probe_latencies); + var start_gate: std.atomic.Value(bool) = .init(false); + + for (threads) |*thread| { + thread.* = try std.Thread.spawn(.{}, firehosePollWorker, .{ polls_per_subscriber, &start_gate }); + } + const probe_thread = try std.Thread.spawn(.{}, accountProbeWorker, .{ + state.account.did, + probe_latencies, + &start_gate, + }); + const start = nowNs(); + start_gate.store(true, .release); + for (threads) |thread| thread.join(); + probe_thread.join(); + + const elapsed = nowNs() - start; + concurrentResult( + "resident reads", + .{ .callers = 1, .ops_per_caller = probe_latencies.len }, + probe_latencies.len, + elapsed, + probe_latencies, + ).print(); + std.debug.print( + "firehose polling subscribers={d} polls={d} elapsed={d:.3}ms\n", + .{ subscribers, subscribers * polls_per_subscriber, nsToMs(elapsed) }, + ); +} + +fn firehosePollWorker(iterations: usize, start: *std.atomic.Value(bool)) void { + var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); + defer arena.deinit(); + while (!start.load(.acquire)) std.atomic.spinLoopHint(); + for (0..iterations) |_| { + _ = arena.reset(.retain_capacity); + _ = zds.storage.store.listSeqEvents(arena.allocator(), 0, 100) catch return; + } +} + fn benchRepoCarSince( allocator: std.mem.Allocator, account: zds.auth.tokens.Account, @@ -1287,6 +1339,7 @@ fn parseScenario(value: []const u8) !Scenario { if (std.mem.eql(u8, value, "list-records")) return .list_records; if (std.mem.eql(u8, value, "repo-size")) return .repo_size; if (std.mem.eql(u8, value, "get-repo")) return .get_repo; + if (std.mem.eql(u8, value, "firehose")) return .firehose; if (std.mem.eql(u8, value, "space")) return .space; if (std.mem.eql(u8, value, "write-profile")) return .write_profile; return error.UnknownScenario; @@ -1294,7 +1347,7 @@ fn parseScenario(value: []const u8) !Scenario { fn usage() void { std.debug.print( - \\usage: zds-bench [--scenario all|write|read|blob|repo|metastore|get-cid|get-block|decode-record|render-record|get-record|list-records|repo-size|get-repo|space|write-profile] [--records N] [--blobs N] [--blob-size BYTES] [--callers N --ops-per-caller N] + \\usage: zds-bench [--scenario all|write|read|blob|repo|metastore|get-cid|get-block|decode-record|render-record|get-record|list-records|repo-size|get-repo|firehose|space|write-profile] [--records N] [--blobs N] [--blob-size BYTES] [--callers N --ops-per-caller N] \\ , .{}); } diff --git a/src/storage/store.zig b/src/storage/store.zig index 826d6c4..b261219 100644 --- a/src/storage/store.zig +++ b/src/storage/store.zig @@ -3107,12 +3107,26 @@ pub fn writeLatestCommitJson(allocator: std.mem.Allocator, did: []const u8) ![]c } pub fn listSeqEvents(allocator: std.mem.Allocator, cursor: u64, limit: usize) ![]SeqEvent { - db_mutex.lockUncancelable(store_io); - defer db_mutex.unlock(store_io); - try requireInitialized(); - try backfillSeqEventsLocked(allocator); + const path = database_path orelse return Error.StoreNotInitialized; + if (std.mem.eql(u8, path, ":memory:")) { + db_mutex.lockUncancelable(store_io); + defer db_mutex.unlock(store_io); + try requireInitialized(); + return listSeqEventsFrom(conn, allocator, cursor, limit); + } - var rows = try conn.rows( + const read_conn = try openRepoReadConnection(path); + defer read_conn.close(); + return listSeqEventsFrom(read_conn, allocator, cursor, limit); +} + +fn listSeqEventsFrom( + read_conn: zqlite.Conn, + allocator: std.mem.Allocator, + cursor: u64, + limit: usize, +) ![]SeqEvent { + var rows = try read_conn.rows( \\SELECT seq, evt \\FROM seq_events \\WHERE seq > ? @@ -5560,53 +5574,6 @@ fn latestRootFrom( }; } -fn backfillSeqEventsLocked(allocator: std.mem.Allocator) !void { - var rows = try conn.rows( - \\SELECT c.seq, c.did, c.cid, c.rev, c.prev, p.rev - \\FROM commits c - \\LEFT JOIN commits p ON p.did = c.did AND p.cid = c.prev - \\LEFT JOIN seq_events e ON e.seq = c.seq - \\WHERE e.seq IS NULL - \\ORDER BY c.seq ASC - , .{}); - defer rows.deinit(); - - var missing: std.ArrayList(struct { - seq: u64, - did: []const u8, - cid: []const u8, - rev: []const u8, - prev: ?[]const u8, - since_rev: ?[]const u8, - }) = .empty; - while (rows.next()) |row| { - try missing.append(allocator, .{ - .seq = @intCast(row.int(0)), - .did = try allocator.dupe(u8, row.text(1)), - .cid = try allocator.dupe(u8, row.text(2)), - .rev = try allocator.dupe(u8, row.text(3)), - .prev = if (row.nullableText(4)) |prev| try allocator.dupe(u8, prev) else null, - .since_rev = if (row.nullableText(5)) |rev| try allocator.dupe(u8, rev) else null, - }); - } - if (rows.err) |err| return err; - - for (missing.items) |commit| { - { - var commit_arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); - defer commit_arena.deinit(); - const commit_allocator = commit_arena.allocator(); - const repo_car = writeRepoCarFromLocked(commit_allocator, commit.did, commit.cid) catch continue; - const frame = commitEventFrameFromCar(commit_allocator, commit.seq, commit.did, commit.cid, commit.rev, commit.since_rev, null, repo_car, &.{}) catch continue; - try conn.exec( - \\INSERT INTO seq_events (seq, did, commit_cid, evt) - \\VALUES (?, ?, ?, ?) - \\ON CONFLICT(seq) DO NOTHING - , .{ @as(i64, @intCast(commit.seq)), commit.did, commit.cid, zqlite.blob(frame) }); - } - } -} - fn rebuildSeqEventsSync11Locked(allocator: std.mem.Allocator) !void { var rows = try conn.rows( \\SELECT c.seq, c.did, c.cid, c.rev, c.prev, p.rev @@ -6702,6 +6669,53 @@ test "commit firehose since uses previous repo rev" { try std.testing.expect(found_second); } +test "listing sequenced events does not perform event repair" { + var arena = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena.deinit(); + const allocator = arena.allocator(); + + try init(std.Options.debug_io, ":memory:"); + defer close(); + + const account = try createAccount( + allocator, + "firehose-read.test", + "firehose-read@test.com", + "password", + "did:plc:firehoseread", + true, + ); + const value = try std.json.parseFromSlice( + std.json.Value, + allocator, + "{\"$type\":\"app.bsky.feed.post\",\"text\":\"sequenced at write time\"}", + .{}, + ); + defer value.deinit(); + _ = try applyWritesWithOptions(allocator, account, &.{.{ .create = .{ + .collection = "app.bsky.feed.post", + .rkey = "3jzfcijpj2z2e", + .value = value.value, + } }}, .{}); + + const before = try scalarCountLocked("SELECT COUNT(*) FROM seq_events WHERE did = ?", account.did); + try std.testing.expect(before > 0); + try conn.exec("DELETE FROM seq_events WHERE did = ?", .{account.did}); + + const events = try listSeqEvents(allocator, 0, 100); + for (events) |event| { + const decoded = zat.firehose.decodeFrame(allocator, event.frame) catch continue; + switch (decoded) { + .commit => |commit| try std.testing.expect(!std.mem.eql(u8, commit.repo, account.did)), + else => {}, + } + } + try std.testing.expectEqual( + @as(u64, 0), + try scalarCountLocked("SELECT COUNT(*) FROM seq_events WHERE did = ?", account.did), + ); +} + test "getRepo since filters repo blocks by revision" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); -- 2.51.2