Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623import * as Effect from "effect/Effect";import type { ResearchObservation } from "./research-tool";import type { ToolDiagnostic } from "./diagnostics";import { Think, type Action, 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 CapabilityActivity, 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; capability?: CapabilityActivity | ((input: unknown) => CapabilityActivity | undefined); 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>(); private actionDescriptors: Record<string, ToolActivityDescriptor> = {};
protected applicationToolNames() { return Effect.gen({ self: this }, function* () { const actions = yield* Effect.promise(() => Promise.resolve(this.getActions()), ); // The complete current native registry owns presentation lifetime. A // durable pause remains pending regardless of catalog size or refresh. this.actionDescriptors = Object.fromEntries( Object.entries(actions).flatMap(([name, action]) => { const descriptor = this.getActionActivityDescriptor(action); return descriptor ? [[action.config.name ?? name, descriptor]] : []; }), ); this.actionNames = new Set( Object.entries(actions).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 getActionActivityDescriptor( action: Action, ): ToolActivityDescriptor | undefined { if (action.config.kind !== "durable-pause") return; return { kind: "other", label: action.config.approvalSummary ?? action.config.description, outcome: (output) => output && typeof output === "object" && "status" in output && output.status === "paused" ? "pending" : "succeeded", }; }
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(); }
protected restoreToolApproval(id: string) { const activity = this.readActivity(id); if (activity?.reason === "interrupted") this.writeActivity({ ...activity, status: "pending", reason: undefined, endedAt: undefined, updatedAt: Date.now(), }); }
protected recordToolApproval(id: string, decision: "approve" | "deny") { this.restoreToolApproval(id); const activity = this.readActivity(id); if (!activity) return; const now = Date.now(); this.writeActivity({ ...activity, approvalDecision: decision === "approve" ? "approved" : "denied", status: decision === "approve" ? "running" : "pending", reason: undefined, endedAt: undefined, updatedAt: now, ...(decision === "approve" ? { startedAt: activity.startedAt ?? now } : {}), }); if (decision === "deny") this.finish(id, "cancelled", undefined, "approval-denied"); }
protected completeToolApproval( id: string, decision: "approve" | "deny", result: unknown, ) { this.restoreToolApproval(id); const failed = result !== null && typeof result === "object" && ("error" in result || ("status" in result && result.status === "error")); this.finish( id, decision === "deny" ? "cancelled" : failed ? "failed" : "succeeded", result, decision === "deny" ? "approval-denied" : failed ? "tool-error" : undefined, true, ); }
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) {}
protected async observeResearchCall<T>( call: ResearchObservation, run: () => Promise<T>, ): Promise<T> { this.beforeToolCall({ type: "tool-call", toolCallId: call.id, toolName: call.name, input: call.input, stepNumber: undefined, messages: [], abortSignal: call.signal, }); try { const output = await run(); this.outcome(call.id, call.name, output); return output; } catch (error) { this.finish( call.id, call.signal.aborted ? "cancelled" : "failed", undefined, call.signal.aborted ? "turn-cancelled" : "tool-error", ); throw error; } }
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, capability, } = activity; try { this.onToolDiagnostic({ toolCallId, toolName, kind, status, attempts, createdAt, startedAt, endedAt, updatedAt, reason, capability, }); } 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] ?? this.actionDescriptors[name] ?? fallback ); }
private observe(id: string, name: string, input?: unknown) { const existing = this.readActivity(id); if ( existing && (existing.capability || terminal(existing.status) || input === undefined) ) return existing; const descriptor = this.descriptor(name); // Streaming starts before parsed input exists. Enrich provenance once the // native call supplies it, then retain that identity through later refreshes. const capability = format(() => { if (typeof descriptor.capability !== "function") return descriptor.capability; return input === undefined ? undefined : descriptor.capability(input); }, undefined); const provenance = capability ? { ...capability, name: safeText(capability.name), source: { ...capability.source, name: safeText(capability.source.name), }, } : undefined; if (existing) { if (provenance) { const enriched = { ...existing, capability: provenance }; this.writeActivity(enriched); return enriched; } 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 now = Date.now(); const activity: ToolActivity = { toolCallId: id, toolName: descriptor === fallback ? "unknown" : name, kind: descriptor.kind, ...(provenance ? { capability: provenance } : {}), 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 = reason === "approval-denied" ? "Owner denied this action" : 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( (activity.approvalDecision === "approved" ? "Owner approved once. " : "") + (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"); } }}