From 51ab60f88e72e0c7235ec89c9841db871a7fff30 Mon Sep 17 00:00:00 2001 From: Steven Roussey Date: Thu, 1 Oct 2026 11:56:03 -0700 Subject: [PATCH 1/4] feat(ai): AgentTask turns that end in a checked answer, under budgets MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit What a batch host reading documents needs from the turn loop, measured against thousands of recorded agent sessions: - outputSchema adds a submit_answer tool taking that schema. A passing call ends the turn ("submitted", the answer on `object`); a failing one is answered as a correction, so the model resubmits instead of the host failing the turn. A text reply is reminded twice to submit. - checkSubmission lets the host judge an answer the schema passed — a section left empty that the document has, two figures that disagree — and send it back with a reason, at most twice. - toolConcurrency runs a round's calls at once, results in the order asked; a round with a call put to a person still runs one at a time. - steps[] records each round: timings, attempts, each tool's outcome and size, usage and cost; costUsd totals them. - Budgets (maxInputTokens, maxCostUsd, maxDurationMs) end the turn "budget", never between a tool_use and its result. - maxRoundRetries retries a retryable round failure as the provider's retry-after asks; roundTimeoutMs abandons a round a provider accepted and never answered, as a retryable failure. An unset object port reaches the task as {}, and an empty schema counts as none. Co-Authored-By: Claude Opus 5.5 --- .claude/CLAUDE.md | 22 + packages/ai/src/task/AgentTask.ts | 509 +++++++++++++++-- packages/ai/src/task/AgentTurnRecord.ts | 148 +++++ packages/ai/src/task/index.ts | 1 + .../src/test/ai/AgentTaskTurnControls.test.ts | 529 ++++++++++++++++++ 5 files changed, 1159 insertions(+), 50 deletions(-) create mode 100644 packages/ai/src/task/AgentTurnRecord.ts create mode 100644 packages/test/src/test/ai/AgentTaskTurnControls.test.ts diff --git a/.claude/CLAUDE.md b/.claude/CLAUDE.md index 79175ac33..3427713ed 100644 --- a/.claude/CLAUDE.md +++ b/.claude/CLAUDE.md @@ -228,6 +228,28 @@ approval), so a host draws one lifecycle rather than one per way a call can end. `approval: "never"` turns it off for a headless run, and with no connector registered such a call is refused rather than run. +**A turn that must end in a typed answer** passes `outputSchema`: the loop adds a +`submit_answer` tool (`AGENT_SUBMIT_TOOL_NAME`) taking that schema, ends with +`stopReason: "submitted"` and the answer on `object` when a call passes it, and answers a +failing one as a correction to make, so the model resubmits instead of the host failing the +turn. A model that replies in text is reminded twice, then the turn ends `"answered"` with no +`object`. An unset object port reaches the task as `{}`, and an empty schema counts as none. +`checkSubmission` is where a host enforces what a schema cannot say — a section left empty +that the document plainly has, two figures that must agree: given an answer that passed the +schema and the turn so far, a returned reason goes back to the model as the submit call's +error. It is turned away at most twice, then accepted, so a check can never hold a turn; the +reasons are on `submissionRejections`. +The rest are controls a batch host needs: `toolConcurrency` (opt-in; a round holding a call +put to a person still runs one at a time, and results return in the order asked), +`maxRoundRetries` (default 2, for a `RetryableJobError`, waiting as the provider's +retry-after says within 1–60 s), `roundTimeoutMs` (a provider can accept a request and never answer; +past it the round is abandoned as a retryable failure, so the retries cover it), and budgets — `maxInputTokens`, `maxCostUsd` (refused +without a price card, since a budget it cannot measure never stops anything), +`maxDurationMs` — each ending the turn `"budget"`, never between a `tool_use` and its result. +Every round leaves an `AgentStep` on `steps` (timings, attempts, each tool's outcome and +size, usage, cost) and rides on the `snapshot` beside `messages`; `costUsd` totals them when +every round could be priced. + **The loop and the model call can live in different places.** A host whose tools are closures (they draw on a screen, ask a person, read state only that process holds) but whose model is reachable only through a backend binds `AGENT_ROUND_RUNNER` on the run's registry: every round is diff --git a/packages/ai/src/task/AgentTask.ts b/packages/ai/src/task/AgentTask.ts index 5004a7bce..fa01ff0c7 100644 --- a/packages/ai/src/task/AgentTask.ts +++ b/packages/ai/src/task/AgentTask.ts @@ -11,17 +11,35 @@ import type { StreamEvent, TaskConfig, TaskEntitlements, + TaskOutput, } from "@workglow/task-graph"; -import { CreateWorkflow, Entitlements, Task, Workflow } from "@workglow/task-graph"; -import type { DataPortSchema } from "@workglow/util/schema"; +import { + CreateWorkflow, + Entitlements, + Task, + TaskConfigurationError, + Workflow, +} from "@workglow/task-graph"; +import type { DataPortSchema, JsonSchema } from "@workglow/util/schema"; import type { Capability } from "../capability/Capabilities"; import { createEmitQueue } from "../capability/emitQueue"; +import { readUsage } from "../capability/UsageTelemetry"; import type { ModelConfig } from "../model/ModelSchema"; import type { AgentRoundRunner } from "./AgentRoundRunner"; import { AGENT_ROUND_RUNNER } from "./AgentRoundRunner"; import { assertModelMeetsRequires } from "./base/AiTask"; import type { AgentApprovalMode } from "./AgentToolExecution"; -import { clampToolText, runAgentTool } from "./AgentToolExecution"; +import { clampToolText, runAgentTool, toolCallNeedsApproval } from "./AgentToolExecution"; +import type { AgentStep, AgentSubmissionCheck, AgentToolRecord } from "./AgentTurnRecord"; +import { + AGENT_SUBMIT_TOOL_NAME, + agentModelPricing, + isRetryableRoundError, + promptTokens, + retryDelayMs, + roundCostUsd, + sleepUnlessAborted, +} from "./AgentTurnRecord"; import { DEFAULT_MAX_HISTORY_CHARS, normalizeHistoryForModel, @@ -39,7 +57,8 @@ import { collectToolUseIds, uniquifyToolCallIds } from "./ToolCallIds"; import type { ToolCallingTaskInput, ToolCallingTaskOutput } from "./ToolCallingTask"; import { ToolCallingInputSchema, ToolCallingTask } from "./ToolCallingTask"; import type { ToolCall, ToolDefinition } from "./ToolCallingUtils"; -import { compileToolValidators, sanitizeToolArgs } from "./ToolCallingUtils"; +import { compileToolValidators, sanitizeToolArgs, ToolCallError } from "./ToolCallingUtils"; +import { RetryableJobError } from "@workglow/job-queue"; /** Rounds before the loop gives up on the model reaching an answer. */ const DEFAULT_MAX_ROUNDS = 8; @@ -54,6 +73,31 @@ const DEFAULT_MAX_ROUNDS = 8; */ const DEFAULT_MAX_TOOL_RESULT_CHARS = 20_000; +/** + * Retries of one round's model call after a retryable failure — a rate limit, + * an overloaded provider, a timeout. A turn several rounds in has already paid + * for those rounds, and failing it for one 429 throws them away. + */ +const DEFAULT_MAX_ROUND_RETRIES = 2; + +/** + * Times a turn with an `outputSchema` reminds a model that replied in text to + * submit instead. Some models narrate a final answer out of habit and submit + * when told; one that still will not after this is not going to. + */ +const MAX_SUBMIT_REMINDERS = 2; + +/** + * Times a host's `checkSubmission` may turn an answer away. A check that keeps + * objecting is either wrong or asking for what the filing does not hold, and + * the turn must still end with the answer the model could give. + */ +const MAX_SUBMISSION_REJECTIONS = 2; + +const SUBMIT_REMINDER = + `Call ${AGENT_SUBMIT_TOOL_NAME} with your final answer. ` + + "An answer given as text is not recorded."; + export const AgentInputSchema = { type: "object", properties: { @@ -91,6 +135,66 @@ export const AgentInputSchema = { minimum: 1, "x-ui-group": "Configuration", }, + outputSchema: { + type: "object", + title: "Output Schema", + description: `JSON Schema of the answer. Given one, the turn adds a ${AGENT_SUBMIT_TOOL_NAME} tool taking it, ends on a call that passes it, and reports the answer on the object output`, + additionalProperties: true, + "x-ui-group": "Configuration", + }, + checkSubmission: { + title: "Check Submission", + description: `A host function given an answer that passed outputSchema and the turn so far; a returned reason sends the answer back to the model to fix (at most ${MAX_SUBMISSION_REJECTIONS} times), undefined accepts it`, + "x-ui-hidden": true, + }, + roundTimeoutMs: { + type: "integer", + title: "Round Timeout (ms)", + description: + "Time one model call may take before it is abandoned and retried as a retryable failure", + minimum: 1, + "x-ui-group": "Configuration", + }, + toolConcurrency: { + type: "integer", + title: "Tool Concurrency", + description: + "Tool calls of one round run at once. Results return in the order the model asked; a round with a call that needs approval runs its calls one at a time", + minimum: 1, + "x-ui-group": "Configuration", + }, + maxRoundRetries: { + type: "integer", + title: "Max Round Retries", + description: + "Times one round's model call is retried after a retryable failure (rate limit, overload, timeout), waiting as the provider asks", + minimum: 0, + "x-ui-group": "Configuration", + }, + maxInputTokens: { + type: "integer", + title: "Max Input Tokens", + description: + "Prompt tokens, cached or not, summed over the turn's rounds, after which it stops with stopReason budget", + minimum: 1, + "x-ui-group": "Configuration", + }, + maxCostUsd: { + type: "number", + title: "Max Cost (USD)", + description: + "Estimated spend after which the turn stops with stopReason budget. Needs a price card for the model", + exclusiveMinimum: 0, + "x-ui-group": "Configuration", + }, + maxDurationMs: { + type: "integer", + title: "Max Duration (ms)", + description: + "Wall-clock time after which no further round starts; the turn stops with stopReason budget", + minimum: 1, + "x-ui-group": "Configuration", + }, approval: { type: "string", title: "Approval", @@ -130,11 +234,35 @@ export const AgentOutputSchema = { type: "string", title: "Stop Reason", description: - '"answered" when the model replied without asking for a tool, "max-rounds" when it ran out of rounds first', - enum: ["answered", "max-rounds"], + '"answered" when the model replied without asking for a tool, "submitted" when it submitted an answer that passed outputSchema, "budget" when a token, cost or time budget ran out, "max-rounds" when it ran out of rounds first', + enum: ["answered", "submitted", "budget", "max-rounds"], + }, + object: { + title: "Object", + description: + "The submitted answer, when the turn had an outputSchema and stopped as submitted", + }, + steps: { + type: "array", + items: { type: "object", additionalProperties: true }, + title: "Steps", + description: + "One record per round: when it started, how long the model call and the tools took, retries, each tool's outcome, usage and cost", + }, + submissionRejections: { + type: "array", + items: { type: "string" }, + title: "Submission Rejections", + description: "Each reason checkSubmission gave for sending an answer back, in order", + }, + costUsd: { + type: "number", + title: "Cost (USD)", + description: + "Estimated spend of the turn, when every round reported usage and the model has a USD price card", }, }, - required: ["text", "messages", "rounds", "stopReason"], + required: ["text", "messages", "rounds", "stopReason", "steps", "submissionRejections"], additionalProperties: false, } as const satisfies DataPortSchema; @@ -155,14 +283,31 @@ export type AgentTaskInput = { readonly maxRounds?: number | undefined; readonly maxToolResultChars?: number | undefined; readonly maxHistoryChars?: number | undefined; + readonly outputSchema?: JsonSchema | undefined; + readonly checkSubmission?: AgentSubmissionCheck | undefined; + readonly roundTimeoutMs?: number | undefined; + readonly toolConcurrency?: number | undefined; + readonly maxRoundRetries?: number | undefined; + readonly maxInputTokens?: number | undefined; + readonly maxCostUsd?: number | undefined; + readonly maxDurationMs?: number | undefined; readonly approval?: AgentApprovalMode | undefined; }; +export type AgentStopReason = "answered" | "submitted" | "budget" | "max-rounds"; + export type AgentTaskOutput = { text: string; messages: ChatMessage[]; rounds: number; - stopReason: "answered" | "max-rounds"; + stopReason: AgentStopReason; + /** Present when the turn stopped as `"submitted"`. */ + object?: unknown; + steps: AgentStep[]; + /** Each reason `checkSubmission` gave for sending an answer back. */ + submissionRejections: string[]; + /** Absent when any round's cost could not be estimated. */ + costUsd?: number | undefined; }; export type AgentTaskConfig = TaskConfig; @@ -225,6 +370,54 @@ function toolResult( }; } +interface RunCallOptions { + readonly byName: ReadonlyMap; + readonly validators: ReturnType; + readonly approval: AgentApprovalMode; + readonly maxResultChars: number; + /** The tool that submits the answer, when the turn has an `outputSchema`. */ + readonly submitToolName: string | undefined; +} + +/** + * The tool a turn with an `outputSchema` finishes with: its arguments are the + * answer, checked against that schema like any tool's arguments, so a failing + * answer comes back to the model as an error it can correct. + */ +function submitToolFor( + outputSchema: JsonSchema | undefined, + onSubmit: (value: unknown) => Promise +): ToolDefinition | undefined { + // The runner fills an object port the caller left unset with `{}`, and an + // empty schema constrains nothing: a submit tool built from one would accept + // any arguments at all, so it counts as no schema. + if ( + outputSchema === undefined || + (typeof outputSchema === "object" && Object.keys(outputSchema).length === 0) + ) { + return undefined; + } + return { + name: AGENT_SUBMIT_TOOL_NAME, + description: + "Submit your final answer. It is checked against the required schema: if it is rejected, fix what the error names and call this again.", + inputSchema: outputSchema, + // Recording an answer reaches nothing beyond the turn itself. + requiresApproval: false, + execute: async (input) => { + const reason = await onSubmit(input); + // Thrown as the tool's own words: the model reads the reason itself, not + // a wrapped error, and answers it in its next round. + if (reason !== undefined) { + throw new ToolCallError( + `Answer not recorded yet. ${reason} Then call submit_answer again.` + ); + } + return "Answer recorded."; + }, + }; +} + /** * One conversational turn: call the model, run every tool it asks for, feed the * results back, and repeat until it answers without tools. @@ -241,8 +434,10 @@ function toolResult( * - **Tool-call ids are unique across the whole conversation.** A model that * restarts its numbering at `call_0` each turn would otherwise attach this * turn's result to an earlier turn's call. - * - **Tools run in order.** One may block on a person; the next must not race - * ahead of the answer. + * - **Tools run in order unless the host asks otherwise.** One may block on a + * person, and the next must not race ahead of the answer, so `toolConcurrency` + * is opt-in and a round holding a call that needs approval runs one at a time + * whatever it says. Results go back in the order the model asked either way. * * What it deliberately does NOT do is keep the conversation. `messages` comes * in and goes out, so the host owns the record — which is what lets the same @@ -290,15 +485,68 @@ export class AgentTask extends Task [tool.name, tool])); + + let submitted: { readonly value: unknown } | undefined; + const submissionRejections: string[] = []; + // Read through a getter: the transcript and step list are declared below, + // and a check sees them as they stand when the model submits. + let turnSoFar = (): { messages: readonly ChatMessage[]; steps: readonly AgentStep[] } => ({ + messages: [], + steps: [], + }); + const submitTool = submitToolFor(input.outputSchema, async (value) => { + if ( + input.checkSubmission !== undefined && + submissionRejections.length < MAX_SUBMISSION_REJECTIONS + ) { + const reason = await input.checkSubmission(value, { + ...turnSoFar(), + rejections: submissionRejections.length, + }); + if (reason !== undefined && reason.trim() !== "") { + submissionRejections.push(reason); + return reason; + } + } + submitted = { value }; + return undefined; + }); + if (submitTool !== undefined && input.tools.some((tool) => tool.name === submitTool.name)) { + throw new TaskConfigurationError( + `AgentTask: a tool named "${submitTool.name}" is reserved for submitting against outputSchema` + ); + } + const tools = submitTool === undefined ? input.tools : [...input.tools, submitTool]; + const validators = compileToolValidators(tools); + const byName = new Map(tools.map((tool) => [tool.name, tool])); + const roundTurn: AgentTaskInput = { ...input, tools }; + + const pricing = await agentModelPricing(input.model, context.registry); + if (input.maxCostUsd !== undefined && pricing === undefined) { + // A budget that cannot be measured would never stop the turn, and the + // caller set it precisely so that something would. + throw new TaskConfigurationError( + "AgentTask: maxCostUsd needs a price card for the model, and none was found" + ); + } + + const startedAt = Date.now(); + const steps: AgentStep[] = []; + let spentTokens = 0; + let spentUsd = 0; + let costKnown = true; + let reminders = 0; const messages: ChatMessage[] = [...(input.messages ?? []), promptToUserMessage(input.prompt)]; + turnSoFar = () => ({ messages, steps }); /** * The transcript so far, for a host drawing one while the turn is still * running — a tool card wants to exist from the moment the model asks for - * it, not once the turn is over. + * it, not once the turn is over. The step records ride along, so a host + * persisting a turn round by round has each one as it settles. * * A `snapshot` rather than an `object-delta` on the `messages` port: an * object-delta carrying an array is folded as an upsert list, so successive @@ -306,34 +554,92 @@ export class AgentTask extends Task => - ({ type: "snapshot", data: { messages: [...messages] } }) as StreamEvent; + ({ + type: "snapshot", + data: { messages: [...messages], steps: [...steps] }, + }) as StreamEvent; let text = ""; let rounds = 0; + const finish = (stopReason: AgentStopReason): StreamEvent => ({ + type: "finish", + data: { + text, + messages, + rounds, + stopReason, + steps, + submissionRejections, + ...(submitted === undefined ? {} : { object: submitted.value }), + ...(costKnown && steps.length > 0 ? { costUsd: spentUsd } : {}), + }, + }); + const overBudget = (): boolean => + (input.maxInputTokens !== undefined && spentTokens >= input.maxInputTokens) || + (input.maxCostUsd !== undefined && spentUsd >= input.maxCostUsd); yield transcript(); for (let round = 0; round < maxRounds; round++) { context.signal.throwIfAborted(); + if (input.maxDurationMs !== undefined && Date.now() - startedAt >= input.maxDurationMs) { + yield finish("budget"); + return; + } rounds = round + 1; await context.updateProgress(undefined, "Thinking"); + const roundStart = Date.now(); const captured: { output: ToolCallingTaskOutput | undefined } = { output: undefined }; - for await (const event of this.streamRound( - rounds, - captured, - input, - messages, - maxHistoryChars, - context - )) { - yield event; + let attempts = 0; + for (;;) { + attempts++; + try { + for await (const event of this.streamRound( + rounds, + captured, + roundTurn, + messages, + maxHistoryChars, + context, + input.roundTimeoutMs + )) { + yield event; + } + break; + } catch (err) { + if (attempts > maxRoundRetries || !isRetryableRoundError(err)) throw err; + context.signal.throwIfAborted(); + const wait = retryDelayMs(err, attempts); + await context.updateProgress(undefined, `Retrying in ${Math.ceil(wait / 1000)}s`); + await sleepUnlessAborted(wait, context.signal); + } } + const modelMs = Date.now() - roundStart; const output = captured.output; + const usage = readUsage(output as TaskOutput | undefined); + const costUsd = roundCostUsd(usage, pricing, new Date(roundStart)); // Summed from the round's settled output rather than from the deltas just // forwarded: a provider that reports its text only on the finish event // streams nothing, and counting deltas would report an empty answer while // `messages` carried the real one. text += output?.text ?? ""; + const record = (toolRecords: readonly AgentToolRecord[]): void => { + steps.push({ + round: rounds, + startedAt: new Date(roundStart).toISOString(), + modelMs, + durationMs: Date.now() - roundStart, + attempts, + text: output?.text ?? "", + tools: toolRecords, + usage, + costUsd, + }); + spentTokens += promptTokens(usage); + if (costUsd === undefined) costKnown = false; + else spentUsd += costUsd; + }; + const calls = uniquifyToolCallIds( (output?.toolCalls ?? []).filter(isAnswerable), collectToolUseIds(messages) @@ -347,7 +653,14 @@ export class AgentTask extends Task { + const tool = byName.get(call.name); + return tool !== undefined && toolCallNeedsApproval(tool, approval, context.registry); + }) + ? 1 + : toolConcurrency; const results: ContentBlockToolResult[] = []; - for (const call of calls) { + const toolRecords: AgentToolRecord[] = []; + for await (const event of this.runCalls(calls, width, results, toolRecords, context, { + byName, + validators, + approval, + maxResultChars: maxToolResultChars, + submitToolName: submitTool?.name, + })) { + yield event; + } + messages.push({ role: "tool", content: results }); + record(toolRecords); + yield transcript(); + if (submitted !== undefined) { + yield finish("submitted"); + return; + } + if (overBudget()) { + yield finish("budget"); + return; + } + } + + yield finish("max-rounds"); + } + + /** + * Runs one round's calls, `width` at a time, and reports each one running and + * settled. `results` and `records` are filled in the order the model asked, + * whatever order the calls finish in, because that is the order its next + * round reads them. + */ + private async *runCalls( + calls: readonly ToolCall[], + width: number, + results: ContentBlockToolResult[], + records: AgentToolRecord[], + context: IExecuteContext, + options: RunCallOptions + ): AsyncIterable> { + const queue = createEmitQueue>(); + const settled: Array<{ result: ContentBlockToolResult; record: AgentToolRecord } | undefined> = + calls.map(() => undefined); + let next = 0; + const worker = async (): Promise => { + while (next < calls.length) { context.signal.throwIfAborted(); + const index = next++; + const call = calls[index]!; await context.updateProgress(undefined, `Running ${call.name}`); - yield { type: "tool-call", status: "running", toolCallId: call.id, name: call.name }; - const result = await this.runCall(call, byName, validators, context, { - approval, - maxResultChars: maxToolResultChars, - }); - results.push(result); + queue.push({ type: "tool-call", status: "running", toolCallId: call.id, name: call.name }); + const began = Date.now(); + const result = await this.runCall(call, context, options); + const resultText = toolResultText(result); + settled[index] = { + result, + record: { + id: call.id, + name: call.name, + isError: result.is_error === true, + chars: resultText.length, + durationMs: Date.now() - began, + }, + }; // Read back off the block rather than from a second copy of the text: // what the card says the call produced is then the same string the // model is about to read, clamp included. - yield { + queue.push({ type: "tool-call", status: result.is_error === true ? "failed" : "completed", toolCallId: call.id, name: call.name, - result: toolResultText(result), - }; + result: resultText, + }); } - messages.push({ role: "tool", content: results }); - yield transcript(); + }; + const run = Promise.all( + Array.from({ length: Math.min(width, calls.length) }, () => worker()) + ).finally(() => queue.close()); + // Held so a rejection while the queue is still draining is not reported as + // unhandled; it is re-thrown below, once the events already produced have + // reached the caller. + run.catch(() => {}); + for await (const event of queue.iterable) yield event; + await run; + for (const entry of settled) { + results.push(entry!.result); + records.push(entry!.record); } - - yield { type: "finish", data: { text, messages, rounds, stopReason: "max-rounds" } }; } /** @@ -410,7 +795,8 @@ export class AgentTask extends Task> { const roundInput: ToolCallingTaskInput = { model: input.model, @@ -425,6 +811,12 @@ export class AgentTask extends Task>(); let off = (): void => {}; let runRound: () => Promise; + // A round that never answers holds the turn forever: a provider can accept + // a request and then send nothing back. Abandoning it as a retryable + // failure hands it to the turn's round retries. + const roundAbort = new AbortController(); + let abandon = (): void => roundAbort.abort(); + let timedOut = false; if (hostRunner) { // The owned ToolCallingTask gates the model before it dispatches; a round // handed elsewhere owes the same gate here, or a model that cannot use @@ -435,7 +827,7 @@ export class AgentTask extends Task hostRunner(roundInput, { - signal: context.signal, + signal: AbortSignal.any([context.signal, roundAbort.signal]), onTextDelta: (delta) => queue.push({ type: "text-delta", port: "text", textDelta: delta }), // Progress is advisory: a failed update must not surface as an @@ -447,6 +839,7 @@ export class AgentTask extends Task turn.abort(); off = turn.subscribe("stream_chunk", (event: StreamEvent) => { if (event.type !== "text-delta") return; if ((event.port ?? "text") !== "text") return; @@ -454,10 +847,23 @@ export class AgentTask extends Task turn.run(roundInput); } + const timer = + timeoutMs === undefined + ? undefined + : setTimeout(() => { + timedOut = true; + abandon(); + }, timeoutMs); const run = (async () => { try { captured.output = await runRound(); + } catch (err) { + if (timedOut && !context.signal.aborted) { + throw new RetryableJobError(`Round ${round} had no answer within ${timeoutMs} ms`); + } + throw err; } finally { + if (timer !== undefined) clearTimeout(timer); off(); queue.close(); } @@ -478,11 +884,10 @@ export class AgentTask extends Task, - validators: ReturnType, context: IExecuteContext, - options: { readonly approval: AgentApprovalMode; readonly maxResultChars: number } + options: RunCallOptions ): Promise { + const { byName, validators } = options; const tool = byName.get(call.name); if (!tool) { const known = [...byName.keys()].join(", "); @@ -499,15 +904,19 @@ export class AgentTask extends Task error.message).join("; ") || "invalid arguments"; - return toolResult( - call, - `Invalid arguments for ${call.name}: ${detail}`, - true, - options.maxResultChars - ); + // The answer itself failed its schema: say so as a correction to make, + // which is what gets a model to resubmit rather than give up. + const text = + call.name === options.submitToolName + ? `Answer rejected; nothing was recorded. Fix these and call ${call.name} again: ${detail}` + : `Invalid arguments for ${call.name}: ${detail}`; + return toolResult(call, text, true, options.maxResultChars); } } - const result = await runAgentTool(tool, { ...call, input: sanitized }, context, options); + const result = await runAgentTool(tool, { ...call, input: sanitized }, context, { + approval: options.approval, + maxResultChars: options.maxResultChars, + }); return toolResult(call, result.text, result.isError, options.maxResultChars); } } diff --git a/packages/ai/src/task/AgentTurnRecord.ts b/packages/ai/src/task/AgentTurnRecord.ts new file mode 100644 index 000000000..73e43a659 --- /dev/null +++ b/packages/ai/src/task/AgentTurnRecord.ts @@ -0,0 +1,148 @@ +/** + * @license + * Copyright 2026 Steven Roussey + * SPDX-License-Identifier: Apache-2.0 + */ + +import { isRetryableError, RetryableJobError } from "@workglow/job-queue"; +import type { Usage } from "@workglow/task-graph"; +import type { ServiceRegistry } from "@workglow/util"; +import { estimateCost } from "../capability/CostEstimate"; +import type { ModelPricing } from "../model/ModelPricing"; +import { mergeModelPricing } from "../model/ModelPricing"; +import { getGlobalModelRepository } from "../model/ModelRegistry"; +import type { ModelConfig } from "../model/ModelSchema"; +import type { ChatMessage } from "./ChatMessage"; +import { getAiProviderRegistry } from "../provider/AiProviderRegistry"; + +/** The tool an agent given an `outputSchema` finishes its turn with. */ +export const AGENT_SUBMIT_TOOL_NAME = "submit_answer"; + +/** One tool call a round made, as it ran. Its arguments and result are in `messages`. */ +export interface AgentToolRecord { + readonly id: string; + readonly name: string; + readonly isError: boolean; + /** Characters of the result the model was shown. */ + readonly chars: number; + readonly durationMs: number; +} + +/** + * One round of a turn: the model call and the tools it asked for. + * + * Kept beside `messages` rather than inside it, so a host can persist, chart or + * budget a turn without parsing a transcript, and join back to it by tool-use id. + */ +export interface AgentStep { + readonly round: number; + readonly startedAt: string; + /** The model call alone, retries included. */ + readonly modelMs: number; + /** The model call and every tool it asked for. */ + readonly durationMs: number; + /** Model calls this round took: 1, plus each retry after a retryable failure. */ + readonly attempts: number; + readonly text: string; + readonly tools: readonly AgentToolRecord[]; + /** As the provider reported it; undefined when it reported none. */ + readonly usage: Usage | undefined; + /** Undefined when the usage or the model's price card is missing, or the card is not in USD. */ + readonly costUsd: number | undefined; +} + +/** What a host's submission check sees besides the answer. */ +export interface AgentSubmissionContext { + /** The turn's transcript so far: every call the model made and what it read. */ + readonly messages: readonly ChatMessage[]; + readonly steps: readonly AgentStep[]; + /** Answers this check has already sent back this turn. */ + readonly rejections: number; +} + +/** + * Judges an answer that passed `outputSchema`: a reason to send it back to the + * model — something missing, something inconsistent — or undefined to accept. + * Where the invariants a schema cannot state are enforced in code. + */ +export type AgentSubmissionCheck = ( + answer: unknown, + turn: AgentSubmissionContext +) => string | undefined | Promise; + +/** + * Tokens a round sent the model: every slice of the prompt, cached or not. + * What a context window and an input budget both measure. + */ +export function promptTokens(usage: Usage | undefined): number { + if (usage === undefined) return 0; + return ( + (usage.input ?? 0) + + (usage.cached ?? 0) + + (usage.cacheWrite ?? 0) + + (usage.imageInput ?? 0) + + (usage.imageCached ?? 0) + ); +} + +/** + * The rate card a turn's rounds are priced with: the provider's list card, + * overridden field by field by any card on the model's own record. Undefined + * when the model id resolves to nothing or no card exists. + */ +export async function agentModelPricing( + model: string | ModelConfig, + registry: ServiceRegistry +): Promise { + const config = + typeof model === "string" ? await getGlobalModelRepository(registry).findByName(model) : model; + if (config === undefined) return undefined; + const provider = getAiProviderRegistry(registry).getProvider(config.provider); + return mergeModelPricing(config.pricing, provider?.modelPricing(config)); +} + +/** A round's cost in US dollars, priced at the instant it was sent. */ +export function roundCostUsd( + usage: Usage | undefined, + pricing: ModelPricing | undefined, + at: Date +): number | undefined { + if (usage === undefined || pricing === undefined) return undefined; + const estimate = estimateCost(usage, pricing, { at }); + return estimate?.currency === "USD" ? estimate.amount : undefined; +} + +/** Whether a failed model call may be made again: a rate limit, an overload, a timeout. */ +export function isRetryableRoundError(err: unknown): boolean { + return isRetryableError(err); +} + +const MIN_RETRY_MS = 1_000; +const MAX_RETRY_MS = 60_000; + +/** + * How long to wait before retry `attempt` (1-based) of a round: the provider's + * own retry-after when it gave one, else exponential from two seconds; within + * one second and one minute either way. + */ +export function retryDelayMs(err: unknown, attempt: number, now: number = Date.now()): number { + const stated = err instanceof RetryableJobError ? err.retryDate?.getTime() : undefined; + const wanted = stated !== undefined ? stated - now : 2_000 * 2 ** (attempt - 1); + return Math.min(MAX_RETRY_MS, Math.max(MIN_RETRY_MS, wanted)); +} + +/** Resolves after `ms`, or rejects with the signal's reason the moment it aborts. */ +export function sleepUnlessAborted(ms: number, signal: AbortSignal): Promise { + signal.throwIfAborted(); + return new Promise((resolve, reject) => { + const onAbort = (): void => { + clearTimeout(timer); + reject(signal.reason); + }; + const timer = setTimeout(() => { + signal.removeEventListener("abort", onAbort); + resolve(); + }, ms); + signal.addEventListener("abort", onAbort, { once: true }); + }); +} diff --git a/packages/ai/src/task/index.ts b/packages/ai/src/task/index.ts index 107f12d36..4eef43254 100644 --- a/packages/ai/src/task/index.ts +++ b/packages/ai/src/task/index.ts @@ -11,6 +11,7 @@ export { registerAiTasks } from "./registerAiTasks"; export * from "./AgentRoundRunner"; export * from "./AgentTask"; export * from "./AgentToolExecution"; +export * from "./AgentTurnRecord"; export * from "./AiChatTask"; export * from "./AiChatWithKbTask"; export * from "./BackgroundRemovalTask"; diff --git a/packages/test/src/test/ai/AgentTaskTurnControls.test.ts b/packages/test/src/test/ai/AgentTaskTurnControls.test.ts new file mode 100644 index 000000000..cd0f57e92 --- /dev/null +++ b/packages/test/src/test/ai/AgentTaskTurnControls.test.ts @@ -0,0 +1,529 @@ +/** + * @license + * Copyright 2026 Steven Roussey + * SPDX-License-Identifier: Apache-2.0 + */ + +import type { AiProviderRunFn, ModelConfig, ToolDefinition } from "@workglow/ai"; +import { + AGENT_SUBMIT_TOOL_NAME, + AgentTask, + AiProviderRegistry, + DirectExecutionStrategy, + getAiProviderRegistry, + retryDelayMs, + setAiProviderRegistry, +} from "@workglow/ai"; +import { PermanentJobError, RetryableJobError } from "@workglow/job-queue"; +import type { Usage } from "@workglow/task-graph"; +import { TaskConfigurationError } from "@workglow/task-graph"; +import type { IHumanConnector } from "@workglow/util"; +import { Container, HUMAN_CONNECTOR, ServiceRegistry } from "@workglow/util"; +import { afterEach, beforeEach, describe, expect, it } from "vitest"; + +const PROVIDER = "mock-agent-turn-controls"; + +const MODEL: ModelConfig = { + model_id: "mock/turn-controls", + provider: PROVIDER, + title: "", + description: "", + capabilities: ["tool-use"], + provider_config: {}, + metadata: {}, +}; + +/** $1 per million uncached input tokens, 10c cached, $2 output: round numbers to check sums against. */ +const PRICED_MODEL: ModelConfig = { + ...MODEL, + pricing: { currency: "USD", input: 1, cached: 0.1, output: 2 }, +} as ModelConfig; + +interface Round { + readonly text?: string; + readonly calls?: ReadonlyArray<{ id: string; name: string; input: Record }>; + readonly usage?: Usage; + /** Thrown instead of answering. */ + readonly error?: Error; + /** Never answers: waits until the call is abandoned. */ + readonly hang?: boolean; +} + +function usage(input: number, cached: number, output: number): Usage { + return { + input, + cached, + output, + cacheWrite: 0, + reasoning: undefined, + total: undefined, + extra: undefined, + }; +} + +/** One entry per model call, the last repeating. */ +function script(rounds: readonly Round[]): () => number { + let called = 0; + const runFn: AiProviderRunFn = async (_input, _model, signal, emit) => { + const round = rounds[Math.min(called, rounds.length - 1)]!; + called++; + if (round.error) throw round.error; + if (round.hang) { + await new Promise((_resolve, reject) => { + if (signal.aborted) reject(signal.reason); + signal.addEventListener("abort", () => reject(signal.reason), { once: true }); + }); + } + if (round.text) emit({ type: "text-delta", port: "text", textDelta: round.text }); + if (round.calls?.length) { + emit({ type: "object-delta", port: "toolCalls", objectDelta: [...round.calls] }); + } + emit({ type: "finish", data: {}, ...(round.usage ? { usage: round.usage } : {}) }); + }; + getAiProviderRegistry().registerRunFn(PROVIDER, { serves: ["tool-use"], runFn }); + return () => called; +} + +const ANSWER_SCHEMA = { + type: "object", + properties: { trust: { type: "integer" } }, + required: ["trust"], + additionalProperties: false, +} as const; + +/** A function tool that records how many of its calls overlap, and finishes in `ms`. */ +function timedTool( + name: string, + ms: number, + overlap: { active: number; peak: number }, + extra: Partial = {} +): ToolDefinition { + return { + name, + description: name, + inputSchema: { type: "object", properties: {}, additionalProperties: false }, + execute: async () => { + overlap.active++; + overlap.peak = Math.max(overlap.peak, overlap.active); + await new Promise((resolve) => setTimeout(resolve, ms)); + overlap.active--; + return `${name} done`; + }, + ...extra, + }; +} + +const ECHO: ToolDefinition = { + name: "echo", + description: "Echoes", + inputSchema: { type: "object", properties: { text: { type: "string" } } }, + execute: async (input) => String((input as { text?: string }).text ?? ""), +}; + +function toolResults(messages: readonly { role: string; content: readonly unknown[] }[]) { + return messages + .filter((message) => message.role === "tool") + .flatMap((message) => message.content as ReadonlyArray>); +} + +describe("AgentTask turn controls", () => { + let registry: ServiceRegistry; + + beforeEach(() => { + setAiProviderRegistry(new AiProviderRegistry()); + getAiProviderRegistry().setDefaultStrategy(new DirectExecutionStrategy()); + registry = new ServiceRegistry(new Container()); + }); + + afterEach(() => { + getAiProviderRegistry().unregisterProvider(PROVIDER); + }); + + describe("outputSchema", () => { + it("ends the turn on a submission that passes, after correcting one that did not", async () => { + const called = script([ + { calls: [{ id: "s1", name: AGENT_SUBMIT_TOOL_NAME, input: { trust: "lots" } }] }, + { calls: [{ id: "s2", name: AGENT_SUBMIT_TOOL_NAME, input: { trust: 230 } }] }, + ]); + const output = await new AgentTask().run( + { model: MODEL, prompt: "?", tools: [], outputSchema: ANSWER_SCHEMA, approval: "never" }, + { registry } + ); + expect(called()).toBe(2); + expect(output.stopReason).toBe("submitted"); + expect(output.object).toEqual({ trust: 230 }); + const [rejected, accepted] = toolResults(output.messages); + expect(rejected).toMatchObject({ tool_use_id: "s1", is_error: true }); + expect(JSON.stringify(rejected)).toContain("Answer rejected"); + expect(accepted).toMatchObject({ tool_use_id: "s2", is_error: undefined }); + }); + + it("reminds a model that answered in text to submit", async () => { + script([ + { text: "The trust holds 230." }, + { calls: [{ id: "s1", name: AGENT_SUBMIT_TOOL_NAME, input: { trust: 230 } }] }, + ]); + const output = await new AgentTask().run( + { model: MODEL, prompt: "?", tools: [], outputSchema: ANSWER_SCHEMA, approval: "never" }, + { registry } + ); + expect(output.stopReason).toBe("submitted"); + expect(output.messages.map((message) => message.role)).toEqual([ + "user", + "assistant", + "user", + "assistant", + "tool", + ]); + expect(JSON.stringify(output.messages[2])).toContain(AGENT_SUBMIT_TOOL_NAME); + }); + + it("stops reminding after two, and reports the turn answered with no object", async () => { + const called = script([{ text: "230." }]); + const output = await new AgentTask().run( + { model: MODEL, prompt: "?", tools: [], outputSchema: ANSWER_SCHEMA, approval: "never" }, + { registry } + ); + expect(called()).toBe(3); + expect(output.stopReason).toBe("answered"); + expect(output.object).toBeUndefined(); + }); + + it("sends an answer back with checkSubmission's reason, and records the next", async () => { + script([ + { calls: [{ id: "s1", name: AGENT_SUBMIT_TOOL_NAME, input: { trust: 0 } }] }, + { calls: [{ id: "s2", name: AGENT_SUBMIT_TOOL_NAME, input: { trust: 230 } }] }, + ]); + const seen: number[] = []; + const output = await new AgentTask().run( + { + model: MODEL, + prompt: "?", + tools: [], + outputSchema: ANSWER_SCHEMA, + approval: "never", + checkSubmission: (answer, turn) => { + seen.push(turn.messages.length); + return (answer as { trust: number }).trust === 0 + ? "The trust section states the amount; read it." + : undefined; + }, + }, + { registry } + ); + expect(output.stopReason).toBe("submitted"); + expect(output.object).toEqual({ trust: 230 }); + expect(output.submissionRejections).toEqual([ + "The trust section states the amount; read it.", + ]); + const [rejected] = toolResults(output.messages); + expect(rejected).toMatchObject({ tool_use_id: "s1", is_error: true }); + expect(JSON.stringify(rejected)).toContain("read it."); + // The check saw the transcript as it stood, the submitting call included. + expect(seen[0]).toBeGreaterThanOrEqual(2); + }); + + it("accepts an answer after two rejections, whatever the check says", async () => { + script([{ calls: [{ id: "s", name: AGENT_SUBMIT_TOOL_NAME, input: { trust: 1 } }] }]); + const output = await new AgentTask().run( + { + model: MODEL, + prompt: "?", + tools: [], + outputSchema: ANSWER_SCHEMA, + approval: "never", + checkSubmission: () => "Still wrong.", + }, + { registry } + ); + expect(output.stopReason).toBe("submitted"); + expect(output.submissionRejections).toHaveLength(2); + expect(output.rounds).toBe(3); + }); + + it("refuses a tool already named like the submit tool", async () => { + script([{ text: "x" }]); + await expect( + new AgentTask().run( + { + model: MODEL, + prompt: "?", + tools: [{ ...ECHO, name: AGENT_SUBMIT_TOOL_NAME }], + outputSchema: ANSWER_SCHEMA, + approval: "never", + }, + { registry } + ) + ).rejects.toBeInstanceOf(TaskConfigurationError); + }); + }); + + describe("toolConcurrency", () => { + const calls = [ + { id: "a", name: "slow", input: {} }, + { id: "b", name: "fast", input: {} }, + ]; + + it("runs a round's calls at once, and returns their results in the order asked", async () => { + script([{ calls }, { text: "done" }]); + const overlap = { active: 0, peak: 0 }; + const output = await new AgentTask().run( + { + model: MODEL, + prompt: "?", + tools: [timedTool("slow", 40, overlap), timedTool("fast", 5, overlap)], + toolConcurrency: 2, + approval: "never", + }, + { registry } + ); + expect(overlap.peak).toBe(2); + expect(toolResults(output.messages).map((result) => result.tool_use_id)).toEqual(["a", "b"]); + expect(output.steps[0]!.tools.map((tool) => tool.id)).toEqual(["a", "b"]); + }); + + it("runs them one at a time by default", async () => { + script([{ calls }, { text: "done" }]); + const overlap = { active: 0, peak: 0 }; + await new AgentTask().run( + { + model: MODEL, + prompt: "?", + tools: [timedTool("slow", 20, overlap), timedTool("fast", 5, overlap)], + approval: "never", + }, + { registry } + ); + expect(overlap.peak).toBe(1); + }); + + it("runs a round one at a time when any of its calls is put to a person", async () => { + script([{ calls }, { text: "done" }]); + const approve: IHumanConnector = { + send: async (request) => ({ + requestId: request.requestId, + action: "accept", + content: undefined, + done: true, + }), + }; + registry.registerInstance(HUMAN_CONNECTOR, approve); + const overlap = { active: 0, peak: 0 }; + const output = await new AgentTask().run( + { + model: MODEL, + prompt: "?", + tools: [ + timedTool("slow", 20, overlap, { requiresApproval: true }), + timedTool("fast", 5, overlap), + ], + toolConcurrency: 2, + }, + { registry } + ); + // Both ran: the approved one too, so the peak measures ordering, not a refusal. + expect(toolResults(output.messages).map((result) => result.is_error)).toEqual([ + undefined, + undefined, + ]); + expect(overlap.peak).toBe(1); + }); + }); + + describe("steps", () => { + it("records each round's tools, usage and cost, and totals the cost", async () => { + script([ + { + calls: [{ id: "c1", name: "echo", input: { text: "hi" } }], + usage: usage(1_000, 9_000, 100), + }, + { text: "done", usage: usage(500, 10_000, 50) }, + ]); + const output = await new AgentTask().run( + { model: PRICED_MODEL, prompt: "?", tools: [ECHO], approval: "never" }, + { registry } + ); + expect(output.steps).toHaveLength(2); + const [first, second] = output.steps; + expect(first).toMatchObject({ + round: 1, + attempts: 1, + tools: [{ id: "c1", name: "echo", isError: false, chars: 2 }], + }); + expect(first!.usage?.input).toBe(1_000); + // 1,000 × $1 + 9,000 × $0.10 + 100 × $2, per million. + expect(first!.costUsd).toBeCloseTo(0.0021, 8); + expect(second!.costUsd).toBeCloseTo(0.0016, 8); + expect(output.costUsd).toBeCloseTo(0.0037, 8); + }); + + it("leaves the total cost out when a round's could not be priced", async () => { + script([{ text: "done", usage: usage(10, 0, 1) }]); + const output = await new AgentTask().run( + { model: MODEL, prompt: "?", tools: [], approval: "never" }, + { registry } + ); + expect(output.steps[0]!.costUsd).toBeUndefined(); + expect(output.costUsd).toBeUndefined(); + }); + }); + + describe("budgets", () => { + const forever: Round = { + calls: [{ id: "c", name: "echo", input: { text: "again" } }], + usage: usage(1_000, 0, 10), + }; + + it("stops on maxInputTokens once the rounds have spent it, with every call answered", async () => { + const called = script([forever]); + const output = await new AgentTask().run( + { + model: MODEL, + prompt: "?", + tools: [ECHO], + maxInputTokens: 1_500, + maxRounds: 10, + approval: "never", + }, + { registry } + ); + expect(called()).toBe(2); + expect(output.stopReason).toBe("budget"); + const uses = output.messages.flatMap((message) => + message.content.filter((block) => block.type === "tool_use") + ); + expect(toolResults(output.messages)).toHaveLength(uses.length); + }); + + it("stops on maxCostUsd", async () => { + // 1,000 uncached tokens at $1 per million plus 10 output at $2: $0.00102 a round. + const called = script([forever]); + const output = await new AgentTask().run( + { + model: PRICED_MODEL, + prompt: "?", + tools: [ECHO], + maxCostUsd: 0.0025, + maxRounds: 10, + approval: "never", + }, + { registry } + ); + expect(called()).toBe(3); + expect(output.stopReason).toBe("budget"); + }); + + it("refuses maxCostUsd for a model with no price card", async () => { + script([forever]); + await expect( + new AgentTask().run( + { model: MODEL, prompt: "?", tools: [ECHO], maxCostUsd: 1, approval: "never" }, + { registry } + ) + ).rejects.toBeInstanceOf(TaskConfigurationError); + }); + + it("starts no round after maxDurationMs", async () => { + const called = script([forever]); + const slowEcho: ToolDefinition = { + ...ECHO, + execute: async () => { + await new Promise((resolve) => setTimeout(resolve, 30)); + return "late"; + }, + }; + const output = await new AgentTask().run( + { + model: MODEL, + prompt: "?", + tools: [slowEcho], + maxDurationMs: 20, + maxRounds: 10, + approval: "never", + }, + { registry } + ); + expect(called()).toBe(1); + expect(output.stopReason).toBe("budget"); + }); + }); + + describe("round retries", () => { + it("retries a round after a retryable failure, and counts the attempt", async () => { + const called = script([ + { error: new RetryableJobError("429 rate limited", new Date(Date.now())) }, + { text: "done" }, + ]); + const output = await new AgentTask().run( + { model: MODEL, prompt: "?", tools: [], approval: "never" }, + { registry } + ); + expect(called()).toBe(2); + expect(output.stopReason).toBe("answered"); + expect(output.steps[0]!.attempts).toBe(2); + }); + + it("fails at once on a permanent failure, or when retries are off", async () => { + const permanent = script([ + { error: new PermanentJobError("401 bad key") }, + { text: "never" }, + ]); + await expect( + new AgentTask().run( + { model: MODEL, prompt: "?", tools: [], approval: "never" }, + { registry } + ) + ).rejects.toThrow(/401/); + expect(permanent()).toBe(1); + + getAiProviderRegistry().unregisterProvider(PROVIDER); + const off = script([{ error: new RetryableJobError("429") }, { text: "never" }]); + await expect( + new AgentTask().run( + { model: MODEL, prompt: "?", tools: [], maxRoundRetries: 0, approval: "never" }, + { registry } + ) + ).rejects.toThrow(/429/); + expect(off()).toBe(1); + }); + + it("abandons a round with no answer after roundTimeoutMs, and retries it", async () => { + const called = script([{ hang: true }, { text: "done" }]); + const output = await new AgentTask().run( + { model: MODEL, prompt: "?", tools: [], roundTimeoutMs: 50, approval: "never" }, + { registry } + ); + expect(called()).toBe(2); + expect(output.stopReason).toBe("answered"); + expect(output.steps[0]!.attempts).toBe(2); + }); + + it("fails the turn on a round with no answer when retries are off", async () => { + script([{ hang: true }]); + await expect( + new AgentTask().run( + { + model: MODEL, + prompt: "?", + tools: [], + roundTimeoutMs: 50, + maxRoundRetries: 0, + approval: "never", + }, + { registry } + ) + ).rejects.toThrow(/no answer within 50 ms/); + }); + + it("waits as the provider asks, within a second and a minute", () => { + const now = 1_000_000; + expect(retryDelayMs(new RetryableJobError("x", new Date(now + 5_000)), 1, now)).toBe(5_000); + expect(retryDelayMs(new RetryableJobError("x", new Date(now + 10)), 1, now)).toBe(1_000); + expect(retryDelayMs(new RetryableJobError("x", new Date(now + 600_000)), 1, now)).toBe( + 60_000 + ); + expect(retryDelayMs(new RetryableJobError("x"), 1, now)).toBe(2_000); + expect(retryDelayMs(new RetryableJobError("x"), 3, now)).toBe(8_000); + }); + }); +}); From a8abd4ed8e3493218edebc94f4495211248d305d Mon Sep 17 00:00:00 2001 From: Steven Roussey Date: Thu, 1 Oct 2026 12:19:38 -0700 Subject: [PATCH 2/4] fix(deepseek): bill every model off-peak in the published windows; opt out of turbo agent guidance DeepSeek's v4-pro card still carried the retired 16:30-00:30 discount; every model is now billed at half rate outside 01:00-04:00 and 06:00-10:00 UTC, as the published table states. Checked against billing: 2,933 deepseek-flash calls priced at $6.07 on this card, against $6.08 off the account balance. turbo 2.11 rewrites AGENTS.md with a guidance block when it detects an AI agent, which replaces the symlink to .claude/CLAUDE.md with a regular file. agentGuidance: false turns that off. Co-Authored-By: Claude Opus 5.5 --- packages/ai/src/model/ModelPricing.ts | 2 +- .../test/ai-provider/DeepSeek_Pricing.test.ts | 46 +++++++++++++++++++ .../src/ai/common/DeepSeek_Pricing.ts | 36 +++++++-------- turbo.json | 3 +- 4 files changed, 67 insertions(+), 20 deletions(-) create mode 100644 packages/test/src/test/ai-provider/DeepSeek_Pricing.test.ts diff --git a/packages/ai/src/model/ModelPricing.ts b/packages/ai/src/model/ModelPricing.ts index eef345614..2d482629a 100644 --- a/packages/ai/src/model/ModelPricing.ts +++ b/packages/ai/src/model/ModelPricing.ts @@ -53,7 +53,7 @@ export interface ModelUsageTier { /** * A rate card that replaces the base one inside a daily clock window, which is - * how time-of-day discounts are published (DeepSeek's runs 16:30-00:30 UTC). + * how time-of-day discounts are published (DeepSeek's off-peak hours are an example). * * `start` and `end` are `HH:MM` in **UTC** — providers publish these windows in * UTC and a local-time reading would silently misprice by the host's offset. diff --git a/packages/test/src/test/ai-provider/DeepSeek_Pricing.test.ts b/packages/test/src/test/ai-provider/DeepSeek_Pricing.test.ts new file mode 100644 index 000000000..969a50bb9 --- /dev/null +++ b/packages/test/src/test/ai-provider/DeepSeek_Pricing.test.ts @@ -0,0 +1,46 @@ +/** + * @license + * Copyright 2026 Steven Roussey + * SPDX-License-Identifier: Apache-2.0 + */ + +import { estimateCost, resolveEffectiveRates } from "@workglow/ai"; +import { getDeepSeekModelPricing } from "@workglow/deepseek/ai"; +import { describe, expect, it } from "vitest"; + +const at = (iso: string): Date => new Date(iso); + +describe("DeepSeek price card", () => { + it.each(["deepseek-flash", "deepseek-v4-pro"])( + "bills %s at half rate outside the weekday peak hours", + (model) => { + const card = getDeepSeekModelPricing(model)!; + const peak = resolveEffectiveRates(card, { at: at("2026-10-01T02:00:00Z") }); + const morningPeak = resolveEffectiveRates(card, { at: at("2026-10-01T08:00:00Z") }); + const gap = resolveEffectiveRates(card, { at: at("2026-10-01T05:00:00Z") }); + const evening = resolveEffectiveRates(card, { at: at("2026-10-01T17:00:00Z") }); + expect(morningPeak.input).toBe(peak.input); + expect(gap.input).toBe(peak.input! / 2); + expect(evening.output).toBe(peak.output! / 2); + } + ); + + it("prices a measured batch as the account billed it", () => { + // 2,933 calls between 15:38 and 18:48 UTC came to $6.07 on this card and + // $6.08 off the account balance; one representative call at 17:00: + const cost = estimateCost( + { + input: 68_243, + cached: 1_257_216, + output: 32_836, + cacheWrite: 0, + reasoning: undefined, + total: undefined, + extra: undefined, + }, + getDeepSeekModelPricing("deepseek-flash"), + { at: at("2026-10-01T17:00:00Z") } + )!; + expect(cost.amount).toBeCloseTo((68_243 * 0.15 + 1_257_216 * 0.003 + 32_836 * 0.6) / 1e6, 8); + }); +}); diff --git a/providers/deepseek/src/ai/common/DeepSeek_Pricing.ts b/providers/deepseek/src/ai/common/DeepSeek_Pricing.ts index a8b2ebf9a..53fa4d0ef 100644 --- a/providers/deepseek/src/ai/common/DeepSeek_Pricing.ts +++ b/providers/deepseek/src/ai/common/DeepSeek_Pricing.ts @@ -8,12 +8,20 @@ import type { ModelPricing, ModelTimingTier } from "@workglow/ai"; import { resolveModelPricingFromTable } from "@workglow/ai"; /** - * DeepSeek's nightly discount, published as 16:30-00:30 UTC. Shared by every - * model that gets it so the window is stated once. + * The windows DeepSeek bills at half its peak rates, for every model. Peak is + * 01:00-04:00 and 06:00-10:00 UTC on weekdays; everything else is off-peak. + * These two daily windows cover the weekday gaps; a clock window cannot drop + * weekend mornings, which the published table also treats as off-peak, so a + * weekend call in a peak hour is priced at the peak rate. */ -const DEEPSEEK_OFF_PEAK: ModelTimingTier[] = [ - { start: "16:30", end: "00:30", pricing: { input: 0.66, output: 1.98, cached: 0.022 } }, -]; +function offPeak(rates: { input: number; output: number; cached: number }): ModelTimingTier[] { + return [ + { start: "10:00", end: "01:00", pricing: { ...rates } }, + { start: "04:00", end: "06:00", pricing: { ...rates } }, + ]; +} + +const DEEPSEEK_OFF_PEAK: ModelTimingTier[] = offPeak({ input: 0.66, output: 1.98, cached: 0.022 }); const DEEPSEEK_PRO: ModelPricing = { currency: "USD", @@ -23,19 +31,11 @@ const DEEPSEEK_PRO: ModelPricing = { timingTiers: DEEPSEEK_OFF_PEAK, }; -/** Half the peak Flash rates, stated once for both off-peak windows. */ -const DEEPSEEK_FLASH_OFF_PEAK_RATES = { input: 0.15, output: 0.6, cached: 0.003 } as const; - -/** - * DeepSeek Flash off-peak windows. Peak is 01:00–04:00 and 06:00–10:00 UTC - * weekdays; everything else is half those rates. These two clock windows cover - * the daily peak gaps; they cannot drop weekend mornings, which the published - * table also treats as off-peak. - */ -const DEEPSEEK_FLASH_OFF_PEAK: ModelTimingTier[] = [ - { start: "10:00", end: "01:00", pricing: { ...DEEPSEEK_FLASH_OFF_PEAK_RATES } }, - { start: "04:00", end: "06:00", pricing: { ...DEEPSEEK_FLASH_OFF_PEAK_RATES } }, -]; +const DEEPSEEK_FLASH_OFF_PEAK: ModelTimingTier[] = offPeak({ + input: 0.15, + output: 0.6, + cached: 0.003, +}); /** Peak list rates for DeepSeek Flash; the dated and versioned ids share this card. */ const DEEPSEEK_FLASH: ModelPricing = { diff --git a/turbo.json b/turbo.json index e8eb7ece1..3933fe4ef 100644 --- a/turbo.json +++ b/turbo.json @@ -1,5 +1,6 @@ { "$schema": "https://turbo.build/schema.json", + "agentGuidance": false, "globalDependencies": [ "tsconfig.json", "vitest.config.ts", @@ -105,4 +106,4 @@ ] } } -} \ No newline at end of file +} From 6004c476c0084d10807a77729f2b95b31359ea2f Mon Sep 17 00:00:00 2001 From: Steven Roussey Date: Thu, 1 Oct 2026 12:52:14 -0700 Subject: [PATCH 3/4] fix(ai): retry a provider's statusless server error; ask for the whole answer back MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit OpenAI reports a server failure mid-stream as a server_error event with no HTTP status ("An error occurred while processing the request."), which fell through to the unknown-error default and failed the turn outright; it is classified retryable like a 5xx now. A rejected submission now asks for the complete answer again. Told only what to fix, a model resubmitted that part and dropped sections it had already filled — thirteen underwriters to none. Co-Authored-By: Claude Opus 5.5 --- packages/ai/src/job/AiJob.ts | 26 +++++++++++++++++++ packages/ai/src/task/AgentTask.ts | 7 +++-- .../ai/AiJob_classifyProviderError.test.ts | 16 ++++++++++++ 3 files changed, 47 insertions(+), 2 deletions(-) diff --git a/packages/ai/src/job/AiJob.ts b/packages/ai/src/job/AiJob.ts index a78b1a2e6..ec8472cc5 100644 --- a/packages/ai/src/job/AiJob.ts +++ b/packages/ai/src/job/AiJob.ts @@ -202,6 +202,15 @@ export function classifyProviderError(err: unknown, taskType: string, provider: ); } + // A provider's own server failure reported in the body or mid-stream rather + // than as an HTTP status — OpenAI's streaming "server_error" event, "An error + // occurred while processing the request." — is as transient as a 5xx. + if (isServerErrorReport(err, message)) { + return new RetryableJobError( + withJobErrorDiagnostics(`Server error from ${provider} for ${taskType}: ${message}`, err) + ); + } + if ( message.includes("ECONNREFUSED") || message.includes("ECONNRESET") || @@ -227,6 +236,23 @@ export function classifyProviderError(err: unknown, taskType: string, provider: ); } +/** + * Whether an error is a provider reporting its own server failure without an + * HTTP status: a `server_error` code or type on the error or the body it + * carries, or the message providers send with one. + */ +function isServerErrorReport(err: unknown, message: string): boolean { + const fields = (value: unknown): unknown[] => { + if (value === null || typeof value !== "object") return []; + const record = value as { code?: unknown; type?: unknown }; + return [record.code, record.type]; + }; + const body = + err !== null && typeof err === "object" ? (err as { error?: unknown }).error : undefined; + if ([...fields(err), ...fields(body)].includes("server_error")) return true; + return /an error occurred while processing (the|your) request/i.test(message); +} + export class AiJob< Input extends AiJobInput = AiJobInput, Output extends TaskOutput = TaskOutput, diff --git a/packages/ai/src/task/AgentTask.ts b/packages/ai/src/task/AgentTask.ts index fa01ff0c7..208ebec8a 100644 --- a/packages/ai/src/task/AgentTask.ts +++ b/packages/ai/src/task/AgentTask.ts @@ -409,8 +409,11 @@ function submitToolFor( // Thrown as the tool's own words: the model reads the reason itself, not // a wrapped error, and answers it in its next round. if (reason !== undefined) { + // The whole answer again, not a patch: a model resubmitting only what the + // reason named drops every section it had already filled. throw new ToolCallError( - `Answer not recorded yet. ${reason} Then call submit_answer again.` + `Answer not recorded yet. ${reason} Then call submit_answer again with the complete ` + + "answer — every field, including those this does not mention, as you had them." ); } return "Answer recorded."; @@ -908,7 +911,7 @@ export class AgentTask extends Task { + it("retries a provider's own server failure reported without an HTTP status", () => { + const streamed = Object.assign(new Error("An error occurred while processing the request."), { + code: "server_error", + }); + expect(classifyProviderError(streamed, "ToolCallingTask", "OPENAI")).toBeInstanceOf( + RetryableJobError + ); + const nested = Object.assign(new Error("stream failed"), { error: { type: "server_error" } }); + expect(classifyProviderError(nested, "ToolCallingTask", "OPENAI")).toBeInstanceOf( + RetryableJobError + ); + expect( + classifyProviderError(new Error("something unexplained"), "ToolCallingTask", "OPENAI") + ).toBeInstanceOf(PermanentJobError); + }); + it("maps ProviderUnsupportedFeatureError to PermanentJobError", () => { const err = new ProviderUnsupportedFeatureError("mask", "m", "not supported"); const classified = classifyProviderError(err, "ImageGenerateTask", "TEST_PROVIDER"); From a07972d0d331d61d351c4965640a20223e80782e Mon Sep 17 00:00:00 2001 From: Steven Roussey Date: Thu, 1 Oct 2026 13:06:29 -0700 Subject: [PATCH 4/4] test(eval): expect DeepSeek v4-pro's published off-peak windows The card now bills every DeepSeek model at half rate outside the weekday peak hours (01:00-04:00, 06:00-10:00 UTC); this test still pinned the retired 16:30-00:30 window. Co-Authored-By: Claude Opus 5.5 --- examples/eval/src/test/models.test.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/examples/eval/src/test/models.test.ts b/examples/eval/src/test/models.test.ts index 8ec381375..aeb24ecfa 100644 --- a/examples/eval/src/test/models.test.ts +++ b/examples/eval/src/test/models.test.ts @@ -27,8 +27,10 @@ describe("resolveModelConfig", () => { input: 1.32, output: 3.96, cached: 0.044, + // Half rate outside DeepSeek's weekday peak hours, 01:00-04:00 and 06:00-10:00 UTC. timingTiers: [ - { start: "16:30", end: "00:30", pricing: { input: 0.66, output: 1.98, cached: 0.022 } }, + { start: "10:00", end: "01:00", pricing: { input: 0.66, output: 1.98, cached: 0.022 } }, + { start: "04:00", end: "06:00", pricing: { input: 0.66, output: 1.98, cached: 0.022 } }, ], }); });