import http from "node:http"; import fs from "node:fs/promises"; import path from "node:path"; import { afterEach, describe, expect, test } from "vitest"; import { ConsumerScheduler } from "../src/agents/scheduler.js"; import { TelegramBotClient } from "../src/connectors/telegram-bot.js"; import type { EventCandidate, ThoughtEvent } from "../src/events/types.js"; import { IncidentLedger } from "../src/incidents/ledger.js"; import { appendSchedulerExhaustedIncident, listOperationalIncidents, OperationalIncidentProjector, } from "../src/incidents/projector.js"; import { IncidentTelegramDispatcher } from "../src/incidents/telegram-alerts.js"; import type { JazzThoughtStore } from "../src/jazz/store.js"; import type { AgentRun, AgentRunStatus } from "../src/store/types.js"; import { temporaryProject, testStore } from "./helpers.js"; const SECRET = "PRIVATE_SOURCE_PROVIDER_TOOL_CREDENTIAL_SENTINEL"; const stores: JazzThoughtStore[] = []; const roots: string[] = []; const servers: http.Server[] = []; afterEach(async () => { for (const server of servers.splice(0)) await new Promise((resolve) => server.close(() => resolve())); for (const store of stores.splice(0)) await store.close(); for (const root of roots.splice(0)) await fs.rm(root, { recursive: true, force: true }); }); describe("operational incidents", () => { test("normalizes failure families into content-dark rows and one private restart-safe ledger", async () => { const { project, store } = await fixtureStore(); const connectorFailure = await store.appendEvent(connectorEvent( "stream.thought.connector.failed", "connector-failed", { status: "failed", phase: SECRET, error: SECRET }, )); await store.appendEvent(connectorEvent( "stream.thought.connector.recovered", "connector-recovered", { status: "recovered", priorError: SECRET }, )); await store.appendEvent(connectorEvent( "stream.thought.connector.subscription.stopped", "connector-terminal", { status: "failed", reason: SECRET }, )); const trigger = await appendTrigger(store, "failed-one", SECRET); await appendTerminalRun(store, trigger.event, "failed", 1, SECRET); const blockedTrigger = await appendTrigger(store, "blocked-one", SECRET, "2026-07-22T01:00:04.000Z"); await appendTerminalRun(store, blockedTrigger.event, "blocked", 1, SECRET); const abandonedTrigger = await appendTrigger(store, "abandoned-one", SECRET, "2026-07-22T01:00:05.000Z"); await appendTerminalRun(store, abandonedTrigger.event, "abandoned", 1, SECRET); await store.appendEvent({ type: "stream.thought.action.telegram.send.failed", schemaVersion: 1, source: "telegram-dispatcher:fixture", sourceKind: "system", externalId: "delivery-one", idempotencyKey: "delivery-one:failed", occurredAt: "2026-07-22T01:00:06.000Z", actor: "telegram-dispatcher:fixture", rootEventId: trigger.event.rootEventId, parentEventId: trigger.event.id, correlationId: "delivery-one", privacy: "sensitive", payload: { status: "failed", messageKind: "observation", error: SECRET }, }); await appendSchedulerExhaustedIncident(store, { operationKey: `consumer:${SECRET}`, occurredAt: "2026-07-22T01:00:07.000Z", source: "jetstream:cameron-bluesky", agentId: "resident-fixture", agentVersion: 1, }); const projection = await new OperationalIncidentProjector().project(store); expect(projection.inserted).toBe(7); const values = await listOperationalIncidents(store); expect(new Set(values.map(({ incident }) => incident.category))).toEqual(new Set([ "connector-failure", "connector-recovered", "connector-terminal", "scheduler-exhausted", "agent-run-failed", "agent-run-blocked", "agent-run-abandoned", "telegram-delivery-failed", ])); const durable = JSON.stringify(values); expect(durable).not.toContain(SECRET); expect(values.every(({ event }) => event.privacy === "sensitive")).toBe(true); expect(values.find(({ incident }) => incident.category === "connector-failure")?.incident.stage) .toBe("connector-operation"); expect(values.find(({ incident }) => incident.category === "agent-run-blocked")?.incident) .toMatchObject({ retryable: false, progress: "advanced" }); expect(values.find(({ incident }) => incident.category === "agent-run-abandoned")?.incident) .toMatchObject({ retryable: false, progress: "unchanged" }); expect(values.find(({ event }) => event.parentEventId === connectorFailure.event.id)).toBeDefined(); const ledger = await IncidentLedger.open(project, ".thoughtstream/error-ledger.jsonl"); const first = await ledger.append(values.map(({ incident }) => incident)); expect(first).toEqual({ appended: 8, unchanged: 0 }); const ledgerPath = path.join(project, ".thoughtstream", "error-ledger.jsonl"); const stat = await fs.stat(ledgerPath); expect(stat.mode & 0o777).toBe(0o600); const text = await fs.readFile(ledgerPath, "utf8"); expect(text).not.toContain(SECRET); expect(text.trim().split("\n")).toHaveLength(8); const reopened = await IncidentLedger.open(project, ".thoughtstream/error-ledger.jsonl"); expect(await reopened.append(values.map(({ incident }) => incident))).toEqual({ appended: 0, unchanged: 8 }); await expect(IncidentLedger.open(project, "../escape.jsonl")).rejects.toThrow("escapes the runtime root"); }); test("classifies deferred budget blocks as retryable with progress unchanged", async () => { const { store } = await fixtureStore(); const trigger = await appendTrigger(store, "deferred-budget-block", SECRET); await appendTerminalRun(store, trigger.event, "blocked", 1, SECRET, true); expect((await new OperationalIncidentProjector().project(store)).inserted).toBe(1); const values = await listOperationalIncidents(store); expect(values).toHaveLength(1); expect(values[0]!.incident).toMatchObject({ category: "agent-run-blocked", code: "agent-run-blocked", retryable: true, progress: "unchanged", }); }); test("keeps the first projected progress classification authoritative after progress advances", async () => { const { store } = await fixtureStore(); const trigger = await appendTrigger(store, "projection-stability", SECRET); await appendTerminalRun(store, trigger.event, "failed", 1, SECRET); const projector = new OperationalIncidentProjector(); expect(await projector.project(store)).toMatchObject({ offered: 1, inserted: 1, unchanged: 0 }); const first = await listOperationalIncidents(store); expect(first[0]?.incident).toMatchObject({ retryable: true, progress: "unchanged" }); await store.initializeConsumerProgress({ id: "consumer-progress:resident-fixture:1:jetstream:cameron-bluesky", consumerId: "resident-fixture", consumerVersion: 1, source: trigger.event.source, lastSequence: trigger.event.sourceSequence, lastEventId: trigger.event.id, updatedAt: "2026-07-22T01:00:10.000Z", }); expect(await projector.project(store)).toMatchObject({ offered: 1, inserted: 0, unchanged: 1 }); const replayed = await listOperationalIncidents(store); expect(replayed).toHaveLength(1); expect(replayed[0]?.incident).toEqual(first[0]?.incident); }); test("classifies an unprojected failed attempt before its later successful retry", async () => { const { store } = await fixtureStore(); const trigger = await appendTrigger(store, "later-attempt", SECRET); await appendTerminalRun(store, trigger.event, "failed", 1, SECRET); await appendCompletedRun(store, trigger.event, 2); await store.initializeConsumerProgress({ id: "consumer-progress:resident-fixture:1:jetstream:cameron-bluesky", consumerId: "resident-fixture", consumerVersion: 1, source: trigger.event.source, lastSequence: trigger.event.sourceSequence, lastEventId: trigger.event.id, updatedAt: "2026-07-22T01:00:10.000Z", }); expect(await new OperationalIncidentProjector().project(store)) .toMatchObject({ offered: 1, inserted: 1, unchanged: 0 }); const incidents = await listOperationalIncidents(store); expect(incidents).toHaveLength(1); expect(incidents[0]?.incident).toMatchObject({ attempt: 1, retryable: true, progress: "unchanged", }); }); test("alerts on a Jetstream-rooted failure without treating successful Bluesky work as notification input", async () => { const { store } = await fixtureStore(); const telegram = await telegramFixture("success"); const dispatcher = new IncidentTelegramDispatcher({ id: "incident-dispatcher:fixture", client: new TelegramBotClient({ token: "fixture-token", baseUrl: telegram.baseUrl }), chatId: "123456789", categories: ["agent-run-failed"], cooldownMs: 15 * 60_000, maxMessagesPerWindow: 3, windowMs: 15 * 60_000, }); const since = await dispatcher.activate(store, new Date("2026-07-22T01:00:00.000Z")); const failedTrigger = await appendTrigger(store, "failed-alert", SECRET, "2026-07-22T01:00:01.000Z"); await appendTerminalRun(store, failedTrigger.event, "failed", 1, SECRET); const completedTrigger = await appendTrigger(store, "successful-bluesky", SECRET, "2026-07-22T01:00:03.000Z"); await appendCompletedRun(store, completedTrigger.event); await new OperationalIncidentProjector().project(store); const first = await dispatcher.sendPending(store, { since, now: new Date("2026-07-22T01:00:10.000Z"), }); expect(first).toMatchObject({ pending: 1, eligible: 1, delivered: 1, failed: 0 }); expect(telegram.messages).toHaveLength(1); expect(telegram.messages[0]).toContain("Agent run failed"); expect(telegram.messages[0]).not.toContain("jetstream:cameron-bluesky"); expect(JSON.stringify(await listOperationalIncidents(store))).toContain("jetstream:cameron-bluesky"); expect(telegram.messages[0]).not.toContain(SECRET); expect(JSON.stringify(await store.listEvents({ source: "incident-dispatcher:fixture" }))).not.toContain(SECRET); const repeatedTrigger = await appendTrigger(store, "failed-alert-two", SECRET, "2026-07-22T01:01:00.000Z"); await appendTerminalRun(store, repeatedTrigger.event, "failed", 1, SECRET); await new OperationalIncidentProjector().project(store); const repeated = await dispatcher.sendPending(store, { since, now: new Date("2026-07-22T01:01:10.000Z"), }); expect(repeated).toMatchObject({ delivered: 0, cooldownDeferred: 1 }); expect(telegram.messages).toHaveLength(1); }); test("suppresses a connector alert when later recovery evidence has the same normalized fingerprint", async () => { const { store } = await fixtureStore(); const telegram = await telegramFixture("success"); const dispatcher = new IncidentTelegramDispatcher({ id: "incident-dispatcher:recovery", client: new TelegramBotClient({ token: "fixture-token", baseUrl: telegram.baseUrl }), chatId: "123456789", categories: ["connector-terminal"], }); const since = await dispatcher.activate(store, new Date("2026-07-22T00:59:00.000Z")); await store.appendEvent(connectorEvent( "stream.thought.connector.subscription.stopped", "connector-terminal", { status: "failed", reason: SECRET }, )); await store.appendEvent(connectorEvent( "stream.thought.connector.recovered", "connector-recovered", { status: "recovered" }, "2026-07-22T01:00:03.000Z", )); await new OperationalIncidentProjector().project(store); const result = await dispatcher.sendPending(store, { since, now: new Date("2026-07-22T01:05:00.000Z"), }); expect(result).toMatchObject({ pending: 0, delivered: 0, failed: 0 }); expect(telegram.calls).toBe(0); }); test("records a failed incident alert without persisting the Bot API body or recursively alerting it", async () => { const { project, store } = await fixtureStore(); const telegram = await telegramFixture("failure"); const dispatcher = new IncidentTelegramDispatcher({ id: "incident-dispatcher:failure", client: new TelegramBotClient({ token: "fixture-token", baseUrl: telegram.baseUrl }), chatId: "123456789", categories: ["agent-run-failed", "telegram-delivery-failed"], }); const since = await dispatcher.activate(store, new Date("2026-07-22T01:00:00.000Z")); const trigger = await appendTrigger(store, "failed-delivery", SECRET, "2026-07-22T01:00:01.000Z"); await appendTerminalRun(store, trigger.event, "failed", 1, SECRET); const projector = new OperationalIncidentProjector(); await projector.project(store); const first = await dispatcher.sendPending(store, { since, now: new Date("2026-07-22T01:00:10.000Z") }); expect(first).toMatchObject({ delivered: 0, failed: 1 }); expect(telegram.calls).toBe(1); const dispatcherEvents = await store.listEvents({ source: "incident-dispatcher:failure" }); expect(JSON.stringify(dispatcherEvents)).not.toContain(SECRET); const failedReceipt = dispatcherEvents.find((event) => event.type === "stream.thought.action.telegram.send.failed"); expect(failedReceipt?.payload).toMatchObject({ status: "failed", errorCode: "telegram-send-failed", errorClass: "TelegramBotApiError", }); expect(failedReceipt?.payload).not.toHaveProperty("error"); await projector.project(store); const incidents = await listOperationalIncidents(store); expect(incidents.some(({ incident }) => incident.category === "telegram-delivery-failed")).toBe(true); const second = await dispatcher.sendPending(store, { since, now: new Date("2026-07-22T01:20:00.000Z") }); expect(second.delivered).toBe(0); expect(second.failed).toBe(0); expect(telegram.calls).toBe(1); const ledger = await IncidentLedger.open(project, ".thoughtstream/error-ledger.jsonl"); await ledger.append(incidents.map(({ incident }) => incident)); expect(await fs.readFile(ledger.filePath, "utf8")).not.toContain(SECRET); }); test("records one content-dark scheduler incident only after all retries are exhausted", async () => { const { store } = await fixtureStore(); let attempts = 0; const scheduler = new ConsumerScheduler({ concurrency: 1, attempts: 3, retryDelayMs: () => 0, onError: async (_error, operationKey, context) => { const value = context as { source: string; agentId: string; agentVersion: number }; await appendSchedulerExhaustedIncident(store, { operationKey, occurredAt: "2026-07-22T01:00:00.000Z", source: value.source, agentId: value.agentId, agentVersion: value.agentVersion, }); }, }); scheduler.enqueue(`scheduler:${SECRET}`, async () => { attempts += 1; throw new Error(SECRET); }, { source: "jetstream:cameron-bluesky", agentId: "resident-fixture", agentVersion: 1, }); await expect(scheduler.drain()).rejects.toThrow("Consumer cycle failed"); expect(attempts).toBe(3); const incidents = await listOperationalIncidents(store); expect(incidents).toHaveLength(1); expect(incidents[0]?.incident).toMatchObject({ category: "scheduler-exhausted", retryable: true, progress: "unchanged", source: "jetstream:cameron-bluesky", }); expect(JSON.stringify(incidents)).not.toContain(SECRET); }); test("alert text includes retry status and receipt without raw field labels", async () => { const { store } = await fixtureStore(); const telegram = await telegramFixture("success"); const dispatcher = new IncidentTelegramDispatcher({ id: "incident-dispatcher:language", client: new TelegramBotClient({ token: "fixture-token", baseUrl: telegram.baseUrl }), chatId: "123456789", categories: ["agent-run-failed", "connector-terminal", "scheduler-exhausted"], }); const since = await dispatcher.activate(store, new Date("2026-07-22T01:00:00.000Z")); const trigger = await appendTrigger(store, "language-test", SECRET, "2026-07-22T01:00:01.000Z"); await appendTerminalRun(store, trigger.event, "failed", 1, SECRET); await new OperationalIncidentProjector().project(store); const result = await dispatcher.sendPending(store, { since, now: new Date("2026-07-22T01:00:10.000Z"), }); expect(result.delivered).toBe(1); expect(telegram.messages).toHaveLength(1); const message = telegram.messages[0]!; // Must contain the accepted The Stream header expect(message).toContain("The Stream ·"); // Must contain retry disposition expect(message).toMatch(/Will retry|Won't retry automatically/); // Must contain a short receipt hash expect(message).toMatch(/Receipt [a-f0-9]/); // Must not contain raw field labels const rawLabels = [ "component:", "source:", "classification:", "retry:", "progress:", ]; for (const term of rawLabels) { expect(message).not.toContain(term); } // Must not contain agent ids, event ids, or run ids expect(message).not.toContain("agent-run-failed"); expect(message).not.toContain("incidentId"); expect(message).not.toContain("fingerprint"); }); test("does not suppress canary agent incidents", async () => { const { store } = await fixtureStore(); const telegram = await telegramFixture("success"); const dispatcher = new IncidentTelegramDispatcher({ id: "incident-dispatcher:canary-alert", client: new TelegramBotClient({ token: "fixture-token", baseUrl: telegram.baseUrl }), chatId: "123456789", categories: ["agent-run-failed"], }); const since = await dispatcher.activate(store, new Date("2026-07-22T01:00:00.000Z")); const trigger = await appendTrigger(store, "canary-failure", SECRET, "2026-07-22T01:00:01.000Z"); await appendTerminalRun(store, trigger.event, "failed", 1, SECRET); await new OperationalIncidentProjector().project(store); const result = await dispatcher.sendPending(store, { since, now: new Date("2026-07-22T01:00:10.000Z"), }); expect(result.delivered).toBe(1); expect(telegram.messages).toHaveLength(1); expect(telegram.calls).toBe(1); }); test("does not suppress retryable failures with unchanged progress", async () => { const { store } = await fixtureStore(); const telegram = await telegramFixture("success"); const dispatcher = new IncidentTelegramDispatcher({ id: "incident-dispatcher:retryable-alert", client: new TelegramBotClient({ token: "fixture-token", baseUrl: telegram.baseUrl }), chatId: "123456789", categories: ["agent-run-failed"], }); const since = await dispatcher.activate(store, new Date("2026-07-22T01:00:00.000Z")); const trigger = await appendTrigger(store, "retryable-unchanged", SECRET, "2026-07-22T01:00:01.000Z"); await appendTerminalRun(store, trigger.event, "failed", 1, SECRET); await new OperationalIncidentProjector().project(store); const incidents = await listOperationalIncidents(store); const failedIncident = incidents.find(({ incident }) => incident.category === "agent-run-failed"); expect(failedIncident?.incident).toMatchObject({ retryable: true, progress: "unchanged" }); const result = await dispatcher.sendPending(store, { since, now: new Date("2026-07-22T01:00:10.000Z"), }); expect(result.delivered).toBe(1); expect(telegram.messages).toHaveLength(1); expect(telegram.calls).toBe(1); }); test("alerts for terminal connector failures with clear format", async () => { const { store } = await fixtureStore(); const telegram = await telegramFixture("success"); const dispatcher = new IncidentTelegramDispatcher({ id: "incident-dispatcher:terminal-alert", client: new TelegramBotClient({ token: "fixture-token", baseUrl: telegram.baseUrl }), chatId: "123456789", categories: ["connector-terminal"], }); const since = await dispatcher.activate(store, new Date("2026-07-22T01:00:00.000Z")); await store.appendEvent(connectorEvent( "stream.thought.connector.subscription.stopped", "connector-terminal-alert", { status: "failed", reason: SECRET }, "2026-07-22T01:00:01.000Z", )); await new OperationalIncidentProjector().project(store); const result = await dispatcher.sendPending(store, { since, now: new Date("2026-07-22T01:00:10.000Z"), }); expect(result.delivered).toBe(1); expect(telegram.messages).toHaveLength(1); const message = telegram.messages[0]!; expect(message).toContain("The Stream ·"); expect(message).not.toContain("component:"); expect(message).not.toContain("source:"); expect(message).not.toContain("classification:"); expect(message).toMatch(/Receipt [a-f0-9]/); }); }); async function fixtureStore(): Promise<{ project: string; store: JazzThoughtStore }> { const project = await temporaryProject("thoughtstream-incidents-"); roots.push(project); const store = await testStore(project); stores.push(store); return { project, store }; } function connectorEvent( type: string, id: string, payload: Record, occurredAt?: string, ): EventCandidate { return { type, schemaVersion: 1, source: "jetstream:cameron-bluesky", sourceKind: "jetstream", externalId: id, idempotencyKey: id, occurredAt: occurredAt ?? (id === "connector-failed" ? "2026-07-22T01:00:00.000Z" : id === "connector-recovered" ? "2026-07-22T01:00:01.000Z" : "2026-07-22T01:00:02.000Z"), actor: "jetstream:cameron-bluesky", correlationId: id, privacy: "public-source", payload: payload as EventCandidate["payload"], }; } async function appendTrigger( store: JazzThoughtStore, id: string, sourceSentinel: string, occurredAt = "2026-07-22T01:00:03.000Z", ) { return store.appendEvent({ type: "stream.thought.source.atproto.commit", schemaVersion: 1, source: "jetstream:cameron-bluesky", sourceKind: "jetstream", externalId: id, idempotencyKey: id, occurredAt, actor: "did:plc:fixture", correlationId: id, privacy: "public-source", payload: { collection: "app.bsky.feed.post", text: sourceSentinel }, }); } async function appendTerminalRun( store: JazzThoughtStore, trigger: ThoughtEvent, status: Extract, attempt: number, errorSentinel: string, deferred = false, ): Promise { const runId = `run_${trigger.externalId}_${status}_${attempt}`; const at = new Date(Date.parse(trigger.occurredAt) + 1_000).toISOString(); const run: AgentRun = { id: runId, executionKey: `execution_${trigger.externalId}`, triggerEventId: trigger.id, agentId: "resident-fixture", agentVersion: 1, status, inputEventIds: [trigger.id], outputEventIds: [], attempt, provider: "letta-cloud", model: "fixture-model", promptHash: "fixture-prompt", contextManifest: {}, errorText: errorSentinel, result: { failureDiagnostic: { code: status === "failed" ? "provider-run-failed" : `agent-run-${status}`, stage: status === "failed" ? "provider" : "recovery", ...(deferred ? { progressDisposition: "deferred", retryAt: "2026-07-22T01:05:00.000Z", } : {}), rawProviderBody: errorSentinel, }, }, createdAt: at, startedAt: at, completedAt: at, updatedAt: at, }; await store.upsertRun(run); await store.appendEvent({ type: `stream.thought.agent.run.${status}`, schemaVersion: 1, source: "agent:resident-fixture", sourceKind: "agent", externalId: runId, idempotencyKey: `${runId}:${status}`, occurredAt: at, actor: "resident-fixture", rootEventId: trigger.rootEventId, parentEventId: trigger.id, correlationId: runId, privacy: "sensitive", payload: { runId, agentId: run.agentId, agentVersion: run.agentVersion, inputEventIds: run.inputEventIds, attempt, status, error: errorSentinel, failureDiagnostic: run.result!.failureDiagnostic!, }, }); return run; } async function appendCompletedRun(store: JazzThoughtStore, trigger: ThoughtEvent, attempt = 1): Promise { const at = new Date(Date.parse(trigger.occurredAt) + 1_000).toISOString(); await store.upsertRun({ id: `run_${trigger.externalId}_completed_${attempt}`, executionKey: `execution_${trigger.externalId}`, triggerEventId: trigger.id, agentId: "resident-fixture", agentVersion: 1, status: "completed", inputEventIds: [trigger.id], outputEventIds: [], attempt, provider: "letta-cloud", model: "fixture-model", promptHash: "fixture-prompt", contextManifest: {}, result: { summary: "Successful private observation", tags: [], importance: "normal", confidence: 1 }, createdAt: at, startedAt: at, completedAt: at, updatedAt: at, }); } async function telegramFixture(mode: "success" | "failure"): Promise<{ baseUrl: string; messages: string[]; calls: number; }> { const messages: string[] = []; const state = { calls: 0 }; const server = http.createServer(async (request, response) => { state.calls += 1; const chunks: Buffer[] = []; for await (const chunk of request) chunks.push(Buffer.from(chunk)); const body = JSON.parse(Buffer.concat(chunks).toString("utf8")) as { text?: string }; if (body.text) messages.push(body.text); response.setHeader("content-type", "application/json"); if (mode === "failure") { response.statusCode = 500; response.end(JSON.stringify({ ok: false, error_code: 500, description: SECRET })); return; } response.statusCode = 200; response.end(JSON.stringify({ ok: true, result: { message_id: 42, date: Math.floor(Date.now() / 1_000), chat: { id: 123456789, type: "private", first_name: "Fixture" }, text: body.text ?? "", }, })); }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); servers.push(server); const address = server.address(); if (!address || typeof address === "string") throw new Error("Fixture server did not obtain a port"); return { baseUrl: `http://127.0.0.1:${address.port}`, messages, get calls() { return state.calls; }, }; }