import fs from "node:fs/promises"; import path from "node:path"; import { afterEach, describe, expect, test } from "vitest"; import { loadAgentDeclarations } from "../src/agents/declarations.js"; import { ThoughtAgentRuntime } from "../src/agents/runtime.js"; import { JetstreamConnector } from "../src/connectors/jetstream.js"; import type { JazzThoughtStore } from "../src/jazz/store.js"; import { buildRootActivity } from "../src/projections/activity.js"; import { temporaryProject, testDeclarationEnvironment, testStore } from "./helpers.js"; const stores: JazzThoughtStore[] = []; const roots: string[] = []; afterEach(async () => { await Promise.all(stores.splice(0).map((store) => store.close())); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true }))); }); describe("Bluesky activity observer", () => { test("turns a Jetstream post into a legible rule result and activity lineage", async () => { const project = await temporaryProject(); roots.push(project); await fs.cp(path.join(process.cwd(), "agents"), path.join(project, "agents"), { recursive: true }); await fs.cp(path.join(process.cwd(), "prompts"), path.join(project, "prompts"), { recursive: true }); const store = await testStore(project); stores.push(store); const declarations = await loadAgentDeclarations(path.join(project, "agents"), testDeclarationEnvironment); const runtime = new ThoughtAgentRuntime(store); const consumers = await runtime.startConsumers(declarations); const connector = new JetstreamConnector({ id: "jetstream:public-posts", collections: ["app.bsky.feed.post"] }); const ingest = await connector.ingestBatch(store, [{ did: "did:plc:alicefixture", time_us: 1784042401000000, kind: "commit", commit: { rev: "3mfixture001", operation: "create", collection: "app.bsky.feed.post", rkey: "post-one", cid: "bafyfixturepostone", record: { $type: "app.bsky.feed.post", text: "A public post with a link", facets: [{ index: { byteStart: 19, byteEnd: 23 }, features: [] }], }, }, }]); await waitFor(async () => (await store.listRuns()).some((run) => run.agentId === "bluesky-activity-observer")); await consumers.drain(); await consumers.stop(); expect(ingest.inserted).toBe(1); const [run] = await store.listRuns(); expect(run).toMatchObject({ agentId: "bluesky-activity-observer", status: "completed", result: { summary: "Create post: A public post with a link", tags: ["atproto", "app.bsky.feed.post", "create", "post", "facets"], }, }); const activity = await buildRootActivity(store); expect(activity.items.find((item) => item.id === ingest.events[0]?.id)).toMatchObject({ summary: "Create post: A public post with a link", presentation: { title: "Posted on Bluesky", body: "A public post with a link", objectLabel: "Bluesky post", url: "https://bsky.app/profile/did:plc:alicefixture/post/post-one", urlLabel: "Open post", }, descendantEventCount: 4, consumerRuns: [{ agentId: "bluesky-activity-observer", kind: "rule", status: "completed", outputCount: 1, }], }); }); test("projects a Bluesky like as a useful list label without exposing the raw payload", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new JetstreamConnector({ id: "jetstream:cameron-atproto", collections: ["app.bsky.feed.like"] }); const ingest = await connector.ingestBatch(store, [{ did: "did:plc:cameronfixture", time_us: 1786405008000000, kind: "commit", commit: { rev: "3mfixturelike", operation: "create", collection: "app.bsky.feed.like", rkey: "like-one", cid: "bafyfixturelikeone", record: { $type: "app.bsky.feed.like", subject: { uri: "at://did:plc:authorfixture/app.bsky.feed.post/post-one", cid: "bafyfixturepost" }, createdAt: "2026-08-10T22:36:48.000Z", }, }, }]); const activity = await buildRootActivity(store); expect(activity.items.find((item) => item.id === ingest.events[0]?.id)).toMatchObject({ presentation: { title: "Liked a post", body: "", objectLabel: "Bluesky post", url: "https://bsky.app/profile/did:plc:authorfixture/post/post-one", urlLabel: "Open liked post", }, }); }); }); async function waitFor(predicate: () => Promise, timeoutMs = 2_000): Promise { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (await predicate()) return; await new Promise((resolve) => setTimeout(resolve, 10)); } throw new Error("Timed out waiting for the ATProto observer"); }