Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419import { 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<Env> { 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<TurnConfig> { 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<SessionConnection>, context: ConnectionContext, ) { return Effect.runPromise( connectSession(this, this.env, connection, context), ); }
onClose(connection: Connection<SessionConnection>) { return Effect.runPromise(closeSession(this, connection)); }
expireSession(payload: SessionExpiry) { expireSession(this, payload); }}