Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
7.3 kB · 190 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process";import fs from "node:fs/promises";import path from "node:path";import { afterEach, describe, expect, test } from "vitest";import { AGENT_MESSAGE_RESPONSE_EVENT_TYPE, AGENT_MESSAGE_SOURCE_EVENT_TYPE,} from "../src/agents/agent-messages.js";import { declarationFingerprint } from "../src/agents/declarations.js";import { OBSERVATION_OUTPUT_CONTRACT } from "../src/agents/output-contracts.js";import { ThoughtAgentRuntime } from "../src/agents/runtime.js";import type { AgentRunner, ThoughtAgentDeclaration } from "../src/agents/types.js";import { JazzThoughtStore } from "../src/jazz/store.js";import { temporaryProject } from "./helpers.js";
const SECRET = "PRIVATE_CO_TO_STREAM_MESSAGE";const RESPONSE = "PRIVATE_STREAM_TO_CO_RESPONSE";const roots: string[] = [];const stores: JazzThoughtStore[] = [];
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("thought stream agent-send command", () => { test("appends an idempotent body-dark Co message and rejects widened routes", async () => { const root = await temporaryProject("thoughtstream-agent-send-cli-"); roots.push(root); const messageFile = path.join(root, "message.txt"); await fs.writeFile(messageFile, SECRET, { mode: 0o600 }); const args = [ "--silent", "thought", "stream", "agent-send", "--from", "co", "--to", "stream-agent-conversation", "--thread", "co-stream-cli", "--message-id", "co-cli-1", "--file", messageFile, ]; const first = await run(args, { THOUGHTSTREAM_ROOT: root }); const repeated = await run(args, { THOUGHTSTREAM_ROOT: root }); const widened = await run(args.map((value) => value === "stream-agent-conversation" ? "other-agent" : value), { THOUGHTSTREAM_ROOT: root, });
expect(first).toMatchObject({ code: 0, stderr: "" }); expect(repeated).toMatchObject({ code: 0, stderr: "" }); expect(widened.code).toBe(1); expect(first.stdout).not.toContain(SECRET); expect(repeated.stdout).not.toContain(SECRET); expect(JSON.parse(first.stdout)).toMatchObject({ agentMessage: { inserted: true, messageId: "co-cli-1", threadId: "co-stream-cli", senderAgentId: "co", recipientAgentId: "stream-agent-conversation", }, }); expect(JSON.parse(repeated.stdout)).toMatchObject({ agentMessage: { inserted: false } }); }, 20_000);
test("waits for the exact typed response and reveals it only after double opt-in", async () => { const root = await temporaryProject("thoughtstream-agent-send-wait-"); roots.push(root); const messageFile = path.join(root, "message.txt"); await fs.writeFile(messageFile, SECRET, { mode: 0o600 }); const store = await JazzThoughtStore.open({ projectRoot: root }); stores.push(store); const child = start([ "--silent", "thought", "stream", "agent-send", "--from", "co", "--to", "stream-agent-conversation", "--thread", "co-stream-cli-wait", "--message-id", "co-cli-wait-1", "--file", messageFile, "--wait-seconds", "15", "--show-response", "--acknowledge-sensitive-private", ], { THOUGHTSTREAM_ROOT: root }); await waitFor(async () => (await store.listEvents({ types: [AGENT_MESSAGE_SOURCE_EVENT_TYPE], source: "agent-message:co", })).length === 1, 10_000); const runner: AgentRunner = { mode: "deterministic", run: async () => ({ summary: RESPONSE, tags: ["conversation"], importance: "normal", confidence: 0.5, }), }; const runtime = new ThoughtAgentRuntime(store, [runner]); await runtime.consumeBacklog([agentMessageDeclaration()]); const result = await collect(child); expect(result).toMatchObject({ code: 0, stderr: "" }); expect(result.stdout).toContain(RESPONSE); expect(result.stdout).not.toContain(SECRET); expect(JSON.parse(result.stdout)).toMatchObject({ run: { status: "completed", agentId: "stream-agent-conversation" }, response: { inReplyToMessageId: "co-cli-wait-1", threadId: "co-stream-cli-wait", senderAgentId: "stream-agent-conversation", recipientAgentId: "co", text: RESPONSE, }, });
const denied = await run([ "--silent", "thought", "stream", "agent-send", "--from", "co", "--to", "stream-agent-conversation", "--thread", "co-stream-cli-wait", "--message-id", "co-cli-denied", "--file", messageFile, "--wait-seconds", "1", "--show-response", ], { THOUGHTSTREAM_ROOT: root }); expect(denied.code).toBe(1); expect(denied.stdout).not.toContain(SECRET); expect(denied.stderr).not.toContain(SECRET); }, 30_000);});
function agentMessageDeclaration(): ThoughtAgentDeclaration { const declaration: ThoughtAgentDeclaration = { id: "stream-agent-conversation", version: 1, name: "Stream", description: "fixture private agent conversation", mode: "deterministic", role: "standard", outputContract: { ...OBSERVATION_OUTPUT_CONTRACT.identity }, outputMode: "conversation-text", eventTypes: [AGENT_MESSAGE_SOURCE_EVENT_TYPE], compiledEventTypes: [AGENT_MESSAGE_SOURCE_EVENT_TYPE], sourcePatterns: ["agent-message:co"], acceptedPrivacy: ["sensitive"], initialReplay: "beginning", outputEventType: AGENT_MESSAGE_RESPONSE_EVENT_TYPE, emit: [AGENT_MESSAGE_RESPONSE_EVENT_TYPE], promptRef: "fixture-agent-message.md", systemPrompt: "Reply to Co.", enabled: true, maxEvents: 20, maxInputChars: 8_000, contextStrategy: "agent-conversation", maxOutputTokens: 2_000, timeoutMs: 60_000, tools: [], proposals: [], externalActions: false, }; declaration.declarationFingerprint = declarationFingerprint(declaration); return declaration;}
function start(args: string[], extraEnv: Record<string, string>): ChildProcessWithoutNullStreams { return spawn("pnpm", args, { cwd: process.cwd(), env: { ...process.env, ...extraEnv }, stdio: ["pipe", "pipe", "pipe"], });}
async function run(args: string[], extraEnv: Record<string, string>) { return collect(start(args, extraEnv));}
async function collect(child: ChildProcessWithoutNullStreams): Promise<{ code: number | null; stdout: string; stderr: string }> { return new Promise((resolve, reject) => { let stdout = ""; let stderr = ""; child.stdout.on("data", (chunk) => { stdout += String(chunk); }); child.stderr.on("data", (chunk) => { stderr += String(chunk); }); child.once("error", reject); child.once("exit", (code) => resolve({ code, stdout, stderr })); });}
async function waitFor(predicate: () => Promise<boolean>, timeoutMs: number): Promise<void> { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (await predicate()) return; await new Promise((resolve) => setTimeout(resolve, 50)); } throw new Error("Timed out waiting for agent-send source event");}