Something went wrong. Try again.
A Zig library for Bluesky activity and feed analytics, built on Zat.
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229const std = @import("std");const tend = @import("root.zig");const paths_module = @import("paths.zig");const store = @import("store.zig");const usage = \\tend — evidence and analysis for Bluesky \\tend paths \\tend probe HOST ACTOR \\tend whoami \\tend follows \\tend mutes \\tend timeline [PAGES] \\tend feed-report FILE \\tend archive-report FORMAT FILE DID UNTIL_MICROS [DAYS] \\tend archive-evidence FORMAT FILE DID UNTIL_MICROS [DAYS] \\tend archive-follows FORMAT FILE DID \\tend mirror-fetch HOST DID [DIRECTORY] \\tend survey HOST FOLLOWS.json [DAYS [CONCURRENCY]] \\tend survey-resume MANIFEST.json [CONCURRENCY] \\tend survey-report HOST FOLLOWS.json UNTIL_MICROS [DAYS] \\ \\FORMAT is car or star (STAR v1). Hubble collection uses CAR. \\Authenticated commands use ATPROTO_HANDLE, ATPROTO_PASSWORD and TEND_APPVIEW. \\TEND_REQUESTS_PER_SECOND optionally sets shared request pacing. \\Reports are offline. No commands change your account. \\;const FollowSet = struct { did: []const u8, revision: []const u8, subjects: []const []const u8 };
pub fn main(init: std.process.Init) u8 { run(init) catch |err| { std.debug.print("tend: {t}\n", .{err}); return 1; }; return 0;}
fn run(init: std.process.Init) !void { const a = init.arena.allocator(); const io = init.io; var args = std.process.Args.Iterator.init(init.minimal.args); _ = args.skip(); const command = args.next() orelse return std.Io.File.stdout().writeStreamingAll(io, usage); if (eq(command, "--help") or eq(command, "help")) return std.Io.File.stdout().writeStreamingAll(io, usage); if (eq(command, "paths")) return emit(io, a, try localPaths(init)); if (eq(command, "archive-report") or eq(command, "archive-evidence") or eq(command, "archive-follows")) { const format = try parseFormat(args.next() orelse return error.MissingFormat); const path = args.next() orelse return error.MissingPath; const did = args.next() orelse return error.MissingDid; const follows = eq(command, "archive-follows"); const observed = std.Io.Timestamp.now(io, .real).toMicroseconds(); const until = if (follows) observed else try std.fmt.parseInt(i64, args.next() orelse return error.MissingWindowEnd, 10); const days = if (follows) 90 else try positive(args.next() orelse "90"); const window = try windowFor(until, days); if (args.next() != null) return error.UnexpectedArgument; const file = try std.Io.Dir.cwd().openFile(io, path, .{}); defer file.close(io); var buffer: [64 * 1024]u8 = undefined; var input = file.reader(io, &buffer); var analysis = try tend.archive.read(init.gpa, &input.interface, .{ .format = format, .did = did, .window = window, .observed_at = observed, .collect_follows = follows, .retain_evidence = eq(command, "archive-evidence") }); defer analysis.deinit(); const value = analysis.value; if (follows) return emit(io, a, FollowSet{ .did = value.did, .revision = value.revision, .subjects = value.follows }); return emit(io, a, value); } if (eq(command, "mirror-fetch")) { const host = args.next() orelse return error.MissingHost; const did = args.next() orelse return error.MissingDid; const directory = args.next() orelse try mirrorDirectory(init); if (args.next() != null) return error.UnexpectedArgument; try secure(host); _ = try std.Io.Dir.cwd().createDirPathStatus(io, directory, .fromMode(0o700)); var dir = try std.Io.Dir.cwd().openDir(io, directory, .{}); defer dir.close(io); var client = tend.client.Client.init(io, init.gpa, host); defer client.deinit(); try configureRate(init, &client); var result = try tend.hubble.refresh(&client, dir, did, .{ .window = try windowFor(std.Io.Timestamp.now(io, .real).toMicroseconds(), 90) }); defer result.deinit(); return emit(io, a, .{ .directory = directory, .result = result.value }); } if (eq(command, "survey") or eq(command, "survey-report")) { const host = args.next() orelse return error.MissingHost; const path = args.next() orelse return error.MissingPath; const reporting = eq(command, "survey-report"); const until = if (reporting) try std.fmt.parseInt(i64, args.next() orelse return error.MissingWindowEnd, 10) else std.Io.Timestamp.now(io, .real).toMicroseconds(); const window = try windowFor(until, try positive(args.next() orelse "90")); const concurrency = if (reporting) 3 else try positive(args.next() orelse "3"); if (args.next() != null) return error.UnexpectedArgument; const follows = try load(FollowSet, init, path); const directory = try mirrorDirectory(init); if (!reporting) _ = try std.Io.Dir.cwd().createDirPathStatus(io, directory, .fromMode(0o700)); var dir = try std.Io.Dir.cwd().openDir(io, directory, .{}); defer dir.close(io); if (reporting) { var result = try tend.survey.report(io, init.gpa, init.gpa, host, follows.subjects, dir, window); defer result.deinit(); return emit(io, a, .{ .account = follows.did, .follow_revision = follows.revision, .source = host, .window = window, .rows = result.value }); } try secure(host); const manifest = try std.fmt.allocPrint(a, "survey-{d}.json", .{until}); const pending: tend.survey.Run = .{ .account = follows.did, .follow_revision = follows.revision, .source = host, .window = window, .subjects = follows.subjects }; try store.save(io, a, dir, manifest, pending, .create); return executeSurvey(init, dir, directory, manifest, pending, concurrency); } if (eq(command, "survey-resume")) { const path = args.next() orelse return error.MissingPath; const concurrency = try positive(args.next() orelse "3"); if (args.next() != null) return error.UnexpectedArgument; var pending = try load(tend.survey.Run, init, path); try secure(pending.source); try pending.window.validate(); const directory = std.fs.path.dirname(path) orelse "."; var dir = try std.Io.Dir.cwd().openDir(io, directory, .{}); defer dir.close(io); const manifest = std.fs.path.basename(path); pending.finished = false; pending.outcomes = null; try store.save(io, a, dir, manifest, pending, .replace); return executeSurvey(init, dir, directory, manifest, pending, concurrency); } if (eq(command, "feed-report")) { const path = args.next() orelse return error.MissingPath; if (args.next() != null) return error.UnexpectedArgument; const sample = try load(tend.feed.Sample, init, path); const introducers = try tend.feed.rank(init.gpa, sample.entries, .introducer); defer introducers.deinit(init.gpa); const authors = try tend.feed.rank(init.gpa, sample.entries, .author); defer authors.deinit(init.gpa); return emit(io, a, .{ .account = sample.account, .appview = sample.appview, .observed_at = sample.observed_at, .received = sample.received, .unclassified = sample.unclassified, .failure = sample.failure, .introducers = introducers, .authors = authors }); } if (eq(command, "probe")) { const host = args.next() orelse return error.MissingHost; const actor = args.next() orelse return error.MissingActor; if (args.next() != null) return error.UnexpectedArgument; try secure(host); var client = tend.client.Client.init(io, init.gpa, host); defer client.deinit(); try configureRate(init, &client); const response = try client.query(tend.client.Profile, "app.bsky.actor.getProfile", &.{.{ .name = "actor", .value = actor }}); defer response.deinit(); return emit(io, a, .{ .profile = response.value, .response = client.last_response }); } const timeline = eq(command, "timeline"); if (!timeline and !eq(command, "whoami") and !eq(command, "follows") and !eq(command, "mutes")) return error.UnknownCommand; const pages = if (timeline) try positive(args.next() orelse "10") else 0; if (pages > 1000) return error.PageLimitTooLarge; if (args.next() != null) return error.UnexpectedArgument; const handle = init.environ_map.get("ATPROTO_HANDLE") orelse return error.MissingHandle; const password = init.environ_map.get("ATPROTO_PASSWORD") orelse return error.MissingPassword; var identity = try tend.client.resolve(io, init.gpa, handle); defer identity.deinit(); var client = tend.client.Client.init(io, init.gpa, identity.value.pds); defer client.deinit(); try configureRate(init, &client); client.appview = init.environ_map.get("TEND_APPVIEW"); if (!eq(command, "whoami")) try tend.client.validateAppview(client.appview orelse return error.MissingAppview); try client.login(identity.value.did, password); if (eq(command, "whoami")) return emit(io, a, identity.value); if (timeline) { var sample = try tend.feed.sample(&client, pages); defer sample.deinit(); const paths = try localPaths(init); const path = try paths.snapshot(a, identity.value.did, client.appview.?, sample.value.observed_at); _ = try std.Io.Dir.cwd().createDirPathStatus(io, std.fs.path.dirname(path).?, .fromMode(0o700)); try store.save(io, a, std.Io.Dir.cwd(), path, sample.value, .create); return emit(io, a, .{ .saved = path, .entries = sample.value.entries.len, .failure = sample.value.failure }); } var result = if (eq(command, "follows")) try client.follows() else try client.mutes(); defer result.deinit(); return emit(io, a, .{ .account = identity.value.did, .appview = client.appview, .observed_at = std.Io.Timestamp.now(io, .real).toMicroseconds(), .profiles = result.value });}
fn executeSurvey(init: std.process.Init, dir: std.Io.Dir, directory: []const u8, manifest: []const u8, pending: tend.survey.Run, concurrency: usize) !void { var result = pending; var outcomes = try tend.survey.collect(init.io, init.gpa, pending.source, pending.subjects, dir, .{ .window = pending.window,
.concurrency = concurrency, .requests_per_second = if (init.environ_map.get("TEND_REQUESTS_PER_SECOND")) |raw| try std.fmt.parseFloat(f64, raw) else null, }, .{ .on_finished = surveyProgress }); defer outcomes.deinit(); result.outcomes = outcomes.value; result.finished = true; try store.save(init.io, init.gpa, dir, manifest, result, .replace); try emit(init.io, init.arena.allocator(), .{ .saved = try std.fs.path.join(init.arena.allocator(), &.{ directory, manifest }), .result = result });}fn surveyProgress(_: ?*anyopaque, outcome: tend.survey.Outcome) void { std.debug.print("{s}: saved={any}, reused={any}, failure={s}\n", .{ outcome.did, outcome.saved, outcome.reused, outcome.failure orelse "none" });}fn load(comptime T: type, init: std.process.Init, path: []const u8) !T { const a = init.arena.allocator(); const bytes = try std.Io.Dir.cwd().readFileAlloc(init.io, path, a, .limited(256 * 1024 * 1024)); return std.json.parseFromSliceLeaky(T, a, bytes, .{ .ignore_unknown_fields = true });}fn localPaths(init: std.process.Init) !paths_module.Paths { return paths_module.Paths.init(init.arena.allocator(), init.environ_map.get("HOME"), init.environ_map.get("XDG_CONFIG_HOME"), init.environ_map.get("XDG_CACHE_HOME"), init.environ_map.get("XDG_STATE_HOME"));}fn mirrorDirectory(init: std.process.Init) ![]const u8 { return std.fs.path.join(init.arena.allocator(), &.{ (try localPaths(init)).cache, "mirrors" });}fn configureRate(init: std.process.Init, client: *tend.client.Client) !void { if (init.environ_map.get("TEND_REQUESTS_PER_SECOND")) |raw| try client.governor.setRate(try std.fmt.parseFloat(f64, raw));}fn eq(a: []const u8, b: []const u8) bool { return std.mem.eql(u8, a, b);}fn positive(raw: []const u8) !usize { const value = try std.fmt.parseInt(usize, raw, 10); if (value == 0) return error.ExpectedPositiveNumber; return value;}fn parseFormat(raw: []const u8) !tend.archive.Format { return std.meta.stringToEnum(tend.archive.Format, raw) orelse error.InvalidFormat;}fn secure(host: []const u8) !void { if (!std.mem.startsWith(u8, host, "https://")) return error.InsecureHost;}fn windowFor(until: i64, days: usize) !tend.activity.Window { if (days == 0 or days > 36500) return error.InvalidDays; return .{ .since = std.math.sub(i64, until, @as(i64, @intCast(days)) * std.time.us_per_day) catch return error.InvalidWindow, .until = until };}fn emit(io: std.Io, a: std.mem.Allocator, value: anytype) !void { const json = try std.json.Stringify.valueAlloc(a, value, .{ .whitespace = .indent_2 }); defer a.free(json); try std.Io.File.stdout().writeStreamingAll(io, json); try std.Io.File.stdout().writeStreamingAll(io, "\n");}