Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342import 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<Env> { 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<TurnConfig> { 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<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); }}