///
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);
}
}
}