atproto pds in zig pds.zat.dev
pds atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623const 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;}