Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
9.6 kB · 246 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247import fs from "node:fs/promises";import path from "node:path";import { afterEach, describe, expect, test } from "vitest";import { buildActivityBatchContextPacket } from "../src/agents/context.js";import { loadAgentDeclarations } from "../src/agents/declarations.js";import { DeterministicBatcher } from "../src/batches/runtime.js";import type { EventCandidate } from "../src/events/types.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import type { BatchDeclarationManifest } from "../src/runtime/manifest.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("broad activity batch context", () => { test("reconstructs exact mixed-source members under one most-private snapshot", async () => { const root = await temporaryProject("thoughtstream-activity-context-"); roots.push(root); const store = await testStore(root); stores.push(store); const first = (await store.appendEvent(sourceEvent("jetstream:cameron-atproto", "jetstream", "one", "public-source"))).event; const second = (await store.appendEvent(sourceEvent("x:cameron-private", "x-webhook", "two", "sensitive"))).event; const batchDeclaration: BatchDeclarationManifest = { id: "stream-activity-ten-minute", version: 1, enabled: true, input: { eventTypes: ["stream.thought.source.atproto.commit"], sourceIds: [first.source, second.source], }, output: { sourceId: "batch:stream-activity", eventType: "stream.thought.derived.event.batch" }, quietWindowMs: 100, maxAgeMs: 100, maxItems: 200, privacy: "most-private", replay: "beginning", pollIntervalMs: 100, }; const result = await new DeterministicBatcher(store, () => new Date("2026-08-11T00:00:01.000Z")).cycle(batchDeclaration); const batch = await store.getEvent(result.batchEventId!); if (!batch) throw new Error("Missing activity batch"); const declarations = await loadAgentDeclarations(path.join(process.cwd(), "agents"), testDeclarationEnvironment); const resident = declarations.find((declaration) => declaration.id === "resident-letta-conversation"); if (!resident) throw new Error("Missing resident Stream declaration");
const packet = await buildActivityBatchContextPacket(store, resident, batch); expect(packet.manifest).toMatchObject({ contextStrategy: "activity-batch", inputEventIds: [batch.id], includedEventIds: [first.id, second.id], memberCount: 2, privacy: "sensitive", }); expect(packet.text).toContain("PUBLIC MEMBER BODY"); expect(packet.text).toContain("SENSITIVE MEMBER BODY"); expect(packet.text).toContain("privateEvidence"); expect(packet.text).toContain('authority="untrusted-data"'); const replay = await buildActivityBatchContextPacket(store, resident, batch); expect(replay).toEqual(packet); });
test("admits only exact completed source-listener outputs to social synthesis", async () => { const root = await temporaryProject("thoughtstream-social-context-"); roots.push(root); const store = await testStore(root); stores.push(store); const bluesky = await appendListenerObservation( store, "cameron-bluesky-listener", "bluesky", "public-source", "Cameron's Bluesky posts are clustering around persistent-agent evaluation.", ); const x = await appendListenerObservation( store, "cameron-x-listener", "x", "sensitive", "Cameron's X likes show attention to rollbackable policy releases.", ); const batchDeclaration = socialBatchDeclaration([bluesky.source, x.source]); const result = await new DeterministicBatcher(store, () => new Date("2026-08-11T00:02:00.000Z")) .cycle(batchDeclaration); const batch = await store.getEvent(result.batchEventId!); if (!batch) throw new Error("Missing social observation batch"); const declarations = await loadAgentDeclarations(path.join(process.cwd(), "agents"), testDeclarationEnvironment); const social = declarations.find((declaration) => declaration.id === "cameron-social-listener"); if (!social) throw new Error("Missing social listener declaration");
const packet = await buildActivityBatchContextPacket(store, social, batch); expect(packet.manifest).toMatchObject({ contextStrategy: "activity-batch", includedEventIds: [bluesky.id, x.id], privacy: "sensitive", completedAgentOutputs: [ { agentId: "cameron-bluesky-listener", outputEventId: bluesky.id }, { agentId: "cameron-x-listener", outputEventId: x.id }, ], }); expect(packet.text).toContain("persistent-agent evaluation"); expect(packet.text).toContain("rollbackable policy releases"); });
test("rejects an observation batch member without a completed source-listener run", async () => { const root = await temporaryProject("thoughtstream-social-context-forged-"); roots.push(root); const store = await testStore(root); stores.push(store); const forged = await appendListenerObservation( store, "cameron-x-listener", "forged", "sensitive", "An observation with no completed run.", false, ); const result = await new DeterministicBatcher(store, () => new Date("2026-08-11T00:02:00.000Z")) .cycle(socialBatchDeclaration([forged.source])); const batch = await store.getEvent(result.batchEventId!); if (!batch) throw new Error("Missing forged social batch"); const declarations = await loadAgentDeclarations(path.join(process.cwd(), "agents"), testDeclarationEnvironment); const social = declarations.find((declaration) => declaration.id === "cameron-social-listener"); if (!social) throw new Error("Missing social listener declaration"); await expect(buildActivityBatchContextPacket(store, social, batch)) .rejects.toThrow("run receipt is unavailable"); });});
function socialBatchDeclaration(sourceIds: string[]): BatchDeclarationManifest { return { id: "cameron-social-observation-window", version: 1, enabled: true, input: { eventTypes: ["stream.thought.derived.message.observation"], sourceIds, }, output: { sourceId: "batch:cameron-social-observations", eventType: "stream.thought.derived.event.batch" }, quietWindowMs: 100, maxAgeMs: 100, maxItems: 20, privacy: "most-private", replay: "beginning", pollIntervalMs: 100, };}
async function appendListenerObservation( store: JazzThoughtStore, agentId: "cameron-bluesky-listener" | "cameron-x-listener", key: string, privacy: "public-source" | "sensitive", summary: string, completed = true,) { const trigger = (await store.appendEvent(sourceEvent( `source:${key}`, "system", key, privacy, ))).event; const runId = `run-${key}`; const executionKey = `execution:${agentId}:${key}`; const result = { summary, tags: [key, "social-listener"], importance: "normal", confidence: 0.8, }; const output = (await store.appendEvent({ type: "stream.thought.derived.message.observation", schemaVersion: 1, source: `agent:${agentId}`, sourceKind: "agent", externalId: runId, idempotencyKey: `${runId}:output`, occurredAt: "2026-08-11T00:01:00.000Z", observedAt: "2026-08-11T00:01:00.000Z", actor: agentId, rootEventId: trigger.rootEventId, parentEventId: trigger.id, correlationId: trigger.correlationId, privacy, payload: { runId, executionKey, inputEventId: trigger.id, inputSourceSequence: trigger.sourceSequence, ...result, }, traceId: runId, })).event; if (completed) { await store.upsertRun({ id: runId, executionKey, triggerEventId: trigger.id, agentId, agentVersion: 1, status: "completed", inputEventIds: [trigger.id], outputEventIds: [output.id], attempt: 1, provider: "letta-cloud", model: "chatgpt-plus-pro/gpt-5.6-luna", promptHash: `prompt-${key}`, contextManifest: {}, result, createdAt: "2026-08-11T00:00:30.000Z", startedAt: "2026-08-11T00:00:30.000Z", completedAt: "2026-08-11T00:01:00.000Z", updatedAt: "2026-08-11T00:01:00.000Z", }); } return output;}
function sourceEvent(source: string, sourceKind: EventCandidate["sourceKind"], key: string, privacy: "public-source" | "sensitive"): EventCandidate { return { type: "stream.thought.source.atproto.commit", schemaVersion: 1, source, sourceKind, externalId: `at://did:plc:fixture/app.bsky.feed.post/${key}`, idempotencyKey: `${source}:${key}`, occurredAt: `2026-08-11T00:00:00.${key === "one" ? "000" : "100"}Z`, observedAt: `2026-08-11T00:00:00.${key === "one" ? "000" : "100"}Z`, actor: "did:plc:fixture", correlationId: `activity-${key}`, privacy, payload: { atUri: `at://did:plc:fixture/app.bsky.feed.post/${key}`, collection: "app.bsky.feed.post", operation: "create", record: { text: key === "one" ? "PUBLIC MEMBER BODY" : "SENSITIVE MEMBER BODY", privateEvidence: key === "two" ? "bounded private context" : "none", }, }, };}