Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568import { DiagnosticStore, diagnosticId, modelDiagnostic, toolDiagnostic, type ToolDiagnostic,} from "./diagnostics";import type { PrepareStepContext, StepContext } from "@cloudflare/think";import type { TaskRun } from "../shared/tasks";import { submissionProjection } from "./task-execution";import type { ThinkSubmissionInspection } from "@cloudflare/think";import { requestedWebReader, WEB_INSTRUCTIONS } from "../shared/web";import { SHELL_INSTRUCTIONS } from "../shared/shell";import { scheduleActionInput, scheduleActionResult, SCHEDULE_INSTRUCTIONS,} from "./schedule-action";import { createShellTool, shellActivity } from "./shell-tool";import { createWebTools, webActivityDescriptors } from "./web-tools";import { ActivityThink, type ToolActivityDescriptor } from "./tool-activity";import { action, type ThinkModel, type TurnConfig, type TurnContext, type ActionContext,} from "@cloudflare/think";import { tool } from "ai";import { z } from "zod";import { MAX_MEMORY_LENGTH, MAX_MEMORY_QUERY_LENGTH, memoryContext,} from "../shared/memory";import { callable, type Connection, type ConnectionContext } from "agents";import { PersonalAgent, type Env } from "./personal-agent";import { Secret } from "../configuration/secrets";import { DEFAULT_MODEL, parseModelConfiguration, type ModelConfiguration, type ConversationModelSettings,} from "./model-settings";import { createConfiguredModel, modelProviderOptions, protectModel,} from "./model-provider";import type { DiagnosticEvent } from "../shared/diagnostics";import { isCredentialProvider, missingProviderKey,} from "../shared/model-providers";import { connectSession, closeSession, expireSession, type SessionConnection, type SessionExpiry,} from "./socket-session";
// One Think instance owns one transcript, queue and resumable stream. Native// child facets give each conversation independent SQLite and execution state.export class Conversation extends ActivityThink<Env> { protected readonly diagnostics = new DiagnosticStore(this.sql.bind(this)); observability = this.diagnostics.receiver();
diagnosticSnapshot() { return this.diagnostics.snapshot(); }
onStepEnd(ctx: StepContext) { const inherited = super.onStepEnd(ctx); try { const event = modelDiagnostic(ctx); // Clear removes the admitted start before native cancellation; late steps // cannot restore the cleared history, even if a new turn is now running. if (this.diagnostics.hasStart("turn", event.requestId)) this.diagnostics.record(event); } catch { /* Optional metadata cannot fail an awaited Think hook. */ } return inherited; }
protected onToolDiagnostic(activity: ToolDiagnostic) { this.diagnostics.record(toolDiagnostic(activity)); }
maxSteps = 8; chatStreamStallTimeoutMs = 60_000; chatRecovery = { maxAttempts: 3, terminalMessage: "The response was interrupted. Please try again.", }; includeMcpTools = false; workspaceBash = false; storeMessages = false; storeTools = false; private requiredWebReader: "read_url" | "browser_read" | undefined;
private modelOverride(): ModelConfiguration | null { this.sql`CREATE TABLE IF NOT EXISTS flarebot_conversation_model ( singleton INTEGER PRIMARY KEY CHECK (singleton = 1), configuration TEXT NOT NULL )`; const row = this.sql<{ configuration: string; }>`SELECT configuration FROM flarebot_conversation_model WHERE singleton = 1`[0]; return row ? parseModelConfiguration(JSON.parse(row.configuration)) : null; }
@callable() async getConversationModelSettings(): Promise<ConversationModelSettings> { const parent = await this.parentAgent(PersonalAgent); return { ...(await parent.getModelSettings()), override: this.modelOverride(), }; }
async submitTaskRun(run: TaskRun) { const parent = await this.parentAgent(PersonalAgent); const task = await parent.authorizeTaskRun(run); if (!task || task.conversationId !== this.name) return null; const { configuration } = await parent.getModelSettings(); const receipt = await this.submitMessages( [ { id: run.submissionId, role: "user", parts: [{ type: "text", text: task.instructions }], metadata: { scheduledModelConfiguration: configuration }, }, ], { submissionId: run.submissionId, idempotencyKey: run.submissionId, metadata: { taskRun: run, modelConfiguration: configuration }, }, ); if (!(await parent.authorizeTaskRun(run))) await this.cancelSubmission( run.submissionId, "Task changed or was deleted", ); return receipt; }
inspectTaskRun(id: string) { return this.inspectSubmission(id); } cancelTaskRun(id: string) { return this.cancelSubmission(id, "Task changed or was deleted"); }
protected async onSubmissionStatus(submission: ThinkSubmissionInspection) { await super.onSubmissionStatus(submission); const run = submission.metadata?.taskRun as TaskRun | undefined; if ( !run || run.submissionId !== submission.submissionId || run.conversationId !== this.name ) return; if (submission.status === "pending" || submission.status === "running") { let allowed = false; try { const parent = await this.parentAgent(PersonalAgent); allowed = Boolean(await parent.authorizeTaskRun(run)); } catch { /* Fail closed. */ } if (!allowed) { // Running is emitted before Think applies queued messages. Explicit native // cancellation also closes that queued-start race; throwing would not. await this.cancelSubmission( run.submissionId, "Task changed or unavailable", ); return; } } const parent = await this.parentAgent(PersonalAgent); await parent.projectTaskRun(submissionProjection(run, submission)); }
getModel(): ThinkModel { // Think resolves a synchronous default before it invokes beforeTurn. return DEFAULT_MODEL.model; }
protected createModel(configuration: ModelConfiguration, key?: Secret) { return createConfiguredModel( this.env.AI, configuration, key, this.sessionAffinity, ); }
async beforeTurn(ctx: TurnContext): Promise<TurnConfig> { const diagnosticGeneration = this.diagnostics.generation(); const parent = await this.parentAgent(PersonalAgent); // Native history is already persisted here; unlike model-facing ctx.messages // it contains neither synthetic continuation prompts nor private context. const firstUser = this.messages.find((message) => message.role === "user"); if (firstUser) { const text = firstUser.parts .filter((part) => part.type === "text") .map((part) => part.text) .join(" "); this.ctx.waitUntil( parent .nameConversationFromFirstMessage(this.name, text) .catch(() => {}), ); } const lastUser = ctx.messages.findLast( (message) => message.role === "user", ); const query = ( typeof lastUser?.content === "string" ? lastUser.content : (lastUser?.content .filter((part) => part.type === "text") .map((part) => part.text) .join(" ") ?? "") ).slice(0, MAX_MEMORY_QUERY_LENGTH); // Native Think durably captures the custom body with the submitted turn and // restores it for recovery. Never resolve a continuation from mutable UI state. // Submission metadata belongs to Think's queue ledger, not activeTurnMetadata. // Carry the task snapshot on its durable user message as well, so recovery // cannot accidentally reuse the last interactive WebSocket request body. const metadata = this.messages.findLast( (message) => message.role === "user", )?.metadata as { scheduledModelConfiguration?: unknown } | undefined; const scheduled = metadata?.scheduledModelConfiguration !== undefined; const selected = scheduled ? metadata?.scheduledModelConfiguration : ctx.body?.modelConfiguration; const override = selected === undefined ? scheduled ? undefined : (this.modelOverride() ?? undefined) : parseModelConfiguration(selected); if ( !scheduled && ctx.body?.useDefaultModel !== undefined && typeof ctx.body.useDefaultModel !== "boolean" ) throw new Error("Invalid model selection"); const [{ configuration, apiKey }, instructions, memories] = await Promise.all([ parent.readModelConfiguration(override), parent.readInstructions(), parent.searchMemories(query), ]); if (isCredentialProvider(configuration.provider) && !apiKey) throw new Error(missingProviderKey(configuration.provider)); if (!scheduled && !ctx.continuation && selected !== undefined) { this.modelOverride(); // Initialize this facet's table, including legacy chats. if (ctx.body?.useDefaultModel === true) this.sql`DELETE FROM flarebot_conversation_model WHERE singleton = 1`; else this .sql`INSERT INTO flarebot_conversation_model VALUES (1, ${JSON.stringify(configuration)}) ON CONFLICT(singleton) DO UPDATE SET configuration = excluded.configuration`; } const model = protectModel( this.createModel(configuration, apiKey ? new Secret(apiKey) : undefined), configuration.provider, (event) => { const diagnostic: DiagnosticEvent = { kind: "model-attempt", timestamp: Date.now(), id: diagnosticId(event.attemptId), conversationId: null, requestId: null, phase: event.phase, status: event.status, durationMs: event.durationMs, details: { provider: configuration.provider, model: configuration.model, cost: null, ...(event.status === "error" ? { httpStatus: event.httpStatus, failureStage: event.failureStage, errorCategory: event.errorCategory, gatewayErrorCode: event.gatewayErrorCode, cfRay: event.cfRay, } : {}), }, }; this.diagnostics.record(diagnostic, diagnosticGeneration); if (event.status === "error") console.error("flarebot.model-error", JSON.stringify(diagnostic)); }, ); // Think also assembles workspace/context/client tools by default. Only // explicitly supplied application tools are enabled at this stage. this.requiredWebReader = requestedWebReader(ctx.messages, ctx.continuation); return { model, // A complete native override replaces the frozen fallback prompt. Never // append to ctx.system: it can contain an obsolete/default instruction set. instructions: instructions + memoryContext(memories) + WEB_INSTRUCTIONS + SHELL_INSTRUCTIONS + SCHEDULE_INSTRUCTIONS + `\nCurrent UTC time: ${new Date().toISOString()}\n`, activeTools: await this.applicationToolNames(), providerOptions: modelProviderOptions(configuration), // One bounded total-output allowance for explicit reasoning, not an // effort-to-budget mapping. Leave existing provider-default turns unchanged. maxOutputTokens: configuration.effort ? 16384 : 4096, }; }
beforeStep(ctx: PrepareStepContext) { if (ctx.stepNumber !== 0 || !this.requiredWebReader) return; return { toolChoice: { type: "tool" as const, toolName: this.requiredWebReader, }, }; }
private turnGeneration = 0;
protected resetTurnState() { this.turnGeneration++; this.diagnostics.clear(); super.resetTurnState(); }
private async memoryParent(ctx: ActionContext) { const generation = this.turnGeneration; ctx.signal.throwIfAborted(); const parent = await this.parentAgent(PersonalAgent); ctx.signal.throwIfAborted(); if (generation !== this.turnGeneration) throw new Error("Memory operation interrupted"); return parent; }
getActions() { const content = z.string().min(1).max(MAX_MEMORY_LENGTH); const id = z .string() .regex(/^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/); const version = z.number().int().positive().max(Number.MAX_SAFE_INTEGER); return { createSchedule: action({ description: "Create a durable task in this conversation only at the user's explicit request. Clarify an exact time and timezone first: recurring schedules are UTC only; one-off instants require an explicit offset. Returns the saved task, actual next run, status and editable task URL.", inputSchema: scheduleActionInput, idempotencyKey: ({ ctx }) => ctx.toolCallId, execute: async (input, ctx) => { const generation = this.turnGeneration; ctx.signal.throwIfAborted(); // Native action settlement and parent persistence are separate commits. // Replays of the same tool call must reach the same task even after a // lost parent reply; prompt text is intentionally not the identity. const digest = await crypto.subtle.digest( "SHA-256", new TextEncoder().encode( JSON.stringify(["schedule", this.name, ctx.toolCallId]), ), ); ctx.signal.throwIfAborted(); if (generation !== this.turnGeneration) throw new Error("Schedule operation interrupted"); const hex = Array.from(new Uint8Array(digest)) .slice(0, 16) .map((byte) => byte.toString(16).padStart(2, "0")) .join(""); const taskId = `${hex.slice(0, 8)}-${hex.slice(8, 12)}-${hex.slice(12, 16)}-${hex.slice(16, 20)}-${hex.slice(20)}`; const parent = await this.parentAgent(PersonalAgent); ctx.signal.throwIfAborted(); if (generation !== this.turnGeneration) throw new Error("Schedule operation interrupted"); return scheduleActionResult( await parent.createTaskForConversation(this.name, taskId, input), ); }, }), remember: action({ description: "Save one fact only when the user explicitly asks you to remember it. Never automatically harvest conversation history.", inputSchema: z.object({ content }).strict(), idempotencyKey: ({ ctx }) => ctx.toolCallId, execute: async ({ content }, ctx) => { const generation = this.turnGeneration; // The parent write and native action ledger are separate commits. A // deterministic server-derived ID closes the lost-reply duplicate gap. const digest = await crypto.subtle.digest( "SHA-256", new TextEncoder().encode( JSON.stringify([this.name, ctx.toolCallId]), ), ); const hex = Array.from(new Uint8Array(digest)) .slice(0, 16) .map((byte) => byte.toString(16).padStart(2, "0")) .join(""); const factId = `${hex.slice(0, 8)}-${hex.slice(8, 12)}-${hex.slice(12, 16)}-${hex.slice(16, 20)}-${hex.slice(20)}`; const parent = await this.memoryParent(ctx); ctx.signal.throwIfAborted(); if (generation !== this.turnGeneration) throw new Error("Memory operation interrupted"); return parent.rememberForConversation(this.name, factId, content); }, }), updateMemory: action({ description: "Update a saved fact at the user's request. Recall the current fact ID and version first; stale edits fail.", inputSchema: z.object({ id, content, version }).strict(), execute: async ({ id, content, version }, ctx) => { const parent = await this.memoryParent(ctx); ctx.signal.throwIfAborted(); return parent.updateMemoryForConversation( this.name, id, content, version, ); }, }), forget: action({ description: "Delete a saved fact at the user's request using its current ID and version. Does not erase historical messages.", inputSchema: z.object({ id, version }).strict(), execute: async ({ id, version }, ctx) => { const parent = await this.memoryParent(ctx); ctx.signal.throwIfAborted(); return parent.deleteMemoryForConversation(this.name, id, version); }, }), }; }
protected browserOptions() { return { lease: () => { const generation = this.turnGeneration; const parentPromise = this.parentAgent(PersonalAgent); parentPromise.catch(() => {}); return { create: async (expiresAt: number) => (await parentPromise).createResearchBrowser(this.name, expiresAt), acquired: async () => { if (generation !== this.turnGeneration) throw new Error("Browser operation interrupted"); }, closed: async (sessionId: string) => { const parent = await parentPromise; await parent.releaseResearchBrowser(sessionId); }, }; }, progress: (id: string, stage: number) => this.reportToolProgress(id, stage), }; }
getTools() { return { shell: createShellTool((timeoutMs) => { const parent = this.parentAgent(PersonalAgent); return { reserve: async () => (await parent).reserveShellWorkspace(this.name, timeoutMs), launch: async (id, command) => (await parent).launchShellWorkspace(id, command), close: async (id) => (await parent).closeShellWorkspace(id), }; }), ...createWebTools( this.env.BROWSER, this.env.AI, (promise) => this.ctx.waitUntil(promise), this.browserOptions(), ), recall: tool({ description: "Search saved facts by relevant words when factual context is needed. Returns current fact IDs and versions for explicit edits or deletion.", inputSchema: z .object({ query: z.string().min(1).max(MAX_MEMORY_QUERY_LENGTH) }) .strict(), execute: async ({ query }, { abortSignal }) => { const parent = await this.parentAgent(PersonalAgent); abortSignal?.throwIfAborted(); return parent.searchMemories(query); }, }), }; }
protected getToolActivityDescriptors(): Record< string, ToolActivityDescriptor > { return { ...webActivityDescriptors, shell: shellActivity, createSchedule: { kind: "schedule", label: "Create scheduled task", outcome: (output) => (output as { status?: string }).status === "unavailable" ? "failed" : "succeeded", outputSummary: (output) => (output as { status?: string }).status === "unavailable" ? "Task saved; scheduling unavailable" : "Scheduled task saved", }, remember: { kind: "memory", label: "Remember fact", outputSummary: () => "Fact saved", }, updateMemory: { kind: "memory", label: "Update memory", outputSummary: () => "Fact updated", }, forget: { kind: "memory", label: "Forget fact", outputSummary: () => "Fact deleted", }, recall: { kind: "memory", label: "Recall memories", outputSummary: () => "Memory search completed", }, }; }
validateStateChange(_state: unknown, source: Connection | "server") { if (source !== "server") throw new Error("State is server managed"); }
onConnect( connection: Connection<SessionConnection>, context: ConnectionContext, ) { return connectSession(this, this.env, connection, context); }
onClose(connection: Connection<SessionConnection>) { return closeSession(this, connection); }
expireSession(payload: SessionExpiry) { expireSession(this, payload); }}