Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
12 kB · 351 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352import { canonicalJson, sha256, type JsonObject } from "../core/json.js";import { stableKey } from "../core/ids.js";import type { EventCandidate, ThoughtEvent } from "../events/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import type { AgentRun, ConsumerProgress } from "../store/types.js";import { incidentPayloadJson, OPERATIONAL_INCIDENT_ACTOR, OPERATIONAL_INCIDENT_EVENT_TYPE, OPERATIONAL_INCIDENT_SOURCE, operationalIncidentPayloadSchema, parseOperationalIncidentEvent, type OperationalIncident, type OperationalIncidentCategory,} from "./types.js";
const CONNECTOR_STAGES = new Set([ "live-subscribe", "captured-batch-ingest", "telegram-webhook-ingest", "telegram-spool-ingest", "rss-poll", "fastmail-capture", "filesystem-scan",]);const AGENT_STAGES = new Set([ "runner", "provider", "sandbox-execution", "broker", "final-output-validation", "accounting-reservation", "sdk-stream", "sdk-setup", "configuration", "recovery",]);const AGENT_CODES = new Set([ "invalid-final-output", "semantic-output-invalid", "timeout", "provider-run-failed", "broker-start-failed", "inference-budget-exhausted", "unclassified-run-failure", "credential-missing", "rate-limited",]);
const PROJECTED_EVENT_TYPES = [ "stream.thought.connector.failed", "stream.thought.connector.recovered", "stream.thought.connector.subscription.stopped", "stream.thought.agent.run.failed", "stream.thought.agent.run.blocked", "stream.thought.agent.run.abandoned", "stream.thought.action.telegram.send.failed",];
export interface OperationalIncidentProjectionResult { offered: number; inserted: number; unchanged: number; incidents: ThoughtEvent[];}
export interface SchedulerExhaustionEvidence { operationKey: string; occurredAt?: string | undefined; source?: string | undefined; agentId?: string | undefined; agentVersion?: number | undefined;}
export class OperationalIncidentProjector { async project(store: JazzThoughtStore): Promise<OperationalIncidentProjectionResult> { const [events, progress, runs, existingIncidents] = await Promise.all([ store.listEvents({ types: PROJECTED_EVENT_TYPES }), store.listConsumerProgress(), store.listRuns(), listOperationalIncidents(store), ]); const existingByOrigin = new Map(existingIncidents.flatMap(({ event, incident }) => { const originEventId = stringField(incident.references.originEventId); return originEventId ? [[originEventId, { event, incident }] as const] : []; })); let offered = 0; let inserted = 0; let unchanged = 0; const incidents: ThoughtEvent[] = []; const pending: Array<{ incident: OperationalIncident; origin: ThoughtEvent }> = []; for (const event of events) { const existing = existingByOrigin.get(event.id); if (existing) { offered += 1; unchanged += 1; incidents.push(existing.event); continue; } const incident = await incidentForEvent(store, event, progress, runs); if (!incident) continue; offered += 1; pending.push({ incident, origin: event }); } if (pending.length > 0) { const result = await store.appendProducerBatch(pending.map(({ incident, origin }) => ( incidentCandidate(incident, origin) ))); inserted += result.inserted.length; unchanged += result.unchanged.length; incidents.push(...result.events); } return { offered, inserted, unchanged, incidents }; }}
export async function appendSchedulerExhaustedIncident( store: JazzThoughtStore, evidence: SchedulerExhaustionEvidence,): Promise<ThoughtEvent> { const occurredAt = evidence.occurredAt ?? new Date().toISOString(); const component = evidence.agentId ? `consumer:${bounded(evidence.agentId, 400)}` : "consumer-scheduler"; const source = evidence.source ? bounded(evidence.source, 500) : undefined; const incidentId = stableKey( "operational-incident", "scheduler-exhausted", sha256(evidence.operationKey), occurredAt, ); const incident = incidentValue({ incidentId, state: "open", category: "scheduler-exhausted", severity: "error", component, code: "consumer-cycle-exhausted", stage: "consumer-scheduler", occurredAt, retryable: true, progress: "unchanged", ...(source ? { source } : {}), ...(evidence.agentId ? { agentId: bounded(evidence.agentId, 500) } : {}), ...(evidence.agentVersion ? { agentVersion: evidence.agentVersion } : {}), references: {}, }); const result = await store.appendEvent(incidentCandidate(incident)); return result.event;}
export async function listOperationalIncidents(store: JazzThoughtStore): Promise<Array<{ event: ThoughtEvent; incident: OperationalIncident;}>> { const events = await store.listEvents({ types: [OPERATIONAL_INCIDENT_EVENT_TYPE] }); return events.flatMap((event) => { const incident = parseOperationalIncidentEvent(event); return incident ? [{ event, incident }] : []; });}
async function incidentForEvent( store: JazzThoughtStore, event: ThoughtEvent, progress: ConsumerProgress[], runs: AgentRun[],): Promise<OperationalIncident | undefined> { if (event.type === "stream.thought.connector.failed") { return connectorIncident(event, "connector-failure", "warning", "open", "connector-failed"); } if (event.type === "stream.thought.connector.recovered") { return connectorIncident(event, "connector-recovered", "info", "recovered", "connector-recovered"); } if (event.type === "stream.thought.connector.subscription.stopped") { if (event.payload.status !== "failed") return undefined; return connectorIncident(event, "connector-terminal", "error", "open", "connector-terminal"); } if (event.type === "stream.thought.action.telegram.send.failed") { return incidentValue({ incidentId: incidentIdFor(event, "telegram-delivery-failed"), state: "open", category: "telegram-delivery-failed", severity: "error", component: bounded(event.source, 500), code: "telegram-send-failed", stage: "telegram-delivery", occurredAt: event.occurredAt, retryable: false, progress: "not-applicable", source: bounded(event.source, 500), references: { originEventId: event.id, deliveryId: bounded(event.externalId, 500), }, }); } if (event.type.startsWith("stream.thought.agent.run.")) { return agentRunIncident(store, event, progress, runs); } return undefined;}
function connectorIncident( event: ThoughtEvent, category: Extract<OperationalIncidentCategory, "connector-failure" | "connector-terminal" | "connector-recovered">, severity: OperationalIncident["severity"], state: OperationalIncident["state"], code: string,): OperationalIncident { return incidentValue({ incidentId: incidentIdFor(event, category), state, category, severity, component: bounded(event.source, 500), code, stage: connectorStage(event.payload.phase), occurredAt: event.occurredAt, retryable: state === "open", progress: "not-applicable", source: bounded(event.source, 500), references: { originEventId: event.id }, }, "connector-health", { code: "connector-health" });}
async function agentRunIncident( store: JazzThoughtStore, event: ThoughtEvent, progress: ConsumerProgress[], runs: AgentRun[],): Promise<OperationalIncident | undefined> { const runId = stringField(event.payload.runId); if (!runId) return undefined; const run = runs.find((candidate) => candidate.id === runId); if (!run) return undefined; const trigger = await store.getEvent(run.triggerEventId); if (!trigger) return undefined; const advanced = progress.some((item) => ( item.consumerId === run.agentId && item.consumerVersion === run.agentVersion && item.source === trigger.source && item.lastSequence >= trigger.sourceSequence )); const status = run.status; if (status !== "failed" && status !== "blocked" && status !== "abandoned") return undefined; const category = `agent-run-${status}` as Extract<OperationalIncidentCategory, "agent-run-failed" | "agent-run-blocked" | "agent-run-abandoned">; // A later attempt can exist only when this terminal attempt left the trigger // unconsumed. Current progress may already include that later attempt, so it // cannot be used alone when rebuilding the earlier incident. const retried = runs.some((candidate) => ( candidate.executionKey === run.executionKey && candidate.attempt > run.attempt )); const diagnostic = objectField(run.result?.failureDiagnostic); const deferred = diagnostic?.progressDisposition === "deferred"; const progressAdvanced = status === "blocked" ? !deferred : (!retried && advanced); const code = agentCode(diagnostic?.code) ?? `agent-run-${status}`; const stage = agentStage(diagnostic?.stage); return incidentValue({ incidentId: incidentIdFor(event, category), state: "open", category, severity: status === "failed" ? "error" : status === "blocked" ? "warning" : "info", component: `agent:${bounded(run.agentId, 494)}`, code, ...(stage ? { stage } : {}), occurredAt: event.occurredAt, retryable: (status === "failed" || deferred) && !progressAdvanced, progress: progressAdvanced ? "advanced" : "unchanged", source: bounded(trigger.source, 500), agentId: bounded(run.agentId, 500), agentVersion: run.agentVersion, attempt: run.attempt, references: { originEventId: event.id, runId: run.id, triggerEventId: trigger.id, }, });}
function incidentValue( input: Omit<OperationalIncident, "incidentVersion" | "fingerprint">, fingerprintCategory: string = input.category, fingerprintClassification: { code: string; stage?: string } = { code: input.code, ...(input.stage ? { stage: input.stage } : {}) },): OperationalIncident { const fingerprint = sha256(canonicalJson({ category: fingerprintCategory, component: input.component, source: input.source ?? null, code: fingerprintClassification.code, stage: fingerprintClassification.stage ?? null, })); return operationalIncidentPayloadSchema.parse({ incidentVersion: 1, ...input, fingerprint, });}
function incidentCandidate(incident: OperationalIncident, origin?: ThoughtEvent): EventCandidate { return { type: OPERATIONAL_INCIDENT_EVENT_TYPE, schemaVersion: 1, source: OPERATIONAL_INCIDENT_SOURCE, sourceKind: "system", externalId: incident.incidentId, idempotencyKey: incident.incidentId, occurredAt: incident.occurredAt, actor: OPERATIONAL_INCIDENT_ACTOR, ...(origin ? { rootEventId: origin.rootEventId, parentEventId: origin.id, correlationId: origin.correlationId, } : { correlationId: incident.incidentId }), privacy: "sensitive", payload: incidentPayloadJson(incident), };}
function incidentIdFor(event: ThoughtEvent, category: OperationalIncidentCategory): string { return stableKey("operational-incident", category, event.id);}
function connectorStage(value: unknown): string { return typeof value === "string" && CONNECTOR_STAGES.has(value) ? value : "connector-operation";}
function agentCode(value: unknown): string | undefined { if (typeof value !== "string") return undefined; if (AGENT_CODES.has(value)) return value; return /^(sandbox|letta)-[a-z0-9.-]{1,100}$/.test(value) ? value : undefined;}
function agentStage(value: unknown): string | undefined { return typeof value === "string" && AGENT_STAGES.has(value) ? value : undefined;}
function stringField(value: unknown): string | undefined { return typeof value === "string" && value.length > 0 ? value : undefined;}
function objectField(value: unknown): JsonObject | undefined { return value && typeof value === "object" && !Array.isArray(value) ? value as JsonObject : undefined;}
function bounded(value: string, maximum: number): string { return value.length <= maximum ? value : value.slice(0, maximum);}