From 985ee7c8afed625daeb84a272bc7b73c3041e2ef Mon Sep 17 00:00:00 2001 From: Nathan Beddoe Date: Mon, 7 Sep 2026 15:10:46 +0200 Subject: [PATCH] Add conversation-scoped model and effort controls to the composer Amp-Thread-ID: https://ampcode.com/threads/T-01a078ff-7bea-7528-a6fa-d6fc3b61c60f Co-authored-by: Amp --- README.md | 15 ++ docs/bug-lessons.md | 18 ++ shared/model-providers.ts | 52 +++++ src/components/model-effort-select.tsx | 39 ++++ src/routes/Conversation.tsx | 111 +++++++++- src/routes/ModelSettings.tsx | 23 +- src/runtime/conversation-session.ts | 74 ++++++- src/styles.css | 25 ++- tests/chat-ui.test.mjs | 147 +++++++++++++ tests/execution.test.mjs | 102 ++++++++- tests/fixtures/think-worker.ts | 20 +- tests/model-provider.test.mjs | 278 ++++++++++++++++++++++++- tests/think.test.mjs | 65 +++++- worker/conversation.ts | 78 ++++++- worker/model-provider.ts | 45 +++- worker/model-settings.ts | 18 +- worker/personal-agent.ts | 4 +- 17 files changed, 1074 insertions(+), 40 deletions(-) create mode 100644 src/components/model-effort-select.tsx diff --git a/README.md b/README.md index cf8deb5..07ac28d 100644 --- a/README.md +++ b/README.md @@ -33,6 +33,21 @@ The published sidebar currently requires browser APIs, so the shell uses TanStack Router's `ClientOnly` boundary with server-rendered navigation and page content until hydration. +## Model selection + +The chat composer uses Octane Kumo model and effort selectors. Choices are saved +with the next message or retry, apply only to that conversation, and cannot change +an in-flight response. Changing models resets effort to the provider default. +Only verified model-specific effort levels are offered; there are no token-budget +inputs or synthetic effort-to-budget mappings. + +Settings retains provider credentials and the installation's default model/effort. +Chats using **Installation default** refresh it before sending. Scheduled tasks +snapshot that default when submitted, independently of the chat's override. +Provider default omits effort options; explicit effort uses a bounded 16,384-token +total output allowance (including reasoning), rather than the usual 4,096. Higher +effort can increase latency and cost. + ## Validation ```sh diff --git a/docs/bug-lessons.md b/docs/bug-lessons.md index 3847b34..3747822 100644 --- a/docs/bug-lessons.md +++ b/docs/bug-lessons.md @@ -558,3 +558,21 @@ A ready installation is not evidence that an unsubmitted upgrade succeeded. The - **Resolution:** Run the opt-in paid smoke separately from `tests/fixtures/config.mjs`, using `unstable_startWorker` with a remote AI binding and without `dev.remote: false`. - **Regression signal:** `FLAREBOT_WEB_LIVE_SMOKE=1 node --test tests/web-search-live.test.mjs` changed from an immediate failure to three real provider sources in about 12 seconds, without changing the search request. - **Prevention rule:** Verify that a live inference harness establishes remote bindings before diagnosing model access or request shape. Keep paid smoke authentication separate from credentials-free fixtures. + +## 2026-09-07 — Reasoning signatures are protocol state, not diagnostics + +- **Affected area:** `worker/model-provider.ts`, Anthropic-compatible reasoning and tool follow-ups. +- **Symptom signature:** Enabling adaptive reasoning while stripping all stream `providerMetadata` removes the signature carried by an empty reasoning delta. The next SDK-generated tool-follow-up request loses the signed thinking block; redacted thinking is likewise lost. +- **Root cause:** The privacy boundary treated every provider metadata field as optional diagnostics, but the Anthropic adapter uses `anthropic.signature` and `anthropic.redactedData` to reconstruct required conversation blocks. +- **Resolution:** Preserve only those string fields on reasoning chunks and generated reasoning parts. Continue removing other provider metadata, raw events, request/response envelopes, and unsafe errors. +- **Regression signal:** `pnpm test:providers` exercises a native Anthropic SDK two-step tool round-trip with empty signed thinking and redacted thinking, and separately checks that unrelated metadata is stripped. +- **Prevention rule:** When enabling reasoning, verify the actual follow-up request—not just the initial effort payload. Preserve required opaque protocol state and empty signature-bearing deltas without treating them as user-visible diagnostics. + +## 2026-09-07 — Queued submission metadata is not active turn metadata + +- **Affected area:** `worker/conversation.ts`, Think 0.17 scheduled model selection. +- **Symptom signature:** A scheduled task in a chat with High effort used High even after the installation default became Low; the submission inspection correctly contained Low. +- **Root cause:** `submitMessages(..., { metadata })` stores queue-ledger metadata but does not stamp the user message's reserved `turnMetadata`. `activeTurnMetadata` therefore did not identify the scheduled turn, and `beforeTurn` fell back to the last interactive request body. +- **Resolution:** Also carry the scheduled model snapshot in application-owned metadata on the submitted user message and resolve that snapshot ahead of the interactive body. Keep queue metadata for submission inspection and task lifecycle handling. +- **Regression signal:** `FLAREBOT_EXECUTION_CASE='scheduled turns ignore' node --test --test-reporter=tap tests/execution.test.mjs` failed with High versus Low, then passed with the scheduled tool follow-up using Low, the chat override remaining High, and a subsequent interactive message using High. +- **Prevention rule:** Test each native submission path's actual hook-visible data. Do not assume queue inspection metadata is automatically available through per-turn APIs. diff --git a/shared/model-providers.ts b/shared/model-providers.ts index 52ac511..a84198a 100644 --- a/shared/model-providers.ts +++ b/shared/model-providers.ts @@ -81,9 +81,61 @@ export type ModelConfiguration = { [P in ModelProvider]: { provider: P; model: keyof (typeof MODEL_PROVIDERS)[P]["models"]; + effort?: ModelEffort; }; }[ModelProvider]; +export const EFFORT_LABELS = { + low: "Low", + medium: "Medium", + high: "High", + xhigh: "Extra high", + max: "Maximum", +} as const; +export type ModelEffort = keyof typeof EFFORT_LABELS; + +// Verified 2026-09-07. Native effort only: never synthesize token budgets. +// https://openrouter.ai/api/v1/models (reasoning.supported_efforts) +// https://platform.claude.com/docs/en/build-with-claude/effort +// https://tinker-docs.thinkingmachines.ai/tinker/compatible-apis/anthropic/ +// Empty means no verified effort control on this exact provider route. +export const MODEL_EFFORTS = { + "workers-ai": { + "@cf/meta/llama-3.3-70b-instruct-fp8-fast": [], + "@cf/meta/llama-4-scout-17b-16e-instruct": [], + "thinkingmachines/inkling-256k": ["low", "medium", "high", "xhigh", "max"], + }, + anthropic: { + "claude-sonnet-5": ["low", "medium", "high", "xhigh", "max"], + "claude-haiku-4-5-20251001": [], + }, + openrouter: { + "openai/gpt-5.6-luna": ["low", "medium", "high", "xhigh", "max"], + "tencent/hy4-preview": ["low", "high"], + "z-ai/glm-5.3-flash": ["low", "high", "max"], + "deepseek/deepseek-v4-flash-0731": ["low", "high", "max"], + "minimax/minimax-m3:free": [], + "deepseek/deepseek-v4-flash": ["high", "xhigh"], + "tencent/hy3": ["low", "high"], + "nvidia/nemotron-3-ultra-550b-a55b:free": ["medium", "high"], + "z-ai/glm-5.3": ["low", "high", "max"], + "xiaomi/mimo-v2.5": [], + }, +} as const satisfies { + [P in ModelProvider]: Record< + keyof (typeof MODEL_PROVIDERS)[P]["models"], + readonly ModelEffort[] + >; +}; + +export function modelEfforts( + configuration: ModelConfiguration, +): readonly ModelEffort[] { + const models: Record = + MODEL_EFFORTS[configuration.provider]; + return models[configuration.model]; +} + export const MODEL_CATALOG = Object.fromEntries( Object.entries(MODEL_PROVIDERS).map(([id, provider]) => [ id, diff --git a/src/components/model-effort-select.tsx b/src/components/model-effort-select.tsx new file mode 100644 index 0000000..57ea057 --- /dev/null +++ b/src/components/model-effort-select.tsx @@ -0,0 +1,39 @@ +import { Select } from "octane-kumo/components/select"; +import { + EFFORT_LABELS, + modelEfforts, + type ModelConfiguration, +} from "../../shared/model-providers"; + +export function ModelEffortSelect({ + configuration, + disabled, + onChange, +}: { + configuration: ModelConfiguration; + disabled: boolean; + onChange: (configuration: ModelConfiguration) => void; +}) { + const efforts = modelEfforts(configuration); + if (!efforts.length) return null; + return ( + + value === "default" + ? "Installation default" + : modelChoices + .find((choice) => choice.value === value) + ?.configuration.model.split("/") + .at(-1) + } + onValueChange={(value) => { + if (value === "default") session.selectModel(null); + else { + const choice = modelChoices.find( + (choice) => choice.value === value, + ); + if (choice && providerReady(choice.configuration.provider)) + session.selectModel(choice.configuration); + } + }} + > + + Installation default + + {(Object.keys(MODEL_PROVIDERS) as ModelProvider[]).map( + (provider) => ( + + + {MODEL_PROVIDERS[provider].name} + {providerReady(provider) ? "" : " · Set up in Settings"} + + {modelChoices + .filter( + (choice) => + choice.configuration.provider === provider, + ) + .map((choice) => ( + + {choice.configuration.model} + + ))} + + ), + )} + + {configuration && ( + + )} + {busy ? ( )} +

+ {modelReady + ? "This conversation only · Saved with your next message or retry." + : "Set up the selected provider before sending."}{" "} + Provider settings + {configuration?.effort + ? " · Higher effort may take longer and cost more." + : ""} +

) : null} diff --git a/src/routes/ModelSettings.tsx b/src/routes/ModelSettings.tsx index 1254f6e..e3fb690 100644 --- a/src/routes/ModelSettings.tsx +++ b/src/routes/ModelSettings.tsx @@ -2,6 +2,7 @@ import { useEffect, useRef, useState } from "octane"; import { Button } from "octane-kumo/components/button"; import { Input } from "octane-kumo/components/input"; import { Select } from "octane-kumo/components/select"; +import { ModelEffortSelect } from "../components/model-effort-select"; import type { MODEL_CATALOG, ModelConfiguration, @@ -144,7 +145,8 @@ export function ModelSettings({ saved && draft && (saved.configuration.provider !== draft.provider || - saved.configuration.model !== draft.model); + saved.configuration.model !== draft.model || + saved.configuration.effort !== draft.effort); const openRouterEnabled = saved?.providers?.openrouter === true; const draftProviderEnabled = draft?.provider !== "openrouter" || openRouterEnabled; @@ -154,8 +156,9 @@ export function ModelSettings({

Model and provider

- Choose how Flarebot responds. Changes apply from the next turn in - every conversation and scheduled task. + Choose the installation default for new chats, conversations using the + default, and scheduled tasks. Use the chat composer to override the + model and effort for one conversation.

{loading && connected &&

Loading model settings…

} @@ -205,7 +208,19 @@ export function ModelSettings({ onValueChange={(value) => { if (!value) return; dirty.current = true; - setDraft({ ...draft, model: value } as ModelConfiguration); + setDraft({ + provider: draft.provider, + model: value, + } as ModelConfiguration); + setNotice(""); + }} + /> + { + dirty.current = true; + setDraft(value); setNotice(""); }} /> diff --git a/src/runtime/conversation-session.ts b/src/runtime/conversation-session.ts index f758502..b0a2d75 100644 --- a/src/runtime/conversation-session.ts +++ b/src/runtime/conversation-session.ts @@ -18,6 +18,8 @@ import type { ToolActivityPage, ToolActivityState, } from "../../shared/tool-activity"; +import type { ModelConfiguration } from "../../shared/model-providers"; +import type { ConversationModelSettings } from "../../worker/model-settings"; type Connection = "loading" | "connected" | "offline" | "signed-out" | "missing"; @@ -35,6 +37,8 @@ export interface ConversationView { activities: ReadonlyMap; activityCursor: number | null | undefined; activitiesLoading: boolean; + modelSettings?: ConversationModelSettings; + modelSelection?: ModelConfiguration | null; } class NativeChat extends AbstractChat { constructor(state: ChatState, options: ChatInit) { @@ -84,7 +88,7 @@ export class ConversationSession { private reconnectQueued = false; private idleConfirmed = false; private probeTimer?: ReturnType; - private operation?: { generation: number }; + private operation?: { generation: number; cancelled: boolean }; private streamEpoch = 0; private pendingInput?: { id: string; text: string }; constructor(readonly id: string) {} @@ -353,6 +357,17 @@ export class ConversationSession { .finally(() => signal.removeEventListener("abort", abort)); }); if (!this.current(generation)) return; + const modelSettings = await client.call( + "getConversationModelSettings", + ); + if (!this.current(generation)) return; + this.publish({ + modelSettings, + modelSelection: + this.view.modelSelection === undefined + ? modelSettings.override + : this.view.modelSelection, + }); this.publish({ connection: "connected", notice: undefined }); await this.history(generation); if (!this.current(generation)) return; @@ -716,6 +731,24 @@ export class ConversationSession { if (this.current(generation)) this.publish({ activitiesLoading: false }); } }; + selectModel = (configuration: ModelConfiguration | null) => { + this.publish({ modelSelection: configuration }); + }; + private async modelBody(generation: number) { + const selection = this.view.modelSelection; + if (selection) + return { modelConfiguration: selection, useDefaultModel: false }; + // Refresh inherited defaults just before submission, not only when opening + // the chat. The resulting snapshot travels with the native request. + const settings = await this.client!.call( + "getConversationModelSettings", + ); + if (this.current(generation)) this.publish({ modelSettings: settings }); + return { + modelConfiguration: settings.configuration, + useDefaultModel: true, + }; + } send = async (text: string) => { if ( !this.chat || @@ -726,7 +759,7 @@ export class ConversationSession { ) return false; const generation = this.generation; - const operation = { generation }; + const operation = { generation, cancelled: false }; this.operation = operation; const id = crypto.randomUUID(); this.idleConfirmed = false; @@ -735,13 +768,25 @@ export class ConversationSession { notice: undefined, error: undefined, unsentText: undefined, + status: "submitted", }); try { - await this.chat.sendMessage({ - id, - role: "user", - parts: [{ type: "text", text }], - }); + const body = await this.modelBody(generation); + if (!this.current(generation) || operation.cancelled) return false; + await this.chat.sendMessage( + { + id, + role: "user", + parts: [{ type: "text", text }], + }, + { body }, + ); + } catch { + if (this.current(generation)) + this.publish({ + error: new Error("Could not submit the message"), + unsentText: text, + }); } finally { if (this.operation === operation) this.operation = undefined; if (this.current(generation)) { @@ -762,6 +807,7 @@ export class ConversationSession { ) return; this.publish({ stopping: true }); + if (this.operation) this.operation.cancelled = true; const cancelled = this.transport.cancelActiveServerTurn(); await this.chat.stop(); if (!cancelled) { @@ -790,11 +836,19 @@ export class ConversationSession { const retained = this.view.messages; this.idleConfirmed = false; const generation = this.generation; - const operation = { generation }; + const operation = { generation, cancelled: false }; this.operation = operation; - this.publish({ error: undefined, notice: undefined }); + this.publish({ error: undefined, notice: undefined, status: "submitted" }); try { - await this.chat.regenerate({ messageId: last.id }); + const body = await this.modelBody(generation); + if (!this.current(generation) || operation.cancelled) return; + await this.chat.regenerate({ + messageId: last.id, + body, + }); + } catch { + if (this.current(generation)) + this.publish({ error: new Error("Could not retry the response") }); } finally { if (this.operation === operation) this.operation = undefined; if (this.current(generation)) { diff --git a/src/styles.css b/src/styles.css index 3ffe89d..bee7369 100644 --- a/src/styles.css +++ b/src/styles.css @@ -761,10 +761,33 @@ body, .chat-composer-actions { display: flex; justify-content: space-between; - align-items: center; + align-items: flex-end; + flex-wrap: wrap; gap: 0.75rem; margin-top: 0.5rem; } +.chat-model-controls { + display: flex; + flex: 1; + flex-wrap: wrap; + align-items: flex-end; + gap: 0.5rem; + min-width: 0; +} +.chat-model-controls > * { + min-width: 0; + max-width: 100%; +} +.chat-model-controls > :first-child { + width: 18rem; +} +.chat-model-controls > :nth-child(2) { + width: 10rem; +} +.chat-model-hint { + font-size: 12px; + margin: 0.5rem 0 0; +} .chat-jump { height: 0; position: relative; diff --git a/tests/chat-ui.test.mjs b/tests/chat-ui.test.mjs index f5c434b..4fe09d8 100644 --- a/tests/chat-ui.test.mjs +++ b/tests/chat-ui.test.mjs @@ -260,6 +260,153 @@ test( .click(); }; const first = await make("Private chat fixture"); + await run( + "composer model and effort are scoped, native Kumo controls", + async () => { + const conversation = await make("Model selection"); + const other = await make("Independent model selection"); + const defaults = await call("getModelSettings"); + const choose = async (label, option) => { + await page + .getByRole("combobox", { name: label, exact: true }) + .click(); + await page + .getByRole("option", { name: option, exact: true }) + .click(); + }; + await open(conversation.id); + assert.equal( + await page + .getByRole("combobox", { name: "Effort", exact: true }) + .count(), + 0, + ); + await page + .getByRole("combobox", { name: "Model", exact: true }) + .click(); + assert.equal( + await page + .getByRole("option", { name: "claude-sonnet-5", exact: true }) + .getAttribute("aria-disabled"), + "true", + ); + await page.keyboard.press("Escape"); + await choose("Model", "thinkingmachines/inkling-256k"); + await choose("Effort", "High"); + await mkdir(screenshots, { recursive: true }); + await page.screenshot({ + path: join(screenshots, "composer-effort-desktop.png"), + }); + await send("stream"); + await page.getByRole("button", { name: "Stop response" }).waitFor(); + assert.ok( + await page + .getByRole("combobox", { name: "Model", exact: true }) + .isDisabled(), + ); + assert.ok( + await page + .getByRole("combobox", { name: "Effort", exact: true }) + .isDisabled(), + ); + await ready(); + const client = await leaf(conversation.id); + assert.deepEqual( + (await client.call("getConversationModelSettings")).override, + { + provider: "workers-ai", + model: "thinkingmachines/inkling-256k", + effort: "high", + }, + ); + assert.deepEqual( + (await call("getModelSettings")).configuration, + defaults.configuration, + ); + await page.reload(); + await ready(); + assert.match( + await page + .getByRole("combobox", { name: "Effort", exact: true }) + .innerText(), + /High/, + ); + await page.setViewportSize({ width: 390, height: 844 }); + await page.screenshot({ + path: join(screenshots, "composer-effort-mobile.png"), + }); + assert.ok( + await page.evaluate( + () => document.documentElement.scrollWidth <= innerWidth, + ), + ); + await page + .getByRole("combobox", { name: "Effort", exact: true }) + .click(); + assert.deepEqual(await page.getByRole("option").allTextContents(), [ + "Provider default", + "Low", + "Medium", + "High", + "Extra high", + "Maximum", + ]); + await page.screenshot({ + path: join(screenshots, "composer-effort-menu.png"), + }); + await page.keyboard.press("Escape"); + await page.setViewportSize({ width: 1280, height: 900 }); + await choose("Model", "@cf/meta/llama-4-scout-17b-16e-instruct"); + assert.equal( + await page + .getByRole("combobox", { name: "Effort", exact: true }) + .count(), + 0, + ); + await choose("Model", "thinkingmachines/inkling-256k"); + assert.match( + await page + .getByRole("combobox", { name: "Effort", exact: true }) + .innerText(), + /Provider default/, + ); + await choose("Model", "Installation default"); + await send("fast"); + await ready(); + assert.equal( + (await client.call("getConversationModelSettings")).override, + null, + ); + await open(other.id); + assert.match( + await page + .getByRole("combobox", { name: "Model", exact: true }) + .innerText(), + /Installation default/, + ); + // Defaults changed in another tab are refreshed on submission, not pinned + // forever to the configuration that was loaded when opening this chat. + await call("updateModelSettings", { + provider: "workers-ai", + model: "thinkingmachines/inkling-256k", + effort: "low", + }); + await send("fast"); + await ready(); + assert.match( + await page + .getByRole("combobox", { name: "Effort", exact: true }) + .innerText(), + /Low/, + ); + assert.equal( + (await (await leaf(other.id)).call("getConversationModelSettings")) + .override, + null, + ); + await call("updateModelSettings", defaults.configuration); + }, + ); await run( "first-message title updates the sidebar without reload", async () => { diff --git a/tests/execution.test.mjs b/tests/execution.test.mjs index cde4c80..bd1772a 100644 --- a/tests/execution.test.mjs +++ b/tests/execution.test.mjs @@ -127,13 +127,21 @@ test( assert.fail(`${message}: ${JSON.stringify(value)}`); } const scenario = (name, body) => - t.test(name, async () => { - try { - await body(); - } finally { - if (!clients.includes(owner)) await connect(); - } - }); + t.test( + name, + { + skip: process.env.FLAREBOT_EXECUTION_CASE + ? !name.includes(process.env.FLAREBOT_EXECUTION_CASE) + : false, + }, + async () => { + try { + await body(); + } finally { + if (!clients.includes(owner)) await connect(); + } + }, + ); try { worker = await start(); await connect(); @@ -633,6 +641,86 @@ test( }, ); + await scenario( + "scheduled turns ignore the chat override and retain their default snapshot", + async () => { + const defaults = await call("getModelSettings"); + const conversation = await call( + "createConversation", + "Scheduled model isolation", + ); + const client = await connect(conversation.id); + const selection = { + provider: "workers-ai", + model: "thinkingmachines/inkling-256k", + effort: "high", + }; + const transport = new WebSocketChatTransport({ agent: client }); + const stream = await transport.sendMessages({ + chatId: conversation.id, + trigger: "submit-message", + messages: [ + { + id: crypto.randomUUID(), + role: "user", + parts: [{ type: "text", text: "configuration" }], + }, + ], + body: { modelConfiguration: selection }, + abortSignal: new AbortController().signal, + }); + for await (const _ of stream) { + /* Persist the real chat override. */ + } + const taskDefault = { ...selection, effort: "low" }; + await call("updateModelSettings", taskDefault); + const task = await create(conversation.id, { + instructions: "tool", + schedule: { kind: "once", at: instant(180_000) }, + }); + const run = await call("runTaskNow", task.id, crypto.randomUUID()); + await complete(task.id); + const calls = await client.call("fixtureModelCalls"); + assert.deepEqual(calls[0].configuration, selection); + assert.ok(calls.length >= 3); + for (const call of calls.slice(1)) { + assert.deepEqual(call.configuration, taskDefault); + assert.deepEqual(call.options.anthropic, { effort: "low" }); + } + assert.deepEqual( + (await client.call("getConversationModelSettings")).override, + selection, + ); + const submission = (await snapshot(conversation.id)).submissions.find( + (submission) => submission.submissionId === run.submissionId, + ); + assert.deepEqual(submission.metadata.modelConfiguration, taskDefault); + const followUp = await transport.sendMessages({ + chatId: conversation.id, + trigger: "submit-message", + messages: [ + ...(await snapshot(conversation.id)).messages, + { + id: crypto.randomUUID(), + role: "user", + parts: [{ type: "text", text: "configuration" }], + }, + ], + body: { modelConfiguration: selection }, + abortSignal: new AbortController().signal, + }); + for await (const _ of followUp) { + /* Scheduled metadata must not leak to chat. */ + } + assert.deepEqual( + (await client.call("fixtureModelCalls")).at(-1).configuration, + selection, + ); + await call("updateModelSettings", defaults.configuration); + await call("deleteTask", task.id, task.version); + }, + ); + await scenario( "accepted scheduled turn reaches native recovery or safe terminal error after restart", async () => { diff --git a/tests/fixtures/think-worker.ts b/tests/fixtures/think-worker.ts index 60def2f..9e428c3 100644 --- a/tests/fixtures/think-worker.ts +++ b/tests/fixtures/think-worker.ts @@ -5,7 +5,7 @@ import { type Env, } from "../../worker/personal-agent"; import { Conversation as RuntimeConversation } from "../../worker/conversation"; -import { getAgentByName } from "agents"; +import { callable, getAgentByName } from "agents"; import { MockLanguageModelV3 } from "ai/test"; import { tool } from "ai"; import { action } from "@cloudflare/think"; @@ -155,6 +155,19 @@ type ModelChunk = // Only this test entry replaces inference. No environment flag or development // authentication/model bypass is present in the customer runtime artifact. export class Conversation extends RuntimeConversation { + @callable() + fixtureModelCalls() { + this + .sql`CREATE TABLE IF NOT EXISTS fixture_model_calls (configuration TEXT, options TEXT)`; + return this.sql<{ + configuration: string; + options: string; + }>`SELECT * FROM fixture_model_calls`.map((row) => ({ + configuration: JSON.parse(row.configuration), + options: JSON.parse(row.options), + })); + } + fixtureDiagnostics(fail: boolean) { this.diagnosticSnapshot(); if (fail) { @@ -165,7 +178,10 @@ export class Conversation extends RuntimeConversation { } protected createModel(configuration: ModelConfiguration, key?: Secret) { return new MockLanguageModelV3({ - doStream: async ({ prompt, abortSignal, tools }) => { + doStream: async ({ prompt, abortSignal, tools, providerOptions }) => { + this.fixtureModelCalls(); + this + .sql`INSERT INTO fixture_model_calls VALUES (${JSON.stringify(configuration)}, ${JSON.stringify(providerOptions ?? {})})`; // Native recovery appends a synthetic continuation instruction. Select // the last fixture scenario from the real user history, not that prompt. const user = prompt.findLast( diff --git a/tests/model-provider.test.mjs b/tests/model-provider.test.mjs index 8649e13..7cdccf7 100644 --- a/tests/model-provider.test.mjs +++ b/tests/model-provider.test.mjs @@ -1,10 +1,16 @@ import assert from "node:assert/strict"; import { test } from "node:test"; import { MockLanguageModelV3 } from "ai/test"; +import { streamText, stepCountIs, tool, jsonSchema } from "ai"; import { Secret } from "../configuration/secrets.ts"; -import { MODEL_PROVIDERS } from "../shared/model-providers.ts"; +import { + MODEL_PROVIDERS, + MODEL_EFFORTS, + EFFORT_LABELS, +} from "../shared/model-providers.ts"; import { createConfiguredModel, + modelProviderOptions, protectModel, } from "../worker/model-provider.ts"; import { @@ -29,6 +35,276 @@ const openrouter = { model: "openai/gpt-5.6-luna", }; +test("effort validation is exhaustive per model and never accepts budgets or toggles", () => { + for (const [provider, definition] of Object.entries(MODEL_PROVIDERS)) { + assert.deepEqual( + Object.keys(MODEL_EFFORTS[provider]).sort(), + Object.keys(definition.models).sort(), + ); + for (const model of Object.keys(definition.models)) { + const configuration = { provider, model }; + assert.deepEqual(parseModelConfiguration(configuration), configuration); + assert.equal(modelProviderOptions(configuration), undefined); + for (const effort of [ + ...Object.keys(EFFORT_LABELS), + "default", + "off", + "on", + "none", + "adaptive", + 2048, + null, + undefined, + ]) { + if (MODEL_EFFORTS[provider][model].includes(effort)) + assert.deepEqual( + parseModelConfiguration({ ...configuration, effort }), + { ...configuration, effort }, + ); + else + assert.throws( + () => parseModelConfiguration({ ...configuration, effort }), + /supported effort/, + ); + } + for (const extra of [ + { budget: 2048 }, + { thinking: true }, + { providerOptions: {} }, + ]) + assert.throws( + () => parseModelConfiguration({ ...configuration, ...extra }), + /supported provider/, + ); + } + } + assert.deepEqual(MODEL_EFFORTS.openrouter["tencent/hy4-preview"], [ + "low", + "high", + ]); + assert.deepEqual(MODEL_EFFORTS.openrouter["deepseek/deepseek-v4-flash"], [ + "high", + "xhigh", + ]); + assert.deepEqual(MODEL_EFFORTS.anthropic[anthropic.model], []); +}); + +test("native SDKs serialize effort on both stream and generate without budget mappings", async (t) => { + let captured; + const response = () => + Response.json({ error: { message: "wire captured" } }, { status: 400 }); + t.mock.method(globalThis, "fetch", async (_url, init) => { + captured = JSON.parse(init.body); + return response(); + }); + const binding = { + run: async (_model, body) => { + captured = body; + return response(); + }, + gateway: () => ({ + run: async (request) => { + captured = request.query; + return response(); + }, + }), + }; + for (const [provider, models] of Object.entries(MODEL_EFFORTS)) { + for (const [model, efforts] of Object.entries(models)) { + // Workers AI's ordinary non-reasoning transport is covered separately. + if (provider === "workers-ai" && !efforts.length) continue; + for (const effort of [undefined, ...efforts]) { + const configuration = parseModelConfiguration({ + provider, + model, + ...(effort ? { effort } : {}), + }); + const languageModel = createConfiguredModel( + binding, + configuration, + new Secret(credential), + ); + for (const operation of ["doGenerate", "doStream"]) { + captured = undefined; + await assert.rejects( + languageModel[operation]({ + ...params, + providerOptions: modelProviderOptions(configuration), + }), + ); + assert.ok( + captured, + `${provider}/${model}/${operation} must reach transport`, + ); + assert.equal(captured.reasoning_effort, undefined); + if (provider === "openrouter") { + assert.deepEqual( + captured.reasoning, + effort ? { effort } : undefined, + ); + assert.equal(captured.output_config, undefined); + } else { + assert.deepEqual( + captured.output_config, + effort ? { effort } : undefined, + ); + assert.deepEqual( + captured.thinking, + effort && provider === "anthropic" + ? { type: "adaptive" } + : undefined, + ); + } + assert.ok(!JSON.stringify(captured).includes("budget_tokens")); + } + } + } + } +}); + +test("adaptive reasoning preserves signed and redacted blocks across a real SDK tool round-trip", async (t) => { + const configuration = { + provider: "anthropic", + model: "claude-sonnet-5", + effort: "high", + }; + const requests = []; + t.mock.method(globalThis, "fetch", async (_url, init) => { + const body = JSON.parse(init.body); + requests.push(body); + const first = requests.length === 1; + const blocks = first + ? [ + { type: "thinking", thinking: "", signature: "" }, + { type: "redacted_thinking", data: "opaque-redacted-block" }, + { type: "tool_use", id: "call-1", name: "echo", input: {} }, + ] + : [{ type: "text", text: "Done" }]; + const events = [ + { + type: "message_start", + message: { + id: `msg-${requests.length}`, + type: "message", + role: "assistant", + model: configuration.model, + content: [], + stop_reason: null, + stop_sequence: null, + usage: { input_tokens: 1, output_tokens: 0 }, + }, + }, + ]; + for (const [index, block] of blocks.entries()) { + events.push({ type: "content_block_start", index, content_block: block }); + if (block.type === "thinking") + events.push({ + type: "content_block_delta", + index, + delta: { + type: "signature_delta", + signature: "opaque-thinking-signature", + }, + }); + events.push({ type: "content_block_stop", index }); + } + events.push( + { + type: "message_delta", + delta: { + stop_reason: first ? "tool_use" : "end_turn", + stop_sequence: null, + }, + usage: { output_tokens: 1 }, + }, + { type: "message_stop" }, + ); + return new Response(anthropicMessagesStream(events), { + headers: { "content-type": "text/event-stream" }, + }); + }); + const result = streamText({ + model: protectModel( + createConfiguredModel({}, configuration, new Secret(credential)), + "anthropic", + ), + providerOptions: modelProviderOptions(configuration), + prompt: "Use the echo tool", + tools: { + echo: tool({ + inputSchema: jsonSchema({ type: "object", properties: {} }), + execute: async () => "ok", + }), + }, + stopWhen: stepCountIs(2), + }); + await result.consumeStream(); + assert.equal(requests.length, 2); + for (const request of requests) + assert.deepEqual(request.output_config, { effort: "high" }); + const assistant = requests[1].messages.find( + (message) => message.role === "assistant", + ); + assert.deepEqual(assistant.content.slice(0, 2), [ + { type: "thinking", thinking: "", signature: "opaque-thinking-signature" }, + { type: "redacted_thinking", data: "opaque-redacted-block" }, + ]); +}); + +test("reasoning protocol metadata is allowlisted, not arbitrary provider diagnostics", async () => { + const metadata = { + anthropic: { signature: "", redactedData: "opaque", secret: credential }, + other: { secret: credential }, + }; + const model = protectModel( + new MockLanguageModelV3({ + doGenerate: async () => ({ + content: [{ type: "reasoning", text: "", providerMetadata: metadata }], + finishReason: { unified: "stop" }, + usage: { inputTokens: {}, outputTokens: {} }, + }), + doStream: async () => ({ + stream: new ReadableStream({ + start(controller) { + controller.enqueue({ + type: "reasoning-start", + id: "r", + providerMetadata: metadata, + }); + controller.enqueue({ + type: "reasoning-delta", + id: "r", + delta: "", + providerMetadata: metadata, + }); + controller.enqueue({ + type: "reasoning-end", + id: "r", + providerMetadata: metadata, + }); + controller.close(); + }, + }), + }), + }), + "anthropic", + ); + const generated = await model.doGenerate(params); + assert.deepEqual(generated.content[0].providerMetadata, { + anthropic: { signature: "", redactedData: "opaque" }, + }); + const chunks = []; + for await (const chunk of (await model.doStream(params)).stream) { + chunks.push(chunk); + assert.deepEqual( + chunk.providerMetadata, + generated.content[0].providerMetadata, + ); + } + assert.equal(chunks.length, 3); + assert.ok(!JSON.stringify({ generated, chunks }).includes(credential)); +}); + test("OpenRouter model selection and credentials are bounded by the provider registry", () => { for (const model of Object.keys(MODEL_PROVIDERS.openrouter.models)) assert.deepEqual(parseModelConfiguration({ ...openrouter, model }), { diff --git a/tests/think.test.mjs b/tests/think.test.mjs index 2f3542f..5a859aa 100644 --- a/tests/think.test.mjs +++ b/tests/think.test.mjs @@ -350,7 +350,7 @@ test( const first = await connect(firstId); const observer = await connect(firstId); const second = await connect(secondId); - async function send(connection, id, text) { + async function send(connection, id, text, body) { const messages = [ ...(await history(id)), { @@ -363,6 +363,7 @@ test( const stream = await connection.transport.sendMessages({ chatId: id, messages, + body, trigger: "submit-message", abortSignal: new AbortController().signal, }); @@ -617,6 +618,68 @@ test( `Reply ${configuration.provider}/${configuration.model} complete`, ); }; + // The composer sends configuration in the native request, not a separate + // settings mutation. Each step sees that same model and effort. + const scopedId = (await create("Scoped model and effort")).id; + const scoped = await connect(scopedId); + const selected = { ...external, effort: "low" }; + assert.equal( + (await scoped.client.call("getConversationModelSettings")).override, + null, + ); + await ( + await send(scoped, scopedId, "tool", { modelConfiguration: selected }) + ).done; + const calls = await scoped.client.call("fixtureModelCalls"); + assert.ok(calls.length >= 2, "tool follow-up must invoke the model"); + for (const call of calls) { + assert.deepEqual(call.configuration, selected); + assert.deepEqual(call.options.anthropic, { + effort: "low", + thinking: { type: "adaptive" }, + }); + } + assert.deepEqual( + (await scoped.client.call("getConversationModelSettings")).override, + selected, + ); + assert.equal( + (await second.client.call("getConversationModelSettings")).override, + null, + ); + assert.deepEqual( + (await owner.call("getModelSettings")).configuration, + external, + ); + const invalidEffort = await send(scoped, scopedId, "configuration", { + modelConfiguration: { ...selected, effort: "off" }, + }); + await assert.rejects(invalidEffort.done, /supported effort/); + assert.deepEqual( + (await scoped.client.call("getConversationModelSettings")).override, + selected, + ); + scoped.client.close(); + const reopened = await connect(scopedId); + assert.deepEqual( + (await reopened.client.call("getConversationModelSettings")).override, + selected, + ); + await assertConfiguration(reopened, scopedId, selected); + // Reset applies an explicit snapshot to this turn while clearing the + // conversation override for future default-following turns. + await ( + await send(reopened, scopedId, "configuration", { + modelConfiguration: defaultSettings.configuration, + useDefaultModel: true, + }) + ).done; + assert.equal( + (await reopened.client.call("getConversationModelSettings")).override, + null, + ); + reopened.client.close(); + await owner.call("deleteConversation", [scopedId]); const openRouterConfiguration = { provider: "openrouter", model: catalog["openrouter"][0], diff --git a/worker/conversation.ts b/worker/conversation.ts index c431481..3dbd66f 100644 --- a/worker/conversation.ts +++ b/worker/conversation.ts @@ -33,11 +33,20 @@ import { MAX_MEMORY_QUERY_LENGTH, memoryContext, } from "../shared/memory"; -import type { Connection, ConnectionContext } from "agents"; +import { callable, type Connection, type ConnectionContext } from "agents"; import { PersonalAgent, type Env } from "./personal-agent"; import { Secret } from "../configuration/secrets"; -import { DEFAULT_MODEL, type ModelConfiguration } from "./model-settings"; -import { createConfiguredModel, protectModel } from "./model-provider"; +import { + DEFAULT_MODEL, + parseModelConfiguration, + type ModelConfiguration, + type ConversationModelSettings, +} from "./model-settings"; +import { + createConfiguredModel, + modelProviderOptions, + protectModel, +} from "./model-provider"; import type { DiagnosticEvent } from "../shared/diagnostics"; import { isCredentialProvider, @@ -91,22 +100,43 @@ export class Conversation extends ActivityThink { storeTools = false; private requiredWebReader: "read_url" | "browser_read" | undefined; + private modelOverride(): ModelConfiguration | null { + 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; + } + + @callable() + async getConversationModelSettings(): Promise { + const parent = await this.parentAgent(PersonalAgent); + return { + ...(await parent.getModelSettings()), + override: this.modelOverride(), + }; + } + async submitTaskRun(run: TaskRun) { const parent = await this.parentAgent(PersonalAgent); const task = await parent.authorizeTaskRun(run); if (!task || task.conversationId !== this.name) return null; + const { configuration } = await parent.getModelSettings(); const receipt = await this.submitMessages( [ { id: run.submissionId, role: "user", parts: [{ type: "text", text: task.instructions }], + metadata: { scheduledModelConfiguration: configuration }, }, ], { submissionId: run.submissionId, idempotencyKey: run.submissionId, - metadata: { taskRun: run }, + metadata: { taskRun: run, modelConfiguration: configuration }, }, ); if (!(await parent.authorizeTaskRun(run))) @@ -197,14 +227,47 @@ export class Conversation extends ActivityThink { .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 = this.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) + : parseModelConfiguration(selected); + if ( + !scheduled && + ctx.body?.useDefaultModel !== undefined && + typeof ctx.body.useDefaultModel !== "boolean" + ) + throw new Error("Invalid model selection"); const [{ configuration, apiKey }, instructions, memories] = await Promise.all([ - parent.readModelConfiguration(), + parent.readModelConfiguration(override), parent.readInstructions(), parent.searchMemories(query), ]); if (isCredentialProvider(configuration.provider) && !apiKey) throw new Error(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( this.createModel(configuration, apiKey ? new Secret(apiKey) : undefined), configuration.provider, @@ -253,7 +316,10 @@ export class Conversation extends ActivityThink { SCHEDULE_INSTRUCTIONS + `\nCurrent UTC time: ${new Date().toISOString()}\n`, activeTools: await this.applicationToolNames(), - maxOutputTokens: 4096, + 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, }; } diff --git a/worker/model-provider.ts b/worker/model-provider.ts index f5beee1..1a3a354 100644 --- a/worker/model-provider.ts +++ b/worker/model-provider.ts @@ -12,6 +12,34 @@ import { createGatewayModel } from "./gateway-model.ts"; const INKLING_MODEL = "thinkingmachines/inkling-256k"; +// Keep wire-specific options at the provider boundary. Think passes these to +// every step; omission preserves the provider's behavior for existing settings. +export function modelProviderOptions(configuration: ModelConfiguration) { + if (!configuration.effort) return undefined; + if (configuration.provider === "openrouter") + return { openrouter: { reasoning: { effort: configuration.effort } } }; + return { + anthropic: { + effort: configuration.effort, + ...(configuration.provider === "anthropic" + ? { thinking: { type: "adaptive" as const } } + : {}), + }, + }; +} + +function reasoningMetadata(metadata: unknown) { + if (!metadata || typeof metadata !== "object" || !("anthropic" in metadata)) + return undefined; + const value = metadata.anthropic; + if (!value || typeof value !== "object") return undefined; + const fields = value as Record; + const anthropic: Record = {}; + for (const key of ["signature", "redactedData"] as const) + if (typeof fields[key] === "string") anthropic[key] = fields[key]; + return Object.keys(anthropic).length ? { anthropic } : undefined; +} + function createCloudflareAnthropicModel( binding: Ai, model: typeof INKLING_MODEL, @@ -230,6 +258,13 @@ export function protectModel( finish("completed"); return { ...result, + content: result.content.map((part) => ({ + ...part, + providerMetadata: + part.type === "reasoning" + ? reasoningMetadata(part.providerMetadata) + : undefined, + })), request: undefined, response: undefined, warnings: [], @@ -285,7 +320,15 @@ export function protectModel( return; } const chunk = { ...value }; - if ("providerMetadata" in chunk) delete chunk.providerMetadata; + if ("providerMetadata" in chunk) { + // Signed/redacted reasoning is protocol state needed by the + // next tool step, not telemetry. Keep only those exact fields. + const metadata = chunk.type.startsWith("reasoning-") + ? reasoningMetadata(chunk.providerMetadata) + : undefined; + delete chunk.providerMetadata; + if (metadata) chunk.providerMetadata = metadata; + } controller.enqueue(chunk); } catch (error) { finish( diff --git a/worker/model-settings.ts b/worker/model-settings.ts index dce04ea..c5a6b6a 100644 --- a/worker/model-settings.ts +++ b/worker/model-settings.ts @@ -4,7 +4,9 @@ import { MODEL_PROVIDERS, isModelProvider, isCredentialProvider, + modelEfforts, type CredentialProvider, + type ModelEffort, type ModelConfiguration, type ModelProvider, } from "../shared/model-providers.ts"; @@ -20,6 +22,10 @@ export interface ModelSettings { providerSetupUrl: string | null; } +export interface ConversationModelSettings extends ModelSettings { + override: ModelConfiguration | null; +} + export const DEFAULT_MODEL: ModelConfiguration = { provider: "workers-ai", model: MODEL_CATALOG["workers-ai"][0], @@ -30,7 +36,9 @@ export function parseModelConfiguration(value: unknown): ModelConfiguration { throw new Error("Select a supported provider and model"); const record = value as Record; if ( - Object.keys(record).length !== 2 || + Object.keys(record).some( + (key) => !["provider", "model", "effort"].includes(key), + ) || !Object.hasOwn(record, "provider") || !Object.hasOwn(record, "model") || !isModelProvider(record.provider) || @@ -40,9 +48,17 @@ export function parseModelConfiguration(value: unknown): ModelConfiguration { ) ) throw new Error("Select a supported provider and model"); + if ( + Object.hasOwn(record, "effort") && + !modelEfforts(record as ModelConfiguration).includes( + record.effort as ModelEffort, + ) + ) + throw new Error("Select a supported effort for this model"); return { provider: record.provider, model: record.model, + ...(record.effort === undefined ? {} : { effort: record.effort }), } as ModelConfiguration; } diff --git a/worker/personal-agent.ts b/worker/personal-agent.ts index ce226d9..5287463 100644 --- a/worker/personal-agent.ts +++ b/worker/personal-agent.ts @@ -834,13 +834,13 @@ export class PersonalAgent extends Agent { // Internal parent RPC only. Read both fields before awaiting crypto, so a turn // cannot combine one settings version with a concurrent credential replacement. // No Secret instance crosses RPC (custom prototypes are not serializable). - async readModelConfiguration(): Promise<{ + async readModelConfiguration(override?: ModelConfiguration): Promise<{ configuration: ModelConfiguration; apiKey?: string; }> { const row = this.modelSettingsRow(); const configuration = parseModelConfiguration( - JSON.parse(row.configuration), + override ?? JSON.parse(row.configuration), ); this.requireEnabledProvider(configuration.provider); const encrypted = this.providerKeyRow(configuration.provider); -- 2.51.2