Something went wrong. Try again.
A Zig library for Bluesky activity and feed analytics, built on Zat.
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240const std = @import("std");const zat = @import("zat");const Owned = @import("owned.zig").Owned;const client = @import("client.zig");
pub const Entry = struct { uri: []const u8, author: client.Profile, introduced_by: client.Profile, indexed_at: i64, is_reply: bool, is_repost: bool,};
pub const Attribution = enum { introducer, author };pub const Rank = struct { account: client.Profile, entries: usize = 0, posts: usize = 0, replies: usize = 0, reposts: usize = 0, share: f64 = 0,};pub const Report = struct { total_entries: usize, ranking: []const Rank, pub fn deinit(self: Report, allocator: std.mem.Allocator) void { allocator.free(self.ranking); }};
pub fn parse(value: std.json.Value) !Entry { const uri = zat.json.getString(value, "post.uri") orelse return error.MissingPost; const author = try profile(value, "post.author"); var introducer = author; var repost = false; var timestamp = zat.json.getString(value, "post.indexedAt") orelse return error.MissingTimestamp; if (zat.json.getPath(value, "reason")) |reason| { const tag = zat.json.getString(reason, "$type") orelse return error.InvalidReason; if (!std.mem.eql(u8, tag, "app.bsky.feed.defs#reasonRepost")) return error.UnsupportedReason; introducer = try profile(reason, "by"); timestamp = zat.json.getString(reason, "indexedAt") orelse return error.MissingTimestamp; repost = true; } const stamp = zat.Datetime.parse(timestamp) orelse return error.InvalidTimestamp; const record = zat.json.getPath(value, "post.record") orelse return error.MissingRecord; const kind = zat.json.getString(record, "$type") orelse return error.MissingType; if (!std.mem.eql(u8, kind, "app.bsky.feed.post")) return error.UnsupportedRecord; var reply = false; if (zat.json.getPath(record, "reply")) |ref| { if (zat.json.getString(ref, "parent.uri") == null or zat.json.getString(ref, "root.uri") == null) return error.InvalidReply; reply = true; } return .{ .uri = uri, .author = author, .introduced_by = introducer, .indexed_at = stamp.micros, .is_reply = reply, .is_repost = repost };}
fn profile(value: std.json.Value, path: []const u8) !client.Profile { const object = zat.json.getPath(value, path) orelse return error.MissingProfile; return .{ .did = zat.json.getString(object, "did") orelse return error.MissingDid, .handle = zat.json.getString(object, "handle") orelse return error.MissingHandle, .displayName = zat.json.getString(object, "displayName") orelse "", };}
pub fn rank(allocator: std.mem.Allocator, entries: []const Entry, attribution: Attribution) !Report { var rows: std.ArrayList(Rank) = .empty; errdefer rows.deinit(allocator); var indices = std.StringHashMap(usize).init(allocator); defer indices.deinit(); for (entries) |entry| { const account = switch (attribution) { .introducer => entry.introduced_by, .author => entry.author, }; const slot = try indices.getOrPut(account.did); if (!slot.found_existing) { slot.value_ptr.* = rows.items.len; try rows.append(allocator, .{ .account = account }); } const row = &rows.items[slot.value_ptr.*]; row.entries += 1; if (entry.is_repost) { row.reposts += 1; } else if (entry.is_reply) { row.replies += 1; } else { row.posts += 1; } } for (rows.items) |*row| row.share = @as(f64, @floatFromInt(row.entries)) / @as(f64, @floatFromInt(entries.len)); std.mem.sort(Rank, rows.items, {}, struct { fn less(_: void, a: Rank, b: Rank) bool { if (a.entries != b.entries) return a.entries > b.entries; return std.mem.lessThan(u8, a.account.did, b.account.did); } }.less); return .{ .total_entries = entries.len, .ranking = try rows.toOwnedSlice(allocator) };}
test "repost attribution separates introducer from author and retains reply dimension" { const raw = \\{"post":{"uri":"at://alice/post/1","author":{"did":"did:plc:alice","handle":"alice.test"},"indexedAt":"2026-09-01T00:00:00Z","record":{"$type":"app.bsky.feed.post","reply":{"root":{"uri":"at://root"},"parent":{"uri":"at://parent"}}}},"reason":{"$type":"app.bsky.feed.defs#reasonRepost","by":{"did":"did:plc:bob","handle":"bob.test"},"indexedAt":"2026-09-02T00:00:00Z"}} ; const value = try std.json.parseFromSlice(std.json.Value, std.testing.allocator, raw, .{}); defer value.deinit(); const repost = try parse(value.value); try std.testing.expect(repost.is_reply and repost.is_repost); try std.testing.expectEqual(zat.Datetime.parse("2026-09-02T00:00:00Z").?.micros, repost.indexed_at); var original = repost; original.introduced_by = original.author; original.is_repost = false; const entries = [_]Entry{ repost, original }; const by_introducer = try rank(std.testing.allocator, &entries, .introducer); defer std.testing.allocator.free(by_introducer.ranking); try std.testing.expectEqual(@as(usize, 2), by_introducer.ranking.len); try std.testing.expectEqual(@as(f64, 0.5), by_introducer.ranking[0].share); const by_author = try rank(std.testing.allocator, &entries, .author); defer std.testing.allocator.free(by_author.ranking); try std.testing.expectEqual(@as(usize, 1), by_author.ranking.len); try std.testing.expectEqual(@as(usize, 1), by_author.ranking[0].reposts); try std.testing.expectEqual(@as(usize, 1), by_author.ranking[0].replies); try std.testing.expectEqual(@as(f64, 1), by_author.ranking[0].share);}
test "empty feed has no rankings" { const report = try rank(std.testing.allocator, &.{}, .introducer); defer std.testing.allocator.free(report.ranking); try std.testing.expectEqual(@as(usize, 0), report.total_entries); try std.testing.expectEqual(@as(usize, 0), report.ranking.len);}
pub const Sample = struct { account: []const u8, appview: []const u8, started_at: i64, observed_at: i64, entries: []const Entry, pages: usize, received: usize, duplicates: usize, unclassified: usize, exhausted: bool, cursor: ?[]const u8, failure: ?[]const u8,};
pub fn sample(rpc: *client.Client, max_pages: usize) !Owned(Sample) { if (max_pages == 0) return error.InvalidPageLimit; const session = rpc.session orelse return error.NotAuthenticated; const appview = rpc.appview orelse return error.MissingAppview; try client.validateAppview(appview); const started = std.Io.Timestamp.now(rpc.io, .real).toMicroseconds(); var arena = std.heap.ArenaAllocator.init(rpc.allocator); errdefer arena.deinit(); const a = arena.allocator(); var entries: std.ArrayList(Entry) = .empty; var cursors = std.StringHashMap(void).init(a); defer cursors.deinit(); var keys = std.StringHashMap(void).init(a); defer keys.deinit(); var result: Sample = .{ .account = try a.dupe(u8, session.did), .appview = try a.dupe(u8, appview), .started_at = started, .observed_at = started, .entries = &.{}, .pages = 0, .received = 0, .duplicates = 0, .unclassified = 0, .exhausted = false, .cursor = null, .failure = null }; for (0..max_pages) |_| { var params: std.ArrayList(client.Param) = .empty; defer params.deinit(a); try params.append(a, .{ .name = "limit", .value = "100" }); if (result.cursor) |cursor| try params.append(a, .{ .name = "cursor", .value = cursor }); const page = rpc.query(struct { feed: []std.json.Value, cursor: ?[]const u8 = null }, "app.bsky.feed.getTimeline", params.items) catch |err| { result.failure = @errorName(err); break; }; defer page.deinit(); result.pages += 1; result.received += page.value.feed.len; for (page.value.feed) |value| { var entry = parse(value) catch { result.unclassified += 1; continue; }; const key = try std.fmt.allocPrint(a, "{s}\n{s}\n{d}\n{any}", .{ entry.uri, entry.introduced_by.did, entry.indexed_at, entry.is_repost }); if (keys.contains(key)) { result.duplicates += 1; continue; } try keys.put(key, {}); entry.uri = try a.dupe(u8, entry.uri); entry.author = try client.copyProfile(a, entry.author); entry.introduced_by = try client.copyProfile(a, entry.introduced_by); try entries.append(a, entry); } result.cursor = if (page.value.cursor) |cursor| try a.dupe(u8, cursor) else null; if (result.cursor) |cursor| { if (cursors.contains(cursor)) { result.failure = "RepeatedCursor"; break; } try cursors.put(cursor, {}); } else { result.exhausted = true; break; } } result.entries = try entries.toOwnedSlice(a); result.observed_at = std.Io.Timestamp.now(rpc.io, .real).toMicroseconds(); return .{ .value = result, .arena = arena };}
test "timeline pagination preserves omissions and stops on a repeated cursor" { const origin_module = @import("test_origin.zig"); const io = std.testing.io; var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const a = arena.allocator(); const entry = \\{"post":{"uri":"at://did:plc:author/app.bsky.feed.post/one","author":{"did":"did:plc:author","handle":"author.test"},"indexedAt":"2026-09-01T00:00:00Z","record":{"$type":"app.bsky.feed.post"}}} ; const first = try std.fmt.allocPrint(a, "{{\"feed\":[{s}],\"cursor\":\"same\"}}", .{entry}); const second = try std.fmt.allocPrint(a, "{{\"feed\":[{s},{{\"unknown\":true}}],\"cursor\":\"same\"}}", .{entry}); const replies: []const []const u8 = &.{ try origin_module.response(a, first), try origin_module.response(a, second) }; var origin = try origin_module.Origin.init(io, replies); defer origin.deinit(); var serving = try io.concurrent(origin_module.Origin.serve, .{&origin}); defer _ = serving.cancel(io) catch {}; var rpc = client.Client.init(io, a, try origin.url(a)); defer rpc.deinit(); rpc.session = .{ .did = "did:plc:viewer", .handle = "viewer.test", .accessJwt = "test-only", .refreshJwt = "test-only" }; rpc.appview = "did:web:appview.test#bsky_appview"; var owned = try sample(&rpc, 5); defer owned.deinit(); const result = owned.value; try serving.await(io); try std.testing.expectEqual(@as(usize, 1), result.entries.len); try std.testing.expectEqual(@as(usize, 3), result.received); try std.testing.expectEqual(@as(usize, 1), result.duplicates); try std.testing.expectEqual(@as(usize, 1), result.unclassified); try std.testing.expectEqualStrings("RepeatedCursor", result.failure.?); try std.testing.expect(!result.exhausted); try std.testing.expectEqualStrings(rpc.appview.?, result.appview);}