Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
15 kB · 382 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383import type { LettaCodeClientSessionOptions, LettaCodeSession, ListMessagesResult, SDKMessage, SendMessage,} from "@letta-ai/letta-agent-sdk";import fs from "node:fs/promises";import { afterEach, describe, expect, test, vi } from "vitest";import { LETTA_AGENT_SDK_ADAPTER_REVISION, LettaAgentSdkRunner, type LettaAgentSdkClient, type LettaRunUsageClient,} from "../src/agents/letta-agent-sdk.js";import { ThoughtAgentRuntime } from "../src/agents/runtime.js";import { DeterministicBatcher } from "../src/batches/runtime.js";import type { AgentOutput, AgentRunInput, AgentRunner, RunnerTrace, ThoughtAgentDeclaration } from "../src/agents/types.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { temporaryProject, 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("Letta Agent SDK runtime integration", () => { test("settles event, output, provenance, traces, progress, and accounting through Jazz", async () => { const project = await temporaryProject("thoughtstream-letta-sdk-runtime-"); roots.push(project); const store = await testStore(project); stores.push(store); const sourceBody = "PRIVATE-NEW-EVENT-BODY"; await store.appendEvent({ type: "stream.thought.source.telegram.message", schemaVersion: 1, source: "telegram:letta-runtime-fixture", sourceKind: "telegram", externalId: "message-1", idempotencyKey: "message-1", occurredAt: "2026-07-21T00:00:00.000Z", actor: "telegram-user", correlationId: "chat-fixture", privacy: "sensitive", payload: { text: sourceBody }, }); const session = new RuntimeFakeSession([ { type: "assistant", content: "PRIVATE-MODEL-REPLY", uuid: "assistant-1", runId: "run-sdk-1" }, { type: "result", success: true, result: "Runtime reply.", durationMs: 250, conversationId: "conv-runtime-fixture", runIds: ["run-sdk-1"], }, ]); const client: LettaAgentSdkClient = { resumeSession: (_id, options) => { session.options = options; return session as unknown as LettaCodeSession; }, }; const runUsageClient: LettaRunUsageClient = { retrieve: async () => ({ inputTokens: 321, outputTokens: 45 }), }; const declaration = runtimeDeclaration(); const runtime = new ThoughtAgentRuntime(store, [new LettaAgentSdkRunner({ client, runUsageClient, reconciliationDelayMs: 0, })]);
const [result] = await runtime.consumeBacklog([declaration]);
expect(result).toMatchObject({ output: { summary: "Runtime reply.", model: { provider: "letta-cloud" } } }); expect(session.sent).toHaveLength(1); expect(session.sent[0]).toContain(sourceBody); expect(session.closed).toBe(true); expect(session.options).toMatchObject({ permissionMode: "unrestricted", sandbox: { ttlMinutes: 5, terminateOnClose: false }, }); const [run] = await store.listRuns(); expect(run).toMatchObject({ status: "completed", provider: "letta-cloud", model: "agent-default", checkpointRevision: expect.stringContaining(LETTA_AGENT_SDK_ADAPTER_REVISION), executionAdapterRevision: LETTA_AGENT_SDK_ADAPTER_REVISION, attempt: 1, }); const traces = await store.listTrace(run!.id); const durable = JSON.stringify(traces); expect(traces.map((trace) => trace.type)).toEqual(expect.arrayContaining([ "letta.session.opened", "letta.prompt", "letta.turn.sent", "letta.assistant", "letta.result", "letta.usage", ])); expect(durable).not.toContain(sourceBody); expect(durable).not.toContain("PRIVATE-MODEL-REPLY"); expect(durable).not.toContain("Runtime reply."); const [accounting] = await store.listInferenceAccounting({ agentId: declaration.id }); expect(accounting).toMatchObject({ status: "settled", usageStatus: "reported", actualUsage: { inputTokens: 321, outputTokens: 45 }, estimate: { calls: 1, inputTokens: 1_000, outputTokens: 100 }, charged: { calls: 1, inputTokens: 321, outputTokens: 45 }, }); const events = await store.listEvents(); expect(events.filter((event) => event.type === declaration.outputEventType)).toHaveLength(1); expect(events.filter((event) => event.type === "stream.thought.agent.run.completed")).toHaveLength(1); expect(JSON.stringify({ run, traces, accounting })).not.toContain(sourceBody); });
test("runs one resident turn for an ATProto burst and serializes it independently with Telegram", async () => { const project = await temporaryProject("thoughtstream-letta-sdk-multi-source-"); roots.push(project); const store = await testStore(project); stores.push(store); vi.spyOn(store, "subscribeConsumerEvents").mockImplementation(() => () => undefined); const declaration: ThoughtAgentDeclaration = { ...runtimeDeclaration(), version: 2, eventTypes: ["stream.thought.source.telegram.message", "stream.thought.derived.event.batch"], compiledEventTypes: ["stream.thought.source.telegram.message", "stream.thought.derived.event.batch"], sourcePatterns: ["telegram:thoughtstream-bot", "batch:cameron-atproto"], acceptedPrivacy: ["sensitive", "public-source"], initialReplay: "beginning", payloadFields: ["text", "atUri", "cid", "collection", "operation", "record"], contextStrategy: "atproto-batch", maxInputChars: 12_000, atprotoObjectContext: true, }; const probe = new ResidentConcurrencyProbe(); let bskyFetches = 0; const atprotoTargets: string[] = []; const runtime = new ThoughtAgentRuntime(store, [probe], { reconcileIntervalMs: 10, atprotoObjectContext: { fetchAtprotoDocument: async (request) => { atprotoTargets.push(request.target); return { markdown: `SYNTHETIC ${request.target.toUpperCase()} CONTEXT`, details: { atUri: request.atUri }, }; }, fetchBskyDocument: async (request) => { bskyFetches += 1; return { markdown: "SYNTHETIC BLUESKY SOCIAL CONTEXT", details: { atUri: request.atUri }, }; }, }, }); const atproto = (await store.appendEvent({ type: "stream.thought.source.atproto.commit", schemaVersion: 1, source: "jetstream:cameron-bluesky", sourceKind: "jetstream", externalId: "at://did:plc:cameron/app.bsky.feed.post/post-one", idempotencyKey: "post-one", occurredAt: "2026-07-21T00:00:00.000Z", actor: "did:plc:cameron", correlationId: "jetstream-fixture", privacy: "public-source", payload: { atUri: "at://did:plc:cameron/app.bsky.feed.post/post-one", collection: "app.bsky.feed.post", operation: "create", record: { text: "Public post fixture" }, }, })).event; await store.appendEvent({ type: "stream.thought.source.atproto.commit", schemaVersion: 1, source: "jetstream:cameron-bluesky", sourceKind: "jetstream", externalId: "at://did:plc:fixtureowner/network.cosmik.collectionLink/link-runtime-fixture", idempotencyKey: "collection-link-runtime-fixture", occurredAt: "2026-07-21T00:00:00.500Z", actor: "did:plc:fixtureowner", correlationId: "jetstream-fixture", privacy: "public-source", payload: { atUri: "at://did:plc:fixtureowner/network.cosmik.collectionLink/link-runtime-fixture", cid: "bafy-runtime-link", collection: "network.cosmik.collectionLink", operation: "create", record: { $type: "network.cosmik.collectionLink", card: { uri: "at://did:plc:fixturecard/network.cosmik.card/card-runtime-fixture", cid: "bafy-runtime-card", }, collection: { uri: "at://did:plc:fixturecollection/network.cosmik.collection/collection-runtime-fixture", cid: "bafy-runtime-collection", }, }, }, }); const batch = new DeterministicBatcher(store, () => new Date("2026-07-21T00:00:10.000Z")); const batchResult = await batch.cycle({ id: "resident-atproto-fixture", version: 1, enabled: true, input: { eventTypes: ["stream.thought.source.atproto.commit"], sourceIds: ["jetstream:cameron-bluesky"] }, output: { sourceId: "batch:cameron-atproto", eventType: "stream.thought.derived.event.batch" }, quietWindowMs: 100, maxAgeMs: 1_000, maxItems: 2, privacy: "preserve", replay: "beginning", pollIntervalMs: 100, }); expect(batchResult.emitted).toBe(1); await store.appendEvent({ type: "stream.thought.source.telegram.message", schemaVersion: 1, source: "telegram:thoughtstream-bot", sourceKind: "telegram", externalId: "message-one", idempotencyKey: "message-one", occurredAt: "2026-07-21T00:00:01.000Z", actor: "telegram-user", correlationId: "telegram-fixture", privacy: "sensitive", payload: { text: "Private Telegram fixture" }, }); await runtime.consumeBacklog([declaration]); await runtime.consumeBacklog([declaration]);
expect(probe.sources.sort()).toEqual([ "batch:cameron-atproto", "telegram:thoughtstream-bot", ]); expect(probe.contexts.join("\n")).toContain("thoughtstream-atproto-card"); expect(probe.contexts.join("\n")).toContain("SYNTHETIC COLLECTION CONTEXT"); expect(atprotoTargets.sort()).toEqual(["card", "collection", "event", "link"]); expect(bskyFetches).toBe(1); expect(probe.maximumActive).toBe(1); const progress = (await store.listConsumerProgress()) .filter((entry) => entry.consumerId === declaration.id && entry.consumerVersion === declaration.version); expect(progress.map((entry) => entry.source).sort()).toEqual([ "batch:cameron-atproto", "telegram:thoughtstream-bot", ]); expect(progress.every((entry) => entry.lastSequence > 0)).toBe(true); expect((await store.listRuns()).filter((run) => run.agentVersion === declaration.version)).toHaveLength(2); const atprotoBatch = (await store.listEvents({ source: "batch:cameron-atproto" }))[0]!; const atprotoDerived = (await store.listEvents()) .filter((event) => event.parentEventId === atprotoBatch.id); expect(atprotoDerived.filter((event) => event.type === declaration.outputEventType)) .toEqual([expect.objectContaining({ privacy: "sensitive" })]); expect(atprotoDerived.filter((event) => event.type.startsWith("stream.thought.agent.run."))) .toEqual(expect.arrayContaining([expect.objectContaining({ privacy: "sensitive" })])); });});
class ResidentConcurrencyProbe implements AgentRunner { readonly mode = "letta-agent-sdk" as const; readonly sources: string[] = []; readonly contexts: string[] = []; maximumActive = 0; private active = 0;
async run( input: AgentRunInput, _onTrace: (trace: RunnerTrace) => Promise<void>, ): Promise<AgentOutput> { this.active += 1; this.maximumActive = Math.max(this.maximumActive, this.active); this.sources.push(input.event.source); this.contexts.push(input.context.text); await new Promise((resolve) => setTimeout(resolve, 25)); this.active -= 1; return { summary: `Observed ${input.event.source}`, tags: ["conversation"], importance: "normal", confidence: 1, }; }}
async function waitFor(predicate: () => boolean, timeoutMs = 2_000): Promise<void> { const deadline = Date.now() + timeoutMs; while (!predicate()) { if (Date.now() >= deadline) throw new Error("Timed out waiting for resident source reconciliation"); await new Promise((resolve) => setTimeout(resolve, 10)); }}
class RuntimeFakeSession { readonly agentId = "agent-runtime-fixture"; readonly sessionId = "session-runtime-fixture"; readonly conversationId = "conv-runtime-fixture"; readonly sent: string[] = []; options: LettaCodeClientSessionOptions | undefined; closed = false;
constructor(private readonly messages: SDKMessage[]) {}
async send(message: SendMessage): Promise<void> { if (typeof message !== "string") throw new Error("Fixture accepts text only"); this.sent.push(message); }
async *stream(): AsyncGenerator<SDKMessage> { for (const message of this.messages) yield message; }
async abort(): Promise<void> {}
async listMessages(): Promise<ListMessagesResult> { return { messages: [], nextBefore: null, hasMore: false }; }
close(): void { this.closed = true; }}
function runtimeDeclaration(): ThoughtAgentDeclaration { return { id: "letta-runtime-fixture", version: 1, name: "Letta runtime fixture", description: "End-to-end SDK runtime fixture", mode: "letta-agent-sdk", provider: "letta-cloud", lettaAgent: { backend: "cloud", agentIdEnv: "THOUGHTSTREAM_LETTA_RUNTIME_FIXTURE_AGENT_ID", agentId: "agent-runtime-fixture", conversation: "main", responseMode: "conversation-text", outputOnly: false, permissionMode: "unrestricted", dreaming: { trigger: "off" }, sandbox: { ttlMinutes: 5, terminateOnClose: false }, }, eventTypes: ["stream.thought.source.telegram.message"], compiledEventTypes: ["stream.thought.source.telegram.message"], sourcePatterns: ["telegram:letta-runtime-fixture"], acceptedPrivacy: ["sensitive"], initialReplay: "beginning", outputEventType: "stream.thought.derived.message.observation", emit: ["stream.thought.derived.message.observation"], promptRef: "fixture", systemPrompt: "Reply to the current message.", enabled: true, maxEvents: 1, maxInputChars: 2_000, contextStrategy: "single-event", payloadFields: ["text"], maxOutputTokens: 100, timeoutMs: 60_000, accounting: { leaseMs: 180_000, reservation: { inputTokens: 1_000, outputTokens: 100 }, limits: [{ window: "rolling", durationMs: 3_600_000, maxCalls: 100, maxInputTokens: 100_000, maxOutputTokens: 10_000, }], }, tools: [], externalActions: false, };}