import fs from "node:fs/promises"; import { afterEach, describe, expect, test, vi } from "vitest"; import { buildAtprotoBatchContextPacket } from "../src/agents/context.js"; import type { ThoughtAgentDeclaration } from "../src/agents/types.js"; import { DeterministicBatcher } from "../src/batches/runtime.js"; import type { JsonObject } from "../src/core/json.js"; import type { EventCandidate, ThoughtEvent } from "../src/events/types.js"; import type { JazzThoughtStore } from "../src/jazz/store.js"; import type { BatchDeclarationManifest } from "../src/runtime/manifest.js"; import { temporaryProject, testStore } from "./helpers.js"; const stores: JazzThoughtStore[] = []; const roots: string[] = []; afterEach(async () => { vi.restoreAllMocks(); 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("trusted ATProto batch context compiler", () => { test("loads exact members, gives Bluesky dual views, and never sends Semble to bsky", async () => { const store = await setup(); await store.appendEvent(post()); await store.appendEvent(semble()); const batch = await flush(store); const bskyUris: string[] = []; const atprotoTargets: string[] = []; const packet = await buildAtprotoBatchContextPacket(store, declaration(), batch, { fetchAtprotoDocument: async (request) => { atprotoTargets.push(request.target); return { markdown: `ATPROTO ${request.target} ${request.atUri}`, details: { atUri: request.atUri } }; }, fetchBskyDocument: async (request) => { bskyUris.push(request.atUri); return { markdown: `BSKY ${request.atUri}`, details: { atUri: request.atUri } }; }, }); expect(packet.manifest).toMatchObject({ contextStrategy: "atproto-batch", memberCount: 2 }); expect(packet.manifest.includedEventIds).toEqual((batch.payload.members as Array<{ eventId: string }>).map((member) => member.eventId)); expect(packet.text).toContain("thoughtstream-bluesky-social"); expect(packet.text).toContain("thoughtstream-atproto-card"); expect(atprotoTargets.sort()).toEqual(["card", "collection", "event", "link"]); expect(bskyUris).toEqual(["at://did:plc:fixture/app.bsky.feed.post/one"]); expect(packet.text.length).toBeLessThanOrEqual(declaration().maxInputChars); }); test("snapshots once and reuses the exact packet on retry", async () => { const store = await setup(); await store.appendEvent(post()); const batch = await flush(store); let fetches = 0; const options = { fetchAtprotoDocument: async (request: { atUri: string }) => { fetches += 1; return { markdown: "FIRST SOURCE VIEW", details: { atUri: request.atUri } }; }, fetchBskyDocument: async (request: { atUri: string }) => { fetches += 1; return { markdown: "FIRST SOCIAL VIEW", details: { atUri: request.atUri } }; }, }; const first = await buildAtprotoBatchContextPacket(store, declaration(), batch, options as never); const second = await buildAtprotoBatchContextPacket(store, declaration(), batch, { fetchAtprotoDocument: async () => { throw new Error("must not refetch"); }, fetchBskyDocument: async () => { throw new Error("must not refetch"); }, }); expect(fetches).toBe(2); expect(second).toEqual(first); }); test("fails closed for missing or mismatched referenced members", async () => { const store = await setup(); await store.appendEvent(post()); const batch = await flush(store); const missing: ThoughtEvent = { ...batch, id: `${batch.id}-missing`, payload: { ...batch.payload, members: [{ ...(batch.payload.members as Array>)[0]!, eventId: "evt_missing" }], }, }; await expect(buildAtprotoBatchContextPacket(store, declaration(), missing)).rejects.toThrow("missing or mismatched"); const mismatch: ThoughtEvent = { ...batch, id: `${batch.id}-mismatch`, payload: { ...batch.payload, members: [{ ...(batch.payload.members as Array>)[0]!, sourceSequence: 99 }], }, }; await expect(buildAtprotoBatchContextPacket(store, declaration(), mismatch)).rejects.toThrow("missing or mismatched"); }); test("fails before any model boundary when fixed member envelopes do not fit", async () => { const store = await setup(); await store.appendEvent(post()); await store.appendEvent({ ...post(), externalId: "at://did:plc:fixture/app.bsky.feed.post/two", idempotencyKey: "post-two" }); const batch = await flush(store); const tiny = { ...declaration(), maxInputChars: 2_500 }; const fetch = vi.fn(); await expect(buildAtprotoBatchContextPacket(store, tiny, batch, { fetchAtprotoDocument: fetch, fetchBskyDocument: fetch, })).rejects.toThrow("fixed member envelopes exceed"); expect(fetch).not.toHaveBeenCalled(); }); }); async function setup(): Promise { const root = await temporaryProject("thoughtstream-atproto-batch-context-"); roots.push(root); const store = await testStore(root); stores.push(store); return store; } async function flush(store: JazzThoughtStore): Promise { const declaration: BatchDeclarationManifest = { id: "context-batch", version: 1, enabled: true, input: { eventTypes: ["stream.thought.source.atproto.commit"], sourceIds: ["jetstream:fixture"] }, output: { sourceId: "batch:fixture", eventType: "stream.thought.derived.event.batch" }, quietWindowMs: 100, maxAgeMs: 1_000, maxItems: 10, privacy: "preserve", replay: "beginning", pollIntervalMs: 100, }; const result = await new DeterministicBatcher(store, () => new Date("2026-07-21T00:00:10.000Z")).cycle(declaration); return (await store.getEvent(result.batchEventId!))!; } function declaration(): ThoughtAgentDeclaration { return { id: "resident-fixture", version: 1, name: "Resident fixture", description: "Fixture", mode: "letta-agent-sdk", provider: "letta-cloud", lettaAgent: { backend: "cloud", agentIdEnv: "FIXTURE_AGENT_ID", agentId: "agent-fixture", conversation: "main", responseMode: "conversation-text", outputOnly: false, permissionMode: "unrestricted", dreaming: { trigger: "off" }, sandbox: { ttlMinutes: 5, terminateOnClose: false }, }, eventTypes: ["stream.thought.derived.event.batch"], compiledEventTypes: ["stream.thought.derived.event.batch"], sourcePatterns: ["batch:fixture"], acceptedPrivacy: ["public-source"], outputEventType: "stream.thought.derived.message.observation", emit: ["stream.thought.derived.message.observation"], promptRef: "fixture", systemPrompt: "Observe.", enabled: true, maxEvents: 1, maxInputChars: 20_000, contextStrategy: "atproto-batch", atprotoObjectContext: true, maxOutputTokens: 100, timeoutMs: 60_000, tools: [], externalActions: false, }; } function post(): EventCandidate { return commit("post-one", "app.bsky.feed.post", { text: "Synthetic post", }); } function semble(): EventCandidate { return commit("semble-one", "network.cosmik.collectionLink", { card: { uri: "at://did:plc:card/network.cosmik.card/one", cid: "bafy-card" }, collection: { uri: "at://did:plc:collection/network.cosmik.collection/one", cid: "bafy-collection" }, }); } function commit(key: string, collection: string, record: JsonObject): EventCandidate { return { type: "stream.thought.source.atproto.commit", schemaVersion: 1, source: "jetstream:fixture", sourceKind: "jetstream", externalId: `at://did:plc:fixture/${collection}/${key === "post-one" ? "one" : key}`, idempotencyKey: key, occurredAt: key === "post-one" ? "2026-07-21T00:00:00.000Z" : "2026-07-21T00:00:01.000Z", observedAt: key === "post-one" ? "2026-07-21T00:00:00.000Z" : "2026-07-21T00:00:01.000Z", actor: "did:plc:fixture", correlationId: "fixture", privacy: "public-source", payload: { atUri: `at://did:plc:fixture/${collection}/${key === "post-one" ? "one" : key}`, cid: `bafy-${key}`, collection, operation: "create", record, }, }; }