jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424//! bench — hot-path timings for each subsystem.//!//! Local and comparative only: the numbers mean something against another run//! on the same machine, and only as a ratio. This is here to catch//! order-of-magnitude regressions, not to resolve a few percent. See//! docs/benchmarks.md for what each case corresponds to upstream.//!//! Deliberately naive: fixed iteration counts, a short prime, no statistics.//! If a number matters enough to argue about, measure it properly.
const std = @import("std");const archive_mod = @import("internal/storage/archive.zig");const cold = @import("internal/storage/cold.zig");const gloom = @import("internal/storage/gloom.zig");const segment = @import("internal/storage/segment.zig");const tail_mod = @import("internal/serve/tail.zig");const wire = @import("internal/serve/wire.zig");
const Io = std.Io;const Allocator = std.mem.Allocator;
/// Reported per case. `unit_count` is whatever the case counts as one unit of/// useful work (rows, or uncompressed payload bytes) so the derived throughput/// column is honest about what it is dividing.const Result = struct { name: []const u8, subsystem: []const u8, iters: u64, total_ns: u64, unit_count: u64, unit: []const u8,
fn nsPerOp(self: Result) f64 { return @as(f64, @floatFromInt(self.total_ns)) / @as(f64, @floatFromInt(self.iters)); }
fn throughput(self: Result) f64 { const secs = @as(f64, @floatFromInt(self.total_ns)) / std.time.ns_per_s; if (secs == 0) return 0; return @as(f64, @floatFromInt(self.unit_count)) / secs; }};
var results: std.ArrayList(Result) = .empty;
/// Nanoseconds between two Io timestamps. std.time.nanoTimestamp does not/// exist in Zig 0.16; the rest of the codebase measures through Io the same way.fn since(io: Io, start: Io.Timestamp) u64 { return @intCast(@max(start.durationTo(Io.Timestamp.now(io, .awake)).nanoseconds, 0));}
fn record(alloc: Allocator, r: Result) !void { try results.append(alloc, r);}
const payload = "\xa2\x64text\x6bhello world\x65$type\x72app.bsky.feed.post";
fn sampleEvent(i: usize) segment.Event { return .{ .seq = 0, .witnessed_at = 1_700_000_000_000_000 + @as(i64, @intCast(i)), .indexed_at = 0, .kind = .create, .did = "did:plc:benchmarkaccount000", .collection = "app.bsky.feed.post", .rkey = "3l3qo2vuowo2b", .rev = "3l3qo2vutsw2b", .payload = payload, };}
/// Archive append under the two shapes upstream distinguishes: live is one row/// at a time with the flush cadence that implies, backfill is batched.fn benchArchiveAppend(alloc: Allocator, io: Io, comptime batched: bool) !Result { var tmp_buf: [Io.Dir.max_path_bytes]u8 = undefined; const dir = try std.fmt.bufPrint(&tmp_buf, "/tmp/stream-bench-{s}", .{if (batched) "backfill" else "live"}); Io.Dir.cwd().deleteTree(io, dir) catch {}; defer Io.Dir.cwd().deleteTree(io, dir) catch {};
var archive = try archive_mod.Archive.init(alloc, io, dir); defer archive.deinit();
const iters: u64 = if (batched) 20_000 else 5_000; const batch_size = 64;
var events: [batch_size]segment.Event = undefined; const start = Io.Timestamp.now(io, .awake); if (batched) { var done: u64 = 0; while (done < iters) : (done += batch_size) { for (&events, 0..) |*e, j| e.* = sampleEvent(@intCast(done + j)); _ = try archive.appendBatch(&events, 1_700_000_000_000_000); } } else { var i: u64 = 0; while (i < iters) : (i += 1) { _ = try archive.append(sampleEvent(@intCast(i)), 1_700_000_000_000_000); } } try archive.flushBlock(); const elapsed = since(io, start); archive.close();
return .{ .name = if (batched) "archive_append_backfill" else "archive_append_live", .subsystem = "ingest->storage", .iters = iters, .total_ns = elapsed, .unit_count = iters, .unit = "rows/s", };}
/// Seal cost, then the two reader-open paths: checksum-verifying and not./// Splitting them is upstream's ReaderOpen vs ReaderOpenNoVerify.fn benchSealAndParse(alloc: Allocator, io: Io) ![3]Result { var tmp_buf: [Io.Dir.max_path_bytes]u8 = undefined; const dir = try std.fmt.bufPrint(&tmp_buf, "/tmp/stream-bench-seal", .{}); Io.Dir.cwd().deleteTree(io, dir) catch {}; defer Io.Dir.cwd().deleteTree(io, dir) catch {};
var archive = try archive_mod.Archive.init(alloc, io, dir); defer archive.deinit();
const rows: u64 = 20_000; var i: u64 = 0; while (i < rows) : (i += 1) _ = try archive.append(sampleEvent(@intCast(i)), 1_700_000_000_000_000);
const seal_start = Io.Timestamp.now(io, .awake); try archive.rotate(); const seal_ns = since(io, seal_start);
var seg_buf: [Io.Dir.max_path_bytes]u8 = undefined; const seg_path = try std.fmt.bufPrint(&seg_buf, "{s}/segments", .{dir}); var seg_dir = try Io.Dir.cwd().openDir(io, seg_path, .{ .iterate = true }); defer seg_dir.close(io); const bytes = try seg_dir.readFileAlloc(io, "seg_0000000000.jss", alloc, .limited(1 << 30)); defer alloc.free(bytes); archive.close();
const parse_iters: u64 = 2_000; var verified_ns: u64 = 0; { const start = Io.Timestamp.now(io, .awake); var n: u64 = 0; while (n < parse_iters) : (n += 1) { var sealed = try segment.Sealed.parse(alloc, bytes); sealed.deinit(alloc); } verified_ns = since(io, start); } var unchecked_ns: u64 = 0; { const start = Io.Timestamp.now(io, .awake); var n: u64 = 0; while (n < parse_iters) : (n += 1) { var sealed = try segment.Sealed.parseUnchecked(alloc, bytes); sealed.deinit(alloc); } unchecked_ns = since(io, start); }
return .{ .{ .name = "archive_seal", .subsystem = "storage", .iters = 1, .total_ns = seal_ns, .unit_count = rows, .unit = "rows/s", }, .{ .name = "sealed_parse", .subsystem = "storage", .iters = parse_iters, .total_ns = verified_ns, .unit_count = parse_iters, .unit = "parses/s", }, .{ .name = "sealed_parse_unchecked", .subsystem = "storage", .iters = parse_iters, .total_ns = unchecked_ns, .unit_count = parse_iters, .unit = "parses/s", }, };}
/// Decoding one sealed block: the cold reader's inner loop.fn benchSealedReadBlock(alloc: Allocator, io: Io) !Result { var tmp_buf: [Io.Dir.max_path_bytes]u8 = undefined; const dir = try std.fmt.bufPrint(&tmp_buf, "/tmp/stream-bench-block", .{}); Io.Dir.cwd().deleteTree(io, dir) catch {}; defer Io.Dir.cwd().deleteTree(io, dir) catch {};
var archive = try archive_mod.Archive.init(alloc, io, dir); defer archive.deinit(); const rows: u64 = 20_000; var i: u64 = 0; while (i < rows) : (i += 1) _ = try archive.append(sampleEvent(@intCast(i)), 1_700_000_000_000_000); try archive.rotate();
var seg_buf: [Io.Dir.max_path_bytes]u8 = undefined; const seg_path = try std.fmt.bufPrint(&seg_buf, "{s}/segments", .{dir}); var seg_dir = try Io.Dir.cwd().openDir(io, seg_path, .{ .iterate = true }); defer seg_dir.close(io); const bytes = try seg_dir.readFileAlloc(io, "seg_0000000000.jss", alloc, .limited(1 << 30)); defer alloc.free(bytes); archive.close();
var sealed = try segment.Sealed.parse(alloc, bytes); defer sealed.deinit(alloc);
const iters: u64 = 2_000; var decoded_rows: u64 = 0; const start = Io.Timestamp.now(io, .awake); var n: u64 = 0; while (n < iters) : (n += 1) { var block = try sealed.readBlock(alloc, n % sealed.block_index.len); decoded_rows += block.events.len; block.deinit(alloc); } const elapsed = since(io, start);
return .{ .name = "sealed_read_block", .subsystem = "storage", .iters = iters, .total_ns = elapsed, .unit_count = decoded_rows, .unit = "rows/s", };}
/// Per-block DID bloom lookup, consulted once per block per planSnapshot.fn benchBloom(alloc: Allocator, io: Io) !Result { var filter = try gloom.Filter.init(alloc, 10_000, 0.01); defer filter.deinit(alloc); var key_buf: [64]u8 = undefined; var i: usize = 0; while (i < 10_000) : (i += 1) { const key = try std.fmt.bufPrint(&key_buf, "did:plc:bloomkey{d:0>12}", .{i}); filter.add(key); } const blob = try filter.marshal(alloc); defer alloc.free(blob);
const iters: u64 = 200_000; var hits: u64 = 0; const start = Io.Timestamp.now(io, .awake); var n: u64 = 0; while (n < iters) : (n += 1) { const key = try std.fmt.bufPrint(&key_buf, "did:plc:bloomkey{d:0>12}", .{n % 20_000}); if (try gloom.containsMarshaled(blob, key)) hits += 1; } const elapsed = since(io, start); std.mem.doNotOptimizeAway(&hits);
return .{ .name = "bloom_lookup", .subsystem = "storage", .iters = iters, .total_ns = elapsed, .unit_count = iters, .unit = "lookups/s", };}
/// Wire encoding, run once per row per subscriber wire.fn benchWire(alloc: Allocator, io: Io, comptime v2: bool) !Result { const cid = try @import("zat").cbor.Cid.forDagCbor(alloc, "bench"); defer alloc.free(cid.raw); const rec: @import("zat").cbor.Value = .{ .map = &.{ .{ .key = "$type", .value = .{ .text = "app.bsky.feed.post" } }, .{ .key = "text", .value = .{ .text = "hello world" } }, } }; const op: wire.CommitOp = .{ .did = "did:plc:benchmarkaccount000", .time_us = 1_700_000_000_000_000, .rev = "3l3qo2vutsw2b", .operation = "create", .collection = "app.bsky.feed.post", .rkey = "3l3qo2vuowo2b", .record = rec, .cid = cid, };
const iters: u64 = 100_000; var bytes_out: u64 = 0; const start = Io.Timestamp.now(io, .awake); var n: u64 = 0; while (n < iters) : (n += 1) { var buf: std.Io.Writer.Allocating = .init(alloc); defer buf.deinit(); if (v2) { try wire.encodeCommitV2(alloc, &buf.writer, op, n); } else { try wire.encodeCommit(alloc, &buf.writer, op); } bytes_out += buf.written().len; } const elapsed = since(io, start);
return .{ .name = if (v2) "wire_encode_v2" else "wire_encode_v1", .subsystem = "serve", .iters = iters, .total_ns = elapsed, .unit_count = iters, .unit = "frames/s", };}
/// The cold reader's batch loop: what a subscriber behind the hot tail pays.fn benchColdReplay(alloc: Allocator, io: Io) !Result { var tmp_buf: [Io.Dir.max_path_bytes]u8 = undefined; const dir = try std.fmt.bufPrint(&tmp_buf, "/tmp/stream-bench-cold", .{}); Io.Dir.cwd().deleteTree(io, dir) catch {}; defer Io.Dir.cwd().deleteTree(io, dir) catch {};
var archive = try archive_mod.Archive.init(alloc, io, dir); defer archive.deinit(); const rows: u64 = 20_000; var i: u64 = 0; while (i < rows) : (i += 1) _ = try archive.append(sampleEvent(@intCast(i)), 1_700_000_000_000_000); try archive.rotate();
var reader = cold.ColdReader.init(alloc, io, 64 << 20); defer reader.deinit();
const Sink = struct { delivered: u64 = 0, fn wants(_: *@This(), _: wire.Kind, _: []const u8, _: []const u8, _: tail_mod.SkipV1) bool { return true; } fn write(_: *@This(), _: cold.CachedFrame) anyerror!cold.WriteResult { return .sent; } fn observe(self: *@This(), _: u64, result: cold.ReplayResult) anyerror!void { if (result == .sent) self.delivered += 1; } }; var sink: Sink = .{};
const start = Io.Timestamp.now(io, .awake); var from: u64 = 1; while (from <= rows) { const batch = try reader.readSeqBatch( &archive, &sink, Sink.wants, &sink, Sink.write, Sink.observe, false, false, from, rows + 1, 1024, ); if (batch.next_seq <= from) break; from = batch.next_seq; } const elapsed = since(io, start); archive.close();
return .{ .name = "cold_replay_batch", .subsystem = "serve", .iters = sink.delivered, .total_ns = elapsed, .unit_count = sink.delivered, .unit = "rows/s", };}
fn humanize(v: f64, buf: []u8) []const u8 { if (v >= 1_000_000_000) return std.fmt.bufPrint(buf, "{d:.2}G", .{v / 1e9}) catch "?"; if (v >= 1_000_000) return std.fmt.bufPrint(buf, "{d:.2}M", .{v / 1e6}) catch "?"; if (v >= 1_000) return std.fmt.bufPrint(buf, "{d:.2}K", .{v / 1e3}) catch "?"; return std.fmt.bufPrint(buf, "{d:.0}", .{v}) catch "?";}
pub fn main() !void { // Production runs on c_allocator (src/main.zig). Benching on // DebugAllocator would inflate every allocation-heavy path -- wire encode // most of all -- and measure the harness rather than the system. const alloc = std.heap.c_allocator; defer results.deinit(alloc);
var threaded: Io.Threaded = .init(alloc, .{}); defer threaded.deinit(); const io = threaded.io();
try record(alloc, try benchArchiveAppend(alloc, io, false)); try record(alloc, try benchArchiveAppend(alloc, io, true)); const seal = try benchSealAndParse(alloc, io); for (seal) |r| try record(alloc, r); try record(alloc, try benchSealedReadBlock(alloc, io)); try record(alloc, try benchBloom(alloc, io)); try record(alloc, try benchWire(alloc, io, false)); try record(alloc, try benchWire(alloc, io, true)); try record(alloc, try benchColdReplay(alloc, io));
var out_buf: [4096]u8 = undefined; var stdout = Io.File.stdout().writer(io, &out_buf); const w = &stdout.interface; try w.print("{s:<26} {s:<16} {s:>12} {s:>14}\n", .{ "bench", "subsystem", "ns/op", "throughput" }); for (results.items) |r| { var tbuf: [32]u8 = undefined; var nbuf: [32]u8 = undefined; try w.print("{s:<26} {s:<16} {s:>12} {s:>10} {s}\n", .{ r.name, r.subsystem, humanize(r.nsPerOp(), &nbuf), humanize(r.throughput(), &tbuf), r.unit, }); } try w.flush();}