From cf554edee73f04e5a1962bfc7961807723d9f8a7 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sat, 29 Aug 2026 01:03:36 -0500 Subject: [PATCH] =?UTF-8?q?examples:=20segment-export=20=E2=80=94=20one=20?= =?UTF-8?q?local=20jss=20segment=20to=20TSV=20rows?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Thin wrapper on ArchiveBackfill.replaySegment with the onRow protocol so marker rows and seq are kept. Columns: seq, time_us (witnessed), did, kind, op, collection, rkey. Measured 3.18M rows from a 269 MB segment in 10 s on the stream box. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01B927dNwNKNYsbQdUoNJwHS --- build.zig | 18 +++++++++++ examples/segment_export.zig | 60 +++++++++++++++++++++++++++++++++++++ 2 files changed, 78 insertions(+) create mode 100644 examples/segment_export.zig diff --git a/build.zig b/build.zig index 003160d..e851a1b 100644 --- a/build.zig +++ b/build.zig @@ -111,6 +111,24 @@ pub fn build(b: *std.Build) void { }), }); b.installArtifact(cutover); + + const segment_export = b.addExecutable(.{ + .name = "segment-export", + .root_module = b.createModule(.{ + .root_source_file = b.path("examples/segment_export.zig"), + .target = target, + .optimize = optimize, + .link_libc = true, + .imports = &.{ + .{ .name = "jetstream", .module = mod }, + }, + }), + }); + b.installArtifact(segment_export); + const run_segment_export = b.addRunArtifact(segment_export); + if (b.args) |args| run_segment_export.addArgs(args); + const segment_export_step = b.step("segment-export", "one local jss segment -> TSV rows on stdout"); + segment_export_step.dependOn(&run_segment_export.step); const run_cutover = b.addRunArtifact(cutover); if (b.args) |args| run_cutover.addArgs(args); const cutover_step = b.step("example-cutover-smoke", "live smoke: archive sweep -> gapless cutover -> live tail"); diff --git a/examples/segment_export.zig b/examples/segment_export.zig new file mode 100644 index 0000000..aa7c7af --- /dev/null +++ b/examples/segment_export.zig @@ -0,0 +1,60 @@ +//! segment export: one sealed jss segment on local disk → one TSV row per +//! event on stdout, for loading into a columnar store (duckdb, parquet). +//! +//! columns: seq, time_us, did, kind, op, collection, rkey. time_us is the +//! archive's witnessed time (jetstream v2 semantics), not record creation +//! time. record bodies are not emitted. +//! +//! run: zig build segment-export -- + +const std = @import("std"); +const jetstream_sdk = @import("jetstream"); + +const Exporter = struct { + out: *std.Io.Writer, + rows: u64 = 0, + failed: ?anyerror = null, + + pub fn onRow(self: *Exporter, event: jetstream_sdk.Event) bool { + self.write(event) catch |err| { + self.failed = err; + return false; + }; + return true; + } + + fn write(self: *Exporter, event: jetstream_sdk.Event) !void { + switch (event.payload) { + .commit => |c| try self.out.print("{d}\t{d}\t{s}\tcommit\t{s}\t{s}\t{s}\n", .{ + event.seq, event.time_us, event.did, @tagName(c.operation), c.collection, c.rkey, + }), + inline else => |_, tag| try self.out.print("{d}\t{d}\t{s}\t{s}\t\t\t\n", .{ + event.seq, event.time_us, event.did, @tagName(tag), + }), + } + self.rows += 1; + } +}; + +pub fn main(init: std.process.Init) !void { + const allocator = init.gpa; + const args = try init.minimal.args.toSlice(init.arena.allocator()); + if (args.len != 2) { + std.debug.print("usage: zig build segment-export -- \n", .{}); + return error.MissingSegment; + } + const body = try std.Io.Dir.cwd().readFileAlloc(init.io, args[1], allocator, .limited(1 << 31)); + defer allocator.free(body); + + var buffer: [1 << 16]u8 = undefined; + var stdout = std.Io.File.stdout().writer(init.io, &buffer); + var exporter: Exporter = .{ .out = &stdout.interface }; + const started = std.Io.Timestamp.now(init.io, .awake).nanoseconds; + _ = jetstream_sdk.ArchiveBackfill.replaySegment(allocator, body, .{ .host = "" }, &exporter) catch |err| switch (err) { + error.Stopped => return exporter.failed orelse err, + else => return err, + }; + try stdout.interface.flush(); + const elapsed_ms = @divTrunc(std.Io.Timestamp.now(init.io, .awake).nanoseconds - started, std.time.ns_per_ms); + std.debug.print("{d} rows from {d} bytes in {d} ms\n", .{ exporter.rows, body.len, elapsed_ms }); +} -- 2.51.2