/// import { EXTENSION_OPERATIONS, EXTENSION_FAILURES, type ExtensionDiagnosticInput, } from "../shared/extension-diagnostics.ts"; import { createHash } from "node:crypto"; import { safeMcpObservability } from "./mcp-observability.ts"; import type { Agent } from "agents"; import { genericObservability, type Observability, type ObservabilityEvent, } from "agents/observability"; import type { StepContext } from "@cloudflare/think"; import { z } from "zod"; import { DIAGNOSTIC_LIMITS, type DiagnosticEvent, type DiagnosticSnapshot, } from "../shared/diagnostics.ts"; import { MODEL_CATALOG } from "./model-settings.ts"; import { MODEL_PROVIDERS, isModelProvider, type ModelProvider, } from "../shared/model-providers.ts"; import type { ToolActivity } from "../shared/tool-activity"; import type { TaskRun } from "../shared/tasks"; const count = z .number() .finite() .nonnegative() .max(Number.MAX_SAFE_INTEGER) .nullable(); const id = z .string() .regex(/^[a-f0-9]{64}$/) .nullable(); const modelNames: string[] = Object.values(MODEL_CATALOG).flat(); const toolNames = [ "remember", "updateMemory", "forget", "recall", "web_search", "read_url", "browser_read", "research", "shell", "createSchedule", "activate_skill", "read_skill_resource", ]; const detailsSchema = z .object({ extensionOperation: z.enum(EXTENSION_OPERATIONS).optional(), extensionFailure: z.enum(EXTENSION_FAILURES).optional(), capabilityVersion: id.optional(), approvalResolution: z.enum(["approved", "denied"]).optional(), capabilityId: id.optional(), capabilityFingerprint: id.optional(), capabilitySourceId: id.optional(), capabilitySourceKind: z.enum(["mcp", "skill"]).optional(), approvalDecision: z .enum(["allow", "ask", "never", "unavailable"]) .optional(), policyRevision: z .number() .int() .nonnegative() .max(Number.MAX_SAFE_INTEGER) .optional(), policySource: z.enum(["capability", "source", "default"]).optional(), provider: z .enum([ ...(Object.keys(MODEL_PROVIDERS) as ModelProvider[]), "unknown" as const, ]) .optional(), model: z .string() .refine((v) => modelNames.includes(v as (typeof modelNames)[number])) .nullable() .optional(), httpStatus: z.number().int().min(400).max(599).nullable().optional(), failureStage: z.enum(["request", "stream-read", "stream-event"]).optional(), errorCategory: z .enum([ "authentication", "billing", "rate-limit", "not-found", "upstream-unavailable", "invalid-request", "unknown", ]) .optional(), gatewayErrorCode: z .number() .int() .min(1000) .max(9999) .nullable() .optional(), cfRay: z .string() .regex(/^[a-f0-9]{16}-[A-Z]{3}$/) .nullable() .optional(), step: count.optional(), responseTimeMs: count.optional(), stepTimeMs: count.optional(), timeToFirstOutputMs: count.optional(), inputTokens: count.optional(), outputTokens: count.optional(), totalTokens: count.optional(), cacheReadTokens: count.optional(), cacheWriteTokens: count.optional(), reasoningTokens: count.optional(), cost: z.null().optional(), tool: z .string() .refine((v) => toolNames.includes(v)) .nullable() .optional(), toolKind: z .enum(["memory", "web", "browser", "shell", "schedule", "other"]) .optional(), attempts: count.optional(), reason: z .enum([ "tool-error", "turn-cancelled", "interrupted", "incomplete", "approval-denied", ]) .nullable() .optional(), taskId: id.optional(), submissionId: id.optional(), source: z.enum(["scheduled", "manual"]).optional(), claimedAt: count.optional(), completedAt: count.optional(), claimToCompletionMs: count.optional(), code: count.optional(), trigger: z .enum([ "ws-chat", "rpc", "submission", "programmatic", "agent-tool", "unknown", ]) .optional(), }) .strict(); const eventSchema = z .object({ kind: z.enum([ "extension", "approval", "turn", "model", "model-attempt", "tool", "task", "connection", "recovery", ]), timestamp: z.number().finite().nonnegative().max(Number.MAX_SAFE_INTEGER), id, conversationId: id, requestId: id, phase: z.enum(["start", "finish", "transition"]), status: z.enum([ "started", "completed", "error", "aborted", "skipped", "pending", "running", "succeeded", "failed", "cancelled", "dispatching", "dispatch_error", "detected", "connected", "disconnected", "unknown", ]), durationMs: count, details: detailsSchema, }) .strict(); /** IDs can be client/provider supplied. Export fixed hashes, never raw strings. */ export function diagnosticId(value: unknown): string | null { return typeof value === "string" && value.length > 0 && value.length <= 512 ? createHash("sha256").update(value).digest("hex") : null; } export function metric(value: unknown): number | null { return typeof value === "number" && Number.isFinite(value) && value >= 0 && value <= Number.MAX_SAFE_INTEGER ? value : null; } const elapsed = (start: number | null, end: number | null) => start !== null && end !== null ? metric(end - start) : null; export function validateDiagnostic(value: unknown): DiagnosticEvent { return eventSchema.parse(value); } function base( kind: DiagnosticEvent["kind"], rawId: unknown, timestamp = Date.now(), ): DiagnosticEvent { return { kind, timestamp, id: diagnosticId(rawId), conversationId: null, requestId: null, phase: "transition", status: "unknown", durationMs: null, details: {}, }; } export function modelDiagnostic(ctx: StepContext): DiagnosticEvent { const event = base("model", `${ctx.callId}:${ctx.stepNumber}`); event.requestId = diagnosticId( ctx.runtimeContext?.["cloudflare.agents.turn.request_id"], ); event.phase = "finish"; event.status = "completed"; const provider = isModelProvider(ctx.model.provider) ? ctx.model.provider : ctx.model.provider.startsWith("workers-ai") ? "workers-ai" : ctx.model.provider.startsWith("anthropic") ? "anthropic" : "unknown"; const tokens = (value: unknown) => { const n = metric(value); return (provider === "workers-ai" || provider === "unknown") && n === 0 ? null : n; }; event.durationMs = metric(ctx.performance.responseTimeMs); event.details = { provider, model: modelNames.includes(ctx.model.modelId as (typeof modelNames)[number]) ? ctx.model.modelId : null, step: metric(ctx.stepNumber), responseTimeMs: metric(ctx.performance.responseTimeMs), stepTimeMs: metric(ctx.performance.stepTimeMs), timeToFirstOutputMs: metric(ctx.performance.timeToFirstOutputMs), inputTokens: tokens(ctx.usage.inputTokens), outputTokens: tokens(ctx.usage.outputTokens), totalTokens: tokens(ctx.usage.totalTokens), cacheReadTokens: tokens(ctx.usage.inputTokenDetails.cacheReadTokens), cacheWriteTokens: tokens(ctx.usage.inputTokenDetails.cacheWriteTokens), reasoningTokens: tokens(ctx.usage.outputTokenDetails.reasoningTokens), cost: null, }; return event; } export type ToolDiagnostic = Pick< ToolActivity, | "toolCallId" | "toolName" | "kind" | "status" | "attempts" | "createdAt" | "startedAt" | "endedAt" | "updatedAt" | "reason" | "capability" | "approvalDecision" >; export function toolDiagnostic(activity: ToolDiagnostic): DiagnosticEvent { return { ...base("tool", activity.toolCallId, activity.updatedAt), status: activity.status, durationMs: elapsed(metric(activity.startedAt), metric(activity.endedAt)), details: { tool: toolNames.includes(activity.toolName) ? activity.toolName : null, toolKind: activity.kind, attempts: metric(activity.attempts), reason: activity.reason ?? null, ...(activity.approvalDecision ? { approvalResolution: activity.approvalDecision } : {}), ...(activity.capability ? { capabilityVersion: diagnosticId(activity.capability.version), capabilityId: diagnosticId(activity.capability.id), capabilityFingerprint: diagnosticId( activity.capability.fingerprint, ), capabilitySourceId: diagnosticId(activity.capability.source.id), capabilitySourceKind: activity.capability.source.kind, } : {}), }, }; } export function extensionDiagnostic( input: ExtensionDiagnosticInput, ): DiagnosticEvent { const now = Date.now(); return { ...base("extension", null, now), phase: "finish", status: input.status, durationMs: elapsed(metric(input.startedAt), now), details: { extensionOperation: input.operation, ...(input.failure ? { extensionFailure: input.failure } : {}), ...(input.source ? { capabilitySourceKind: input.source.kind, capabilitySourceId: diagnosticId(input.source.id), } : {}), capabilityVersion: diagnosticId(input.version), capabilityFingerprint: diagnosticId(input.fingerprint), }, }; } export type TaskDiagnostic = Pick< TaskRun, | "id" | "taskId" | "conversationId" | "submissionId" | "source" | "status" | "startedAt" | "completedAt" >; export function taskDiagnostic(run: TaskDiagnostic): DiagnosticEvent { const claimedAt = run.startedAt ? metric(Date.parse(run.startedAt)) : null; const completedAt = run.completedAt ? metric(Date.parse(run.completedAt)) : null; return { ...base("task", run.id), conversationId: diagnosticId(run.conversationId), status: run.status, details: { taskId: diagnosticId(run.taskId), submissionId: diagnosticId(run.submissionId), source: run.source, claimedAt, completedAt, claimToCompletionMs: elapsed(claimedAt, completedAt), }, }; } interface Row { sequence: number; event: string; } const bytes = (value: string) => new TextEncoder().encode(value).length; /** This store owns only its two named app tables. No SDK retention/idempotency writes. */ export class DiagnosticStore { private ready = false; private readonly sql: Agent["sql"]; constructor(sql: Agent["sql"]) { this.sql = sql; } private initialize() { if (this.ready) return; this .sql`CREATE TABLE IF NOT EXISTS flarebot_diagnostics (sequence INTEGER PRIMARY KEY AUTOINCREMENT, timestamp INTEGER NOT NULL, kind TEXT NOT NULL, id TEXT, conversation_id TEXT, bytes INTEGER NOT NULL, event TEXT NOT NULL)`; this .sql`CREATE TABLE IF NOT EXISTS flarebot_diagnostic_meta (singleton INTEGER PRIMARY KEY CHECK(singleton = 1), generation INTEGER NOT NULL, pruned INTEGER NOT NULL)`; this.sql`INSERT OR IGNORE INTO flarebot_diagnostic_meta VALUES (1, 0, 0)`; this.ready = true; } generation(): number { try { this.initialize(); return this.sql<{ generation: number; }>`SELECT generation FROM flarebot_diagnostic_meta WHERE singleton = 1`[0] .generation; } catch { return -1; } } clear() { try { this.initialize(); this .sql`UPDATE flarebot_diagnostic_meta SET generation = generation + 1 WHERE singleton = 1`; this.sql`DELETE FROM flarebot_diagnostics`; } catch { /* Best-effort metadata cannot fail the native lifecycle. */ } } deleteConversation(rawId: string) { try { this.initialize(); this .sql`DELETE FROM flarebot_diagnostics WHERE conversation_id = ${diagnosticId(rawId)}`; } catch { /* Same failure isolation as collection. */ } } deleteTask(rawId: string) { try { this.initialize(); this .sql`DELETE FROM flarebot_diagnostics WHERE kind = 'task' AND json_extract(event, '$.details.taskId') = ${diagnosticId(rawId)}`; } catch { /* Optional history cleanup. */ } } private prune(now: number) { const [{ count: before }] = this.sql<{ count: number; }>`SELECT count(*) AS count FROM flarebot_diagnostics`; this .sql`DELETE FROM flarebot_diagnostics WHERE timestamp < ${now - DIAGNOSTIC_LIMITS.maxAgeMs}`; this .sql`DELETE FROM flarebot_diagnostics WHERE sequence NOT IN (SELECT sequence FROM flarebot_diagnostics ORDER BY sequence DESC LIMIT ${DIAGNOSTIC_LIMITS.maxRows})`; this .sql`DELETE FROM flarebot_diagnostics WHERE sequence IN (SELECT sequence FROM (SELECT sequence, sum(bytes) OVER (ORDER BY sequence DESC) AS running_bytes FROM flarebot_diagnostics) WHERE running_bytes > ${DIAGNOSTIC_LIMITS.maxStoreBytes})`; const [{ count: after }] = this.sql<{ count: number; }>`SELECT count(*) AS count FROM flarebot_diagnostics`; this .sql`UPDATE flarebot_diagnostic_meta SET pruned = pruned + ${before - after} WHERE singleton = 1`; } hasStart(kind: "turn" | "connection", key: string | null): boolean { try { this.initialize(); return ( key !== null && this .sql`SELECT 1 FROM flarebot_diagnostics WHERE kind = ${kind} AND id = ${key} AND json_extract(event, '$.phase') = 'start' LIMIT 1` .length > 0 ); } catch { return false; } } record(event: DiagnosticEvent, generation?: number) { try { this.initialize(); if (generation !== undefined && generation !== this.generation()) return; const safe = validateDiagnostic(event); const serialized = JSON.stringify(safe); const size = bytes(serialized); if (size > DIAGNOSTIC_LIMITS.maxEventBytes) return; this .sql`INSERT INTO flarebot_diagnostics (timestamp, kind, id, conversation_id, bytes, event) VALUES (${safe.timestamp}, ${safe.kind}, ${safe.id}, ${safe.conversationId}, ${size}, ${serialized})`; this.prune(Date.now()); } catch { /* An optional diagnostic observer must never break a turn/tool/task. */ } } snapshot(conversationId: string | null = null): DiagnosticSnapshot { try { this.initialize(); this.prune(Date.now()); const rows = this .sql`SELECT sequence, event FROM flarebot_diagnostics WHERE conversation_id IS NULL OR conversation_id = ${conversationId} ORDER BY sequence DESC LIMIT ${DIAGNOSTIC_LIMITS.maxRows}`; const events: DiagnosticSnapshot["events"] = []; let omittedRows = 0; for (const row of rows) { try { events.push({ ...validateDiagnostic(JSON.parse(row.event)), sequence: row.sequence, }); } catch { omittedRows++; } } const [{ pruned }] = this.sql<{ pruned: number; }>`SELECT pruned FROM flarebot_diagnostic_meta WHERE singleton = 1`; return { events, oldestAt: events.at(-1)?.timestamp ?? null, newestAt: events[0]?.timestamp ?? null, prunedRows: pruned, omittedRows, available: true, }; } catch { return { events: [], oldestAt: null, newestAt: null, prunedRows: 0, omittedRows: 0, available: false, }; } } receiver(): Observability { return { emit: (event) => { try { const safe = safeMcpObservability(event); if (safe) genericObservability.emit(safe); } catch { /* Independent native sink. */ } try { this.native(event); } catch { /* Public SDK emit does not isolate observers. */ } }, }; } private native(native: ObservabilityEvent) { if ( native.type === "chat:turn:start" || native.type === "chat:turn:finish" ) { const e = base("turn", native.payload.requestId, native.timestamp); const start = native.type === "chat:turn:start"; if (!start && !this.hasStart("turn", e.id)) return; e.requestId = e.id; e.phase = start ? "start" : "finish"; e.status = start ? "started" : ["completed", "error", "aborted", "skipped"].includes( native.payload.status, ) ? (native.payload.status as DiagnosticEvent["status"]) : "unknown"; if (native.type === "chat:turn:finish") e.durationMs = metric(native.payload.durationMs); e.details.trigger = [ "ws-chat", "rpc", "submission", "programmatic", "agent-tool", ].includes(native.payload.trigger) ? (native.payload.trigger as NonNullable< DiagnosticEvent["details"]["trigger"] >) : "unknown"; this.record(e); } else if (native.type === "connect" || native.type === "disconnect") { const e = base( "connection", native.payload.connectionId, native.timestamp, ); if (native.type === "disconnect" && !this.hasStart("connection", e.id)) return; e.phase = native.type === "connect" ? "start" : "finish"; e.status = native.type === "connect" ? "connected" : "disconnected"; if (native.type === "disconnect") { e.details.code = metric(native.payload.code); const start = this.sql<{ timestamp: number; }>`SELECT timestamp FROM flarebot_diagnostics WHERE kind = 'connection' AND id = ${e.id} AND json_extract(event, '$.phase') = 'start' ORDER BY sequence DESC LIMIT 1`[0]; e.durationMs = elapsed( metric(start?.timestamp), metric(native.timestamp), ); } this.record(e); } else if (native.type === "chat:recovery:detected") { const e = base("recovery", native.payload.incidentId, native.timestamp); e.requestId = diagnosticId(native.payload.requestId); if (!this.hasStart("turn", e.requestId)) return; e.status = "detected"; e.details.attempts = metric(native.payload.attempt); this.record(e); } } }