diff --git a/packages/agent-core-v2/docs/en/llm.md b/packages/agent-core-v2/docs/en/llm.md index 70e38f00c32..24240c7e858 100644 --- a/packages/agent-core-v2/docs/en/llm.md +++ b/packages/agent-core-v2/docs/en/llm.md @@ -5,10 +5,10 @@ llm is a standalone LLM request library inside the human layer (`src/human/llm/` ## Design Principles 1. **Minimal boundary: llm = "a single request"**. llm only handles request encoding/decoding and event emission. auth, usage accounting, HistoryMessage/meta, compaction, switch, the media file system, and Tool Message assembly are all out of scope — they either move up to the turn/agent layer or plug in as contribution points. -2. **Streaming-native; events are the contract**. The only outward surface is a single, purely serializable event stream (requester level: `llm.sent / streaming.headers / streaming.part / streaming.usage / streaming.finish / streaming.message_id / failed.syntax / failed.remote / done`; the machine level adds `llm.retrying / llm.recovering`, and `llm.sent` carries the most recent recovery record). Streaming and non-streaming are isomorphic (non-streaming also accumulates over the stream, just without deltas). Events are emitted as they arrive — no caching, no fallback. -3. **format masks inter-protocol differences; traits express provider customizations**. format lives at the protocol layer and handles encoding/decoding of requests, responses, errors, usage, and finish. Each protocol owns a typed trait interface (`OpenAITrait` / `OpenAIResponsesTrait` / `AnthropicTrait` / `GoogleGenAITrait`) exposing only the customization points that protocol actually consumes — a hook a protocol ignores is unrepresentable, never silently dead. format and trait never import each other: both speak only the neutral wire/chunk types in the protocol's `contract.ts`. The requester is the composition root — `generate` runs a fixed per-protocol pipeline (`planOpenAIRequest` and friends) that alternates pure format stages (lower → assemble → encode → stream parser) with trait hooks (cacheKey/thinking → convertMessage → mergeHistory → convertTool → buildParams → extractUsage), so customization is explicit data flow instead of a closure captured inside format. Endpoint/env resolution and default headers form the provider `connection`, error classification is a requester option, and model capability is a provider-variant field — none of them are format business. Each base's public seam is contract + trait + requester; format, lower, and patterns are internal to the requester pipeline — only bases code and tests may import them (lint-enforced). Protocol differences must not leak into the machine or into requester decorators. +2. **Streaming-native; events are the contract**. The only outward surface is a single, purely serializable event stream (requester level: `llm.sent / streaming.headers / streaming.part / streaming.usage / streaming.finish / streaming.message_id / failed.syntax / failed.remote / done`; the turn level adds `llm.retrying / llm.recovering`, and `llm.sent` carries the most recent recovery record). Streaming and non-streaming are isomorphic (non-streaming also accumulates over the stream, just without deltas). Events are emitted as they arrive — no caching, no fallback. +3. **format masks inter-protocol differences; traits express provider customizations**. format lives at the protocol layer and handles encoding/decoding of requests, responses, errors, usage, and finish. Each protocol owns a typed trait interface (`OpenAITrait` / `OpenAIResponsesTrait` / `AnthropicTrait` / `GoogleGenAITrait`) exposing only the customization points that protocol actually consumes — a hook a protocol ignores is unrepresentable, never silently dead. format and trait never import each other: both speak only the neutral wire/chunk types in the protocol's `contract.ts`. The requester is the composition root — `generate` runs a fixed per-protocol pipeline (`planOpenAIRequest` and friends) that alternates pure format stages (lower → assemble → encode → stream parser) with trait hooks (cacheKey/thinking → convertMessage → mergeHistory → convertTool → buildParams → extractUsage), so customization is explicit data flow instead of a closure captured inside format. Endpoint/env resolution and default headers form the provider `connection`, error classification is a requester option, and model capability is a provider-variant field — none of them are format business. Each base's public seam is contract + trait + requester; format, lower, and patterns are internal to the requester pipeline — only bases code and tests may import them (lint-enforced). Protocol differences must not leak into the turn or into requester decorators. 4. **Two-layer error model**. Internally, code throws the SDK's native errors; local request validation throws the shared `SyntaxRequestFormatError` (`llm/syntax-errors.ts`), which the requester converts uniformly via `toLlmSyntaxErrorMessage`, with no intermediate layer. Externally there are only `llm.failed.syntax` (local message syntax errors, never retried) and `llm.failed.remote` (remote streaming errors, subdivided into connection / timeout / rate_limit / quota_exhausted / context_overflow / request_structure, etc.), converted by format at the boundary. -5. **Stateless core + state machine shell**. `generate(config, content, control)` is a stateless function; errors are delivered via onEvent, never thrown. The llm machine wraps a single request (messageResolvers, abort scope, event forwarding) and drives retry and recovery through the pure policy functions in retry.ts / recovery.ts: recovery re-sends with replacement messages produced by the pure `propose` function (attempt resets to 1), retry backs off in the `retrying` state (honoring Retry-After), and the machine emits `llm.recovering / llm.retrying` for each. Empty response is judged by `withEmptyResponseGuard` at the requester boundary and raised as `llm.failed.remote`, entering the same retry path. Abort is carried by an AbortController owned by the turn: the controller is passed into the machine and the request actor via `LlmInput.signal`, and the turn aborts it directly on `turn.abort`, with the request ending as `llm.failed.remote`; the request actor neither creates its own controller nor touches any signal on teardown, so a finished request can never abort a shared signal. The accumulator is held by the turn and fed by the event stream; on `llm.retrying / llm.recovering` the turn rolls it back and recreates it, so every attempt accumulates from zero while as much interrupted state as possible is preserved (the turn finishes the complete message out of the accumulator at `llm.done`). +5. **Stateless core + turn-driven orchestration**. `generate(config, content, control)` is a stateless function; errors are delivered via onEvent, never thrown. The turn machine invokes the request actor (`createRequestActor`) directly: the actor wraps a single request (messageResolvers, abort scope, event sendBack), and the turn drives retry and recovery through the pure policy functions in retry.ts / recovery.ts: recovery re-sends with replacement messages produced by the pure `propose` function (attempt resets to 1), retry backs off in the `retrying` state (honoring Retry-After), and the turn emits `llm.recovering / llm.retrying` for each. Empty response is judged by `withEmptyResponseGuard` at the requester boundary and raised as `llm.failed.remote`, entering the same retry path. Abort is carried by an AbortController owned by the turn: the controller is passed into the request actor via `LlmInput.signal`, and the turn aborts it directly on `turn.abort`, with the request ending as `llm.failed.remote`; the request actor neither creates its own controller nor touches any signal on teardown, so a finished request can never abort a shared signal. The accumulator is held by the turn and fed by the event stream; on `llm.retrying / llm.recovering` the turn rolls it back and recreates it, so every attempt accumulates from zero while as much interrupted state as possible is preserved (the turn finishes the complete message out of the accumulator at `llm.done`). 6. **No silent fallback**. Configuration is taken exactly as given. For beta features, thinking, empty response, and similar scenarios, define explicit error conditions first, fail at request time, and guide the user to fix the configuration — never fall back silently. 7. **Every variable capability is a contribution point**. Providers, media upload/degradation, usage, traceId, and error recovery (compaction / media degradation) all plug in through extension points; the llm core contains none of these concepts. 8. **Data is data**. A model is pure, function-free data (endpoint url + model uniquely identifies a model), serializable and directly usable as generate input. The catalog is a derived `provider -> models` cache; the dependency direction only goes from models-dev into llm internals, never the reverse. @@ -34,9 +34,9 @@ llm/ ├── requester/ │ ├── requester.ts LlmRequester.generate(config, content, control); │ │ ExtraParams typed per protocol {openai?, responses?, anthropic?, googleGenai?} -│ ├── machine.ts llm state machine (single request + retry/recovery + empty response -│ │ judgment; emits llm.retrying / llm.recovering) -│ ├── retry.ts / recovery.ts pure retry/recovery policy functions (driven by the llm machine; propose is pure) +│ ├── actor.ts request actor: a fromCallback wrapping a single request +│ │ (messageResolvers, abort scope, event sendBack); invoked by the turn +│ ├── retry.ts / recovery.ts pure retry/recovery policy functions (driven by the turn machine; propose is pure) │ ├── empty-response.ts withEmptyResponseGuard: judges empty responses at finish and raises llm.failed.remote │ └── bases/ four protocol bases: openai / openai-responses / anthropic / google-genai │ each with contract / format / lower / patterns / capability / extra-params / trait / requester @@ -53,11 +53,12 @@ llm/ └── media/ media contribution points: cache / degrade / ref / resolver / store / upload ``` -Request lifecycle: `generate` receives (config, content, control) → the requester's `plan*` function composes pure format stages with trait hooks into protocol requestParams (format lowers the generic Message[] through the Pattern Rewriter; trait adjusts kwargs, converted messages, history, tools, and final params in between) → internalGenerate calls the official SDK → streaming chunks are converted by the stateless parser callbacks into `llm.streaming.part / streaming.usage / streaming.finish / streaming.message_id` events → errors are converted by format into `llm.failed.*`; on success the requester emits `llm.done`, on failure it ends with `llm.failed.syntax / llm.failed.remote` and never emits `llm.done`. `withEmptyResponseGuard` judges empty responses at finish and raises `llm.failed.remote`; the llm machine first tries recovery on `llm.failed.remote` (replacement messages from the pure `propose` function, emitting `llm.recovering`), then retries with backoff (honoring Retry-After, emitting `llm.retrying`), and only lands in the failed final state once attempts are exhausted. The upper-layer turn holds the HistoryAccumulator, fed by the event stream, rolls it back and recreates it on `llm.retrying / llm.recovering`, and finishes the complete message at `llm.done`; usage accounting, tracing, compaction, and media degradation all attach to the event stream as plugins/contribution points. +Request lifecycle: `generate` receives (config, content, control) → the requester's `plan*` function composes pure format stages with trait hooks into protocol requestParams (format lowers the generic Message[] through the Pattern Rewriter; trait adjusts kwargs, converted messages, history, tools, and final params in between) → internalGenerate calls the official SDK → streaming chunks are converted by the stateless parser callbacks into `llm.streaming.part / streaming.usage / streaming.finish / streaming.message_id` events → errors are converted by format into `llm.failed.*`; on success the requester emits `llm.done`, on failure it ends with `llm.failed.syntax / llm.failed.remote` and never emits `llm.done`. At `llm.done` the turn judges empty responses via `emptyResponseError` and re-raises them as `llm.failed.remote`; the turn machine first tries recovery on `llm.failed.remote` (replacement messages from the pure `propose` function, emitting `llm.recovering`), then retries with backoff (honoring Retry-After, emitting `llm.retrying`), and only fails the turn once attempts are exhausted. The turn holds the HistoryAccumulator, fed by the event stream, rolls it back and recreates it on `llm.retrying / llm.recovering`, and finishes the complete message at `llm.done`; usage accounting, tracing, compaction, and media degradation all attach to the event stream as plugins/contribution points. ## Rejected Schemes (do not reintroduce) -- Splitting llmActor / llmStreamActor into two actors — a single machine; non-streaming also accumulates over the stream. +- Splitting the request actor into llmActor / llmStreamActor — one actor per request; non-streaming also accumulates over the stream. +- A dedicated llm state machine wrapping the request actor — the turn machine invokes the actor directly and owns retry/recovery; the extra machine layer carried no state anyone consumed. - DDD domain-method wrapping (Generation Domain, etc.) — use the format/trait/provider layering instead. - A single cross-protocol trait bag holding every vendor hook (the old ProtocolTrait) — per-protocol typed traits, composed by the requester's request pipeline. - Binding the trait into the format (a `createOpenAIFormat(trait)` closure, or trait hooks passed as formatRequest options) — the requester pipeline alternates format stages and trait hooks explicitly; the two sides only share the neutral `contract.ts` types. diff --git a/packages/agent-core-v2/docs/zh/llm.md b/packages/agent-core-v2/docs/zh/llm.md index 97ce6dc78c4..344e3acf71c 100644 --- a/packages/agent-core-v2/docs/zh/llm.md +++ b/packages/agent-core-v2/docs/zh/llm.md @@ -5,10 +5,10 @@ llm 是 human 层内一个独立的 LLM 请求库(`src/human/llm/`),提供 ## 设计原则 1. **边界极简:llm = 「一次请求」**。llm 只负责请求编解码与事件回传。auth、usage 统计、HistoryMessage/meta、compaction、switch、媒体文件系统、Tool Message 拼装全部不属于 llm——要么上移到 turn/agent 层,要么以贡献点接入。 -2. **流式原生、事件即契约**。对外只暴露一条纯可序列化的事件流(requester 层:`llm.sent / streaming.headers / streaming.part / streaming.usage / streaming.finish / streaming.message_id / failed.syntax / failed.remote / done`;machine 层补充 `llm.retrying / llm.recovering`,`llm.sent` 携带最近一次 recovery 记录),流式与非流式同构(非流式也走流式累积,只是不发 delta);事件收到即发,不缓存、不兜底。 -3. **format 屏蔽协议间差异,trait 表达 provider 定制**。format 位于 protocol 层,负责请求、响应、错误、usage 和 finish 的编解码。每种协议拥有自己的类型化 trait 接口(`OpenAITrait` / `OpenAIResponsesTrait` / `AnthropicTrait` / `GoogleGenAITrait`),只暴露该协议实际消费的定制点——协议不支持的 hook 在类型上无法表达,而不是配了却静默无效。format 与 trait 互不 import:双方只共享协议 `contract.ts` 里的中立 wire/chunk 类型。requester 是组合根——`generate` 执行每个协议固定的流水线(`planOpenAIRequest` 等),交替调用纯 format 阶段(lower → assemble → encode → stream parser)与 trait hooks(cacheKey/thinking → convertMessage → mergeHistory → convertTool → buildParams → extractUsage),定制逻辑是显式的数据流,而不是捕获在 format 闭包里。endpoint/环境变量解析与默认 headers 属于 provider `connection`,错误归类是 requester 选项,模型能力是 provider variant 字段——都不是 format 的职责。每个 base 的公开接缝是 contract + trait + requester;format、lower、patterns 是 requester 流水线的内部模块——只有 bases 内代码和测试可以 import(lint 强制)。协议差异不允许泄漏到 machine 或 requester 的装饰层。 +2. **流式原生、事件即契约**。对外只暴露一条纯可序列化的事件流(requester 层:`llm.sent / streaming.headers / streaming.part / streaming.usage / streaming.finish / streaming.message_id / failed.syntax / failed.remote / done`;turn 层补充 `llm.retrying / llm.recovering`,`llm.sent` 携带最近一次 recovery 记录),流式与非流式同构(非流式也走流式累积,只是不发 delta);事件收到即发,不缓存、不兜底。 +3. **format 屏蔽协议间差异,trait 表达 provider 定制**。format 位于 protocol 层,负责请求、响应、错误、usage 和 finish 的编解码。每种协议拥有自己的类型化 trait 接口(`OpenAITrait` / `OpenAIResponsesTrait` / `AnthropicTrait` / `GoogleGenAITrait`),只暴露该协议实际消费的定制点——协议不支持的 hook 在类型上无法表达,而不是配了却静默无效。format 与 trait 互不 import:双方只共享协议 `contract.ts` 里的中立 wire/chunk 类型。requester 是组合根——`generate` 执行每个协议固定的流水线(`planOpenAIRequest` 等),交替调用纯 format 阶段(lower → assemble → encode → stream parser)与 trait hooks(cacheKey/thinking → convertMessage → mergeHistory → convertTool → buildParams → extractUsage),定制逻辑是显式的数据流,而不是捕获在 format 闭包里。endpoint/环境变量解析与默认 headers 属于 provider `connection`,错误归类是 requester 选项,模型能力是 provider variant 字段——都不是 format 的职责。每个 base 的公开接缝是 contract + trait + requester;format、lower、patterns 是 requester 流水线的内部模块——只有 bases 内代码和测试可以 import(lint 强制)。协议差异不允许泄漏到 turn 或 requester 的装饰层。 4. **错误两层模型**。内部 throw SDK 原生错误;本地请求校验抛共享的 `SyntaxRequestFormatError`(`llm/syntax-errors.ts`),由 requester 经 `toLlmSyntaxErrorMessage` 统一转换,不加中间层。对外只有 `llm.failed.syntax`(本地消息语法错误,不重试)与 `llm.failed.remote`(远程流式错误,细分为 connection/timeout/rate_limit/quota_exhausted/context_overflow/request_structure 等),由 format 在边界完成转换。 -5. **无状态内核 + 状态机外壳**。`generate(config, content, control)` 是无状态函数,错误走 onEvent 不 throw;llm machine 包装单次请求(messageResolvers、abort 作用域、事件转发),并借助 retry.ts / recovery.ts 的纯策略函数驱动重试与 recovery:recovery 由纯函数 `propose` 产出替换消息直接重发(attempt 重置为 1),重试走 `retrying` 状态的 backoff(尊重 Retry-After),两者分别对外补发 `llm.recovering / llm.retrying` 事件;empty response 由 `withEmptyResponseGuard` 在 requester 边界判定并转为 `llm.failed.remote`,进入同一重试路径;abort 由 turn 持有的 AbortController 承载:controller 经 `LlmInput.signal` 传入 machine 与 request actor,turn 在 `turn.abort` 时直接 abort 它,请求随即以 `llm.failed.remote` 收尾;request actor 不自建 controller、回收时不触碰任何 signal,正常完成的请求绝不可能误 abort 共享 signal。累积器由 turn 持有并随事件流喂入,在 `llm.retrying / llm.recovering` 时 rollback 并重建,每次 attempt 从零累积,从而尽可能保留中断现场(turn 在 `llm.done` 时从累加器 finish 出完整消息)。 +5. **无状态内核 + turn 驱动的编排**。`generate(config, content, control)` 是无状态函数,错误走 onEvent 不 throw;turn machine 直接 invoke 请求 actor(`createRequestActor`):actor 包装单次请求(messageResolvers、abort 作用域、事件 sendBack),turn 借助 retry.ts / recovery.ts 的纯策略函数驱动重试与 recovery:recovery 由纯函数 `propose` 产出替换消息直接重发(attempt 重置为 1),重试走 `retrying` 状态的 backoff(尊重 Retry-After),两者分别由 turn 对外补发 `llm.recovering / llm.retrying` 事件;empty response 由 `withEmptyResponseGuard` 在 requester 边界判定并转为 `llm.failed.remote`,进入同一重试路径;abort 由 turn 持有的 AbortController 承载:controller 经 `LlmInput.signal` 传入 request actor,turn 在 `turn.abort` 时直接 abort 它,请求随即以 `llm.failed.remote` 收尾;request actor 不自建 controller、回收时不触碰任何 signal,正常完成的请求绝不可能误 abort 共享 signal。累积器由 turn 持有并随事件流喂入,在 `llm.retrying / llm.recovering` 时 rollback 并重建,每次 attempt 从零累积,从而尽可能保留中断现场(turn 在 `llm.done` 时从累加器 finish 出完整消息)。 6. **不兜底**。配置是什么就是什么;beta 特性、thinking、empty response 等场景先定义明确报错条件,在请求阶段报错并引导用户修正,而不是静默兜底。 7. **一切可变能力都是贡献点**。provider、媒体上传/降级、usage、traceId、错误恢复(compaction/媒体降级)都通过扩展点接入,llm 内核不含这些概念。 8. **数据即数据**。model 是无函数的纯数据(endpoint url + model 唯一标识一个模型),可序列化、可直接作为 generate 输入;catalog 是 `provider -> models` 的派生缓存,依赖方向只能从 models-dev 指向 llm 内部,不能反向依赖。 @@ -34,9 +34,9 @@ llm/ ├── requester/ │ ├── requester.ts LlmRequester.generate(config, content, control); │ │ ExtraParams 按协议带类型 {openai?, responses?, anthropic?, googleGenai?} -│ ├── machine.ts llm 状态机(单次请求 + 重试/恢复 + empty response 判定; -│ │ 对外补发 llm.retrying / llm.recovering) -│ ├── retry.ts / recovery.ts 重试/恢复策略纯函数(由 llm machine 驱动;propose 为纯函数) +│ ├── actor.ts 请求 actor:包装单次请求的 fromCallback +│ │ (messageResolvers、abort 作用域、事件 sendBack),由 turn invoke +│ ├── retry.ts / recovery.ts 重试/恢复策略纯函数(由 turn machine 驱动;propose 为纯函数) │ ├── empty-response.ts withEmptyResponseGuard:finish 时判定空响应并转为 llm.failed.remote │ └── bases/ 四个协议基座:openai / openai-responses / anthropic / google-genai │ 各自含 contract / format / lower / patterns / capability / extra-params / trait / requester @@ -53,11 +53,12 @@ llm/ └── media/ 媒体贡献点:cache / degrade / ref / resolver / store / upload ``` -请求生命周期:`generate` 收到 (config, content, control) → requester 的 `plan*` 函数将纯 format 阶段与 trait hooks 组合为协议 requestParams(format 将通用 Message[] 经 Pattern Rewriter 降低,trait 在其间调整 kwargs、转换消息、合并历史、转换 tools 并收尾 params)→ internalGenerate 调用官方 SDK → 流式 chunk 经无状态 parser 回调转换为 `llm.streaming.part / streaming.usage / streaming.finish / streaming.message_id` 事件 → 错误由 format 转换为 `llm.failed.*`;成功时 requester 发出 `llm.done`,失败时以 `llm.failed.syntax / llm.failed.remote` 收尾、不再发 `llm.done`。`withEmptyResponseGuard` 在 finish 时判定空响应并转为 `llm.failed.remote`;llm machine 对 `llm.failed.remote` 先尝试 recovery(纯函数 `propose` 产出替换消息,发 `llm.recovering`),再按策略 backoff 重试(尊重 Retry-After,发 `llm.retrying`),耗尽后才以 failed 终态收尾。上层的 turn 持有 HistoryAccumulator 随事件流累积,在 `llm.retrying / llm.recovering` 时 rollback 并重建累加器,`llm.done` 时 finish 出完整消息;usage 统计、trace、compaction、媒体降级均以插件/贡献点身份挂接在事件流上。 +请求生命周期:`generate` 收到 (config, content, control) → requester 的 `plan*` 函数将纯 format 阶段与 trait hooks 组合为协议 requestParams(format 将通用 Message[] 经 Pattern Rewriter 降低,trait 在其间调整 kwargs、转换消息、合并历史、转换 tools 并收尾 params)→ internalGenerate 调用官方 SDK → 流式 chunk 经无状态 parser 回调转换为 `llm.streaming.part / streaming.usage / streaming.finish / streaming.message_id` 事件 → 错误由 format 转换为 `llm.failed.*`;成功时 requester 发出 `llm.done`,失败时以 `llm.failed.syntax / llm.failed.remote` 收尾、不再发 `llm.done`。turn 在 `llm.done` 时经 `emptyResponseError` 判定空响应并转为 `llm.failed.remote`;turn machine 对 `llm.failed.remote` 先尝试 recovery(纯函数 `propose` 产出替换消息,发 `llm.recovering`),再按策略 backoff 重试(尊重 Retry-After,发 `llm.retrying`),耗尽后才将 turn 置为失败。turn 持有 HistoryAccumulator 随事件流累积,在 `llm.retrying / llm.recovering` 时 rollback 并重建累加器,`llm.done` 时 finish 出完整消息;usage 统计、trace、compaction、媒体降级均以插件/贡献点身份挂接在事件流上。 ## 已被否决的方案(不要再引入) -- 拆分 llmActor / llmStreamActor 两个 actor —— 单一 machine,非流式也走流式累积。 +- 拆分 llmActor / llmStreamActor 两个 actor —— 每次请求一个 actor,非流式也走流式累积。 +- 给请求 actor 再包一层专用 llm 状态机 —— turn machine 直接 invoke actor 并持有重试/recovery,额外的 machine 层没有任何被消费的状态。 - DDD 领域方法包装(Generation Domain 等)—— 用 format/trait/provider 分层。 - 用一个跨协议 trait 大包承载所有厂商 hooks(旧的 ProtocolTrait)—— 按协议拆分的类型化 trait,由 requester 的请求流水线组合。 - 把 trait 绑定进 format(`createOpenAIFormat(trait)` 闭包,或把 trait hooks 作为 formatRequest 选项传入)—— requester 流水线显式交替调用 format 阶段与 trait hooks,双方只共享中立的 `contract.ts` 类型。 diff --git a/packages/agent-core-v2/src/agent/loop/machine/engine.ts b/packages/agent-core-v2/src/agent/loop/machine/engine.ts index 1287783b40a..a4f021392ec 100644 --- a/packages/agent-core-v2/src/agent/loop/machine/engine.ts +++ b/packages/agent-core-v2/src/agent/loop/machine/engine.ts @@ -14,7 +14,6 @@ import type { LlmErrorMessage } from '#human/llm/errors'; import type { FinishInfo } from '#human/llm/finish-reason'; import type { StreamedMessagePart, UserMessage } from '#human/llm/message'; import type { LlmModel } from '#human/llm/model'; -import { createLlmMachine } from '#human/llm/requester/machine'; import type { LlmRecovery, LlmRecoveryRecord } from '#human/llm/requester/recovery'; import { resolveMaxAttempts } from '#human/llm/requester/retry'; import type { ToolResult as MachineToolResult, ToolUpdate } from '#human/tool/executor'; @@ -283,15 +282,10 @@ export function createMachineEngine(options: CreateMachineEngineOptions): Machin const actor = createActor( createAgentMachine({ tools: tools.tools, - turnActor: createTurnMachine( - createLlmMachine({ - requester: requester.requester, - }), - { - retry: { maxAttemptsPerStep: options.maxAttemptsPerStep }, - recovery: options.recovery, - }, - ), + turnActor: createTurnMachine(requester.requester, { + retry: { maxAttemptsPerStep: options.maxAttemptsPerStep }, + recovery: options.recovery, + }), abortTimeoutMs: options.abortTimeoutMs, }), { input: { request: { model: options.model, systemPrompt: options.systemPrompt }, store } }, diff --git a/packages/agent-core-v2/src/human/agent/turn.ts b/packages/agent-core-v2/src/human/agent/turn.ts index 630e8b81c14..8b8ed8088e4 100644 --- a/packages/agent-core-v2/src/human/agent/turn.ts +++ b/packages/agent-core-v2/src/human/agent/turn.ts @@ -16,14 +16,14 @@ import { type UserMessage, } from '#/llm/message'; import type { LlmModel } from '#/llm/model'; -import type { createLlmMachine, LlmEvent } from '#/llm/requester/machine'; +import { createRequestActor, type LlmEvent, type MessageResolver } from '#/llm/requester/actor'; import type { LlmRecovery, LlmRecoveryContext, LlmRecoveryProposal, LlmRecoveryRecord, } from '#/llm/requester/recovery'; -import type { LlmRequestConfig } from '#/llm/requester/requester'; +import type { LlmRequestConfig, LlmRequester } from '#/llm/requester/requester'; import { readRetryAfterMs, resolveMaxAttempts, @@ -347,10 +347,11 @@ export interface CreateTurnMachineOptions { readonly recovery?: LlmRecovery; readonly retry?: LlmRetryOptions; readonly abortGraceMs?: number; + readonly messageResolvers?: readonly MessageResolver[]; } export function createTurnMachine( - llmActor: ReturnType, + requester: LlmRequester, options?: CreateTurnMachineOptions, ) { const recovery = options?.recovery; @@ -364,7 +365,7 @@ export function createTurnMachine( output: {} as TurnOutput, }, actors: { - llmActor, + llmActor: createRequestActor(requester, options?.messageResolvers), }, actions: { forwardToParent: ({ self, event }) => { diff --git a/packages/agent-core-v2/src/human/index.ts b/packages/agent-core-v2/src/human/index.ts index d91ecdc0190..39acdd4a2fb 100644 --- a/packages/agent-core-v2/src/human/index.ts +++ b/packages/agent-core-v2/src/human/index.ts @@ -18,7 +18,7 @@ export * from './llm/protocol/patterns'; export * from './llm/media'; export * from './llm/requester/requester'; export * from './llm/empty-response'; -export * from './llm/requester/machine'; +export * from './llm/requester/actor'; export * from './llm/requester/recovery'; export * from './llm/requester/retry'; export * from './llm/requester/bases/openai/contract'; diff --git a/packages/agent-core-v2/src/human/llm/media/resolver.ts b/packages/agent-core-v2/src/human/llm/media/resolver.ts index 73974c1c61e..a3364e207d5 100644 --- a/packages/agent-core-v2/src/human/llm/media/resolver.ts +++ b/packages/agent-core-v2/src/human/llm/media/resolver.ts @@ -1,7 +1,7 @@ import type { ModelCapability } from '#/llm/capability'; import type { ContentPart, Message, VideoURLPart } from '#/llm/message'; import type { Provider } from '#/llm/provider/definition'; -import type { MessageResolveContext, MessageResolver } from '#/llm/requester/machine'; +import type { MessageResolveContext, MessageResolver } from '#/llm/requester/actor'; import type { MediaUploadCache } from './cache'; import { mediaKindForMime, mediaMimeForPath, type MediaKind } from './mime'; diff --git a/packages/agent-core-v2/src/human/llm/requester/actor.ts b/packages/agent-core-v2/src/human/llm/requester/actor.ts new file mode 100644 index 00000000000..d33c7d57652 --- /dev/null +++ b/packages/agent-core-v2/src/human/llm/requester/actor.ts @@ -0,0 +1,78 @@ +import { fromCallback } from '#/xstate2'; + +import type { Message } from '#/llm/message'; +import type { LlmModel } from '#/llm/model'; + +import type { + LlmRequestConfig, + LlmRequestContent, + LlmRequestEvent, + LlmRequester, +} from './requester'; +import type { LlmRecoveryRecord } from './recovery'; + +export interface LlmInput { + readonly config: LlmRequestConfig; + readonly content: LlmRequestContent; + readonly signal: AbortSignal; +} + +export interface MessageResolveContext { + readonly model: LlmModel; + readonly signal: AbortSignal; +} + +export interface MessageResolver { + readonly id: string; + resolve( + messages: readonly Message[], + ctx: MessageResolveContext, + ): Promise; +} + +export type LlmEvent = + | Exclude + | { type: 'llm.sent'; recovery?: LlmRecoveryRecord } + | { + type: 'llm.retrying'; + failedAttempt: number; + nextAttempt: number; + maxAttempts: number; + delayMs: number; + errorName: string; + errorMessage: string; + statusCode?: number; + } + | { + type: 'llm.recovering'; + strategy: string; + action: string; + errorName: string; + errorMessage: string; + statusCode?: number; + }; + +export function createRequestActor( + requester: LlmRequester, + messageResolvers: readonly MessageResolver[] = [], +) { + return fromCallback(({ input, sendBack }) => { + void (async () => { + let messages = input.content.messages; + for (const resolver of messageResolvers) { + messages = await resolver.resolve(messages, { + model: input.config.model, + signal: input.signal, + }); + } + await requester.generate( + input.config, + { ...input.content, messages }, + { + signal: input.signal, + onEvent: sendBack, + }, + ); + })(); + }); +} diff --git a/packages/agent-core-v2/src/human/llm/requester/machine.ts b/packages/agent-core-v2/src/human/llm/requester/machine.ts deleted file mode 100644 index b14af143a0f..00000000000 --- a/packages/agent-core-v2/src/human/llm/requester/machine.ts +++ /dev/null @@ -1,194 +0,0 @@ -import { assign, emit, fromCallback, setup } from '#/xstate2'; - -import type { LlmErrorMessage } from '#/llm/errors'; -import type { Message } from '#/llm/message'; -import type { LlmModel } from '#/llm/model'; - -import type { - LlmRequestConfig, - LlmRequestContent, - LlmRequester, - LlmRequestEvent, -} from './requester'; -import type { LlmRecoveryRecord } from './recovery'; - -export interface LlmInput { - readonly config: LlmRequestConfig; - readonly content: LlmRequestContent; - readonly signal: AbortSignal; -} - -export interface MessageResolveContext { - readonly model: LlmModel; - readonly signal: AbortSignal; -} - -export interface MessageResolver { - readonly id: string; - resolve( - messages: readonly Message[], - ctx: MessageResolveContext, - ): Promise; -} - -export type LlmEvent = - | Exclude - | { type: 'llm.sent'; recovery?: LlmRecoveryRecord } - | { - type: 'llm.retrying'; - failedAttempt: number; - nextAttempt: number; - maxAttempts: number; - delayMs: number; - errorName: string; - errorMessage: string; - statusCode?: number; - } - | { - type: 'llm.recovering'; - strategy: string; - action: string; - errorName: string; - errorMessage: string; - statusCode?: number; - }; - -export type LlmOutput = { type: 'succeeded' } | { type: 'failed'; error: LlmErrorMessage }; - -export interface LlmMachineContext { - input: LlmInput; - outcome?: 'succeeded' | 'failed'; - error?: LlmErrorMessage; -} - -function createRequestActor( - requester: LlmRequester, - messageResolvers: readonly MessageResolver[], -) { - return fromCallback(({ input, sendBack }) => { - void (async () => { - let messages = input.content.messages; - for (const resolver of messageResolvers) { - messages = await resolver.resolve(messages, { - model: input.config.model, - signal: input.signal, - }); - } - await requester.generate( - input.config, - { ...input.content, messages }, - { - signal: input.signal, - onEvent: sendBack, - }, - ); - })(); - }); -} - -export interface CreateLlmMachineOptions { - requester: LlmRequester; - messageResolvers?: readonly MessageResolver[]; -} - -export function createLlmMachine(options: CreateLlmMachineOptions) { - const requestActor = createRequestActor(options.requester, options.messageResolvers ?? []); - return setup({ - types: { - input: {} as LlmInput, - context: {} as LlmMachineContext, - events: {} as LlmEvent, - emitted: {} as LlmEvent, - output: {} as LlmOutput, - }, - actors: { requestActor }, - actions: { - forwardToParent: ({ self, event }) => { - self._parent?.send(event); - }, - sendToParent: ({ self }, params: LlmEvent) => { - self._parent?.send(params); - }, - }, - }).createMachine({ - id: 'llm', - initial: 'generating', - context: ({ input }) => ({ input }), - states: { - generating: { - invoke: { - src: 'requestActor', - input: ({ context }) => context.input, - }, - on: { - 'llm.sent': { - actions: [ - emit({ type: 'llm.sent' as const }), - { type: 'sendToParent', params: { type: 'llm.sent' as const } }, - ], - }, - 'llm.streaming.headers': { - actions: [ - emit(({ event }) => ({ type: 'llm.streaming.headers' as const, headers: event.headers })), - 'forwardToParent', - ], - }, - 'llm.streaming.part': { - actions: [ - emit(({ event }) => ({ type: 'llm.streaming.part' as const, part: event.part })), - 'forwardToParent', - ], - }, - 'llm.streaming.usage': { - actions: [ - emit(({ event }) => ({ type: 'llm.streaming.usage' as const, usage: event.usage })), - 'forwardToParent', - ], - }, - 'llm.streaming.finish': { - actions: [ - emit(({ event }) => ({ type: 'llm.streaming.finish' as const, finish: event.finish })), - 'forwardToParent', - ], - }, - 'llm.streaming.message_id': { - actions: [ - emit(({ event }) => ({ type: 'llm.streaming.message_id' as const, messageId: event.messageId })), - 'forwardToParent', - ], - }, - 'llm.done': { - target: 'succeeded', - actions: [ - assign({ outcome: 'succeeded' as const }), - emit({ type: 'llm.done' as const }), - 'forwardToParent', - ], - }, - 'llm.failed.syntax': { - target: 'failed', - actions: [ - assign({ outcome: 'failed' as const, error: ({ event }) => event.error }), - emit(({ event }) => ({ type: 'llm.failed.syntax' as const, error: event.error })), - 'forwardToParent', - ], - }, - 'llm.failed.remote': { - target: 'failed', - actions: [ - assign({ outcome: 'failed' as const, error: ({ event }) => event.error }), - emit(({ event }) => ({ type: 'llm.failed.remote' as const, error: event.error })), - 'forwardToParent', - ], - }, - }, - }, - succeeded: { type: 'final' }, - failed: { type: 'final' }, - }, - output: ({ context }): LlmOutput => - context.outcome === 'failed' - ? { type: 'failed', error: context.error as LlmErrorMessage } - : { type: 'succeeded' }, - }); -} diff --git a/packages/agent-core-v2/src/human/test/agent/machine.test.ts b/packages/agent-core-v2/src/human/test/agent/machine.test.ts index f8fa34b7b98..1d6ca159dfe 100644 --- a/packages/agent-core-v2/src/human/test/agent/machine.test.ts +++ b/packages/agent-core-v2/src/human/test/agent/machine.test.ts @@ -11,7 +11,7 @@ import { type ToolCall, } from '#/llm/message'; import type { LlmModel } from '#/llm/model'; -import { createLlmMachine, type LlmEvent } from '#/llm/requester/machine'; +import type { LlmEvent } from '#/llm/requester/actor'; import type { LlmRequester, LlmRequestEvent } from '#/llm/requester/requester'; import type { LlmRetryOptions } from '#/llm/requester/retry'; import { emptyUsage, type TokenUsage } from '#/llm/usage'; @@ -79,7 +79,7 @@ function createTestAgentMachine( ) { return createAgentMachine({ tools, - turnActor: createTurnMachine(createLlmMachine({ requester }), { retry }), + turnActor: createTurnMachine(requester, { retry }), abortTimeoutMs, }); } @@ -1377,7 +1377,7 @@ describe('agent machine max steps', () => { const actor = createActor( createAgentMachine({ tools, - turnActor: createTurnMachine(createLlmMachine({ requester })), + turnActor: createTurnMachine(requester), maxStepsPerTurn: 2, }), { input: { request: { model }, store } }, diff --git a/packages/agent-core-v2/src/human/test/agent/turn.test.ts b/packages/agent-core-v2/src/human/test/agent/turn.test.ts index dc91094f2c1..7ada3304dd7 100644 --- a/packages/agent-core-v2/src/human/test/agent/turn.test.ts +++ b/packages/agent-core-v2/src/human/test/agent/turn.test.ts @@ -6,7 +6,7 @@ import type { LlmErrorMessage } from '#/llm/errors'; import type { ContentPart, Message, UserMessage } from '#/llm/message'; import { createMediaDegradeRecovery } from '#/llm/media/degrade'; import type { LlmModel } from '#/llm/model'; -import { createLlmMachine, type LlmEvent } from '#/llm/requester/machine'; +import type { LlmEvent } from '#/llm/requester/actor'; import type { LlmRecovery } from '#/llm/requester/recovery'; import type { LlmRequester } from '#/llm/requester/requester'; import type { LlmRetryOptions } from '#/llm/requester/retry'; @@ -92,7 +92,7 @@ function startTurnActor( events: {} as TurnEvent, emitted: {} as TurnLlmEvent, }, - actors: { turn: createTurnMachine(createLlmMachine({ requester }), options) }, + actors: { turn: createTurnMachine(requester, options) }, }).createMachine({ id: 'harness', initial: 'running', diff --git a/packages/agent-core-v2/src/human/test/media/tool.test.ts b/packages/agent-core-v2/src/human/test/media/tool.test.ts index a55a6ff6de9..55c7f346876 100644 --- a/packages/agent-core-v2/src/human/test/media/tool.test.ts +++ b/packages/agent-core-v2/src/human/test/media/tool.test.ts @@ -27,7 +27,6 @@ import { createMediaRefResolver } from '#/llm/media/resolver'; import { createMemoryMediaStore } from '#/llm/media/store'; import type { LlmModel } from '#/llm/model'; import { createProvider } from '#/llm/provider/definition'; -import { createLlmMachine } from '#/llm/requester/machine'; import type { LlmRequester } from '#/llm/requester/requester'; import { openAIBase, planOpenAIRequest } from '#/llm/requester/bases/openai/requester'; import { createReadMediaFileTool } from '#/media/tool'; @@ -179,18 +178,15 @@ describe('media stack wiring', () => { const actor = createActor( createAgentMachine({ tools, - turnActor: createTurnMachine( - createLlmMachine({ - requester, - messageResolvers: [ - createMediaRefResolver({ - providers: [provider], - source: store, - cache: createMemoryMediaUploadCache(), - }), - ], - }), - ), + turnActor: createTurnMachine(requester, { + messageResolvers: [ + createMediaRefResolver({ + providers: [provider], + source: store, + cache: createMemoryMediaUploadCache(), + }), + ], + }), }), { input: { request: { model }, store: agentStore } }, ); diff --git a/packages/agent-core-v2/src/human/test/session/machine.test.ts b/packages/agent-core-v2/src/human/test/session/machine.test.ts index 947c3490aba..9154ea6a359 100644 --- a/packages/agent-core-v2/src/human/test/session/machine.test.ts +++ b/packages/agent-core-v2/src/human/test/session/machine.test.ts @@ -8,7 +8,6 @@ import { extractText, } from '#/llm/message'; import type { LlmModel } from '#/llm/model'; -import { createLlmMachine } from '#/llm/requester/machine'; import type { LlmRequester } from '#/llm/requester/requester'; import { emptyUsage } from '#/llm/usage'; import { createAgentMachine } from '#/agent/machine'; @@ -53,7 +52,7 @@ function createTestSession(requester: LlmRequester): SessionActor { createSessionMachine({ agent: createAgentMachine({ tools: [], - turnActor: createTurnMachine(createLlmMachine({ requester })), + turnActor: createTurnMachine(requester), }), }), { input: { request: { model } } }, diff --git a/packages/agent-core-v2/src/human/test/session/migrate-v2.test.ts b/packages/agent-core-v2/src/human/test/session/migrate-v2.test.ts index 54592a6902e..a0e3d4c74ff 100644 --- a/packages/agent-core-v2/src/human/test/session/migrate-v2.test.ts +++ b/packages/agent-core-v2/src/human/test/session/migrate-v2.test.ts @@ -9,7 +9,6 @@ import { createActor, waitFor } from '#/xstate2'; import { UNKNOWN_CAPABILITY } from '#/llm/capability'; import { createUserMessage, extractText } from '#/llm/message'; import type { LlmModel } from '#/llm/model'; -import { createLlmMachine } from '#/llm/requester/machine'; import type { LlmRequester } from '#/llm/requester/requester'; import { createAgentMachine } from '#/agent/machine'; import type { StateUpdated } from '#/agent/events'; @@ -44,7 +43,7 @@ function createTestSession() { createSessionMachine({ agent: createAgentMachine({ tools: [], - turnActor: createTurnMachine(createLlmMachine({ requester: createEchoRequester() })), + turnActor: createTurnMachine(createEchoRequester()), }), }), { input: { request: { model } } }, @@ -192,7 +191,7 @@ async function loadAgent(stores: SessionStores, agentId: string): Promise { const roster = (await stores.session()).getState().roster.agents; const agents: LoadedAgent[] = []; - for (const agentId of Object.keys(roster).sort()) { + for (const agentId of Object.keys(roster).toSorted()) { agents.push(await loadAgent(stores, agentId)); } return agents; diff --git a/packages/agent-core-v2/src/human/test/session/stores.test.ts b/packages/agent-core-v2/src/human/test/session/stores.test.ts index 35c7674caf3..a6473c6e136 100644 --- a/packages/agent-core-v2/src/human/test/session/stores.test.ts +++ b/packages/agent-core-v2/src/human/test/session/stores.test.ts @@ -4,7 +4,6 @@ import { createActor, waitFor, type ActorRefFrom } from '#/xstate2'; import { UNKNOWN_CAPABILITY } from '#/llm/capability'; import { createUserMessage, extractText } from '#/llm/message'; import type { LlmModel } from '#/llm/model'; -import { createLlmMachine } from '#/llm/requester/machine'; import type { LlmRequester } from '#/llm/requester/requester'; import { createAgentMachine } from '#/agent/machine'; import { createTurnMachine } from '#/agent/turn'; @@ -55,7 +54,7 @@ function startAgent(store: AgentEventStore, requester: LlmRequester = createEcho const actor = createActor( createAgentMachine({ tools: [], - turnActor: createTurnMachine(createLlmMachine({ requester })), + turnActor: createTurnMachine(requester), }), { input: { request: { model }, store } }, ); diff --git a/packages/agent-core-v2/src/human/test/tool-select/plugin.test.ts b/packages/agent-core-v2/src/human/test/tool-select/plugin.test.ts index 731bca0d5e2..ed8900ccfd8 100644 --- a/packages/agent-core-v2/src/human/test/tool-select/plugin.test.ts +++ b/packages/agent-core-v2/src/human/test/tool-select/plugin.test.ts @@ -11,7 +11,6 @@ import { type UserMessage, } from '#/llm/message'; import type { LlmModel } from '#/llm/model'; -import { createLlmMachine } from '#/llm/requester/machine'; import type { LlmRequestConfig, LlmRequester, LlmRequestEvent } from '#/llm/requester/requester'; import { connectPlugins, type AgentPluginTarget } from '#/plugin'; import { createAgentMachine, type AgentEmitted } from '#/agent/machine'; @@ -305,7 +304,7 @@ describe('tool select agent flow', () => { const actor = createActor( createAgentMachine({ tools: [createSelectToolsTool(state), deferred], - turnActor: createTurnMachine(createLlmMachine({ requester })), + turnActor: createTurnMachine(requester), }), { input: { request: { model }, store } }, ); diff --git a/packages/agent-core-v2/src/human/test/usage/machine.test.ts b/packages/agent-core-v2/src/human/test/usage/machine.test.ts index 1e5271fd565..1f1c1e3cbe1 100644 --- a/packages/agent-core-v2/src/human/test/usage/machine.test.ts +++ b/packages/agent-core-v2/src/human/test/usage/machine.test.ts @@ -5,7 +5,6 @@ import { connectPlugins } from '#/plugin'; import { UNKNOWN_CAPABILITY } from '#/llm/capability'; import { createUserMessage } from '#/llm/message'; import type { LlmModel } from '#/llm/model'; -import { createLlmMachine } from '#/llm/requester/machine'; import type { LlmRequester } from '#/llm/requester/requester'; import type { TokenUsage } from '#/llm/usage'; import { createAgentMachine } from '#/agent/machine'; @@ -133,7 +132,7 @@ describe('usage plugin', () => { const store = await testStore(); const actor = createActor( createAgentMachine({ - turnActor: createTurnMachine(createLlmMachine({ requester })), + turnActor: createTurnMachine(requester), }), { input: { request: { model }, store } }, ); diff --git a/packages/agent-core-v2/src/human/tool-select/resolver.ts b/packages/agent-core-v2/src/human/tool-select/resolver.ts index 9455c6c6cec..462ebcd024c 100644 --- a/packages/agent-core-v2/src/human/tool-select/resolver.ts +++ b/packages/agent-core-v2/src/human/tool-select/resolver.ts @@ -1,5 +1,5 @@ import type { Message } from '#/llm/message'; -import type { MessageResolver } from '#/llm/requester/machine'; +import type { MessageResolver } from '#/llm/requester/actor'; import type { ToolSelectState } from './state';