diff --git a/build.zig b/build.zig index d5c98ab..877ba7b 100644 --- a/build.zig +++ b/build.zig @@ -138,6 +138,22 @@ pub fn build(b: *std.Build) void { const example_read_a_record_step = b.step("example-read-a-record", "read one record from a PDS via getRecord"); example_read_a_record_step.dependOn(&run_example_read_a_record.step); + const example_follow_the_firehose = b.addExecutable(.{ + .name = "example-follow-the-firehose", + .root_module = b.createModule(.{ + .root_source_file = b.path("examples/04_follow_the_firehose.zig"), + .target = target, + .optimize = optimize, + .link_libc = true, + .imports = &.{.{ .name = "zat", .module = mod }}, + }), + }); + b.installArtifact(example_follow_the_firehose); + + const run_example_follow_the_firehose = b.addRunArtifact(example_follow_the_firehose); + const example_follow_the_firehose_step = b.step("example-follow-the-firehose", "consume the firehose into a live language tally"); + example_follow_the_firehose_step.dependOn(&run_example_follow_the_firehose.step); + const example_resolve_identity = b.addExecutable(.{ .name = "example-resolve-identity", .root_module = b.createModule(.{ diff --git a/examples/04_follow_the_firehose.md b/examples/04_follow_the_firehose.md new file mode 100644 index 0000000..6ef06fa --- /dev/null +++ b/examples/04_follow_the_firehose.md @@ -0,0 +1,117 @@ +# follow the firehose + +The firehose is the network's event log — every write, streamed live. You +consume it by folding events into a **materialized view**: a running query you +keep in sync. This one tallies the languages people are posting in, right now. + +The [AT Protocol for distributed systems engineers](https://atproto.com/articles/atproto-for-distsys-engineers) +article is the mental model here — the database turned inside-out, writes flowing +through an event log into views you build yourself. + +```zig +const std = @import("std"); +const zat = @import("zat"); + +// the view we maintain from the stream: a post count per language. +const tracked = [_][]const u8{ "en", "ja", "pt", "es", "de", "fr", "ko" }; + +const LangView = struct { + total: u64 = 0, + counts: [tracked.len]u64 = [_]u64{0} ** tracked.len, + other: u64 = 0, + + pub fn onEvent(self: *LangView, event: zat.JetstreamEvent) void { + const commit = switch (event) { + .commit => |c| c, + else => return, + }; + if (commit.operation != .create) return; + const record = commit.record orelse return; + self.total += 1; + + // a post's declared languages, e.g. ["en"]. many posts omit it. + const langs = zat.json.getArray(record, "langs") orelse { + if (self.total % 500 == 0) self.report(); + return; + }; + if (langs.len > 0) { + switch (langs[0]) { + .string => |primary| { + for (tracked, 0..) |code, i| { + if (std.mem.eql(u8, code, primary)) { + self.counts[i] += 1; + break; + } + } else self.other += 1; + }, + else => {}, + } + } + + if (self.total % 500 == 0) self.report(); + } + + fn report(self: *const LangView) void { + std.debug.print("\n{d} posts\n", .{self.total}); + for (tracked, self.counts) |code, n| { + std.debug.print(" {s:<6}{d}\n", .{ code, n }); + } + std.debug.print(" other {d}\n", .{self.other}); + } + + pub fn onError(_: *LangView, err: anyerror) void { + std.debug.print("stream error: {s}\n", .{@errorName(err)}); + } +}; + +pub fn main() !void { + var da: std.heap.DebugAllocator(.{}) = .init; + defer _ = da.deinit(); + const allocator = da.allocator(); + + var view = LangView{}; + var client = zat.JetstreamClient.init(std.Options.debug_io, allocator, .{ + .hosts = &.{"jetstream2.us-east.bsky.network"}, + .wanted_collections = &.{"app.bsky.feed.post"}, + }); + defer client.deinit(); + + // subscribe blocks, folding each event into `view` until interrupted. + try client.subscribe(&view); +} +``` + +``` +500 posts + en 161 + ja 61 + pt 2 + es 15 + de 27 + fr 10 + ko 10 + other 61 +``` + +This is `examples/04_follow_the_firehose.zig` verbatim — copy it and +`zig build example-follow-the-firehose` (Ctrl-C to stop). + +The view lives entirely in memory and holds no borrowed data — each event's +record is freed after `onEvent` returns, so a durable view would copy out or +persist what it needs. Point `wanted_collections` at any lexicon (or leave it +empty for the whole firehose) and fold it into whatever query you care about. + +`zat.JetstreamClient` reads the JSON firehose. For the raw CBOR event stream +with commit blocks to verify, use `zat.FirehoseClient` — see the commit +verification recipe. + +
+wiring zat into your own build.zig + +```zig +const zat = b.dependency("zat", .{}).module("zat"); +exe.root_module.addImport("zat", zat); +``` + +after `zig fetch --save https://tangled.org/zat.dev/zat/archive/main`. +
diff --git a/examples/04_follow_the_firehose.zig b/examples/04_follow_the_firehose.zig new file mode 100644 index 0000000..a641f25 --- /dev/null +++ b/examples/04_follow_the_firehose.zig @@ -0,0 +1,79 @@ +//! consume the firehose and fold it into a live view. +//! +//! the network's writes stream past as an event log; you keep a "materialized +//! view" — a running query — in sync by folding each event in. here: which +//! languages people are posting in right now. +//! https://atproto.com/articles/atproto-for-distsys-engineers +//! +//! run: zig build example-follow-the-firehose (Ctrl-C to stop) + +const std = @import("std"); +const zat = @import("zat"); + +// the view we maintain from the stream: a post count per language. +const tracked = [_][]const u8{ "en", "ja", "pt", "es", "de", "fr", "ko" }; + +const LangView = struct { + total: u64 = 0, + counts: [tracked.len]u64 = [_]u64{0} ** tracked.len, + other: u64 = 0, + + pub fn onEvent(self: *LangView, event: zat.JetstreamEvent) void { + const commit = switch (event) { + .commit => |c| c, + else => return, + }; + if (commit.operation != .create) return; + const record = commit.record orelse return; + self.total += 1; + + // a post's declared languages, e.g. ["en"]. many posts omit it. + const langs = zat.json.getArray(record, "langs") orelse { + if (self.total % 500 == 0) self.report(); + return; + }; + if (langs.len > 0) { + switch (langs[0]) { + .string => |primary| { + for (tracked, 0..) |code, i| { + if (std.mem.eql(u8, code, primary)) { + self.counts[i] += 1; + break; + } + } else self.other += 1; + }, + else => {}, + } + } + + if (self.total % 500 == 0) self.report(); + } + + fn report(self: *const LangView) void { + std.debug.print("\n{d} posts\n", .{self.total}); + for (tracked, self.counts) |code, n| { + std.debug.print(" {s:<6}{d}\n", .{ code, n }); + } + std.debug.print(" other {d}\n", .{self.other}); + } + + pub fn onError(_: *LangView, err: anyerror) void { + std.debug.print("stream error: {s}\n", .{@errorName(err)}); + } +}; + +pub fn main() !void { + var da: std.heap.DebugAllocator(.{}) = .init; + defer _ = da.deinit(); + const allocator = da.allocator(); + + var view = LangView{}; + var client = zat.JetstreamClient.init(std.Options.debug_io, allocator, .{ + .hosts = &.{"jetstream2.us-east.bsky.network"}, + .wanted_collections = &.{"app.bsky.feed.post"}, + }); + defer client.deinit(); + + // subscribe blocks, folding each event into `view` until interrupted. + try client.subscribe(&view); +} diff --git a/examples/_overview.md b/examples/_overview.md new file mode 100644 index 0000000..de1b7c5 --- /dev/null +++ b/examples/_overview.md @@ -0,0 +1 @@ +Each recipe is a complete, runnable program — copy it and `zig build example-`. They run in order from client-side basics (parse an identifier, resolve one, read a record) toward the network's core (follow the firehose, walk repos, verify commits). diff --git a/scripts/build-site.mjs b/scripts/build-site.mjs index d17e641..0a358cf 100644 --- a/scripts/build-site.mjs +++ b/scripts/build-site.mjs @@ -216,7 +216,9 @@ async function main() { const exampleEntries = []; for (const rel of exampleFiles) { - if (rel === "index.md") continue; // generated below + // index.md is generated below; _-prefixed files (e.g. _overview.md) are + // editable copy for the index, not recipes. + if (rel === "index.md" || rel.startsWith("_")) continue; const src = path.join(examplesDir, rel); const dst = path.join(outDocsDir, "examples", rel); await mkdir(path.dirname(dst), { recursive: true }); @@ -234,9 +236,17 @@ async function main() { // no example list to maintain here. Only the index goes in the top nav. if (exampleEntries.length > 0) { exampleEntries.sort((a, b) => a.path.localeCompare(b.path)); + + // Optional editable overview prose sits above the generated list. + const overviewPath = path.join(examplesDir, "_overview.md"); + const overview = (await exists(overviewPath)) + ? (await readFile(overviewPath, "utf8")).trim() + : ""; + const indexMd = [ "# examples", "", + ...(overview ? [overview, ""] : []), ...exampleEntries.map((e) => `- [${e.title}](${e.path})`), "", ].join("\n");