import fs from "node:fs/promises"; import path from "node:path"; import { afterEach, describe, expect, test } from "vitest"; import { buildSubscribedTelegramConversationContextPacket } from "../src/agents/context.js"; import { loadAgentDeclarations } from "../src/agents/declarations.js"; import { OBSERVATION_OUTPUT_CONTRACT, outputContractIdentityJson, } from "../src/agents/output-contracts.js"; import { CORRECTION_TARGET_LATEST, proposalCapabilitiesSchema } from "../src/agents/proposals.js"; import { ThoughtAgentRuntime } from "../src/agents/runtime.js"; import type { AgentRunner, AgentRunInput, ThoughtAgentDeclaration } from "../src/agents/types.js"; import { CORRECTION_PROPOSAL_EVENT_TYPE, MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE, MEMORY_MATERIALIZED_EVENT_TYPE, MEMORY_PROPOSAL_EVENT_TYPE, PROPOSAL_DECISION_EVENT_TYPE, } from "../src/agent-proposals/contracts.js"; import { materializeMemoryDecision } from "../src/agent-proposals/memory-materializer.js"; import { recordProposalDecision } from "../src/agent-proposals/review.js"; import { FilesystemConnector } from "../src/connectors/filesystem.js"; import { stableKey } from "../src/core/ids.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 { AgentRun, ConsumerProgress } from "../src/store/types.js"; import { projectTrainingExamples, projectPrivateTrainingExamples, recordJudgment } from "../src/training/judgments.js"; import { exportPrivateTrainingDataset } from "../src/training/private-training.js"; import { temporaryProject, testDeclarationEnvironment, testStore } from "./helpers.js"; import { FOCUS_DECLARATION_PROPOSED_EVENT_TYPE, focusDeclarationProposalPayloadSchema, } from "../src/focuses/types.js"; const stores: JazzThoughtStore[] = []; const roots: string[] = []; const telegramSource = "telegram:thoughtstream-bot-webhook"; const memorySource = "filesystem:telegram-agent-context"; const chatId = "fixture-chat"; const senderId = "fixture-user"; const baseMemory = "---\nid: stream-memory\n---\n# Stream memory\n\nExisting fact.\n"; class ProposalFixtureRunner implements AgentRunner { readonly mode = "pi" as const; constructor( private readonly invalidEvidence = false, private readonly memoryOperation: "append" | "replace-document" = "append", private readonly useLatestCorrectionTarget = false, ) {} async run(input: AgentRunInput) { const capabilities = proposalCapabilitiesSchema.parse(input.context.manifest.proposalCapabilities); const target = capabilities.correctionTargets.at(-1); if (!capabilities.memoryTarget || !target) throw new Error("Fixture requires both proposal capabilities"); const evidence = this.invalidEvidence ? ["evt-outside-snapshot"] : [input.event.id]; return { summary: "I understand.", tags: ["conversation"], importance: "normal" as const, confidence: 0.5, model: { provider: "tinker", id: input.declaration.model! }, proposals: [ { toolCallId: "call-memory", kind: "memory-change" as const, arguments: { operation: this.memoryOperation, proposed_text: "## Response preference\n\nKeep technical replies compact.", reason: "Cameron stated a durable response preference.", evidence_event_ids: evidence, }, }, { toolCallId: "call-correction", kind: "self-correction" as const, arguments: { target_output: this.useLatestCorrectionTarget ? CORRECTION_TARGET_LATEST : target.outputEventId, replacement: "The corrected prior reply.", reason: "The prior delivered reply misstated the fact.", evidence_event_ids: evidence, }, }, ], }; } } class FocusProposalFixtureRunner implements AgentRunner { readonly mode = "pi" as const; constructor(private readonly enforceCapabilities = true) {} async run(input: AgentRunInput) { const capabilities = proposalCapabilitiesSchema.parse(input.context.manifest.proposalCapabilities); if (this.enforceCapabilities && !capabilities.enabled.includes("focus-declaration")) { throw new Error("Fixture requires focus proposal capability"); } return { summary: "I saved that as a focus proposal.", tags: ["conversation"], importance: "normal" as const, confidence: 0.5, model: { provider: "tinker", id: input.declaration.model! }, proposals: [{ toolCallId: "call-focus", kind: "focus-declaration" as const, arguments: { declaration: focusDeclarationFixture() }, }], }; } } interface ProposalFixture { root: string; contextRoot: string; store: JazzThoughtStore; declaration: ThoughtAgentDeclaration; trigger: ThoughtEvent; priorRunId: string; priorOutputEventId: string; memoryProposal: ThoughtEvent; correctionProposal: ThoughtEvent; } 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("agent-originated proposals", () => { test("settles one inert focus proposal atomically and reconstructs its native tool history", async () => { const setup = await setupConversation("/focus AI news that changes how I understand agents and training"); const runtime = new ThoughtAgentRuntime(setup.store, [new FocusProposalFixtureRunner()]); const results = await runtime.consumeBacklog([setup.declaration]); expect(results).toHaveLength(1); expect(results[0]?.error).toBeUndefined(); const run = (await setup.store.listRuns()).find((candidate) => candidate.triggerEventId === setup.trigger.id)!; const output = await setup.store.getEvent(run.outputEventIds[0]!); const focusEvent = (await setup.store.listEvents({ types: [FOCUS_DECLARATION_PROPOSED_EVENT_TYPE] }))[0]!; const focus = focusDeclarationProposalPayloadSchema.parse(focusEvent.payload); const completed = (await setup.store.listEvents({ types: ["stream.thought.agent.run.completed"] })) .find((event) => event.payload.runId === run.id)!; expect(focusEvent).toMatchObject({ privacy: "sensitive", rootEventId: setup.trigger.rootEventId, parentEventId: output?.id, payload: { authority: { activatesFocus: false, expandsDataAccess: false, grantsExternalActions: false, startsTraining: false, promotesAdapter: false, }, }, }); expect(focus.declaration).toEqual(focusDeclarationFixture()); expect(focus.proposer).toMatchObject({ kind: "agent", runId: run.id, outputEventId: output?.id, triggerEventId: setup.trigger.id, agentId: setup.declaration.id, agentVersion: setup.declaration.version, }); expect(completed.payload.proposalEventIds).toEqual([focusEvent.id]); await setup.store.appendEvent({ type: "stream.thought.action.telegram.send.delivered", schemaVersion: 1, source: `telegram-dispatcher:${telegramSource}:${chatId}`, sourceKind: "system", externalId: "delivery-focus", idempotencyKey: "delivery-focus", occurredAt: new Date().toISOString(), actor: `telegram-dispatcher:${telegramSource}:${chatId}`, rootEventId: setup.trigger.rootEventId, parentEventId: output!.id, correlationId: setup.trigger.correlationId, privacy: "sensitive", payload: { chatId, messageId: "out-focus", runIds: [run.id] }, }); const next = (await setup.store.appendEvent(telegramMessage( "after-focus", "/focus AI news that changes how I understand agents and training", ))).event; const packet = await buildSubscribedTelegramConversationContextPacket(setup.declaration, next, setup.store); expect(packet.messages).toEqual(expect.arrayContaining([ expect.objectContaining({ role: "assistant", toolCalls: [expect.objectContaining({ name: "propose_focus", arguments: { declaration: focusDeclarationFixture() }, })], }), expect.objectContaining({ role: "toolResult", toolName: "propose_focus", isError: false, }), ])); const repeated = await runtime.consumeBacklog([setup.declaration]); expect(repeated).toHaveLength(1); expect(repeated[0]?.error).toBeUndefined(); const focusEvents = await setup.store.listEvents({ types: [FOCUS_DECLARATION_PROPOSED_EVENT_TYPE] }); expect(focusEvents).toHaveLength(2); expect(new Set(focusEvents.map((event) => event.id)).size).toBe(2); expect(focusEvents.map((event) => event.payload.declarationFingerprint)) .toEqual([focus.declarationFingerprint, focus.declarationFingerprint]); }); test("rejects an alternate runner focus proposal outside an exact /focus turn", async () => { const setup = await setupConversation(); const runtime = new ThoughtAgentRuntime(setup.store, [new FocusProposalFixtureRunner(false)]); const results = await runtime.consumeBacklog([setup.declaration]); const run = (await setup.store.listRuns()).find((candidate) => candidate.triggerEventId === setup.trigger.id)!; expect(results).toHaveLength(1); expect(results[0]?.error).toBe("Agent runner failed"); expect(run.status).toBe("failed"); expect(run.outputEventIds).toEqual([]); expect(await setup.store.listEvents({ types: [FOCUS_DECLARATION_PROPOSED_EVENT_TYPE] })).toEqual([]); expect((await setup.store.listEvents({ types: ["stream.thought.agent.run.completed"] })) .filter((event) => event.payload.runId === run.id)).toEqual([]); const progress = (await setup.store.listConsumerProgress()) .find((candidate) => candidate.consumerVersion === setup.declaration.version && candidate.source === telegramSource)!; expect(progress.lastEventId).toBe(setup.trigger.id); }); test("freezes exact proposal capability evidence and settles output, proposals, lifecycle, run, and progress atomically", async () => { const fixture = await proposalFixture(); const { store, trigger, memoryProposal, correctionProposal } = fixture; const run = (await store.listRuns()).find((candidate) => candidate.triggerEventId === trigger.id)!; const output = await store.getEvent(run.outputEventIds[0]!); const completed = (await store.listEvents({ types: ["stream.thought.agent.run.completed"] })) .find((event) => event.payload.runId === run.id)!; const progress = (await store.listConsumerProgress()).find((item) => item.source === telegramSource)!; expect(run.status).toBe("completed"); expect(run.result).not.toHaveProperty("proposals"); expect(output?.payload.structuredOutput).not.toHaveProperty("proposals"); expect([memoryProposal, correctionProposal].map((event) => event.sourceSequence)).toEqual([ output!.sourceSequence + 1, output!.sourceSequence + 2, ]); expect(completed.sourceSequence).toBe(output!.sourceSequence + 3); expect(completed.payload.proposalEventIds).toEqual([memoryProposal.id, correctionProposal.id]); expect(progress.lastEventId).toBe(trigger.id); expect(memoryProposal).toMatchObject({ type: MEMORY_PROPOSAL_EVENT_TYPE, privacy: "sensitive", parentEventId: output?.id, payload: { proposalState: "agent-proposed", target: expect.objectContaining({ source: memorySource, path: "memory.md" }), publicationEligible: false, }, }); expect(correctionProposal).toMatchObject({ type: CORRECTION_PROPOSAL_EVENT_TYPE, parentEventId: fixture.priorOutputEventId, payload: { proposalState: "agent-proposed", target: expect.objectContaining({ runId: fixture.priorRunId, outputEventId: fixture.priorOutputEventId }), qualityEligible: false, externalExportEligible: false, publicationEligible: false, }, }); const snapshotId = String(run.contextManifest.contextSnapshot && (run.contextManifest.contextSnapshot as JsonObject).id); const snapshotVersion = await store.getDocumentVersion(snapshotId); expect(snapshotVersion).toBeDefined(); const snapshot = JSON.parse(snapshotVersion!.content) as { manifest: JsonObject }; expect(snapshot.manifest.proposalCapabilities).toEqual(run.contextManifest.proposalCapabilities); const before = (await store.listEvents()).length; const replay = await new ThoughtAgentRuntime(store, [new ProposalFixtureRunner()]) .consumeBacklog([fixture.declaration]); expect(replay).toHaveLength(0); expect(await store.listEvents()).toHaveLength(before); }); test("resolves the latest correction shorthand to exact snapshot-bound target evidence", async () => { const fixture = await proposalFixture("append", true); expect(fixture.correctionProposal.payload.target).toMatchObject({ runId: fixture.priorRunId, outputEventId: fixture.priorOutputEventId, }); expect(JSON.stringify(fixture.correctionProposal.payload)).not.toContain(CORRECTION_TARGET_LATEST); }); test("allows exactly one concurrent human decision and rejects fabricated proposal lineage", async () => { const fixture = await proposalFixture(); const decisions = await Promise.allSettled([ recordProposalDecision(fixture.store, { proposalEventId: fixture.memoryProposal.id, disposition: "accept", submissionId: "concurrent-accept", }), recordProposalDecision(fixture.store, { proposalEventId: fixture.memoryProposal.id, disposition: "reject", submissionId: "concurrent-reject", }), ]); expect(decisions.filter((result) => result.status === "fulfilled")).toHaveLength(1); expect(decisions.filter((result) => result.status === "rejected")).toHaveLength(1); expect((await fixture.store.listEvents({ types: [PROPOSAL_DECISION_EVENT_TYPE] })) .filter((event) => event.payload.proposalEventId === fixture.memoryProposal.id)).toHaveLength(1); const forged = (await fixture.store.appendEvent({ type: MEMORY_PROPOSAL_EVENT_TYPE, schemaVersion: 1, source: fixture.memoryProposal.source, sourceKind: "agent", externalId: "forged-proposal", idempotencyKey: "forged-proposal", occurredAt: new Date().toISOString(), actor: fixture.memoryProposal.actor, rootEventId: fixture.memoryProposal.rootEventId, parentEventId: fixture.memoryProposal.parentEventId!, correlationId: fixture.memoryProposal.correlationId, privacy: "sensitive", payload: fixture.memoryProposal.payload, })).event; await expect(recordProposalDecision(fixture.store, { proposalEventId: forged.id, disposition: "accept", submissionId: "forged-decision", })).rejects.toThrow("atomic completed-run receipt"); }); test("keeps invalid proposal settlement from leaking a partial output or proposal", async () => { const setup = await setupConversation(); const declaration = setup.declaration; const runtime = new ThoughtAgentRuntime(setup.store, [new ProposalFixtureRunner(true)]); const results = await runtime.consumeBacklog([declaration]); const run = (await setup.store.listRuns()).find((candidate) => candidate.triggerEventId === setup.trigger.id)!; const outputs = (await setup.store.listEvents()).filter((event) => event.payload.runId === run.id && ( event.type === declaration.outputEventType || event.type === MEMORY_PROPOSAL_EVENT_TYPE || event.type === CORRECTION_PROPOSAL_EVENT_TYPE )); expect(results[0]).toMatchObject({ error: "Agent runner failed" }); expect(run.status).toBe("failed"); expect(run.outputEventIds).toEqual([]); expect(outputs).toEqual([]); expect((await setup.store.listEvents({ types: ["stream.thought.agent.run.failed"] }))) .toEqual([expect.objectContaining({ payload: expect.objectContaining({ runId: run.id }) })]); }); test("records append-only decisions, projects accepted corrections, materializes memory with scan receipts, and exports private quality data only through the private path", async () => { const fixture = await proposalFixture(); const correction = await recordProposalDecision(fixture.store, { proposalEventId: fixture.correctionProposal.id, disposition: "accept", submissionId: "submission-correction", actor: "operator:test", }); const correctionAgain = await recordProposalDecision(fixture.store, { proposalEventId: fixture.correctionProposal.id, disposition: "accept", submissionId: "submission-correction", actor: "operator:test", }); expect(correctionAgain.decision.id).toBe(correction.decision.id); expect(correction.projectedJudgment).toMatchObject({ type: "stream.thought.judgment.training-example", schemaVersion: 2, payload: { runId: fixture.priorRunId, outputEventId: fixture.priorOutputEventId, kind: "correct", criterion: "agent-self-correction", criterionVersion: 1, qualityEligible: true, externalExportEligible: false, }, }); await expect(recordProposalDecision(fixture.store, { proposalEventId: fixture.correctionProposal.id, disposition: "reject", submissionId: "different-decision", })).rejects.toThrow("already has a human decision"); expect(await projectTrainingExamples(fixture.store, { includeSensitivePrivate: true })).toEqual([]); await recordJudgment(fixture.store, { runId: fixture.priorRunId, kind: "accept", criterion: "manual-private-without-feedback-provenance", criterionVersion: 1, qualityEligible: true, externalExportEligible: false, actor: "operator:test", }); const privateExamples = await projectPrivateTrainingExamples(fixture.store); expect(privateExamples).toEqual([ expect.objectContaining({ format: "thoughtstream.private-training-example.v1", privateProvenance: expect.objectContaining({ judgmentEventId: correction.projectedJudgment!.id, runId: fixture.priorRunId, outputEventId: fixture.priorOutputEventId, feedbackSourceEventId: correction.decision.id, }), }), ]); const privateDestination = path.join(fixture.root, "private-training", "stream.jsonl"); const manifest = await exportPrivateTrainingDataset(fixture.store, privateDestination, { acknowledgeSensitivePrivate: true, }); expect(manifest).toMatchObject({ format: "thoughtstream.private-training-dataset-manifest.v1", examples: 1, exactPrivateProvenance: true, externalExportAuthorityRequired: false, }); expect((await fs.stat(privateDestination)).mode & 0o777).toBe(0o600); expect((await fs.stat(`${privateDestination}.manifest.json`)).mode & 0o777).toBe(0o600); const memory = await recordProposalDecision(fixture.store, { proposalEventId: fixture.memoryProposal.id, disposition: "accept", submissionId: "submission-memory", actor: "operator:test", }); const materialized = await materializeMemoryDecision(fixture.store, memory.decision.id, { contextRoot: fixture.contextRoot, }); expect(materialized).toMatchObject({ status: "materialized", event: { type: MEMORY_MATERIALIZED_EVENT_TYPE, payload: { proposalEventId: fixture.memoryProposal.id, decisionEventId: memory.decision.id, result: expect.objectContaining({ filesystemEventId: expect.any(String) }), }, }, }); expect(await fs.readFile(path.join(fixture.contextRoot, "memory.md"), "utf8")).toContain("Keep technical replies compact."); expect((await fs.stat(path.join(fixture.contextRoot, "memory.md"))).mode & 0o777).toBe(0o600); const resultVersionId = String(materialized.event.payload.result && (materialized.event.payload.result as JsonObject).versionId); expect(await fixture.store.getDocumentVersion(resultVersionId)).toMatchObject({ source: memorySource, path: "memory.md", }); }); test("reclaims a verified stale materializer lock while retaining live/invalid-lock fail-closed behavior", async () => { const staleOwner = await proposalFixture(); const decision = await recordProposalDecision(staleOwner.store, { proposalEventId: staleOwner.memoryProposal.id, disposition: "accept", submissionId: "stale-lock", }); const lockPath = path.join(staleOwner.contextRoot, ".thoughtstream-memory-materializer.lock"); await fs.writeFile(lockPath, `${JSON.stringify({ bootId: "definitely-another-boot", pid: 999_999, processStart: "1", token: "stale-lock-token-123456789", createdAt: "2026-01-01T00:00:00.000Z", })}\n`, { mode: 0o600 }); const materialized = await materializeMemoryDecision(staleOwner.store, decision.decision.id, { contextRoot: staleOwner.contextRoot, }); expect(materialized.status).toBe("materialized"); await expect(fs.lstat(lockPath)).rejects.toMatchObject({ code: "ENOENT" }); const invalidOwner = await proposalFixture(); const invalidDecision = await recordProposalDecision(invalidOwner.store, { proposalEventId: invalidOwner.memoryProposal.id, disposition: "accept", submissionId: "invalid-lock", }); const invalidLockPath = path.join(invalidOwner.contextRoot, ".thoughtstream-memory-materializer.lock"); await fs.writeFile(invalidLockPath, "not valid lock evidence\n", { mode: 0o600 }); const failed = await materializeMemoryDecision(invalidOwner.store, invalidDecision.decision.id, { contextRoot: invalidOwner.contextRoot, }); expect(failed).toMatchObject({ status: "failed", event: { payload: { reasonCode: "write-failed" } } }); expect(await fs.readFile(invalidLockPath, "utf8")).toBe("not valid lock evidence\n"); }); test("fails closed on stale bases, symlink swaps, and frontmatter-removing edits without overwriting the file", async () => { const stale = await proposalFixture(); const stalePath = path.join(stale.contextRoot, "memory.md"); const concurrent = `${baseMemory.trimEnd()}\n\nConcurrent operator edit.\n`; await fs.writeFile(stalePath, concurrent); await new FilesystemConnector({ id: memorySource, root: stale.contextRoot }).scan(stale.store); const staleDecision = await recordProposalDecision(stale.store, { proposalEventId: stale.memoryProposal.id, disposition: "accept", submissionId: "stale", }); const staleResult = await materializeMemoryDecision(stale.store, staleDecision.decision.id, { contextRoot: stale.contextRoot }); expect(staleResult).toMatchObject({ status: "failed", event: { payload: { reasonCode: "stale-base", contentRedacted: true } } }); expect(await fs.readFile(stalePath, "utf8")).toBe(concurrent); const symlink = await proposalFixture(); const symlinkPath = path.join(symlink.contextRoot, "memory.md"); const outside = path.join(symlink.root, "outside-memory.md"); await fs.writeFile(outside, "OUTSIDE SENTINEL\n"); const symlinkDecision = await recordProposalDecision(symlink.store, { proposalEventId: symlink.memoryProposal.id, disposition: "accept", submissionId: "symlink", }); const symlinkResult = await materializeMemoryDecision(symlink.store, symlinkDecision.decision.id, { contextRoot: symlink.contextRoot, beforeRename: async () => { await fs.rm(symlinkPath); await fs.symlink(outside, symlinkPath); }, }); expect(symlinkResult).toMatchObject({ status: "failed", event: { payload: { reasonCode: "symlink-refused" } } }); expect(await fs.readFile(outside, "utf8")).toBe("OUTSIDE SENTINEL\n"); const frontmatter = await proposalFixture("replace-document"); const frontmatterPath = path.join(frontmatter.contextRoot, "memory.md"); const before = await fs.readFile(frontmatterPath, "utf8"); const frontmatterDecision = await recordProposalDecision(frontmatter.store, { proposalEventId: frontmatter.memoryProposal.id, disposition: "edit", replacementText: "# No stable identity\n", submissionId: "frontmatter", }); const frontmatterResult = await materializeMemoryDecision(frontmatter.store, frontmatterDecision.decision.id, { contextRoot: frontmatter.contextRoot }); expect(frontmatterResult).toMatchObject({ status: "failed", event: { payload: { reasonCode: "frontmatter-invalid" } } }); expect(await fs.readFile(frontmatterPath, "utf8")).toBe(before); expect((await frontmatter.store.listEvents({ types: [MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE] }))).toHaveLength(1); }); test("rolls back a terminal settlement when a side-effect proposal candidate is invalid", async () => { const root = await temporaryProject("thoughtstream-proposal-rollback-"); roots.push(root); const store = await testStore(root); stores.push(store); const input = (await store.appendEvent(telegramMessage("rollback", "Rollback fixture"))).event; const run = runningRun(input); await store.upsertRun(run); const output: EventCandidate = { type: "stream.thought.derived.message.observation", schemaVersion: 1, source: "agent:telegram-conversation", sourceKind: "agent", externalId: run.id, idempotencyKey: `${run.id}:output`, occurredAt: new Date().toISOString(), actor: "telegram-conversation", rootEventId: input.rootEventId, parentEventId: input.id, correlationId: input.correlationId, privacy: "sensitive", payload: { runId: run.id, summary: "Would have settled" }, }; const invalidProposal: EventCandidate = { ...output, type: MEMORY_PROPOSAL_EVENT_TYPE, externalId: `${run.id}:proposal`, idempotencyKey: `${run.id}:proposal`, payload: { invalid: true }, }; const completed: EventCandidate = { ...output, type: "stream.thought.agent.run.completed", externalId: `${run.id}:completed`, idempotencyKey: `${run.id}:completed`, payload: { runId: run.id, agentId: run.agentId, agentVersion: run.agentVersion, inputEventIds: [input.id], attempt: 1, status: "completed", }, }; const progress: ConsumerProgress = { id: stableKey("consumer-progress", run.agentId, String(run.agentVersion), input.source), consumerId: run.agentId, consumerVersion: run.agentVersion, source: input.source, lastSequence: input.sourceSequence, lastEventId: input.id, updatedAt: new Date().toISOString(), }; await expect(store.settleConsumerSuccess({ run: { ...run, status: "completed", outputEventIds: ["would-be-output"] }, inputEvent: input, output, sideEffects: [invalidProposal], completed, progress, })).rejects.toThrow(); expect((await store.getRun(run.id))?.status).toBe("running"); expect((await store.listEvents()).filter((event) => event.source === output.source)).toEqual([]); expect(await store.getConsumerProgress(progress.id)).toBeUndefined(); }); }); async function proposalFixture( memoryOperation: "append" | "replace-document" = "append", useLatestCorrectionTarget = false, ): Promise { const setup = await setupConversation(); const runtime = new ThoughtAgentRuntime(setup.store, [ new ProposalFixtureRunner(false, memoryOperation, useLatestCorrectionTarget), ]); const results = await runtime.consumeBacklog([setup.declaration]); expect(results).toHaveLength(1); expect(results[0]?.error).toBeUndefined(); const proposals = (await setup.store.listEvents({ types: [MEMORY_PROPOSAL_EVENT_TYPE, CORRECTION_PROPOSAL_EVENT_TYPE], })).filter((event) => event.payload.proposer && (event.payload.proposer as JsonObject).triggerEventId === setup.trigger.id); const memoryProposal = proposals.find((event) => event.type === MEMORY_PROPOSAL_EVENT_TYPE); const correctionProposal = proposals.find((event) => event.type === CORRECTION_PROPOSAL_EVENT_TYPE); if (!memoryProposal || !correctionProposal) throw new Error("Proposal fixture did not settle both proposals"); return { ...setup, memoryProposal, correctionProposal, }; } async function setupConversation(currentText = "Please remember that and correct yourself.") { const root = await temporaryProject("thoughtstream-agent-proposals-"); roots.push(root); const store = await testStore(root); stores.push(store); const contextRoot = path.join(root, "stream-context"); await fs.mkdir(contextRoot); await fs.writeFile(path.join(contextRoot, "identity.md"), "---\nid: stream-identity\n---\n# Stream identity\n"); await fs.writeFile(path.join(contextRoot, "memory.md"), baseMemory); await new FilesystemConnector({ id: memorySource, root: contextRoot }).scan(store); const loadedDeclaration = (await loadAgentDeclarations(path.join(process.cwd(), "agents"), testDeclarationEnvironment)) .find((candidate) => candidate.id === "telegram-conversation")!; const declaration = structuredClone(loadedDeclaration); declaration.enabled = true; delete declaration.conversationCompaction; const prior = (await store.appendEvent(telegramMessage("prior", "The prior user fact."))).event; const priorContext = await buildSubscribedTelegramConversationContextPacket(declaration, prior, store); const priorOutput = (await store.appendEvent({ type: "stream.thought.derived.message.observation", schemaVersion: 1, source: "agent:telegram-conversation", sourceKind: "agent", externalId: "run-prior", idempotencyKey: "run-prior:output", occurredAt: new Date().toISOString(), actor: "telegram-conversation", rootEventId: prior.rootEventId, parentEventId: prior.id, correlationId: prior.correlationId, privacy: "sensitive", payload: { runId: "run-prior", executionKey: "execution-prior", inputEventId: prior.id, inputSourceSequence: prior.sourceSequence, summary: "The incorrect prior reply.", tags: ["conversation"], importance: "normal", confidence: 0.5, outputContract: outputContractIdentityJson(OBSERVATION_OUTPUT_CONTRACT.identity), structuredOutput: { summary: "The incorrect prior reply.", tags: ["conversation"], importance: "normal", confidence: 0.5, }, }, traceId: "run-prior", })).event; const priorRun: AgentRun = { id: "run-prior", executionKey: "execution-prior", triggerEventId: prior.id, agentId: "telegram-conversation", agentVersion: declaration.version, status: "completed", inputEventIds: [prior.id], outputEventIds: [priorOutput.id], attempt: 1, provider: "tinker", model: "thinkingmachines/Inkling-Small", privacy: "sensitive", promptHash: "prompt-prior", contextManifest: priorContext.manifest, result: { summary: "The incorrect prior reply.", tags: ["conversation"], importance: "normal", confidence: 0.5, }, createdAt: prior.observedAt, startedAt: prior.observedAt, completedAt: prior.observedAt, updatedAt: prior.observedAt, }; await store.upsertRun(priorRun); await store.appendEvent({ type: "stream.thought.action.telegram.send.delivered", schemaVersion: 1, source: `telegram-dispatcher:${telegramSource}:${chatId}`, sourceKind: "system", externalId: "delivery-prior", idempotencyKey: "delivery-prior", occurredAt: new Date().toISOString(), actor: `telegram-dispatcher:${telegramSource}:${chatId}`, rootEventId: prior.rootEventId, parentEventId: priorOutput.id, correlationId: prior.correlationId, privacy: "sensitive", payload: { chatId, messageId: "out-prior", runIds: [priorRun.id] }, }); await store.initializeConsumerProgress({ id: stableKey("consumer-progress", declaration.id, String(declaration.version), telegramSource), consumerId: declaration.id, consumerVersion: declaration.version, source: telegramSource, lastSequence: prior.sourceSequence, lastEventId: prior.id, updatedAt: new Date().toISOString(), }); const trigger = (await store.appendEvent(telegramMessage("current", currentText))).event; return { root, contextRoot, store, declaration, trigger, priorRunId: priorRun.id, priorOutputEventId: priorOutput.id, }; } function telegramMessage(id: string, text: string): EventCandidate { return { type: "stream.thought.source.telegram.message", schemaVersion: 1, source: telegramSource, sourceKind: "telegram", externalId: id, idempotencyKey: id, occurredAt: new Date().toISOString(), actor: `telegram-user:${senderId}`, correlationId: `conversation-${chatId}`, privacy: "sensitive", payload: { accountId: "thoughtstream-bot", updateId: id, chatId, messageId: id, senderId, chatType: "direct", text, occurredAt: new Date().toISOString(), transport: "telegram-bot-api-webhook", }, }; } function runningRun(input: ThoughtEvent): AgentRun { const at = new Date().toISOString(); return { id: "run-rollback", executionKey: "execution-rollback", triggerEventId: input.id, agentId: "telegram-conversation", agentVersion: 15, status: "running", inputEventIds: [input.id], outputEventIds: [], attempt: 1, provider: "tinker", model: "thinkingmachines/Inkling-Small", privacy: "sensitive", promptHash: "prompt-rollback", contextManifest: { agentRole: "standard", outputContract: outputContractIdentityJson(OBSERVATION_OUTPUT_CONTRACT.identity), }, createdAt: at, startedAt: at, updatedAt: at, }; } function focusDeclarationFixture() { return { id: "ai-news", version: 1, name: "AI news", description: "Notice public AI news that changes Cameron's model of agents and training.", parentFocusIds: ["news"], scope: { summary: "Public AI research, product, and infrastructure developments.", includeTags: ["ai-news"], excludeTags: [], }, objective: { statement: "Produce concise evidence-grounded observations about meaningful AI developments." }, subscriptions: { eventTypes: ["stream.thought.source.rss.item", "stream.thought.source.atproto.commit"], sourcePatterns: ["rss:*", "jetstream:*"], privacy: ["public-source" as const], replay: "now" as const, }, budgets: { inference: { maxCallsPerDay: 24, maxInputTokensPerDay: 480_000, maxOutputTokensPerDay: 12_288, maxCostMicrousdPerDay: 2_400_000, }, attention: { maxDeliveriesPerDay: 6 }, }, permissions: { requestedReadTools: [], requestedExternalActions: [] }, retirementRule: "Retire or revise after 100 eligible events without accepted value.", }; }