Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205import 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<PersonalAgent>, | "getModelSettings" | "nameConversationFromFirstMessage" | "readModelConfiguration" | "readInstructions" | "searchMemories">;export class ConversationTurn { constructor( private readonly sql: Agent["sql"], private readonly name: () => string, private readonly parent: () => Effect.Effect<TurnParent, AgentFailure>, private readonly diagnostics: DiagnosticStore, private readonly waitUntil: (work: Promise<unknown>) => 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<typeof createConfiguredModel>, toolNames: Effect.Effect<string[], AgentFailure>, ): Effect.Effect<TurnConfig, AgentFailure> { 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, }; }); }}