Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459import type { ToolDiagnostic } from "./diagnostics";import { Think, type ChatResponseResult, type ChunkContext, type ToolCallContext, type ToolCallResultContext,} from "@cloudflare/think";import { callable } from "agents";import { TOOL_ACTIVITY_LIVE_LIMIT, TOOL_ACTIVITY_PAGE_LIMIT, TOOL_ACTIVITY_TEXT_LIMIT, type ToolActivity, type ToolActivityPage, type ToolActivityReason, type ToolActivityState, type ToolActivityStatus, type ToolKind, type ToolProgress,} from "../shared/tool-activity";
/** Trusted server configuration: select allowlisted fields, never stringify raw * input/output/errors. Text limits are a size boundary, not secret redaction. */export interface ToolActivityDescriptor { kind: ToolKind; label: string; inputSummary?: (input: unknown) => string; outputSummary?: (output: unknown) => string; outcome?: (output: unknown) => "succeeded" | "failed" | "pending"; progress?: (value: unknown) => ToolProgress | undefined;}const fallback: ToolActivityDescriptor = { kind: "other", label: "Tool invocation",};const terminal = (status: ToolActivityStatus) => status === "succeeded" || status === "failed" || status === "cancelled";const safeText = (value: string) => value .replace(/[\u0000-\u001f\u007f]/g, " ") .slice(0, TOOL_ACTIVITY_TEXT_LIMIT);function format<T>(formatter: (() => T) | undefined, fallbackValue: T): T { try { return formatter ? formatter() : fallbackValue; } catch { return fallbackValue; }}type Row = { sequence: number; activity: string };
/** Public Think hooks + native Agent persistence/state, with no stream replacement. */export class ActivityThink<Env extends Cloudflare.Env> extends Think< Env, ToolActivityState> { initialState: ToolActivityState = { toolActivityVersion: 1, toolActivities: [], }; private activityReady = false; private actionNames = new Set<string>();
protected async applicationToolNames() { this.actionNames = new Set( Object.entries(await this.getActions()).map( ([name, action]) => action.config.name ?? name, ), ); return [...Object.keys(this.getTools()), ...this.actionNames]; } private signals = new Map<string, () => void>(); private pendingProgress = new Map<string, ToolProgress>();
protected getToolActivityDescriptors(): Record< string, ToolActivityDescriptor > { return {}; }
private initializeActivities() { if (this.activityReady) return; this.sql`CREATE TABLE IF NOT EXISTS flarebot_tool_activity ( sequence INTEGER PRIMARY KEY AUTOINCREMENT, tool_call_id TEXT NOT NULL UNIQUE, status TEXT NOT NULL, activity TEXT NOT NULL )`; this.sql`CREATE INDEX IF NOT EXISTS flarebot_tool_activity_status ON flarebot_tool_activity(status, sequence DESC)`; this.activityReady = true; }
onStart() { this.initializeActivities(); // An isolate restart cannot certify that an external operation stopped or // completed. Close the old observation; let Think own recovery and replay. for (const row of this .sql<Row>`SELECT sequence, activity FROM flarebot_tool_activity WHERE status IN ('pending', 'running')`) { const activity = JSON.parse(row.activity) as ToolActivity; this.finish(activity.toolCallId, "failed", undefined, "interrupted"); } this.publishActivities(); }
private readActivity(id: string): ToolActivity | undefined { this.initializeActivities(); const row = this .sql<Row>`SELECT sequence, activity FROM flarebot_tool_activity WHERE tool_call_id = ${id}`[0]; return row ? ((JSON.parse(row.activity) as ToolActivity | null) ?? undefined) : undefined; }
private publishActivities() { const newest = (statuses: string[], limit: number) => statuses .flatMap( (status) => this.sql<Row>`SELECT sequence, activity FROM flarebot_tool_activity WHERE status = ${status} ORDER BY sequence DESC LIMIT ${limit}`, ) .sort((a, b) => b.sequence - a.sequence) .slice(0, limit); const active = newest(["pending", "running"], TOOL_ACTIVITY_LIVE_LIMIT); const rows = [ ...active, ...newest( ["succeeded", "failed", "cancelled"], TOOL_ACTIVITY_LIVE_LIMIT - active.length, ), ]; this.setState({ ...this.state, toolActivityVersion: 1, toolActivities: rows.map( (row) => JSON.parse(row.activity) as ToolActivity, ), }); }
protected onToolDiagnostic(_activity: ToolDiagnostic) {}
private writeActivity(activity: ToolActivity) { const previous = this.readActivity(activity.toolCallId); this.sql`INSERT INTO flarebot_tool_activity (tool_call_id, status, activity) VALUES (${activity.toolCallId}, ${activity.status}, ${JSON.stringify(activity)}) ON CONFLICT(tool_call_id) DO UPDATE SET status = excluded.status, activity = excluded.activity`; this.publishActivities(); if ( !previous || previous.status !== activity.status || previous.attempts !== activity.attempts ) { const { toolCallId, toolName, kind, status, attempts, createdAt, startedAt, endedAt, updatedAt, reason, } = activity; try { this.onToolDiagnostic({ toolCallId, toolName, kind, status, attempts, createdAt, startedAt, endedAt, updatedAt, reason, }); } catch { /* Optional observer. */ } } }
/** Descending immutable creation cursor, independent of the bounded live window. */ @callable() listToolActivities(before?: number): ToolActivityPage { if (before !== undefined && (!Number.isSafeInteger(before) || before < 1)) throw new Error("Invalid tool activity cursor"); this.initializeActivities(); const rows = this .sql<Row>`SELECT sequence, activity FROM flarebot_tool_activity WHERE status != 'cleared' AND sequence < ${before ?? Number.MAX_SAFE_INTEGER} ORDER BY sequence DESC LIMIT ${TOOL_ACTIVITY_PAGE_LIMIT + 1}`; const page = rows.slice(0, TOOL_ACTIVITY_PAGE_LIMIT); return { activities: page.map((row) => JSON.parse(row.activity) as ToolActivity), nextCursor: rows.length > page.length ? page.at(-1)!.sequence : null, }; }
private descriptor(name: string) { return this.getToolActivityDescriptors()[name] ?? fallback; }
private observe(id: string, name: string, input?: unknown) { const existing = this.readActivity(id); if (existing) return existing; // Clear removes presentation data but keeps an ID-only tombstone so native // callbacks that were already in flight cannot restore cleared history. if ( this .sql`SELECT 1 FROM flarebot_tool_activity WHERE tool_call_id = ${id} AND status = 'cleared'` .length ) return; const descriptor = this.descriptor(name); const now = Date.now(); const activity: ToolActivity = { toolCallId: id, toolName: descriptor === fallback ? "unknown" : name, kind: descriptor.kind, status: "pending", inputSummary: safeText( format( input === undefined ? undefined : () => descriptor.inputSummary?.(input) ?? descriptor.label, descriptor.label, ), ), attempts: 0, createdAt: now, updatedAt: now, }; this.writeActivity(activity); return activity; }
beforeToolCall(ctx: ToolCallContext) { const activity = this.observe(ctx.toolCallId, ctx.toolName, ctx.input); if ( !activity || (terminal(activity.status) && activity.reason !== "interrupted") ) return; const descriptor = this.descriptor(ctx.toolName); const now = Date.now(); this.writeActivity({ ...activity, status: "running", reason: undefined, endedAt: undefined, outputSummary: undefined, attempts: activity.attempts + (activity.status === "running" ? 0 : 1), startedAt: activity.startedAt ?? now, updatedAt: now, inputSummary: safeText( format( () => descriptor.inputSummary?.(ctx.input) ?? descriptor.label, descriptor.label, ), ), }); const abort = () => this.finish(ctx.toolCallId, "cancelled", undefined, "turn-cancelled"); if (ctx.abortSignal?.aborted) abort(); else if (ctx.abortSignal && !this.signals.has(ctx.toolCallId)) { ctx.abortSignal.addEventListener("abort", abort, { once: true }); this.signals.set(ctx.toolCallId, () => ctx.abortSignal!.removeEventListener("abort", abort), ); } }
/** Scalar tools/actions can report progress; generators are also observed via onChunk. * Only descriptor-approved text/counts survive. Coalesce to at most four writes/s/call. */ protected reportToolProgress(id: string, value: unknown) { const activity = this.readActivity(id); if (!activity || terminal(activity.status)) return; const descriptor = this.descriptor(activity.toolName); const progress = format(() => descriptor.progress?.(value), undefined); if (!progress) return; const bounded: ToolProgress = { text: safeText(progress.text) }; for (const key of ["completed", "total"] as const) { const count = progress[key]; if ( typeof count === "number" && Number.isSafeInteger(count) && count >= 0 ) bounded[key] = Math.min(count, 1_000_000_000); } this.pendingProgress.set(id, bounded); if (activity.progress && Date.now() - activity.updatedAt < 250) return; this.writeActivity({ ...activity, progress: bounded, updatedAt: Date.now(), }); this.pendingProgress.delete(id); }
private finish( id: string, status: "succeeded" | "failed" | "cancelled", output?: unknown, reason?: ToolActivityReason, reconcileInterrupted = false, ) { const activity = this.readActivity(id); if ( !activity || (terminal(activity.status) && !(reconcileInterrupted && activity.reason === "interrupted")) ) return; const descriptor = this.descriptor(activity.toolName); const summary = status === "cancelled" ? "Conversation stopped waiting for this tool" : reason === "interrupted" ? "Execution interrupted; outcome unknown" : status === "failed" ? "Tool invocation failed" : "Tool invocation completed"; const now = Date.now(); this.writeActivity({ ...activity, status, endedAt: now, updatedAt: now, outputSummary: safeText( status === "succeeded" ? format(() => descriptor.outputSummary?.(output) ?? summary, summary) : summary, ), reason, ...(reason === "interrupted" ? { interruptedAt: now } : {}), ...(this.pendingProgress.has(id) ? { progress: this.pendingProgress.get(id) } : {}), }); this.signals.get(id)?.(); this.signals.delete(id); this.pendingProgress.delete(id); }
private outcome(id: string, name: string, output: unknown) { if (!this.observe(id, name)) return; // Native actions catch failures into a reserved error envelope. SDK success // describes result delivery, not successful action execution. const actionFailed = this.actionNames.has(name) && typeof output === "object" && output !== null && "error" in output; const status = actionFailed ? "failed" : format( () => this.descriptor(name).outcome?.(output) ?? "succeeded", "failed", ); if (status === "pending") { const activity = this.readActivity(id)!; if (!terminal(activity.status)) this.writeActivity({ ...activity, status: "pending", updatedAt: Date.now(), }); } else this.finish( id, status, output, status === "failed" ? "tool-error" : undefined, true, ); }
afterToolCall(ctx: ToolCallResultContext) { this.observe(ctx.toolCallId, ctx.toolName, ctx.input); if (ctx.toolOutput.type === "tool-error") this.finish(ctx.toolCallId, "failed", undefined, "tool-error"); else this.outcome(ctx.toolCallId, ctx.toolName, ctx.toolOutput.output); }
onChunk({ chunk }: ChunkContext) { switch (chunk.type) { case "tool-input-start": this.observe(chunk.id, chunk.toolName); break; case "tool-call": this.observe(chunk.toolCallId, chunk.toolName, chunk.input); break; case "tool-result": if (chunk.preliminary) this.reportToolProgress(chunk.toolCallId, chunk.output); else this.outcome(chunk.toolCallId, chunk.toolName, chunk.output); break; } }
protected resetTurnState() { // Both native clearMessages and the WebSocket clear path use this public // reset seam. Synchronous tombstones precede cancellation/awaits: subsequent // turns keep their new rows, and old callbacks cannot resurrect summaries. this.initializeActivities(); this .sql`UPDATE flarebot_tool_activity SET status = 'cleared', activity = 'null' WHERE status != 'cleared'`; for (const remove of this.signals.values()) remove(); this.signals.clear(); this.pendingProgress.clear(); this.publishActivities(); super.resetTurnState(); }
onChatResponse(result: ChatResponseResult) { // The hook can run after the next turn starts. Only reconcile IDs belonging // to this native message; never sweep a global current-turn list. for (const part of result.message?.parts ?? []) { if (!("toolCallId" in part) || typeof part.toolCallId !== "string") continue; const activity = this.readActivity(part.toolCallId); if ( !activity || (terminal(activity.status) && activity.reason !== "interrupted") ) continue; if ( "state" in part && part.state === "output-available" && "output" in part && !("preliminary" in part && part.preliminary) ) { this.outcome(part.toolCallId, activity.toolName, part.output); continue; } if (result.status === "aborted") this.finish(part.toolCallId, "cancelled", undefined, "turn-cancelled"); else if ("state" in part && part.state === "output-error") this.finish(part.toolCallId, "failed", undefined, "tool-error"); else if ( "state" in part && (part.state === "approval-requested" || part.state === "approval-responded") ) continue; else this.finish(part.toolCallId, "failed", undefined, "incomplete"); } }}