Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
2.4 kB · 58 lines
TypeScript
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859import { canonicalJson } from "../core/json.js";import type { ThoughtEvent } from "../events/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import type { AgentRun } from "../store/types.js";
export interface CompletedObservationOutput { run: AgentRun; trigger: ThoughtEvent; output: ThoughtEvent;}
export async function requireCompletedObservationOutput( store: JazzThoughtStore, output: ThoughtEvent,): Promise<CompletedObservationOutput> { const runId = output.payload.runId; if (typeof runId !== "string" || runId.length === 0) { throw new Error("Agent observation has no run receipt"); } const run = await store.getRun(runId); if (!run) throw new Error("Agent observation run receipt is unavailable"); const trigger = await store.getEvent(run.triggerEventId); if (!trigger) throw new Error("Agent observation trigger receipt is unavailable"); if (!isCompletedObservationOutput(run, trigger, output)) { throw new Error("Agent observation does not match its completed run receipt"); } return { run, trigger, output };}
export function isCompletedObservationOutput( run: AgentRun, trigger: ThoughtEvent, output: ThoughtEvent,): boolean { return run.status === "completed" && run.inputEventIds.length === 1 && run.inputEventIds[0] === trigger.id && run.outputEventIds.length === 1 && run.outputEventIds[0] === output.id && output.type === "stream.thought.derived.message.observation" && output.source === `agent:${run.agentId}` && output.sourceKind === "agent" && output.actor === run.agentId && output.externalId === run.id && output.traceId === run.id && output.rootEventId === trigger.rootEventId && output.parentEventId === trigger.id && output.correlationId === trigger.correlationId && output.payload.runId === run.id && output.payload.executionKey === run.executionKey && output.payload.inputEventId === trigger.id && output.payload.inputSourceSequence === trigger.sourceSequence && output.payload.summary === run.result?.summary && canonicalJson(output.payload.tags ?? null) === canonicalJson(run.result?.tags ?? null) && output.payload.importance === run.result?.importance && output.payload.confidence === run.result?.confidence && canonicalJson(output.payload.recommendation ?? null) === canonicalJson(run.result?.recommendation ?? null);}