const std = @import("std"); const zds = @import("zds"); const Scenario = enum { all, write, read, blob, blob_gc, repo, metastore, get_cid, get_block, decode_record, render_record, get_record, list_records, repo_size, get_repo, firehose, publish_sync, space, write_profile, }; const bench_collection = "dev.zds.bench.record"; const space_collection = "fm.plyr.stg.track"; const space_type = "fm.plyr.stg.privateMedia"; const valid_bench_did = "did:plc:zdsbenchzdsbenchzdsbenchzd"; const Options = struct { scenario: Scenario = .all, records: usize = 1000, blobs: usize = 128, blob_size: usize = 64 * 1024, callers: ?usize = null, ops_per_caller: ?usize = null, }; const BenchResult = struct { name: []const u8, ops: usize, bytes: usize = 0, elapsed_ns: u64, fn print(self: BenchResult) void { const elapsed_s = @as(f64, @floatFromInt(self.elapsed_ns)) / std.time.ns_per_s; const ops_per_s = @as(f64, @floatFromInt(self.ops)) / elapsed_s; if (self.bytes == 0) { std.debug.print("{s: <20} {d: >8} ops {d: >10.1} ops/s {d: >9.3} ms\n", .{ self.name, self.ops, ops_per_s, elapsed_s * 1000.0, }); return; } const mb = @as(f64, @floatFromInt(self.bytes)) / (1024.0 * 1024.0); const mb_per_s = mb / elapsed_s; std.debug.print("{s: <20} {d: >8} ops {d: >10.1} ops/s {d: >8.1} MB/s {d: >9.3} ms\n", .{ self.name, self.ops, ops_per_s, mb_per_s, elapsed_s * 1000.0, }); } }; const BenchState = struct { arena: std.heap.ArenaAllocator, account: zds.auth.tokens.Account, db_path: []const u8, blob_path: []const u8, fn deinit(self: *BenchState) void { zds.storage.store.close(); cleanupPath(self.db_path); cleanupPath(self.blob_path); self.arena.deinit(); } }; pub fn main(init: std.process.Init) !void { const allocator = init.gpa; const options = try parseOptions(init); std.debug.print("\n=== zds store benchmarks ===\n", .{}); std.debug.print("records={d} blobs={d} blob_size={d}\n\n", .{ options.records, options.blobs, options.blob_size }); switch (options.scenario) { .all => { var state = try initBench(allocator); defer state.deinit(); (try benchWrite(allocator, state.account, options.records)).print(); (try benchRead(allocator, state.account, options.records)).print(); (try benchRepoCar(allocator, state.account, options.records)).print(); (try benchBlob(allocator, state.account, options.blobs, options.blob_size)).print(); }, .write => { var state = try initBench(allocator); defer state.deinit(); (try benchWrite(allocator, state.account, options.records)).print(); }, .read => { var state = try initBench(allocator); defer state.deinit(); try seedRecords(allocator, state.account, options.records); (try benchRead(allocator, state.account, options.records)).print(); }, .repo => { var state = try initBench(allocator); defer state.deinit(); try seedRecords(allocator, state.account, options.records); (try benchRepoCar(allocator, state.account, options.records)).print(); }, .blob => { var state = try initBench(allocator); defer state.deinit(); (try benchBlob(allocator, state.account, options.blobs, options.blob_size)).print(); }, .blob_gc => { var state = try initBench(allocator); defer state.deinit(); (try benchBlobGc(allocator, state.account, options.blobs, options.blob_size)).print(); }, .metastore => try benchMetastore(allocator, options), .get_cid => try benchGetCid(allocator, options), .get_block => try benchGetBlock(allocator, options), .decode_record => try benchDecodeRecord(allocator, options), .render_record => try benchRenderRecord(allocator, options), .get_record => try benchGetRecord(allocator, options), .list_records => try benchListRecords(allocator, options), .repo_size => try benchRepoSizes(allocator, options.records), .get_repo => try benchGetRepo(allocator, options), .firehose => try benchFirehoseContention(allocator, options), .publish_sync => try benchPublishSync(allocator, options), .space => try benchSpace(allocator, options), .write_profile => try benchWriteProfile(allocator, options), } } fn initBench(allocator: std.mem.Allocator) !BenchState { var arena = std.heap.ArenaAllocator.init(allocator); errdefer arena.deinit(); const a = arena.allocator(); const suffix = nowNs(); const db_path = try std.fmt.allocPrint(a, "{s}/zds-bench-{d}.sqlite3", .{ tmpDir(), suffix }); const blob_path = try std.fmt.allocPrint(a, "{s}/zds-bench-blobs-{d}", .{ tmpDir(), suffix }); zds.storage.blobstore.init(std.Options.debug_io, blob_path); try zds.storage.store.init(std.Options.debug_io, db_path); const account = try zds.storage.store.createAccount( a, "bench.test", "bench@test.com", "password", valid_bench_did, true, ); return .{ .arena = arena, .account = account, .db_path = db_path, .blob_path = blob_path, }; } fn benchWrite(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, records: usize) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const start = nowNs(); for (0..records) |i| { _ = arena.reset(.retain_capacity); const value = try postValue(arena.allocator(), i); const result = try zds.storage.store.applyWrites(arena.allocator(), account, &.{.{ .create = .{ .collection = "app.bsky.feed.post", .rkey = null, .value = value, } }}); if (result.records.len != 1) return error.UnexpectedWriteResult; } return .{ .name = "write records", .ops = records, .elapsed_ns = nowNs() - start }; } fn benchRead(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, records: usize) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 500; const limit = @min(records, 100); const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const found = try zds.storage.store.listRecords(arena.allocator(), account.did, "app.bsky.feed.post", limit); if (records > 0 and found.len == 0) return error.MissingRecords; } return .{ .name = "list records", .ops = iterations, .elapsed_ns = nowNs() - start }; } fn benchRepoCar(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, records: usize) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 20; var total_bytes: usize = 0; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const car = try zds.storage.store.writeRepoCar(arena.allocator(), account.did); if (records > 0 and car.len == 0) return error.EmptyRepoCar; total_bytes += car.len; } return .{ .name = "write repo car", .ops = iterations, .bytes = total_bytes, .elapsed_ns = nowNs() - start }; } fn benchBlob(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, blobs: usize, blob_size: usize) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const a = arena.allocator(); const payload = try a.alloc(u8, blob_size); for (payload, 0..) |*byte, i| byte.* = @truncate(i); const start = nowNs(); for (0..blobs) |i| { payload[0] = @truncate(i); const cid = try zds.storage.store.putBlob(a, std.Options.debug_io, account, payload, "application/octet-stream"); const blob = zds.storage.store.getBlob(a, account.did, cid) orelse return error.MissingBlob; if (blob.data.len != payload.len) return error.BlobSizeMismatch; } return .{ .name = "blob put+get", .ops = blobs, .bytes = blobs * blob_size * 2, .elapsed_ns = nowNs() - start }; } fn benchBlobGc(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, blobs: usize, blob_size: usize) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const a = arena.allocator(); const payload = try a.alloc(u8, @max(blob_size, @sizeOf(u64))); @memset(payload, 0xa5); for (0..blobs) |i| { std.mem.writeInt(u64, payload[0..@sizeOf(u64)], i, .little); _ = try zds.storage.store.putBlob(a, std.Options.debug_io, account, payload, "application/octet-stream"); } const start = nowNs(); const result = try zds.storage.store.collectBlobGarbage(a, 0, blobs); if (result.deleted != blobs or result.failed != 0) return error.UnexpectedBlobGcResult; return .{ .name = "blob gc", .ops = result.deleted, .bytes = @intCast(result.bytes), .elapsed_ns = nowNs() - start }; } const ConcurrentStats = struct { p50: u64, p95: u64, p99: u64, max: u64, mean: u64, }; const ConcurrentResult = struct { name: []const u8, callers: usize, ops: usize, elapsed_ns: u64, stats: ConcurrentStats, fn print(self: ConcurrentResult) void { const elapsed_s = @as(f64, @floatFromInt(self.elapsed_ns)) / std.time.ns_per_s; const ops_per_s = @as(f64, @floatFromInt(self.ops)) / elapsed_s; std.debug.print( "{s: <16} callers={d: >4} ops={d: >6} {d: >10.1} ops/s p50={d:.3}ms p95={d:.3}ms p99={d:.3}ms max={d:.3}ms mean={d:.3}ms\n", .{ self.name, self.callers, self.ops, ops_per_s, nsToMs(self.stats.p50), nsToMs(self.stats.p95), nsToMs(self.stats.p99), nsToMs(self.stats.max), nsToMs(self.stats.mean), }, ); } }; const ConcurrencyLevel = struct { callers: usize, ops_per_caller: usize, }; fn benchMetastore(allocator: std.mem.Allocator, options: Options) !void { const default_levels = [_]ConcurrencyLevel{ .{ .callers = 1, .ops_per_caller = 5000 }, .{ .callers = 10, .ops_per_caller = 1000 }, .{ .callers = 100, .ops_per_caller = 200 }, .{ .callers = 1000, .ops_per_caller = 50 }, }; std.debug.print("\n=== zds metastore-shaped benchmarks ===\n", .{}); if (options.callers) |callers| { const level: ConcurrencyLevel = .{ .callers = callers, .ops_per_caller = options.ops_per_caller orelse return error.MissingOpsPerCaller, }; (try benchConcurrentWrites(allocator, level)).print(); (try benchConcurrentGets(allocator, level)).print(); (try benchConcurrentLists(allocator, level)).print(); return; } for (default_levels) |level| { (try benchConcurrentWrites(allocator, level)).print(); (try benchConcurrentGets(allocator, level)).print(); (try benchConcurrentLists(allocator, level)).print(); } } fn benchConcurrentWrites(allocator: std.mem.Allocator, level: ConcurrencyLevel) !ConcurrentResult { var state = try initBench(allocator); defer state.deinit(); var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const a = arena.allocator(); const accounts = try a.alloc(zds.auth.tokens.Account, level.callers); for (accounts, 0..) |*account, i| { account.* = try zds.storage.store.createAccount( a, try std.fmt.allocPrint(a, "bench-{d}.test", .{i}), try std.fmt.allocPrint(a, "bench-{d}@test.com", .{i}), "password", try std.fmt.allocPrint(a, "did:plc:zdsbench{d}", .{i}), true, ); } const total_ops = level.callers * level.ops_per_caller; const latencies = try allocator.alloc(u64, total_ops); defer allocator.free(latencies); const threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(threads); const start = nowNs(); for (threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, writeWorker, .{ accounts[i], i * level.ops_per_caller, latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], }); } for (threads) |thread| thread.join(); return concurrentResult("applyWrites", level, total_ops, nowNs() - start, latencies); } fn benchWriteProfile(allocator: std.mem.Allocator, options: Options) !void { const level: ConcurrencyLevel = .{ .callers = options.callers orelse 10, .ops_per_caller = options.ops_per_caller orelse 100, }; var state = try initBench(allocator); defer state.deinit(); var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const a = arena.allocator(); const accounts = try a.alloc(zds.auth.tokens.Account, level.callers); for (accounts, 0..) |*account, i| { account.* = try zds.storage.store.createAccount( a, try std.fmt.allocPrint(a, "profile-{d}.test", .{i}), try std.fmt.allocPrint(a, "profile-{d}@test.com", .{i}), "password", try std.fmt.allocPrint(a, "did:plc:zdsprofile{d}", .{i}), true, ); } const total_ops = level.callers * level.ops_per_caller; const latencies = try allocator.alloc(u64, total_ops); defer allocator.free(latencies); const profiles = try allocator.alloc(zds.storage.store.WriteProfile, total_ops); defer allocator.free(profiles); const threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(threads); const start = nowNs(); for (threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, profileWorker, .{ accounts[i], i * level.ops_per_caller, latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], profiles[i * level.ops_per_caller ..][0..level.ops_per_caller], }); } for (threads) |thread| thread.join(); const result = concurrentResult("applyWrites profile", level, total_ops, nowNs() - start, latencies); result.print(); printWriteProfile(profiles); } fn benchConcurrentGets(allocator: std.mem.Allocator, level: ConcurrencyLevel) !ConcurrentResult { var state = try initBench(allocator); defer state.deinit(); try seedIndexedRecords(allocator, state.account, 1000); const total_ops = level.callers * level.ops_per_caller; const latencies = try allocator.alloc(u64, total_ops); defer allocator.free(latencies); const threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(threads); const start = nowNs(); for (threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, getWorker, .{ state.account.did, i, latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], }); } for (threads) |thread| thread.join(); return concurrentResult("getRecord", level, total_ops, nowNs() - start, latencies); } fn benchGetCid(allocator: std.mem.Allocator, options: Options) !void { const level: ConcurrencyLevel = .{ .callers = options.callers orelse 10, .ops_per_caller = options.ops_per_caller orelse 1000, }; var state = try initBench(allocator); defer state.deinit(); try seedIndexedRecords(allocator, state.account, 1000); const total_ops = level.callers * level.ops_per_caller; const latencies = try allocator.alloc(u64, total_ops); defer allocator.free(latencies); const threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(threads); const start = nowNs(); for (threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, getCidWorker, .{ state.account.did, i, latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], }); } for (threads) |thread| thread.join(); concurrentResult("getRecordCid", level, total_ops, nowNs() - start, latencies).print(); } fn benchGetBlock(allocator: std.mem.Allocator, options: Options) !void { const level: ConcurrencyLevel = .{ .callers = options.callers orelse 10, .ops_per_caller = options.ops_per_caller orelse 1000, }; var state = try initBench(allocator); defer state.deinit(); try seedIndexedRecords(allocator, state.account, 1000); const total_ops = level.callers * level.ops_per_caller; const latencies = try allocator.alloc(u64, total_ops); defer allocator.free(latencies); const threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(threads); const start = nowNs(); for (threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, getBlockWorker, .{ state.account.did, i, latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], }); } for (threads) |thread| thread.join(); concurrentResult("getRecordBlock", level, total_ops, nowNs() - start, latencies).print(); } fn benchGetRecord(allocator: std.mem.Allocator, options: Options) !void { const level: ConcurrencyLevel = .{ .callers = options.callers orelse 10, .ops_per_caller = options.ops_per_caller orelse 1000, }; const result = try benchConcurrentGets(allocator, level); result.print(); } fn benchListRecords(allocator: std.mem.Allocator, options: Options) !void { const level: ConcurrencyLevel = .{ .callers = options.callers orelse 10, .ops_per_caller = options.ops_per_caller orelse 1000, }; const result = try benchConcurrentLists(allocator, level); result.print(); } fn benchRepoSizes(allocator: std.mem.Allocator, max_records: usize) !void { const tiers = [_]struct { name: []const u8, records: usize, }{ .{ .name = "repo small", .records = 100 }, .{ .name = "repo medium", .records = @min(1000, max_records) }, .{ .name = "repo large", .records = max_records }, }; std.debug.print("\n=== zds repo-size benchmarks ===\n", .{}); var last_records: usize = 0; for (tiers) |tier| { if (tier.records == 0 or tier.records == last_records) continue; last_records = tier.records; std.debug.print("seeding {s} with {d} records\n", .{ tier.name, tier.records }); var state = try initBench(allocator); errdefer state.deinit(); try seedIndexedRecords(allocator, state.account, tier.records); (try benchListRecordsForRepo(allocator, state.account, tier.records)).print(); (try benchRepoCarForRepo(allocator, state.account, tier.records)).print(); state.deinit(); } } fn benchListRecordsForRepo( allocator: std.mem.Allocator, account: zds.auth.tokens.Account, records: usize, ) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = if (records >= 100_000) 50 else 200; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const found = try zds.storage.store.listRecords(arena.allocator(), account.did, bench_collection, 100); if (found.len == 0) return error.MissingRecords; } return .{ .name = "repo tier list", .ops = iterations, .elapsed_ns = nowNs() - start, }; } fn benchRepoCarForRepo( allocator: std.mem.Allocator, account: zds.auth.tokens.Account, records: usize, ) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = if (records >= 100_000) 2 else 5; var total_bytes: usize = 0; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const car = try zds.storage.store.writeRepoCar(arena.allocator(), account.did); if (car.len == 0) return error.EmptyRepoCar; total_bytes += car.len; } return .{ .name = "repo tier CAR", .ops = iterations, .bytes = total_bytes, .elapsed_ns = nowNs() - start, }; } fn benchGetRepo(allocator: std.mem.Allocator, options: Options) !void { const level: ConcurrencyLevel = .{ .callers = options.callers orelse 4, .ops_per_caller = options.ops_per_caller orelse if (options.records >= 100_000) 1 else 5, }; std.debug.print("\n=== zds getRepo benchmarks ===\n", .{}); std.debug.print("seeding repo with {d} records\n", .{options.records}); var state = try initBench(allocator); defer state.deinit(); try seedIndexedRecords(allocator, state.account, options.records); const since_rev = try appendIndexedRecordRev(allocator, state.account, options.records); defer allocator.free(since_rev); const latest_rev = try appendIndexedRecordRev(allocator, state.account, options.records + 1); defer allocator.free(latest_rev); (try benchRepoCarForRepo(allocator, state.account, options.records + 2)).print(); (try benchRepoCarSince(allocator, state.account, since_rev)).print(); (try benchConcurrentRepoCar(allocator, state.account, level)).print(); 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 benchPublishSync(allocator: std.mem.Allocator, options: Options) !void { const level: ConcurrencyLevel = .{ .callers = options.callers orelse 3, .ops_per_caller = options.ops_per_caller orelse 3, }; var state = try initBench(allocator); defer state.deinit(); try seedIndexedRecords(allocator, state.account, @max(options.records, 1)); var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const a = arena.allocator(); const blob_size = @max(options.blob_size, 1); const seed_blob = try a.alloc(u8, blob_size); @memset(seed_blob, 0x5a); const seed_cid = try zds.storage.store.putBlob( a, std.Options.debug_io, state.account, seed_blob, "audio/mpeg", ); const total_ops = level.callers * level.ops_per_caller; const publish_latencies = try allocator.alloc(u64, total_ops); defer allocator.free(publish_latencies); const sync_latencies = try allocator.alloc(u64, total_ops); defer allocator.free(sync_latencies); const probe_latencies = try allocator.alloc(u64, total_ops); defer allocator.free(probe_latencies); const threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(threads); var start: std.atomic.Value(bool) = .init(false); for (threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, publishWorker, .{ state.account, blob_size, i * level.ops_per_caller, publish_latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], &start, }); } const sync_thread = try std.Thread.spawn(.{}, publishSyncWorker, .{ state.account.did, seed_cid, sync_latencies, &start, }); const probe_thread = try std.Thread.spawn(.{}, publishProbeWorker, .{ state.account.did, probe_latencies, &start, }); const bench_start = nowNs(); start.store(true, .release); for (threads) |thread| thread.join(); sync_thread.join(); probe_thread.join(); const elapsed = nowNs() - bench_start; std.debug.print("\n=== mixed publisher/sync benchmark ===\n", .{}); concurrentResult("blob+repo publish", level, total_ops, elapsed, publish_latencies).print(); concurrentResult("blob+repo sync", .{ .callers = 1, .ops_per_caller = total_ops }, total_ops, elapsed, sync_latencies).print(); concurrentResult("metadata probe", .{ .callers = 1, .ops_per_caller = total_ops }, total_ops, elapsed, probe_latencies).print(); } fn benchRepoCarSince( allocator: std.mem.Allocator, account: zds.auth.tokens.Account, since_rev: []const u8, ) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 20; var total_bytes: usize = 0; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const car = try zds.storage.store.writeRepoCarSince(arena.allocator(), account.did, since_rev); if (car.len == 0) return error.EmptyRepoCar; total_bytes += car.len; } return .{ .name = "getRepo since", .ops = iterations, .bytes = total_bytes, .elapsed_ns = nowNs() - start, }; } fn benchConcurrentRepoCar( allocator: std.mem.Allocator, account: zds.auth.tokens.Account, level: ConcurrencyLevel, ) !ConcurrentResult { const total_ops = level.callers * level.ops_per_caller; const latencies = try allocator.alloc(u64, total_ops); defer allocator.free(latencies); const threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(threads); const start = nowNs(); for (threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, repoCarWorker, .{ account.did, latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], }); } for (threads) |thread| thread.join(); return concurrentResult("getRepo full", level, total_ops, nowNs() - start, latencies); } fn benchRepoExportIsolation( allocator: std.mem.Allocator, account: zds.auth.tokens.Account, first_write_index: usize, level: ConcurrencyLevel, ) !void { const export_ops = level.callers * level.ops_per_caller; const probe_ops = @max(export_ops * 10, 20); const write_ops = @max(export_ops, 4); const export_latencies = try allocator.alloc(u64, export_ops); defer allocator.free(export_latencies); const probe_latencies = try allocator.alloc(u64, probe_ops); defer allocator.free(probe_latencies); const write_latencies = try allocator.alloc(u64, write_ops); defer allocator.free(write_latencies); const export_threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(export_threads); var start: std.atomic.Value(bool) = .init(false); for (export_threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, repoCarWorkerStarting, .{ account.did, export_latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], &start, }); } const probe_thread = try std.Thread.spawn(.{}, accountProbeWorker, .{ account.did, probe_latencies, &start, }); const write_thread = try std.Thread.spawn(.{}, repoWriteWorker, .{ account, first_write_index, write_latencies, &start, }); const started = nowNs(); start.store(true, .release); for (export_threads) |thread| thread.join(); probe_thread.join(); write_thread.join(); const elapsed = nowNs() - started; std.debug.print("\n=== getRepo isolation benchmark ===\n", .{}); concurrentResult("full exports", level, export_ops, elapsed, export_latencies).print(); concurrentResult("account probes", .{ .callers = 1, .ops_per_caller = probe_ops }, probe_ops, elapsed, probe_latencies).print(); concurrentResult("concurrent writes", .{ .callers = 1, .ops_per_caller = write_ops }, write_ops, elapsed, write_latencies).print(); } fn benchSpace(allocator: std.mem.Allocator, options: Options) !void { var state = try initBench(allocator); defer state.deinit(); var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const a = arena.allocator(); const records = options.records; const space = try seedSpaceRecords(a, state.account, records); const cid = try zds.storage.store.putBlob(a, std.Options.debug_io, state.account, "hello permissioned audio", "audio/mpeg"); const blob_record_json = try std.fmt.allocPrint( a, "{{\"$type\":\"{s}\",\"title\":\"permissioned audio\",\"audio\":{{\"$type\":\"blob\",\"ref\":{{\"$link\":\"{s}\"}},\"mimeType\":\"audio/mpeg\",\"size\":26}}}}", .{ space_collection, cid }, ); const blob_record = try std.json.parseFromSlice(std.json.Value, a, blob_record_json, .{}); const prepared_blob_record = try zds.storage.store.prepareRecordValue(a, space_collection, "blob-track", blob_record.value); _ = try zds.storage.store.createSpaceRecord(a, space.uri, state.account.did, space_collection, "blob-track", prepared_blob_record); std.debug.print("\n=== zds permissioned-space benchmarks ===\n", .{}); (try benchSpaceListSpaces(allocator, state.account, space.uri)).print(); (try benchSimpleSpaceListMembers(allocator, space.uri)).print(); (try benchSpaceRepoCar(allocator, state.account, space.uri)).print(); (try benchSpaceRepoVerify(allocator, state.account, space.uri)).print(); (try benchSpaceWrite(allocator, state.account, space.uri, records)).print(); (try benchSpaceGetRecord(allocator, state.account, space.uri, records)).print(); (try benchSpaceListRecords(allocator, state.account, space.uri, records)).print(); const repo_state = try zds.storage.store.getSpaceRepoState(a, space.uri, state.account.did); const set_hash = try zds.internal.permissioned_data.LtHash.fromBytes(repo_state.set_hash.?); _ = try zds.storage.store.recordSpaceWriter(a, space.uri, state.account.did, repo_state.rev.?, &set_hash.digest()); (try benchSpaceListRepos(allocator, space.uri)).print(); (try benchSpaceOplog(allocator, state.account, space.uri)).print(); (try benchSpaceListBlobs(allocator, state.account, space.uri)).print(); (try benchSpaceBlob(allocator, state.account, space.uri, cid)).print(); } fn benchSpaceRepoCar(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 20; var total_bytes: usize = 0; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const a = arena.allocator(); const state = try zds.storage.store.getSpaceRepoState(a, space, account.did); const set_hash = state.set_hash orelse return error.MissingRecords; const rev = state.rev orelse return error.MissingRecords; var keypair = try zds.storage.store.signingKeypair(account.did); const commit = try zds.internal.permissioned_data.createCommit(a, std.Options.debug_io, set_hash, .{ .space = space, .author = account.did, .rev = rev, }, &keypair); const records = try zds.storage.store.loadSpaceRepoBlocks(a, space, account.did); const car = try zds.internal.permissioned_data.serializeRepoCar(a, commit, records); total_bytes += car.len; } return .{ .name = "space getRepo", .ops = iterations, .bytes = total_bytes, .elapsed_ns = nowNs() - start }; } fn benchSpaceRepoVerify(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8) !BenchResult { var fixture_arena = std.heap.ArenaAllocator.init(allocator); defer fixture_arena.deinit(); const a = fixture_arena.allocator(); const state = try zds.storage.store.getSpaceRepoState(a, space, account.did); var keypair = try zds.storage.store.signingKeypair(account.did); const commit = try zds.internal.permissioned_data.createCommit(a, std.Options.debug_io, state.set_hash orelse return error.MissingRecords, .{ .space = space, .author = account.did, .rev = state.rev orelse return error.MissingRecords, }, &keypair); const records = try zds.storage.store.loadSpaceRepoBlocks(a, space, account.did); const car = try zds.internal.permissioned_data.serializeRepoCar(a, commit, records); const did_key = try keypair.did(a); const iterations: usize = 20; const start = nowNs(); for (0..iterations) |_| { var verified = try zds.internal.permissioned_data.verifyRepoCarFull( allocator, car, .{ .space = space, .author = account.did }, did_key, ); verified.deinit(); } return .{ .name = "space verifyRepo", .ops = iterations, .bytes = car.len * iterations, .elapsed_ns = nowNs() - start }; } fn seedSpaceRecords( allocator: std.mem.Allocator, account: zds.auth.tokens.Account, records: usize, ) !zds.storage.store.SpaceConfig { const space = try zds.storage.store.createSpace(allocator, .{ .actor_did = account.did, .authority_did = account.did, .space_type = space_type, .skey = "self", .is_authority = true, .read_managing_app = "https://api-stg.plyr.fm", .read_policy = "member-list", .write_managing_app = "https://api-stg.plyr.fm", .write_policy = "member-list", .app_access_json = "{\"type\":\"open\"}", }); for (0..records) |i| { const value = try spaceRecordValue(allocator, i); const prepared = try zds.storage.store.prepareRecordValue( allocator, space_collection, try std.fmt.allocPrint(allocator, "track{d:0>8}", .{i}), value, ); _ = try zds.storage.store.createSpaceRecord( allocator, space.uri, account.did, space_collection, try std.fmt.allocPrint(allocator, "track{d:0>8}", .{i}), prepared, ); } return space; } fn benchSpaceListSpaces(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 500; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const spaces = try zds.storage.store.listSpaces(arena.allocator(), account.did, account.did, space_type, null, 50); if (spaces.len == 0 or !std.mem.eql(u8, spaces[0].uri, space)) return error.MissingSpace; } return .{ .name = "space listSpaces", .ops = iterations, .elapsed_ns = nowNs() - start }; } fn benchSpaceWrite(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8, offset: usize) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 100; const start = nowNs(); for (0..iterations) |i| { _ = arena.reset(.retain_capacity); const idx = offset + i; const rkey = try std.fmt.allocPrint(arena.allocator(), "track{d:0>8}", .{idx}); const value = try spaceRecordValue(arena.allocator(), idx); const prepared = try zds.storage.store.prepareRecordValue(arena.allocator(), space_collection, rkey, value); const result = try zds.storage.store.createSpaceRecord(arena.allocator(), space, account.did, space_collection, rkey, prepared); if (result.cid.len == 0) return error.UnexpectedWriteResult; } return .{ .name = "space createRecord", .ops = iterations, .elapsed_ns = nowNs() - start }; } fn benchSimpleSpaceListMembers(allocator: std.mem.Allocator, space: []const u8) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 500; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const members = try zds.storage.store.listSimpleSpaceMembers(arena.allocator(), space, null, 50); if (members.len == 0) return error.MissingRecords; } return .{ .name = "simplespace listMembers", .ops = iterations, .elapsed_ns = nowNs() - start }; } fn benchSpaceGetRecord(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8, records: usize) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 1000; const start = nowNs(); for (0..iterations) |i| { _ = arena.reset(.retain_capacity); const idx = if (records == 0) 0 else (i * 13) % records; const rkey = try std.fmt.allocPrint(arena.allocator(), "track{d:0>8}", .{idx}); const record = (try zds.storage.store.getSpaceRecord(arena.allocator(), space, account.did, space_collection, rkey)) orelse return error.MissingRecord; if (record.value_json.len == 0) return error.MissingRecord; } return .{ .name = "space getRecord", .ops = iterations, .elapsed_ns = nowNs() - start }; } fn benchSpaceListRecords(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8, records: usize) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 500; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const found = try zds.storage.store.listSpaceRecords(arena.allocator(), space, account.did, space_collection, null, false, 50, true); if (records > 0 and found.len == 0) return error.MissingRecords; if (records > 0 and found[0].value_json == null) return error.MissingRecord; } return .{ .name = "space listRecords", .ops = iterations, .elapsed_ns = nowNs() - start }; } fn benchSpaceListRepos(allocator: std.mem.Allocator, space: []const u8) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 500; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const repos = try zds.storage.store.listSpaceWriters(arena.allocator(), space, null, 50); if (repos.len == 0) return error.MissingRecords; } return .{ .name = "space listRepos", .ops = iterations, .elapsed_ns = nowNs() - start }; } fn benchSpaceOplog(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 500; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const ops = try zds.storage.store.listSpaceRecordOplog(arena.allocator(), space, account.did, null, null, 100, true); if (ops.len == 0) return error.MissingRecords; } return .{ .name = "space listRepoOps", .ops = iterations, .elapsed_ns = nowNs() - start }; } fn benchSpaceBlob(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8, cid: []const u8) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 1000; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); if (!try zds.storage.store.permissionedBlobReferenced(space, account.did, cid)) return error.MissingBlob; const blob = zds.storage.store.getBlob(arena.allocator(), account.did, cid) orelse return error.MissingBlob; if (blob.data.len == 0) return error.MissingBlob; } return .{ .name = "space getBlob", .ops = iterations, .bytes = iterations * "hello permissioned audio".len, .elapsed_ns = nowNs() - start }; } fn benchSpaceListBlobs(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8) !BenchResult { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const iterations: usize = 1000; const start = nowNs(); for (0..iterations) |_| { _ = arena.reset(.retain_capacity); const cids = try zds.storage.store.listPermissionedBlobs(arena.allocator(), space, account.did, null, null, 500); if (cids.len == 0) return error.MissingBlob; } return .{ .name = "space listBlobs", .ops = iterations, .elapsed_ns = nowNs() - start }; } fn benchDecodeRecord(allocator: std.mem.Allocator, options: Options) !void { try benchRecordBlockCpu(allocator, options, .decode); } fn benchRenderRecord(allocator: std.mem.Allocator, options: Options) !void { try benchRecordBlockCpu(allocator, options, .render); } fn benchConcurrentLists(allocator: std.mem.Allocator, level: ConcurrencyLevel) !ConcurrentResult { var state = try initBench(allocator); defer state.deinit(); try seedIndexedRecords(allocator, state.account, 1000); const total_ops = level.callers * level.ops_per_caller; const latencies = try allocator.alloc(u64, total_ops); defer allocator.free(latencies); const threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(threads); const start = nowNs(); for (threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, listWorker, .{ state.account.did, latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], }); } for (threads) |thread| thread.join(); return concurrentResult("listRecords", level, total_ops, nowNs() - start, latencies); } fn writeWorker(account: zds.auth.tokens.Account, offset: usize, latencies: []u64) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies, 0..) |*latency, i| { _ = arena.reset(.retain_capacity); const value = try benchRecordValue(arena.allocator(), offset + i); const start = nowNs(); const result = try zds.storage.store.applyWrites(arena.allocator(), account, &.{.{ .create = .{ .collection = bench_collection, .rkey = try std.fmt.allocPrint(arena.allocator(), "r{d:0>10}", .{offset + i}), .value = value, } }}); if (result.records.len != 1) return error.UnexpectedWriteResult; latency.* = nowNs() - start; } } fn profileWorker( account: zds.auth.tokens.Account, offset: usize, latencies: []u64, profiles: []zds.storage.store.WriteProfile, ) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies, profiles, 0..) |*latency, *profile, i| { _ = arena.reset(.retain_capacity); const value = try benchRecordValue(arena.allocator(), offset + i); const start = nowNs(); const result = try zds.storage.store.applyWritesProfiled(arena.allocator(), account, &.{.{ .create = .{ .collection = bench_collection, .rkey = try std.fmt.allocPrint(arena.allocator(), "r{d:0>10}", .{offset + i}), .value = value, } }}, profile); if (result.records.len != 1) return error.UnexpectedWriteResult; latency.* = nowNs() - start; } } fn printWriteProfile(profiles: []const zds.storage.store.WriteProfile) void { var total: zds.storage.store.WriteProfile = .{}; for (profiles) |profile| { total.validation_ns += profile.validation_ns; total.lock_wait_ns += profile.lock_wait_ns; total.load_repo_ns += profile.load_repo_ns; total.stage_records_ns += profile.stage_records_ns; total.build_commit_ns += profile.build_commit_ns; total.sql_ns += profile.sql_ns; total.event_publish_ns += profile.event_publish_ns; total.total_ns += profile.total_ns; } std.debug.print("\nwrite stage averages across {d} ops:\n", .{profiles.len}); printProfileStage("validation", total.validation_ns, profiles.len, total.total_ns); printProfileStage("lock wait", total.lock_wait_ns, profiles.len, total.total_ns); printProfileStage("load repo", total.load_repo_ns, profiles.len, total.total_ns); printProfileStage("stage records", total.stage_records_ns, profiles.len, total.total_ns); printProfileStage("build commit", total.build_commit_ns, profiles.len, total.total_ns); printProfileStage("sqlite/event", total.sql_ns, profiles.len, total.total_ns); printProfileStage("publish", total.event_publish_ns, profiles.len, total.total_ns); } fn printProfileStage(name: []const u8, ns: u64, ops: usize, total_ns: u64) void { const avg_ms = @as(f64, @floatFromInt(ns)) / @as(f64, @floatFromInt(ops)) / std.time.ns_per_ms; const pct = if (total_ns == 0) 0 else @as(f64, @floatFromInt(ns)) * 100.0 / @as(f64, @floatFromInt(total_ns)); std.debug.print(" {s: <14} {d: >9.3} ms/op {d: >6.1}%\n", .{ name, avg_ms, pct }); } fn getWorker(did: []const u8, task_id: usize, latencies: []u64) !void { for (latencies, 0..) |*latency, i| { const rec_idx = (task_id * 7 + i * 13) % 1000; var rkey_buf: [32]u8 = undefined; const rkey = try std.fmt.bufPrint(&rkey_buf, "rec{d:0>8}", .{rec_idx}); const start = nowNs(); const record = zds.storage.store.get(did, bench_collection, rkey) orelse return error.MissingRecord; std.heap.page_allocator.free(record.did); std.heap.page_allocator.free(record.collection); std.heap.page_allocator.free(record.rkey); std.heap.page_allocator.free(record.cid); std.heap.page_allocator.free(record.value_json); std.heap.page_allocator.free(record.validation_status); std.heap.page_allocator.free(record.rev); latency.* = nowNs() - start; } } fn getCidWorker(did: []const u8, task_id: usize, latencies: []u64) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies, 0..) |*latency, i| { _ = arena.reset(.retain_capacity); const rec_idx = (task_id * 7 + i * 13) % 1000; var rkey_buf: [32]u8 = undefined; const rkey = try std.fmt.bufPrint(&rkey_buf, "rec{d:0>8}", .{rec_idx}); const start = nowNs(); const cid = zds.storage.store.getCid(arena.allocator(), did, bench_collection, rkey) orelse return error.MissingRecord; if (cid.len == 0) return error.MissingRecord; latency.* = nowNs() - start; } } fn getBlockWorker(did: []const u8, task_id: usize, latencies: []u64) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies, 0..) |*latency, i| { _ = arena.reset(.retain_capacity); const rec_idx = (task_id * 7 + i * 13) % 1000; var rkey_buf: [32]u8 = undefined; const rkey = try std.fmt.bufPrint(&rkey_buf, "rec{d:0>8}", .{rec_idx}); const start = nowNs(); const block = zds.storage.store.getRecordBlock(arena.allocator(), did, bench_collection, rkey) orelse return error.MissingRecord; if (block.cid.len == 0 or block.data.len == 0) return error.MissingRecord; latency.* = nowNs() - start; } } fn listWorker(did: []const u8, latencies: []u64) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies) |*latency| { _ = arena.reset(.retain_capacity); const start = nowNs(); const records = try zds.storage.store.listRecords(arena.allocator(), did, bench_collection, 50); if (records.len == 0) return error.MissingRecords; latency.* = nowNs() - start; } } fn repoCarWorker(did: []const u8, latencies: []u64) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies) |*latency| { _ = arena.reset(.retain_capacity); const start = nowNs(); const car = try zds.storage.store.writeRepoCar(arena.allocator(), did); if (car.len == 0) return error.EmptyRepoCar; latency.* = nowNs() - start; } } fn publishWorker( account: zds.auth.tokens.Account, blob_size: usize, offset: usize, latencies: []u64, start: *std.atomic.Value(bool), ) !void { waitForBenchStart(start); var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies, 0..) |*latency, i| { _ = arena.reset(.retain_capacity); const a = arena.allocator(); const payload = try a.alloc(u8, blob_size); @memset(payload, @intCast((offset + i) % 251)); if (payload.len > 0) payload[0] = @intCast((offset + i + 1) % 251); const started = nowNs(); const cid = try zds.storage.store.putBlob(a, std.Options.debug_io, account, payload, "audio/mpeg"); const value = try blobRecordValue(a, offset + i, cid, blob_size); const result = try zds.storage.store.applyWrites(a, account, &.{.{ .create = .{ .collection = bench_collection, .rkey = try std.fmt.allocPrint(a, "publish{d:0>8}", .{offset + i}), .value = value, } }}); if (result.records.len != 1) return error.UnexpectedWriteResult; latency.* = nowNs() - started; } } fn publishSyncWorker( did: []const u8, cid: []const u8, latencies: []u64, start: *std.atomic.Value(bool), ) !void { waitForBenchStart(start); var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies) |*latency| { _ = arena.reset(.retain_capacity); const started = nowNs(); const blob = zds.storage.store.getBlob(arena.allocator(), did, cid) orelse return error.MissingBlob; if (blob.data.len == 0) return error.MissingBlob; const car = try zds.storage.store.writeRepoCar(arena.allocator(), did); if (car.len == 0) return error.EmptyRepoCar; latency.* = nowNs() - started; } } fn publishProbeWorker( did: []const u8, latencies: []u64, start: *std.atomic.Value(bool), ) !void { waitForBenchStart(start); var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies) |*latency| { _ = arena.reset(.retain_capacity); const started = nowNs(); const cid = zds.storage.store.getCid(arena.allocator(), did, bench_collection, "rec00000000") orelse return error.MissingRecord; if (cid.len == 0) return error.MissingRecord; latency.* = nowNs() - started; } } fn waitForBenchStart(start: *std.atomic.Value(bool)) void { while (!start.load(.acquire)) std.Thread.yield() catch {}; } fn repoCarWorkerStarting( did: []const u8, latencies: []u64, start: *std.atomic.Value(bool), ) !void { while (!start.load(.acquire)) std.atomic.spinLoopHint(); return repoCarWorker(did, latencies); } fn accountProbeWorker( did: []const u8, latencies: []u64, 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(); waitForExportsToEnterStore(); for (latencies) |*latency| { _ = arena.reset(.retain_capacity); const started = nowNs(); const account = try zds.storage.store.findAccount(arena.allocator(), did) orelse return error.MissingAccount; if (!std.mem.eql(u8, account.did, did)) return error.MissingAccount; latency.* = nowNs() - started; } } fn repoWriteWorker( account: zds.auth.tokens.Account, first_index: usize, latencies: []u64, 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(); waitForExportsToEnterStore(); for (latencies, 0..) |*latency, offset| { _ = arena.reset(.retain_capacity); const index = first_index + offset; const value = try benchRecordValue(arena.allocator(), index); const started = nowNs(); const result = try zds.storage.store.applyWrites(arena.allocator(), account, &.{.{ .create = .{ .collection = bench_collection, .rkey = try std.fmt.allocPrint(arena.allocator(), "rec{d:0>8}", .{index}), .value = value, } }}); if (result.records.len != 1) return error.UnexpectedWriteResult; latency.* = nowNs() - started; } } fn waitForExportsToEnterStore() void { const until = nowNs() + std.time.ns_per_ms; while (nowNs() < until) std.atomic.spinLoopHint(); } const CpuPath = enum { decode, render }; fn benchRecordBlockCpu(allocator: std.mem.Allocator, options: Options, path: CpuPath) !void { const level: ConcurrencyLevel = .{ .callers = options.callers orelse 10, .ops_per_caller = options.ops_per_caller orelse 1000, }; var state = try initBench(allocator); defer state.deinit(); try seedIndexedRecords(allocator, state.account, 1); var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const block = zds.storage.store.getRecordBlock(arena.allocator(), state.account.did, bench_collection, "rec00000000") orelse return error.MissingRecord; const total_ops = level.callers * level.ops_per_caller; const latencies = try allocator.alloc(u64, total_ops); defer allocator.free(latencies); const threads = try allocator.alloc(std.Thread, level.callers); defer allocator.free(threads); const start = nowNs(); for (threads, 0..) |*thread, i| { thread.* = try std.Thread.spawn(.{}, recordBlockCpuWorker, .{ block.data, path, latencies[i * level.ops_per_caller ..][0..level.ops_per_caller], }); } for (threads) |thread| thread.join(); concurrentResult(switch (path) { .decode => "decodeDagCbor", .render => "renderRecordJson", }, level, total_ops, nowNs() - start, latencies).print(); } fn recordBlockCpuWorker(data: []const u8, path: CpuPath, latencies: []u64) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); for (latencies) |*latency| { _ = arena.reset(.retain_capacity); const start = nowNs(); switch (path) { .decode => { _ = try zds.zat.cbor.decodeAll(arena.allocator(), data); }, .render => { const json = try zds.storage.store.recordJsonFromDagCbor(arena.allocator(), data); if (json.len == 0) return error.EmptyRecordJson; }, } latency.* = nowNs() - start; } } fn seedIndexedRecords(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, target: usize) !void { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); for (0..target) |i| { _ = arena.reset(.retain_capacity); const value = try benchRecordValue(arena.allocator(), i); const result = try zds.storage.store.applyWrites(arena.allocator(), account, &.{.{ .create = .{ .collection = bench_collection, .rkey = try std.fmt.allocPrint(arena.allocator(), "rec{d:0>8}", .{i}), .value = value, } }}); if (result.records.len != 1) return error.UnexpectedWriteResult; } } fn appendIndexedRecordRev(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, index: usize) ![]const u8 { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const value = try benchRecordValue(arena.allocator(), index); const result = try zds.storage.store.applyWrites(arena.allocator(), account, &.{.{ .create = .{ .collection = bench_collection, .rkey = try std.fmt.allocPrint(arena.allocator(), "rec{d:0>8}", .{index}), .value = value, } }}); if (result.records.len != 1) return error.UnexpectedWriteResult; return try allocator.dupe(u8, result.commit.rev); } fn concurrentResult(name: []const u8, level: ConcurrencyLevel, total_ops: usize, elapsed_ns: u64, latencies: []u64) ConcurrentResult { return .{ .name = name, .callers = level.callers, .ops = total_ops, .elapsed_ns = elapsed_ns, .stats = latencyStats(latencies), }; } fn latencyStats(latencies: []u64) ConcurrentStats { std.mem.sort(u64, latencies, {}, std.sort.asc(u64)); var sum: u128 = 0; for (latencies) |latency| sum += latency; const last = latencies.len - 1; return .{ .p50 = latencies[last * 50 / 100], .p95 = latencies[last * 95 / 100], .p99 = latencies[last * 99 / 100], .max = latencies[last], .mean = @intCast(sum / latencies.len), }; } fn nsToMs(ns: u64) f64 { return @as(f64, @floatFromInt(ns)) / std.time.ns_per_ms; } fn seedRecords(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, target: usize) !void { const existing = try zds.storage.store.listRecords(allocator, account.did, "app.bsky.feed.post", 1); if (existing.len > 0) return; _ = try benchWrite(allocator, account, target); } fn postValue(allocator: std.mem.Allocator, index: usize) !std.json.Value { const json = try std.fmt.allocPrint( allocator, "{{\"$type\":\"app.bsky.feed.post\",\"text\":\"zds bench post {d}\",\"createdAt\":\"2026-05-23T00:{d:0>2}:00.000Z\"}}", .{ index, index % 60 }, ); return try std.json.parseFromSliceLeaky(std.json.Value, allocator, json, .{}); } fn benchRecordValue(allocator: std.mem.Allocator, index: usize) !std.json.Value { const json = try std.fmt.allocPrint( allocator, "{{\"$type\":\"{s}\",\"text\":\"zds bench record {d}\",\"createdAt\":\"2026-05-23T00:{d:0>2}:00.000Z\"}}", .{ bench_collection, index, index % 60 }, ); return try std.json.parseFromSliceLeaky(std.json.Value, allocator, json, .{}); } fn blobRecordValue(allocator: std.mem.Allocator, index: usize, cid: []const u8, size: usize) !std.json.Value { const json = try std.fmt.allocPrint( allocator, "{{\"$type\":\"{s}\",\"text\":\"publish {d}\",\"audio\":{{\"$type\":\"blob\",\"ref\":{{\"$link\":{f}}},\"mimeType\":\"audio/mpeg\",\"size\":{d}}}}}", .{ bench_collection, index, std.json.fmt(cid, .{}), size }, ); return try std.json.parseFromSliceLeaky(std.json.Value, allocator, json, .{}); } fn spaceRecordValue(allocator: std.mem.Allocator, index: usize) !std.json.Value { const json = try std.fmt.allocPrint( allocator, "{{\"$type\":\"{s}\",\"title\":\"private track {d}\",\"createdAt\":\"2026-06-15T00:{d:0>2}:00.000Z\",\"audioRef\":\"bench-audio-{d}\",\"size\":{d}}}", .{ space_collection, index, index % 60, index, 1024 + index }, ); return try std.json.parseFromSliceLeaky(std.json.Value, allocator, json, .{}); } fn parseOptions(init: std.process.Init) !Options { var options: Options = .{}; var args = std.process.Args.Iterator.init(init.minimal.args); _ = args.next(); while (args.next()) |arg| { if (std.mem.eql(u8, arg, "--scenario")) { options.scenario = try parseScenario(args.next() orelse return error.MissingScenario); } else if (std.mem.startsWith(u8, arg, "--scenario=")) { options.scenario = try parseScenario(arg["--scenario=".len..]); } else if (std.mem.eql(u8, arg, "--records")) { options.records = try std.fmt.parseInt(usize, args.next() orelse return error.MissingRecords, 10); } else if (std.mem.startsWith(u8, arg, "--records=")) { options.records = try std.fmt.parseInt(usize, arg["--records=".len..], 10); } else if (std.mem.eql(u8, arg, "--blobs")) { options.blobs = try std.fmt.parseInt(usize, args.next() orelse return error.MissingBlobs, 10); } else if (std.mem.startsWith(u8, arg, "--blobs=")) { options.blobs = try std.fmt.parseInt(usize, arg["--blobs=".len..], 10); } else if (std.mem.eql(u8, arg, "--blob-size")) { options.blob_size = try std.fmt.parseInt(usize, args.next() orelse return error.MissingBlobSize, 10); } else if (std.mem.startsWith(u8, arg, "--blob-size=")) { options.blob_size = try std.fmt.parseInt(usize, arg["--blob-size=".len..], 10); } else if (std.mem.eql(u8, arg, "--callers")) { options.callers = try std.fmt.parseInt(usize, args.next() orelse return error.MissingCallers, 10); } else if (std.mem.startsWith(u8, arg, "--callers=")) { options.callers = try std.fmt.parseInt(usize, arg["--callers=".len..], 10); } else if (std.mem.eql(u8, arg, "--ops-per-caller")) { options.ops_per_caller = try std.fmt.parseInt(usize, args.next() orelse return error.MissingOpsPerCaller, 10); } else if (std.mem.startsWith(u8, arg, "--ops-per-caller=")) { options.ops_per_caller = try std.fmt.parseInt(usize, arg["--ops-per-caller=".len..], 10); } else if (std.mem.eql(u8, arg, "--help") or std.mem.eql(u8, arg, "-h")) { usage(); std.process.exit(0); } else { return error.UnknownArgument; } } return options; } fn parseScenario(value: []const u8) !Scenario { if (std.mem.eql(u8, value, "all")) return .all; if (std.mem.eql(u8, value, "write")) return .write; if (std.mem.eql(u8, value, "read")) return .read; if (std.mem.eql(u8, value, "blob")) return .blob; if (std.mem.eql(u8, value, "blob-gc")) return .blob_gc; if (std.mem.eql(u8, value, "repo")) return .repo; if (std.mem.eql(u8, value, "metastore")) return .metastore; if (std.mem.eql(u8, value, "get-cid")) return .get_cid; if (std.mem.eql(u8, value, "get-block")) return .get_block; if (std.mem.eql(u8, value, "decode-record")) return .decode_record; if (std.mem.eql(u8, value, "render-record")) return .render_record; if (std.mem.eql(u8, value, "get-record")) return .get_record; 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, "publish-sync")) return .publish_sync; if (std.mem.eql(u8, value, "space")) return .space; if (std.mem.eql(u8, value, "write-profile")) return .write_profile; return error.UnknownScenario; } fn usage() void { std.debug.print( \\usage: zds-bench [--scenario all|write|read|blob|blob-gc|repo|metastore|get-cid|get-block|decode-record|render-record|get-record|list-records|repo-size|get-repo|firehose|publish-sync|space|write-profile] [--records N] [--blobs N] [--blob-size BYTES] [--callers N --ops-per-caller N] \\ , .{}); } fn nowNs() u64 { return @intCast(std.Io.Clock.awake.now(std.Options.debug_io).toNanoseconds()); } fn tmpDir() []const u8 { return if (std.c.getenv("TMPDIR")) |value| std.mem.span(value) else "/tmp"; } fn cleanupPath(path: []const u8) void { std.Io.Dir.cwd().deleteFile(std.Options.debug_io, path) catch {}; const wal = std.fmt.allocPrint(std.heap.page_allocator, "{s}-wal", .{path}) catch return; defer std.heap.page_allocator.free(wal); const shm = std.fmt.allocPrint(std.heap.page_allocator, "{s}-shm", .{path}) catch return; defer std.heap.page_allocator.free(shm); std.Io.Dir.cwd().deleteFile(std.Options.debug_io, wal) catch {}; std.Io.Dir.cwd().deleteFile(std.Options.debug_io, shm) catch {}; deleteTree(path) catch {}; } fn deleteTree(path: []const u8) !void { const path_z = try std.heap.page_allocator.dupeSentinel(u8, path, 0); defer std.heap.page_allocator.free(path_z); const dir = std.c.opendir(path_z.ptr) orelse return; defer _ = std.c.closedir(dir); while (std.c.readdir(dir)) |entry| { const name = std.mem.sliceTo(&entry.name, 0); if (std.mem.eql(u8, name, ".") or std.mem.eql(u8, name, "..")) continue; const child = try std.fs.path.join(std.heap.page_allocator, &.{ path, name }); defer std.heap.page_allocator.free(child); std.Io.Dir.cwd().deleteFile(std.Options.debug_io, child) catch { deleteTree(child) catch {}; }; } const rc = std.c.rmdir(path_z.ptr); if (rc != 0) return error.DeleteTreeFailed; }