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 { agentCall, runAgent } from "./agent-io"; import { conversationActions } from "./conversation-actions"; import { ConversationTasks } from "./conversation-tasks"; import { ConversationTurn } from "./conversation-turn"; 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 { ActivityThink, type ToolActivityDescriptor } from "./tool-activity"; import { createWebTools, webActivityDescriptors } from "./web-tools"; // 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 { 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), ); 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; @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(), ) .pipe( Effect.tap(() => Effect.sync(() => { this.requiredWebReader = requestedWebReader( ctx.messages, ctx.continuation, ); }), ), ), ); } 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(); } getActions() { return 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() { return { 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)), ), ), ), }; }), ...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: ({ 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 { ...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, 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); } }