import { mcpActions } from "./mcp-actions"; import { skillActions } from "./skill-actions"; import { Workspace } from "@cloudflare/shell"; import { ConversationAttachments } from "./attachments"; import type { PrepareStepContext, StepContext, ThinkSubmissionInspection, } from "@cloudflare/think"; import { type ThinkModel, type TurnConfig, type TurnContext, } from "@cloudflare/think"; import { callable, type Connection, type ConnectionContext } from "agents"; import { tool } from "ai"; import * as Effect from "effect/Effect"; import { z } from "zod"; import { Secret } from "../configuration/secrets"; import { MAX_MEMORY_QUERY_LENGTH } from "../shared/memory"; import type { TaskRun } from "../shared/tasks"; import { requestedWebReader } from "../shared/web"; import { requestedAttachmentReader } from "../shared/attachments"; import { agentCall, runAgent } from "./agent-io"; import { conversationActions } from "./conversation-actions"; import { ConversationTasks } from "./conversation-tasks"; import { ConversationTurn } from "./conversation-turn"; import { ConversationSkills } from "./conversation-skills"; import { DiagnosticStore, modelDiagnostic, toolDiagnostic, type ToolDiagnostic, } from "./diagnostics"; import { createConfiguredModel } from "./model-provider"; import { DEFAULT_MODEL, type ModelConfiguration } from "./model-settings"; import { PersonalAgent, type Env } from "./personal-agent"; import { createShellTool, shellActivity } from "./shell-tool"; import { closeSession, connectSession, expireSession, type SessionConnection, type SessionExpiry, } from "./socket-session"; import type { ToolActivityDescriptor } from "./tool-activity"; import { CapabilityApprovalThink } from "./capability-approvals"; import { createWebTools, webActivityDescriptors } from "./web-tools"; import { createResearchTool } from "./research-tool"; // 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 CapabilityApprovalThink { override workspace = new Workspace({ sql: this.ctx.storage.sql, r2: this.env.ATTACHMENTS, name: () => this.name, }); private readonly attachments = new ConversationAttachments( this.sql.bind(this), this.workspace, this.env.AI, ); uploadAttachment(name: string, bytes: Uint8Array) { return this.attachments.upload(name, bytes); } downloadAttachment(id: string) { return this.attachments.download(id); } removeAttachment(id: string) { if ( this.messages.some((m) => ( m.metadata as { attachments?: { id: string }[] } | undefined )?.attachments?.some((f) => f.id === id), ) ) throw new Error("This file belongs to a saved message."); return this.attachments.remove(id); } protected readonly diagnostics = new DiagnosticStore(this.sql.bind(this)); observability = this.diagnostics.receiver(); private readonly submissions: ConversationTasks = new ConversationTasks( () => this.name, () => agentCall(() => this.parentAgent(PersonalAgent)), this, ); private readonly turn = new ConversationTurn( this.sql.bind(this), () => this.name, () => agentCall(() => this.parentAgent(PersonalAgent)), this.diagnostics, (work) => this.ctx.waitUntil(work), ); private readonly skills = new ConversationSkills(() => agentCall(() => this.parentAgent(PersonalAgent)), ); getSkills() { return [this.skills.source]; } 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 requiredReader: "read_url" | "browser_read" | "read_attachment" | undefined; @callable() getConversationModelSettings() { return runAgent(this.turn.settings()); } submitTaskRun(run: TaskRun) { return runAgent(this.submissions.submit(run)); } inspectTaskRun(id: string) { return this.inspectSubmission(id); } cancelTaskRun(id: string) { return this.cancelSubmission(id, "Task changed or was deleted"); } protected onSubmissionStatus(submission: ThinkSubmissionInspection) { const inherited = () => Promise.resolve(super.onSubmissionStatus(submission)); return runAgent( agentCall(inherited).pipe( Effect.andThen(this.submissions.observe(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, ); } beforeTurn(ctx: TurnContext): Promise { return runAgent( this.turn .prepare( ctx, this.messages, (configuration, key) => this.createModel(configuration, key), this.applicationToolNames(), (configuration) => this.attachments.prepare(this.messages, configuration), ) .pipe( Effect.map((config) => this.skills.configure( config, this.session.getContextBlock("think_skills")?.content, ), ), Effect.tap(() => Effect.sync(() => { this.requiredReader = requestedAttachmentReader( this.messages, ctx.continuation, ) ? "read_attachment" : requestedWebReader(ctx.messages, ctx.continuation); }), ), ), ); } beforeStep(ctx: PrepareStepContext) { if (ctx.stepNumber !== 0 || !this.requiredReader) return; return { toolChoice: { type: "tool" as const, toolName: this.requiredReader, }, }; } private turnGeneration = 0; protected resetTurnState() { this.turnGeneration++; this.diagnostics.clear(); super.resetTurnState(); } async getActions() { const parent = await this.parentAgent(PersonalAgent); const extensions = await mcpActions(parent, this.name, (options) => this.capabilityAction(options), ); const scripts = await skillActions(parent, this.name, (options) => this.capabilityAction(options), ); return { ...extensions, ...scripts, ...conversationActions( () => this.name, () => this.turnGeneration, () => agentCall(() => this.parentAgent(PersonalAgent)), ), }; } protected browserOptions() { return { lease: () => { const generation = this.turnGeneration; const parentPromise = this.parentAgent(PersonalAgent); parentPromise.catch(() => {}); return { create: (expiresAt: number) => runAgent( agentCall(() => parentPromise).pipe( Effect.flatMap((parent) => agentCall(() => parent.createResearchBrowser(this.name, expiresAt), ), ), ), ), acquired: async () => { if (generation !== this.turnGeneration) throw new Error("Browser operation interrupted"); }, closed: (sessionId: string) => runAgent( agentCall(() => parentPromise).pipe( Effect.flatMap((parent) => agentCall(() => parent.releaseResearchBrowser(sessionId)), ), ), ), }; }, progress: (id: string, stage: number) => this.reportToolProgress(id, stage), }; } getTools() { const web = createWebTools( this.env.BROWSER, this.env.AI, (promise) => this.ctx.waitUntil(promise), this.browserOptions(), ); return { research: createResearchTool(web, this.env.LOADER, (call, run) => this.observeResearchCall(call, run), ), shell: createShellTool((timeoutMs) => { const parent = this.parentAgent(PersonalAgent); return { reserve: () => runAgent( agentCall(() => parent).pipe( Effect.flatMap((client) => agentCall(() => client.reserveShellWorkspace(this.name, timeoutMs), ), ), ), ), launch: (id, command) => runAgent( agentCall(() => parent).pipe( Effect.flatMap((client) => agentCall(() => client.launchShellWorkspace(id, command)), ), ), ), close: (id) => runAgent( agentCall(() => parent).pipe( Effect.flatMap((client) => agentCall(() => client.closeShellWorkspace(id)), ), ), ), }; }), ...web, 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: ({ query }, { abortSignal }) => runAgent( Effect.gen({ self: this }, function* () { const parent = yield* agentCall(() => this.parentAgent(PersonalAgent), ); abortSignal?.throwIfAborted(); return yield* agentCall(() => parent.searchMemories(query)); }), abortSignal, ), }), }; } protected getToolActivityDescriptors(): Record< string, ToolActivityDescriptor > { return { ...super.getToolActivityDescriptors(), activate_skill: this.skills.activity, read_skill_resource: this.skills.resources.activity, ...webActivityDescriptors, research: { kind: "web", label: "Research webpages", outputSummary: () => "Research batch completed", }, 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, context: ConnectionContext, ) { return Effect.runPromise( connectSession(this, this.env, connection, context), ); } onClose(connection: Connection) { return Effect.runPromise(closeSession(this, connection)); } expireSession(payload: SessionExpiry) { expireSession(this, payload); } }