import type { TurnConfig, TurnContext } from "@cloudflare/think"; import type { Agent } from "agents"; import type { UIMessage } from "ai"; import * as Effect from "effect/Effect"; import { Secret } from "../configuration/secrets"; import type { DiagnosticEvent } from "../shared/diagnostics"; import { MAX_MEMORY_QUERY_LENGTH, memoryContext } from "../shared/memory"; import { isCredentialProvider, missingProviderKey, } from "../shared/model-providers"; import { SHELL_INSTRUCTIONS } from "../shared/shell"; import { WEB_INSTRUCTIONS } from "../shared/web"; import { agentCall, AgentFailure, agentValidation } from "./agent-io"; import { operationCall } from "./operation-result"; import { diagnosticId, DiagnosticStore } from "./diagnostics"; import { modelProviderOptions, protectModel, type createConfiguredModel, } from "./model-provider"; import { parseModelConfiguration, type ModelConfiguration, } from "./model-settings"; import type { PersonalAgent } from "./personal-agent"; import { SCHEDULE_INSTRUCTIONS } from "./schedule-action"; type TurnParent = Pick< DurableObjectStub, | "getModelSettings" | "nameConversationFromFirstMessage" | "readModelConfiguration" | "readInstructions" | "searchMemories" >; export class ConversationTurn { constructor( private readonly sql: Agent["sql"], private readonly name: () => string, private readonly parent: () => Effect.Effect, private readonly diagnostics: DiagnosticStore, private readonly waitUntil: (work: Promise) => void, ) {} modelOverride() { 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; } settings() { return Effect.gen({ self: this }, function* () { const parent = yield* this.parent(); return { ...(yield* agentCall(() => parent.getModelSettings())), override: this.modelOverride(), }; }); } prepare( ctx: TurnContext, messages: UIMessage[], createModel: ( configuration: ModelConfiguration, key?: Secret, ) => ReturnType, toolNames: Effect.Effect, ): Effect.Effect { return Effect.gen({ self: this }, function* () { const diagnosticGeneration = this.diagnostics.generation(); const parent = yield* this.parent(); // Native history is already persisted here; unlike model-facing ctx.messages // it contains neither synthetic continuation prompts nor private context. const firstUser = messages.find((message) => message.role === "user"); if (firstUser) { const text = firstUser.parts .filter((part) => part.type === "text") .map((part) => part.text) .join(" "); this.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 = 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) : yield* agentValidation(() => parseModelConfiguration(selected)); if ( !scheduled && ctx.body?.useDefaultModel !== undefined && typeof ctx.body.useDefaultModel !== "boolean" ) return yield* Effect.fail( new AgentFailure({ message: "Invalid model selection" }), ); const [{ configuration, apiKey }, instructions, memories] = yield* Effect.all( [ operationCall(() => parent.readModelConfiguration(override)), agentCall(() => parent.readInstructions()), agentCall(() => parent.searchMemories(query)), ], { concurrency: "unbounded" }, ); if (isCredentialProvider(configuration.provider) && !apiKey) return yield* Effect.fail( new AgentFailure({ message: 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( 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. 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: yield* toolNames, 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, }; }); } }