Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 9 additions & 6 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@
"@hono/node-server": "^2.0.4",
"@huggingface/transformers": "^4.1.0",
"@octokit/rest": "^22.0.1",
"aieo": "^0.1.41",
"aieo": "^0.2.0",
"hono": "^4.12.23",
"html-to-text": "^10.0.1",
"js-yaml": "^4.1.1",
Expand All @@ -44,23 +44,26 @@
"ws": "^8.21.3"
},
"devDependencies": {
"@ai-sdk/anthropic": "^3.0.92",
"@ai-sdk/anthropic": "^4.0.56",
"@types/html-to-text": "^9.0.4",
"@types/js-yaml": "^4.0.9",
"@types/node": "^22.19.0",
"@types/uuid": "^10.0.0",
"@types/ws": "^8.18.1",
"ai": "^6.0.196",
"ai": "^7.0.105",
"tsx": "^4.21.0",
"typescript": "^5.9.0",
"zod": "^4.3.6"
"zod": "^4.6.5"
},
"peerDependencies": {
"@ai-sdk/anthropic": "^3.0.92",
"ai": "^6.0.175",
"@ai-sdk/anthropic": "^4.0.56",
"ai": "^7.0.105",
"zod": "^4.3.6"
},
"optionalDependencies": {
"sherpa-onnx-node": "^1.13.7"
},
"engines": {
"node": ">=22"
}
}
15 changes: 8 additions & 7 deletions src/createStrut.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1612,7 +1612,7 @@ export async function createStrut<TServices = unknown>(
});

try {
const { ToolLoopAgent, stepCountIs } = await import("ai");
const { ToolLoopAgent, isStepCount } = await import("ai");
const { buildTools, buildSystem } = await import("./ai/index.js");

// The chat's model (`ChatMeta.model`, set by POST /chat) or the
Expand Down Expand Up @@ -1730,8 +1730,8 @@ export async function createStrut<TServices = unknown>(
instructions: await buildSystem(deps),
tools: buildTools(deps),
maxOutputTokens: llm.maxOutputTokens,
stopWhen: stepCountIs(chatMaxSteps),
onFinish: () => {
stopWhen: isStepCount(chatMaxSteps),
onEnd: () => {
registry = deps.registry;
},
});
Expand All @@ -1740,7 +1740,7 @@ export async function createStrut<TServices = unknown>(

const result = await agent.stream({
messages: modelMessages,
onStepFinish: (step) => {
onStepEnd: (step) => {
const u = step.usage;
console.log(
`[chat ${chatId}] turn ${turn} step ${step.stepNumber} finish=${step.finishReason} tokens=in:${u?.inputTokens ?? "?"}/out:${u?.outputTokens ?? "?"}`,
Expand All @@ -1755,7 +1755,7 @@ export async function createStrut<TServices = unknown>(
},
});

for await (const part of result.fullStream) {
for await (const part of result.stream) {
switch (part.type) {
case "text-delta":
if (part.text) await emit({ type: "text-delta", delta: part.text });
Expand Down Expand Up @@ -1793,8 +1793,9 @@ export async function createStrut<TServices = unknown>(
}
}

const resp = await result.response;
await chatStore.appendMessages(chatId, resp.messages as any);
// Every step's messages (tool calls + results), not just the last:
// v7's `response` is final-step only.
await chatStore.appendMessages(chatId, (await result.responseMessages) as any);
await emit({ type: "chat.end" });
await chatStore.setMeta(chatId, { status: "done" });
} catch (err) {
Expand Down
8 changes: 4 additions & 4 deletions src/pricing.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,11 +28,11 @@ export function emptyUsage(): TokenUsage {
const num = (v: unknown): number => (typeof v === "number" && Number.isFinite(v) ? v : 0);

/**
* Normalize a Vercel AI SDK `LanguageModelUsage` (the `.usage` / `.totalUsage`
* on a generate result) into a flat {@link TokenUsage}. Prefers the v6
* Normalize a Vercel AI SDK `LanguageModelUsage` (the `.usage` on a generate
* result — all steps in v7) into a flat {@link TokenUsage}. Prefers the
* `inputTokenDetails` breakdown (noCache / cacheRead / cacheWrite); falls back
* to the flat `inputTokens` + deprecated `cachedInputTokens` when details are
* absent (other providers), treating the remainder as non-cached input.
* to the flat `inputTokens` + pre-v7 `cachedInputTokens` when details are
* absent, treating the remainder as non-cached input.
*/
export function usageFromResult(usage: unknown): TokenUsage {
if (!usage || typeof usage !== "object") return emptyUsage();
Expand Down
7 changes: 6 additions & 1 deletion src/steps/core/agent.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -836,7 +836,12 @@ describe("mid-stream socket death is resumed, not lost", () => {
});
process.env["ANTHROPIC_BASE_URL"] = `http://127.0.0.1:${s.port}`;
try {
await assert.rejects(run(), /terminated/);
// v7 wraps the socket death (APICallError → TypeError: terminated), so
// match anywhere on the cause chain.
await assert.rejects(run(), (e: any) => {
for (let c = e; c; c = c.cause) if (/terminated/.test(String(c.message))) return true;
return false;
});
// 1 good call + the first sever + MAX_STREAM_ERROR_CONTINUATIONS resumes.
assert.equal(s.calls(), 7);
} finally {
Expand Down
68 changes: 36 additions & 32 deletions src/steps/core/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -732,7 +732,7 @@ export default defineStep({
}),
output: z.any(),
async run(cfg, ctx) {
const { ToolLoopAgent, Output, tool, stepCountIs, hasToolCall, jsonSchema, streamText } = await import("ai");
const { ToolLoopAgent, Output, tool, isStepCount, hasToolCall, jsonSchema, streamText } = await import("ai");

// Model/provider resolution via aieo (shared with mcp) through strut's
// resolver (src/llm.ts): friendly aliases ("sonnet", "grok"), canonical
Expand Down Expand Up @@ -921,8 +921,8 @@ export default defineStep({
wrapToolsWithEmit(tools, ctx);

const stopWhen = !useSchema && cfg.finalAnswer
? [hasToolCall("final_answer"), stepCountIs(cfg.maxSteps)]
: [stepCountIs(cfg.maxSteps)];
? [hasToolCall("final_answer"), isStepCount(cfg.maxSteps)]
: [isStepCount(cfg.maxSteps)];

// Resolved LAST (after all config validation): the key lookup throws when
// no key is configured — a config error should surface before a
Expand Down Expand Up @@ -961,16 +961,17 @@ export default defineStep({
wrapToolsWithEmit(webTools, ctx);
Object.assign(tools, webTools);
// Steps that completed BEFORE a mid-stream failure are unreachable through
// the stream's result promises — `steps`, `response`, `totalUsage` and
// the stream's result promises — `steps`, `responseMessages`, `usage` and
// `text` all reject with the stream error — so the only way to keep that
// work is to bank each step as it finishes. Cumulative across resume
// attempts, which is exactly what a continuation needs to replay.
const bankedSteps: any[] = [];
const bankedMessages: any[] = [];
let bankedUsage = emptyUsage();
// Shared by the main loop and the premature-stop nudge continuation.
const onStepFinish = (sf: any) => {
const onStepEnd = (sf: any) => {
bankedSteps.push(sf);
// Per-step in v7 (v6 made these cumulative, so banking them duplicated history).
bankedMessages.push(...((sf.response?.messages ?? []) as any[]));
bankedUsage = addUsage(bankedUsage, usageFromResult(sf.usage));
// A length finish means the generation was TRUNCATED at the output
Expand Down Expand Up @@ -1006,7 +1007,7 @@ export default defineStep({
...(providerOptions ? { providerOptions } : {}),
...(useSchema ? { output: Output.object({ schema: jsonSchema(cfg.schema) }) } : {}),
prepareStep,
onStepFinish,
onStepEnd,
});

const preamble = buildPreamble(cfg.cwd);
Expand All @@ -1026,9 +1027,8 @@ export default defineStep({
let streamErrorContinuations = 0;
let res!: {
steps: any;
response: any;
totalUsage: any;
usage: unknown;
responseMessages: any[];
usage: any;
text: any;
output: any;
};
Expand All @@ -1044,18 +1044,18 @@ export default defineStep({
maxOutputTokens,
stopWhen:
!useSchema && cfg.finalAnswer
? [hasToolCall("final_answer"), stepCountIs(remaining)]
: [stepCountIs(remaining)],
? [hasToolCall("final_answer"), isStepCount(remaining)]
: [isStepCount(remaining)],
...(providerOptions ? { providerOptions } : {}),
...(useSchema ? { output: Output.object({ schema: jsonSchema(cfg.schema) }) } : {}),
prepareStep,
onStepFinish,
onStepEnd,
})
: agent;
const attempt = resuming
? await runner.stream({
messages: [
// response.messages holds only generated turns, so the task
// responseMessages holds only generated turns, so the task
// itself has to lead the replay.
{ role: "user", content: basePrompt },
...(bankedMessages as any[]),
Expand All @@ -1068,9 +1068,9 @@ export default defineStep({
if (!streamError) {
res = {
steps: await attempt.steps,
response: await attempt.response,
totalUsage: await attempt.totalUsage,
usage: undefined as unknown,
// v7: `response` is final-step only; `responseMessages` spans every step.
responseMessages: await attempt.responseMessages,
usage: await attempt.usage,
text: await attempt.text,
output: useSchema ? await (attempt as any).output : undefined,
};
Expand All @@ -1084,10 +1084,14 @@ export default defineStep({
throw streamError;
}
streamErrorContinuations++;
// v7 wraps the socket fault ("Failed to process successful response");
// the root cause is the useful part of the log line.
let rootCause: any = streamError;
while (rootCause?.cause) rootCause = rootCause.cause;
console.warn(
`[agent] stream severed after ${bankedSteps.length} banked step(s) (${
(streamError as Error).message
}); resuming ${streamErrorContinuations}/${MAX_STREAM_ERROR_CONTINUATIONS}.`,
}${rootCause !== streamError ? `: ${rootCause?.message ?? rootCause}` : ""}); resuming ${streamErrorContinuations}/${MAX_STREAM_ERROR_CONTINUATIONS}.`,
);
}
// A resumed run's final attempt only knows its own segment — the banked
Expand All @@ -1101,15 +1105,15 @@ export default defineStep({
// The full session is HUGE and the runner persists every step's output, so we
// only include it when explicitly asked (a future fork/sub-agent). Off by
// default keeps the explore step's persisted output to `{ result, steps, … }`.
const messages = resumedFromStreamError ? bankedMessages : (res.response?.messages ?? []);
const messages = resumedFromStreamError ? bankedMessages : (res.responseMessages ?? []);
const maybeMessages = cfg.returnMessages ? { messages } : {};

// Token usage + cost across the WHOLE agent loop (totalUsage aggregates every
// step; fall back to the final-step usage). `provider` drives the rate table.
// Token usage + cost across the WHOLE agent loop (v7 `usage` aggregates every
// step). `provider` drives the rate table.
// Mutable so a forced final-answer turn (below) can be folded in.
let usage = resumedFromStreamError
? bankedUsage
: usageFromResult(res.totalUsage ?? res.usage);
: usageFromResult(res.usage);
let cost = costOf(usage);
console.log(
`[agent] tokens in:${usage.inputTokens} cacheRead:${usage.cacheReadTokens} cacheWrite:${usage.cacheWriteTokens} out:${usage.outputTokens} → $${cost.toFixed(4)}`,
Expand Down Expand Up @@ -1138,15 +1142,15 @@ export default defineStep({
tools,
maxOutputTokens,
// At least a few turns even when the stop came near the cap.
stopWhen: [stepCountIs(Math.max(4, cfg.maxSteps - stepsUsed))],
stopWhen: [isStepCount(Math.max(4, cfg.maxSteps - stepsUsed))],
...(providerOptions ? { providerOptions } : {}),
output: Output.object({ schema: jsonSchema(cfg.schema) }),
prepareStep,
onStepFinish,
onStepEnd,
});
const nudged = await nudger.stream({
messages: [
// response.messages holds only generated turns — the task leads.
// responseMessages holds only generated turns — the task leads.
{ role: "user", content: basePrompt },
...(messages as any[]),
{
Expand All @@ -1165,8 +1169,8 @@ export default defineStep({
if (nudgeError) throw nudgeError;
const nudgedSteps = (await nudged.steps) ?? [];
stepsUsed += nudgedSteps.length;
messages.push(...(((await nudged.response)?.messages ?? []) as any[]));
const nu = usageFromResult(await nudged.totalUsage);
messages.push(...(((await nudged.responseMessages) ?? []) as any[]));
const nu = usageFromResult(await nudged.usage);
usage = addUsage(usage, nu);
cost += costOf(nu);
// The continuation may have done real work (a publish) before
Expand Down Expand Up @@ -1231,15 +1235,15 @@ export default defineStep({
hasToolCall("final_answer"),
// At least a few turns even when the stop came near the cap —
// finishing file work takes more than one call.
stepCountIs(Math.max(4, cfg.maxSteps - stepsUsed)),
isStepCount(Math.max(4, cfg.maxSteps - stepsUsed)),
],
...(providerOptions ? { providerOptions } : {}),
prepareStep,
onStepFinish,
onStepEnd,
});
const nudged = await nudger.stream({
messages: [
// The session's own user prompt first — response.messages holds
// The session's own user prompt first — responseMessages holds
// only the generated turns, and the continuation needs the task.
{ role: "user", content: preamble ? `${preamble}\n\n${cfg.prompt}` : cfg.prompt },
...(messages as any[]),
Expand All @@ -1258,14 +1262,14 @@ export default defineStep({
if (nudgeError) throw nudgeError;
const nudgedSteps = (await nudged.steps) ?? [];
stepsUsed += nudgedSteps.length;
messages.push(...(((await nudged.response)?.messages ?? []) as any[]));
messages.push(...(((await nudged.responseMessages) ?? []) as any[]));
for (const step of nudgedSteps) {
for (const item of step.content) {
if (item.type === "text" && item.text?.trim()) lastText = item.text.trim();
}
}
final = extractFinal(nudgedSteps);
const nu = usageFromResult(await nudged.totalUsage);
const nu = usageFromResult(await nudged.usage);
usage = addUsage(usage, nu);
cost += costOf(nu);
} catch (e) {
Expand Down Expand Up @@ -1300,7 +1304,7 @@ export default defineStep({
const ft = ((await forced.text) ?? "").trim();
if (ft) {
final = ft;
const fu = usageFromResult(await forced.totalUsage);
const fu = usageFromResult(await forced.usage);
usage = addUsage(usage, fu);
cost += costOf(fu);
}
Expand Down
2 changes: 1 addition & 1 deletion src/steps/core/llm.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import { jsonSchema } from "ai";
import llm, { toSdkSchema } from "./llm.js";

// OFFLINE: no model call. This covers the normalization a workflow's `schema:`
// goes through before generateObject. The bug it pins: a YAML workflow can only
// goes through before Output.object. The bug it pins: a YAML workflow can only
// write a plain JSON Schema object, and handed one bare the SDK assumes a lazy
// thunk and throws "schema is not a function".
describe("llm step: toSdkSchema", () => {
Expand Down
8 changes: 4 additions & 4 deletions src/steps/core/llm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ export default defineStep({
output: z.any(),
async run(cfg, ctx?: StepContext<StrutCapabilities>) {
// Dynamic import to avoid hard dependency if not using LLM steps
const { generateText, generateObject } = await import("ai");
const { generateText, Output } = await import("ai");

// Provider/model/key via the shared resolver (src/llm.ts → aieo): the key
// comes through the secrets boundary (secret store → env). The output
Expand All @@ -65,13 +65,13 @@ export default defineStep({

if (cfg.schema) {
// Structured output
const result = await generateObject({
const result = await generateText({
model,
prompt: cfg.prompt,
schema: (await toSdkSchema(cfg.schema)) as any,
output: Output.object({ schema: (await toSdkSchema(cfg.schema)) as any }),
maxOutputTokens,
});
return result.object;
return result.output;
} else {
// Free-form text
const result = await generateText({
Expand Down
Loading
Loading