diff --git a/README.md b/README.md index bce743e..ccc45a5 100644 --- a/README.md +++ b/README.md @@ -113,7 +113,9 @@ Start only the local adapter. Press `Ctrl+C` to stop it: agentx proxy ``` -The local API is exposed at `http://127.0.0.1:` and provides `GET /health`, `GET /v1/models`, `POST /v1/messages`, `POST /v1/responses`, and `POST /v1/chat/completions`. On startup, `proxy` prints the full URL of all three client-facing endpoints. +The local API is exposed at `http://127.0.0.1:` and provides `GET /health`, `GET /v1/models`, `GET /v1/models/{id}`, `POST /v1/messages`, `POST /v1/messages/count_tokens`, `POST /v1/responses`, and `POST /v1/chat/completions`. On startup, `proxy` prints the full URL of all three client-facing endpoints. + +Every route except `GET /health` requires the adapter's local token; `/health` stays open so a supervisor can poll it. `POST /v1/messages/count_tokens` answers Claude Code's context-usage query: a native Anthropic upstream is asked for the exact count, and the other protocols — which have no equivalent endpoint — get a local estimate computed from the request itself. ### `doctor` @@ -123,7 +125,7 @@ Inspect the local environment and configuration: agentx doctor ``` -The report includes Node.js, platform/WSL status, architecture, API key presence, supported models, and Claude Code discovery. +The report includes Node.js, platform/WSL status, architecture, API key presence, supported models, Claude Code discovery, and an upstream probe that reports whether the configured endpoint is reachable and whether it *accepts* the key — a rejected key being the most common reason a launch fails. ### `forget` @@ -178,11 +180,14 @@ Print token usage statistics collected from every request the adapter serves: ```bash agentx usage # all time agentx usage --period today # today / week / month / all +agentx usage --session # totals for one client session +agentx usage --json # machine-readable output for scripts and status lines ``` The report groups tokens by provider and model and shows input/output/total counts. Statistics are stored per adapter run; the optional `--period` flag -filters by time range. +filters by time range, `--session ` narrows the report to a single client +session, and `--json` emits the same numbers as JSON instead of a table. `agentx usage --provider ` is a deprecated alias for `agentx quota --provider ` (below); it still works but prints a deprecation notice. @@ -227,7 +232,8 @@ For `claude`/`codex`, `--native` (or **Launch native (skip AgentX)** in the laun | `--provider ` | `AGENTX_PROVIDER` | none | Upstream provider (`opencode`, `deepseek`, `openrouter`) | | `--background-model ` | `AGENTX_BACKGROUND_MODEL` | none | Model for Claude Code's background (haiku) lane | | `--effort ` | `AGENTX_EFFORT` | none | Reasoning effort: `codex` `none`/`minimal`/`low`/`medium`/`high`/`xhigh`/`max`/`ultra`; `claude` `low`/`medium`/`high`/`xhigh`/`ultracode` | -| `--retry ` | `AGENTX_RETRY` | `3` | Retry attempts on upstream 429/502/503/504 (0 disables) | +| `--retry ` | `AGENTX_RETRY` | `3` | Retry attempts on upstream 408/429/502/503/504 (0 disables) | +| `--max-concurrency ` | `AGENTX_MAX_CONCURRENCY` | `0` | Max upstream requests in flight; extra requests queue locally (0 = unlimited) | | `--client-protocol ` | | `anthropic` | `exec` only: env vars to inject for the launched program | | `--verbose` | `AGENTX_LOG_LEVEL` | `info` | Reserved for verbose logging | | | `AGENTX_USAGE_DIR` | `~/.config/agentx` | Directory for token usage statistics | @@ -238,7 +244,7 @@ If the preferred port is already in use, the adapter tries subsequent ports. A n agentx proxy --host 0.0.0.0 ``` -`agentx doctor` accepts `--client ` (default `all`) to limit checks to one client, and `--offline` to skip network-dependent checks. Skipped checks are noted in the report. +`agentx doctor` accepts `--client ` (default `all`) to limit checks to one client, and `--offline` to skip the upstream probe (the only network-dependent check). Skipped checks are noted in the report. ## Credentials and Profiles @@ -279,6 +285,15 @@ Claude Code's reasoning effort is supported the same way: `agentx claude --effor Beyond the three built-in providers, you can register an arbitrary OpenAI- or Anthropic-compatible endpoint — a local model server (Ollama, vLLM, LM Studio), an internal gateway, or any other compatible API. In the interactive launcher's "Change Provider" list, choose **Add custom provider…**, then pick the protocol on a single screen — each option shows the path AgentX appends (`/v1/messages`, `/responses`, or `/chat/completions`), and Chat Completions carries a `legacy` note (it is OpenAI's earlier API, still the shape almost every third-party and local endpoint implements). The Base URL prompt repeats the chosen path, so enter the base URL only — then a display name, and the API key prompt follows. The same picker also offers **Remove custom provider…** once one exists. To do any of this without launching a client, run `agentx config` (see [`config`](#config)). +A gateway that needs extra request headers — attribution headers, or its key under a name of its own — takes them from `--header`, which is repeatable and persists with the provider: + +```bash +agentx config --provider MyGateway --base-url https://gw.example/v1 \ + --header "HTTP-Referer=https://example.com" --header "X-Org-Id=acme" +``` + +Provider headers are merged over AgentX's defaults, so a gateway that wants something other than `Authorization: Bearer` can say so. They are stored in `runtime.json` and must not carry secrets — API keys belong in the credential environment variables. + A custom endpoint exposes no model list, so the first launch asks for its model id directly — the internal placeholder model is never offered as a choice nor exposed to Codex's model picker — and remembers it like any other model afterwards. (Non-interactive runs without `--model` still fall back to the placeholder, so pass `--model` in scripts.) For scripts and non-interactive use, `agentx config --provider --base-url ` registers (and persists) a custom provider without launching a client; `exec`/`claude`/`codex` accept the same flags to define it and run in one step. `--provider` becomes its display name, and `--protocol` selects the upstream shape (`chat-completions` by default, or `responses`/`anthropic`): @@ -422,7 +437,7 @@ Supported translation areas include: - `max_tokens` to the upstream output-token limit - Anthropic text messages and response text - Anthropic streaming events to Anthropic SSE events -- Anthropic `tools`, `tool_use`, and `tool_result` to function tools and function call outputs +- Anthropic `tools`, `tool_use`, and `tool_result` to function tools and function call outputs, including a `tool_result` whose content is a block array: its text and its images are separated rather than serialized together as JSON (see [Media in tool results](#media-in-tool-results)) - Anthropic `thinking` / `output_config.effort` to the upstream's reasoning controls (DeepSeek's `thinking`/`reasoning_effort` over Chat Completions, or `reasoning.effort` over the Responses API) - Anthropic `tool_choice` to the upstream's Chat Completions or Responses tool-choice shape - Responses and Chat Completions usage data to Anthropic usage fields @@ -431,6 +446,18 @@ Supported translation areas include: For DeepSeek specifically, its thinking mode requires every assistant turn's `reasoning_content` to be echoed back anchored to the same message as the tool call it led to; the adapter keeps an assistant message's text, reasoning, and tool calls together instead of splitting them across separate messages, and only forwards `reasoning_content` for DeepSeek (other Chat Completions upstreams do not expect that field). An abnormal upstream stop — `content_filter`, `insufficient_system_resource`, or a stream that ends without either a `finish_reason` or `[DONE]` — surfaces as an error instead of silently reading back as a normal `end_turn`. +### Media in tool results + +A tool can return an image — Claude Code's `Read` does, for an image file — as an `image` block inside `tool_result.content`. Each upstream protocol takes that differently, so the adapter places it where the protocol actually accepts one: + +| Upstream protocol | Where the image goes | +| --- | --- | +| `anthropic` | Unchanged — the request is already Anthropic-shaped and passes through | +| `responses` | Inline, as `input_image` parts of the `function_call_output`'s array `output` | +| `chat-completions` | Lifted out: the tool message keeps the text, and the images follow in their own user message, because a `role:"tool"` message accepts only a string | + +When a model's catalog metadata states that it does not accept image input, the images are dropped and the tool message says so, rather than sending something the upstream would reject. A tool result that contained nothing but images still gets non-empty text, which several upstreams require. + The adapter translates tool protocols only. It does not execute tools and does not persist prompts, tool arguments, or conversation state. ## Token Usage Statistics diff --git a/README.zh-CN.md b/README.zh-CN.md index b4703a6..71f6b09 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -113,7 +113,9 @@ agentx exec --client-protocol openai -- my-openai-compatible-tool agentx proxy ``` -本地 API 地址为 `http://127.0.0.1:`,提供 `GET /health`、`GET /v1/models`、`POST /v1/messages`、`POST /v1/responses`、`POST /v1/chat/completions`。启动 `proxy` 时会打印这三个面向客户端端点各自的完整 URL。 +本地 API 地址为 `http://127.0.0.1:`,提供 `GET /health`、`GET /v1/models`、`GET /v1/models/{id}`、`POST /v1/messages`、`POST /v1/messages/count_tokens`、`POST /v1/responses`、`POST /v1/chat/completions`。启动 `proxy` 时会打印这三个面向客户端端点各自的完整 URL。 + +除 `GET /health` 外的所有路由都需要适配器的本地 token;`/health` 保持开放,便于进程守护轮询。`POST /v1/messages/count_tokens` 回答 Claude Code 的上下文占用查询:原生 Anthropic 上游会被请求精确计数,其余两种协议没有对应端点,则按请求本身在本地估算。 ### `doctor` @@ -123,7 +125,7 @@ agentx proxy agentx doctor ``` -报告包含 Node.js、平台/WSL 状态、CPU 架构、API Key 是否存在、支持的模型,以及 Claude Code 是否可被发现。 +报告包含 Node.js、平台/WSL 状态、CPU 架构、API Key 是否存在、支持的模型、Claude Code 是否可被发现,以及一项上游探测——报告配置的端点是否可达,以及 Key 是否被*接受*(Key 被拒绝是启动失败最常见的原因)。 ### `forget` @@ -171,9 +173,11 @@ AgentX 自身不保存密钥。在交互式管理器里选中未配置的 Provid ```bash agentx usage # 全部时间 agentx usage --period today # today / week / month / all +agentx usage --session # 单个客户端会话的合计 +agentx usage --json # 机器可读输出,便于脚本与状态栏消费 ``` -报告按 Provider 和模型分组显示 Token 数,包含输入/输出/总量。统计按适配器运行保存;`--period` 可按时间范围过滤。 +报告按 Provider 和模型分组显示 Token 数,包含输入/输出/总量。统计按适配器运行保存;`--period` 可按时间范围过滤,`--session ` 可只看单个客户端会话,`--json` 则以 JSON 输出同样的数字而非表格。 `agentx usage --provider ` 是 `agentx quota --provider `(见下文)的过渡期别名,仍然可用,但会打印一行弃用提示。 @@ -217,7 +221,8 @@ agentx version | `--provider ` | `AGENTX_PROVIDER` | 无 | 上游 Provider(`opencode`、`deepseek`、`openrouter`) | | `--background-model ` | `AGENTX_BACKGROUND_MODEL` | 无 | Claude Code 后台(haiku)通道使用的模型 | | `--effort ` | `AGENTX_EFFORT` | 无 | 推理档位:`codex` `none`/`minimal`/`low`/`medium`/`high`/`xhigh`/`max`/`ultra`;`claude` `low`/`medium`/`high`/`xhigh`/`ultracode` | -| `--retry ` | `AGENTX_RETRY` | `3` | 上游 429/502/503/504 的重试次数(0 表示禁用) | +| `--retry ` | `AGENTX_RETRY` | `3` | 上游 408/429/502/503/504 的重试次数(0 表示禁用) | +| `--max-concurrency ` | `AGENTX_MAX_CONCURRENCY` | `0` | 同时在途的上游请求上限,超出的请求在本地排队(0 表示不限) | | `--client-protocol ` | | `anthropic` | 仅 `exec`:决定给被启动程序注入哪种形状的环境变量 | | `--verbose` | `AGENTX_LOG_LEVEL` | `info` | 预留的详细日志选项 | | | `AGENTX_USAGE_DIR` | `~/.config/agentx` | Token 用量统计存储目录 | @@ -228,7 +233,7 @@ agentx version agentx proxy --host 0.0.0.0 ``` -`agentx doctor` 支持 `--client `(默认 `all`)只检查指定客户端,并支持 `--offline` 跳过依赖网络的检查;跳过的项会在报告中标出。 +`agentx doctor` 支持 `--client `(默认 `all`)只检查指定客户端,并支持 `--offline` 跳过上游探测(唯一依赖网络的检查);跳过的项会在报告中标出。 ## 凭据与 Provider Profile @@ -269,6 +274,15 @@ Claude Code 的推理档位同样支持:`agentx claude --effort low|medium|hig 除了三个内置 Provider,你还可以注册任意 OpenAI 或 Anthropic 兼容的端点——本地模型服务(Ollama、vLLM、LM Studio)、内部网关,或其他任何兼容 API。在交互式启动器的「Change Provider」列表里选择 **Add custom provider…**,先在一屏内选择协议——每个选项标出 AgentX 会追加的路径(`/v1/messages`、`/responses` 或 `/chat/completions`),Chat Completions 行尾还带一个 `legacy` 备注(它是 OpenAI 早期的 API,但几乎所有第三方和本地端点仍只实现这一形态)。接着填 Base URL(提示里会再写明该路径,所以只填 base 即可)和显示名称,随后紧接 API key 提示;已经有自定义 Provider 时,同一个列表还会提供 **Remove custom provider…**。不想启动客户端时,可以用 `agentx config` 完成同样的配置(见上文 [`config`](#config))。 +如果网关需要额外的请求头——归因头,或把 Key 放在自己命名的头里——用 `--header` 指定,可重复传入,并会随 Provider 一起持久化: + +```bash +agentx config --provider MyGateway --base-url https://gw.example/v1 \ + --header "HTTP-Referer=https://example.com" --header "X-Org-Id=acme" +``` + +Provider 自定义头会覆盖 AgentX 的默认头,因此需要非 `Authorization: Bearer` 形态的网关可以自行声明。这些头保存在 `runtime.json` 中,不得携带机密——API Key 应放在凭据环境变量里。 + 自定义端点不提供模型列表,因此首次启动会直接询问它的模型 id——内部占位模型不会作为选项出现,也不会出现在 Codex 的模型选择器里——之后会像其他模型一样被记住。(非交互运行未传 `--model` 时仍会回退到占位模型,所以脚本里请显式传 `--model`。) 对于脚本和非交互场景,`agentx config --provider <名称> --base-url ` 可以不启动客户端就注册(并持久化)一个自定义 Provider;`exec`/`claude`/`codex` 也接受同样的参数,在启动的同时完成定义。`--provider` 作为它的显示名,`--protocol` 选择上游协议形状(默认 `chat-completions`,也可以是 `responses`/`anthropic`): @@ -412,7 +426,7 @@ Claude Code 遇到不认识的模型时,会假定它使用默认的约 200k to - `max_tokens` 转换为上游输出 Token 限制 - Anthropic 文本消息与响应文本 - Anthropic 流式事件与上游 SSE 事件转换 -- Anthropic `tools`、`tool_use`、`tool_result` 与 function tool、function call output 转换 +- Anthropic `tools`、`tool_use`、`tool_result` 与 function tool、function call output 转换;`tool_result` 的 content 为 block 数组时,其中的文本与图片会被分开处理,而不是整体序列化成 JSON(见 [工具结果中的多模态内容](#工具结果中的多模态内容)) - Anthropic `thinking` / `output_config.effort` 转换为上游的思考控制参数(Chat Completions 上是 DeepSeek 的 `thinking`/`reasoning_effort`,Responses API 上是 `reasoning.effort`) - Anthropic `tool_choice` 转换为上游 Chat Completions 或 Responses 对应的 tool-choice 结构 - Responses 和 Chat Completions usage 字段转换为 Anthropic usage 字段 @@ -421,6 +435,18 @@ Claude Code 遇到不认识的模型时,会假定它使用默认的约 200k to DeepSeek 的思考模式要求每个 assistant 回合的 `reasoning_content` 必须回传,并且要挂在它所引出的那个 tool call 所在的同一条消息上;适配器会把一条 assistant 消息的文本、reasoning 与 tool call 保持在同一条消息里,而不是拆分成多条,并且只对 DeepSeek 转发 `reasoning_content`(其它 Chat Completions 上游并不期望这个字段)。当上游以异常方式结束——`content_filter`、`insufficient_system_resource`,或者流结束时既没有 `finish_reason` 也没有 `[DONE]`——都会转换为错误返回,而不是被悄悄当成正常的 `end_turn`。 +### 工具结果中的多模态内容 + +工具可以返回图片——Claude Code 的 `Read` 读取图片文件时就会——形式是 `tool_result.content` 里的 `image` block。三种上游协议对此的接受方式不同,适配器会把图片放到各协议真正接受的位置: + +| 上游协议 | 图片的去向 | +| --- | --- | +| `anthropic` | 原样透传,请求本来就是 Anthropic 形态 | +| `responses` | 内联为 `function_call_output` 数组 `output` 中的 `input_image` 部分 | +| `chat-completions` | 抽出:tool 消息保留文本,图片以紧随其后的一条 user 消息发送——因为 `role:"tool"` 消息只接受字符串 | + +当模型的目录元数据明确表示它不接受图片输入时,图片会被丢弃,并在 tool 消息中说明,而不是发送一个上游必定拒绝的请求。只包含图片的工具结果仍会得到非空文本,这是若干上游的硬性要求。 + 适配器只负责工具协议转换,不会执行工具,也不会持久化 prompt、工具参数或对话状态。 ## Token 用量统计 diff --git a/src/cli.ts b/src/cli.ts index ffd8aa7..041e709 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -1,7 +1,7 @@ #!/usr/bin/env node import { readFileSync, realpathSync } from "node:fs"; import { fileURLToPath } from "node:url"; -import { CLAUDE_EFFORT_LEVELS, CODEX_EFFORT_LEVELS, loadConfig, parseCliOptions as options } from "./config.js"; +import { CLAUDE_EFFORT_LEVELS, CODEX_EFFORT_LEVELS, loadConfig, parseCliOptions as options, parseHeaderFlags } from "./config.js"; import { startAdapter } from "./server.js"; import { runCommand, runShellCommand, ClientNotFoundError, CLIENT_INSTALL_COMMANDS, clientEnvironment, codexLaunchArgs, nativeClientEnvironment } from "./process.js"; import { providerEntries, runInteractiveLauncher, runProviderManager, runSavedModelManager, LaunchCancelledError, type ProviderEntry } from "./ui.js"; @@ -40,10 +40,12 @@ const OPTION_LINES = [ " --port Preferred local port (default 8787)", " --host Local bind address (default 127.0.0.1)", " --api-key Upstream API key", - " --retry Retry attempts on upstream 429/5xx (default 3, 0 disables)", + " --retry Retry attempts on upstream 408/429/502/503/504 (default 3, 0 disables)", + " --max-concurrency Max upstream requests in flight; extra requests queue locally (default 0 = unlimited)", " --client-protocol exec only: env vars to inject for the launched program (default anthropic)", " --base-url Define (and persist) a custom provider at this endpoint; --provider names it", " --protocol Upstream protocol for --base-url (default chat-completions)", + " --header Extra request header for --base-url; repeat for several", " --verbose Verbose logging", " --native claude/codex only: launch the real client directly, no adapter or env overrides", ]; @@ -56,9 +58,11 @@ function helpText(command?: string): string { if (command === "auth") { lines.push("Usage: agentx auth --provider "); } else if (command === "usage") { - lines.push("Usage: agentx usage [--period today|week|month|all]"); + lines.push("Usage: agentx usage [--period today|week|month|all] [--session ] [--json]"); lines.push("Options:"); lines.push(" --period Time range for token statistics (default all)"); + lines.push(" --session Totals for one client session id instead of the per-model table"); + lines.push(" --json Emit JSON instead of the table, for scripts and status lines"); lines.push(" --provider Deprecated: use `agentx quota --provider ` instead"); return lines.join("\n"); } else if (command === "quota") { @@ -74,6 +78,7 @@ function helpText(command?: string): string { lines.push(" --base-url Add or update a custom provider at this endpoint, then exit"); lines.push(" --protocol

Upstream protocol for --base-url (default chat-completions)"); lines.push(" --model Model id for a custom provider (default custom-model)"); + lines.push(" --header Extra request header sent to the custom endpoint; repeat for several"); return lines.join("\n"); } else if (command === "doctor") { lines.push("Usage: agentx doctor [options]"); @@ -133,7 +138,7 @@ function isInteractive(): boolean { } /** Flags the adapter consumes itself; never forwarded to the launched client. */ -const ADAPTER_FLAGS = new Set(["--model", "--effort", "--background-model", "--provider", "--port", "--host", "--api-key", "--retry", "--client-protocol", "--base-url", "--protocol", "--verbose", "--native"]); +const ADAPTER_FLAGS = new Set(["--model", "--effort", "--background-model", "--provider", "--port", "--host", "--api-key", "--retry", "--max-concurrency", "--client-protocol", "--base-url", "--protocol", "--header", "--verbose", "--native"]); /** Boolean-ish adapter flags that never consume the following argument as a value. */ const BOOLEAN_ADAPTER_FLAGS = new Set(["--verbose", "--native"]); @@ -334,7 +339,7 @@ function renderProviderList(entries: ProviderEntry[]): string { export async function runConfigCommand(args: string[]): Promise { const opts = options(args); if (opts["base-url"]) { - const definition = await persistCustomProvider(opts); + const definition = await persistCustomProvider(opts, parseHeaderFlags(args)); console.log(`✓ ${definition.name} configured — ${definition.models[0].endpoint}`); console.log(credentialInstructions(definition)); return; @@ -347,7 +352,9 @@ export async function runUsageCommand(args: string[]): Promise { const opts = options(args); if (!opts.provider && !process.env.AGENTX_PROVIDER) { const period = opts.period === "today" || opts.period === "week" || opts.period === "month" || opts.period === "all" ? opts.period : "all"; - console.log(await runUsageStats(period)); + const json = opts.json !== undefined && opts.json !== "false"; + const sessionId = opts.session && opts.session !== "true" ? opts.session : undefined; + console.log(await runUsageStats({ period, json, sessionId })); return; } console.error("Deprecated: use `agentx quota --provider ` instead."); @@ -441,11 +448,11 @@ function resolveLaunchTarget(command: string, args: string[], opts: Record): Promise { +async function persistCustomProvider(opts: Record, headers: Record = {}): Promise { const baseUrl = opts["base-url"] ?? ""; const protocol: ProviderProtocol = opts.protocol === "responses" || opts.protocol === "anthropic" ? opts.protocol : "chat-completions"; - const definition = registerCustomProvider({ name: opts.provider ?? "custom", baseUrl, protocol, model: opts.model }); - await saveCustomProvider(definition.id, { name: definition.name, baseUrl, protocol, model: definition.models[0].model }); + const definition = registerCustomProvider({ name: opts.provider ?? "custom", baseUrl, protocol, model: opts.model, headers }); + await saveCustomProvider(definition.id, { name: definition.name, baseUrl, protocol, model: definition.models[0].model, ...(Object.keys(headers).length ? { headers } : {}) }); return definition; } @@ -471,7 +478,7 @@ export async function runClientLaunch(command: string, args: string[], deps: Cli // for exec/scripts/CI, but not restricted to exec: it works the same way // for claude/codex. --provider doubles as the display name here. if (opts["base-url"]) { - const definition = await persistCustomProvider(opts); + const definition = await persistCustomProvider(opts, parseHeaderFlags(args)); opts.provider = definition.id; } // Native launch is only meaningful for claude/codex: they have their own @@ -609,7 +616,7 @@ export async function runClientLaunch(command: string, args: string[], deps: Cli async function hydrateCustomProviders(): Promise { const saved = await loadCustomProviders(); for (const [, definition] of Object.entries(saved)) { - registerCustomProvider({ name: definition.name, baseUrl: definition.baseUrl, protocol: definition.protocol as ProviderProtocol, model: definition.model }); + registerCustomProvider({ name: definition.name, baseUrl: definition.baseUrl, protocol: definition.protocol as ProviderProtocol, model: definition.model, headers: definition.headers }); } } diff --git a/src/config.ts b/src/config.ts index 70186cf..2f98651 100644 --- a/src/config.ts +++ b/src/config.ts @@ -17,8 +17,10 @@ export interface Config { effort?: ReasoningEffort; apiKey: string; logLevel: string; - /** Retry attempts on upstream network failure or 429/502/503/504; 0 disables retry. */ + /** Retry attempts on upstream network failure or a retryable status; 0 disables retry. */ retry: number; + /** Max upstream requests in flight at once; 0 or unset (the default) disables the local queue. */ + maxConcurrency?: number; } /** Validate the `--effort`/`AGENTX_EFFORT` value against the union of both clients' scales. */ @@ -48,6 +50,28 @@ export function parseCliOptions(args: string[]): Record { + const headers: Record = {}; + const add = (raw: string | undefined) => { + if (!raw) return; + const match = /^\s*([^=:]+?)\s*[=:]\s*(.*)$/.exec(raw); + if (!match) return; + headers[match[1]] = match[2]; + }; + for (let i = 0; i < args.length; i++) { + const arg = args[i]; + if (arg === "--header") { const value = args[i + 1]; if (value !== undefined && !value.startsWith("--")) { add(value); i++; } continue; } + if (arg?.startsWith("--header=")) add(arg.slice("--header=".length)); + } + return headers; +} + interface RememberedModel { provider?: string; model?: string; @@ -73,6 +97,8 @@ export function loadConfig( if (!Number.isInteger(port) || port < 1 || port > 65535) throw new Error("Invalid port"); const retry = Number(options.retry ?? process.env.AGENTX_RETRY ?? 3); if (!Number.isInteger(retry) || retry < 0) throw new Error("Invalid retry count"); + const maxConcurrency = Number(options["max-concurrency"] ?? options.maxConcurrency ?? process.env.AGENTX_MAX_CONCURRENCY ?? 0); + if (!Number.isInteger(maxConcurrency) || maxConcurrency < 0) throw new Error("Invalid max concurrency"); const envModel = process.env.AGENTX_MODEL === "auto" ? undefined : process.env.AGENTX_MODEL; const rememberedModel = remembered.provider && provider && remembered.provider !== provider ? undefined @@ -87,5 +113,6 @@ export function loadConfig( apiKey, logLevel: options.verbose ? "debug" : process.env.AGENTX_LOG_LEVEL ?? "info", retry, + maxConcurrency, }; } diff --git a/src/convert/chat.ts b/src/convert/chat.ts index dc1afb3..94921b7 100644 --- a/src/convert/chat.ts +++ b/src/convert/chat.ts @@ -4,14 +4,16 @@ * Completions protocol. */ import type { AnthropicMessage, AnthropicRequest } from "./shared.js"; -import { chatControlParams, imageDataUri, parse, samplingParams } from "./shared.js"; +import { acceptsImageInput, chatControlParams, imageDataUri, parse, samplingParams, toolResultContent, toolResultText, TOOL_RESULT_MEDIA_PROMPT } from "./shared.js"; +import type { ProviderModel } from "../providers/types.js"; import { isDeepSeekLongContextModel } from "../providers/registry.js"; -export function toChatRequest(input: AnthropicRequest, model: string, provider?: string) { +export function toChatRequest(input: AnthropicRequest, model: string, provider?: ProviderModel) { const messages: any[] = []; if (input.system) messages.push({ role: "system", content: typeof input.system === "string" ? input.system : input.system.map((part) => part.text ?? "").join("\n") }); - const deepSeek = provider === "deepseek" || isDeepSeekLongContextModel(model); - for (const message of input.messages) messages.push(...toChatMessages(message, deepSeek)); + const deepSeek = provider?.provider === "deepseek" || isDeepSeekLongContextModel(model); + const allowImages = acceptsImageInput(provider); + for (const message of input.messages) messages.push(...toChatMessages(message, deepSeek, allowImages)); return { model, messages, ...(input.max_tokens === undefined ? {} : { max_tokens: input.max_tokens }), ...(input.stream ? { stream: true } : {}), ...samplingParams(input as any), ...chatControlParams(input, deepSeek), ...(input.tools ? { tools: input.tools.map((tool) => ({ type: "function", function: { name: tool.name, description: tool.description, parameters: tool.input_schema } })) } : {}) }; } @@ -20,11 +22,13 @@ function toChatImagePart(part: any): any | undefined { return url ? { type: "image_url", image_url: { url } } : undefined; } -function toChatMessages(message: AnthropicMessage, preserveReasoning = false): any[] { +function toChatMessages(message: AnthropicMessage, preserveReasoning = false, allowImages = true): any[] { if (!Array.isArray(message.content)) return [{ role: message.role, content: message.content ?? "" }]; const parts: any[] = []; const toolCalls: any[] = []; const toolResults: any[] = []; + /** Images lifted out of tool results; they follow as their own user message. */ + const toolMedia: string[] = []; let reasoning = ""; const pushText = (text: string) => { if (!text) return; @@ -36,7 +40,9 @@ function toChatMessages(message: AnthropicMessage, preserveReasoning = false): a if (part.type === "tool_use") { toolCalls.push({ id: part.id, type: "function", function: { name: part.name, arguments: JSON.stringify(part.input ?? {}) } }); } else if (part.type === "tool_result") { - toolResults.push({ role: "tool", tool_call_id: part.tool_use_id, content: typeof part.content === "string" ? part.content : JSON.stringify(part.content ?? "") }); + const { text, images } = toolResultContent(part.content); + if (allowImages) toolMedia.push(...images); + toolResults.push({ role: "tool", tool_call_id: part.tool_use_id, content: toolResultText(text, images.length, allowImages) }); } else if (part.type === "text") { pushText(part.text ?? ""); } else if (part.type === "thinking" && preserveReasoning && message.role === "assistant") { @@ -61,6 +67,11 @@ function toChatMessages(message: AnthropicMessage, preserveReasoning = false): a output.push({ role: message.role, content: parts.every((part) => part.type === "text") ? parts.map((part) => part.text).join("") : parts }); } output.push(...toolResults); + // Chat Completions has no place for an image inside a tool message, so the + // media travels as its own user turn right after the results it came from. + if (toolMedia.length) { + output.push({ role: "user", content: [{ type: "text", text: TOOL_RESULT_MEDIA_PROMPT }, ...toolMedia.map((url) => ({ type: "image_url", image_url: { url } }))] }); + } return output; } /** Plain text of a chat message content, ignoring non-text parts (images etc.). */ @@ -70,14 +81,27 @@ function chatText(content: any): string { return ""; } +/** + * Map a Chat Completions `finish_reason` to an Anthropic `stop_reason`. + * Unrecognized reasons used to fall through to "max_tokens", telling the + * client a complete reply had been truncated; "end_turn" is the honest + * default. `chatResponseFailure` still turns the genuinely abnormal reasons + * into an upstream error before this is reached on the proxy path. + */ +function anthropicStopReason(finishReason: unknown, hasToolCalls: boolean): string { + if (hasToolCalls) return "tool_use"; + if (finishReason === "length") return "max_tokens"; + if (finishReason === "content_filter") return "refusal"; + return "end_turn"; +} + export function fromChatResponse(response: any, model: string): Record { const choice = response.choices?.[0] ?? {}; const message = choice.message ?? {}; const content = []; if (message.reasoning_content) content.push({ type: "thinking", thinking: message.reasoning_content }); const text = chatText(message.content); if (text) content.push({ type: "text", text }); for (const call of message.tool_calls ?? []) content.push({ type: "tool_use", id: call.id, name: call.function?.name, input: parse(call.function?.arguments) }); - const finishReason = choice.finish_reason; - return { id: response.id ?? `msg_${crypto.randomUUID()}`, type: "message", role: "assistant", model, content, stop_reason: message.tool_calls?.length ? "tool_use" : finishReason === "length" ? "max_tokens" : finishReason === "stop" || finishReason === undefined ? "end_turn" : "max_tokens", stop_sequence: null, usage: { input_tokens: response.usage?.prompt_tokens ?? 0, output_tokens: response.usage?.completion_tokens ?? 0, cache_creation_input_tokens: 0, cache_read_input_tokens: response.usage?.prompt_tokens_details?.cached_tokens ?? 0 } }; + return { id: response.id ?? `msg_${crypto.randomUUID()}`, type: "message", role: "assistant", model, content, stop_reason: anthropicStopReason(choice.finish_reason, Boolean(message.tool_calls?.length)), stop_sequence: null, usage: { input_tokens: response.usage?.prompt_tokens ?? 0, output_tokens: response.usage?.completion_tokens ?? 0, cache_creation_input_tokens: 0, cache_read_input_tokens: response.usage?.prompt_tokens_details?.cached_tokens ?? 0 } }; } /** Return a failure for a provider completion that is not a normal stop. */ diff --git a/src/convert/responses.ts b/src/convert/responses.ts index 25261f9..2cf341c 100644 --- a/src/convert/responses.ts +++ b/src/convert/responses.ts @@ -1,6 +1,7 @@ /** Anthropic Messages API <-> Responses API: the direction whose upstream speaks the Responses protocol. */ import type { AnthropicMessage, AnthropicRequest } from "./shared.js"; -import { chatThinking, imageDataUri, parse, reasoningEffort, responsesToolChoice } from "./shared.js"; +import { acceptsImageInput, chatThinking, imageDataUri, parse, reasoningEffort, responsesToolChoice, toolResultContent, toolResultText } from "./shared.js"; +import type { ProviderModel } from "../providers/types.js"; import { fromAnthropicResponseToChat, toAnthropicRequestFromChat } from "./anthropic.js"; function convertContent(content: any, role: string): any { @@ -8,18 +9,39 @@ function convertContent(content: any, role: string): any { return content.map((part) => { if (part.type === "text") return { type: role === "assistant" ? "output_text" : "input_text", text: part.text ?? "" }; if (part.type === "image") { const url = imageDataUri(part.source); return url ? { type: "input_image", image_url: url } : null; } - if (part.type === "tool_result") return { type: "function_call_output", call_id: part.tool_use_id, output: typeof part.content === "string" ? part.content : JSON.stringify(part.content) }; - if (part.type === "tool_use") return { type: "function_call", call_id: part.id, name: part.name, arguments: JSON.stringify(part.input ?? {}) }; + // tool_use / tool_result never reach here: toResponsesInput handles those + // blocks itself and only sends single non-tool parts through this map. return part; }).filter((part) => part !== null); } -export function toResponsesRequest(input: AnthropicRequest, model: string): Record { +/** + * Build a Responses `function_call_output` from an Anthropic tool_result. + * Unlike a Chat Completions tool message, `output` accepts either a plain + * string or an array of typed parts, so an image a tool returned survives as + * a real `input_image` instead of being flattened into a JSON string that + * carried its base64 upstream as literal text. + */ +function toFunctionCallOutput(part: any, allowImages: boolean): Record { + const { text, images } = toolResultContent(part.content); + const media = allowImages ? images : []; + if (!media.length) return { type: "function_call_output", call_id: part.tool_use_id, output: toolResultText(text, images.length, false) }; + return { + type: "function_call_output", + call_id: part.tool_use_id, + output: [ + ...(text ? [{ type: "input_text", text }] : []), + ...media.map((image_url) => ({ type: "input_image", image_url })), + ], + }; +} + +export function toResponsesRequest(input: AnthropicRequest, model: string, provider?: ProviderModel): Record { const thinking = chatThinking(input); const effort = reasoningEffort(input); const toolChoice = responsesToolChoice(input.tool_choice); const body: Record = { - model, input: toResponsesInput(input.messages), + model, input: toResponsesInput(input.messages, acceptsImageInput(provider)), ...(input.max_tokens === undefined ? {} : { max_output_tokens: input.max_tokens }), ...(input.stream ? { stream: true } : {}), ...(input.temperature === undefined ? {} : { temperature: input.temperature }), @@ -39,7 +61,7 @@ export function toResponsesRequest(input: AnthropicRequest, model: string): Reco /** Thinking blocks are local-only reasoning echoes from Claude Code; upstream providers never accept them. */ function isThinkingPart(part: any): boolean { return part?.type === "thinking" || part?.type === "redacted_thinking"; } -function toResponsesInput(messages: AnthropicMessage[]): any[] { +function toResponsesInput(messages: AnthropicMessage[], allowImages = true): any[] { const output: any[] = []; for (const message of messages) { if (!Array.isArray(message.content)) { output.push({ ...message, content: convertContent(message.content, message.role) }); continue; } @@ -47,7 +69,7 @@ function toResponsesInput(messages: AnthropicMessage[]): any[] { for (const part of message.content as any[]) { if (isThinkingPart(part)) continue; if (part.type === "tool_use") { output.push({ type: "function_call", call_id: part.id, name: part.name, arguments: JSON.stringify(part.input ?? {}) }); continue; } - if (part.type === "tool_result") { output.push({ type: "function_call_output", call_id: part.tool_use_id, output: typeof part.content === "string" ? part.content : JSON.stringify(part.content ?? "") }); continue; } + if (part.type === "tool_result") { output.push(toFunctionCallOutput(part, allowImages)); continue; } const converted = convertContent([part], message.role)[0]; if (converted) textParts.push(converted); } if (textParts.length) output.push({ role: message.role, content: textParts }); diff --git a/src/convert/shared.ts b/src/convert/shared.ts index 65cd69b..6e65faa 100644 --- a/src/convert/shared.ts +++ b/src/convert/shared.ts @@ -135,3 +135,73 @@ export function chatControlParams(input: any, deepSeek: boolean): Record { + if (!apiKey) return { checked: true, reachable: false, authorized: false, message: "skipped — no API key to test" }; + const url = modelListUrl(provider.endpoint); + try { + const headers: Record = { accept: "application/json", authorization: `Bearer ${apiKey}`, ...provider.headers }; + if (provider.protocol === "anthropic") { headers["x-api-key"] = apiKey; headers["anthropic-version"] = "2023-06-01"; } + const response = await fetcher(url, { headers, signal: AbortSignal.timeout(10_000) }); + if (response.status === 401 || response.status === 403) return { checked: true, reachable: true, authorized: false, message: `rejected the API key (HTTP ${response.status})` }; + if (!response.ok) return { checked: true, reachable: true, authorized: true, message: `reachable, but the model list returned HTTP ${response.status}` }; + return { checked: true, reachable: true, authorized: true }; + } catch (error) { + return { checked: true, reachable: false, authorized: false, message: error instanceof Error ? error.message : "unreachable" }; + } +} + export async function executableExists(command: string): Promise { return new Promise((resolve) => { const child = spawn(command, ["--version"], { stdio: "ignore", shell: process.platform === "win32" }); @@ -54,7 +99,7 @@ function parseClient(value: string | undefined): DoctorClient { } export async function runDoctor(options: DoctorOptions = {}): Promise { - const { client: _client, offline: _offline, ...raw } = options; + const { client: _client, offline: _offline, fetcher, ...raw } = options; const config = loadConfig(raw as Record); const wsl = Boolean(process.env.WSL_INTEROP); const provider = providerById(config.provider ?? "opencode"); @@ -71,6 +116,15 @@ export async function runDoctor(options: DoctorOptions = {}): Promise.`); + } return { nodeVersion: process.version, platform: wsl ? "WSL" : process.platform, @@ -85,6 +139,7 @@ export async function runDoctor(options: DoctorOptions = {}): Promise; } /** @@ -337,7 +339,7 @@ export function registerCustomProvider(input: CustomProviderInput): ProviderDefi let suffix = 2; while (providerRegistry.some((entry) => entry.id === id)) id = `${base}-${suffix++}`; } - const model: ProviderModel = { provider: id, model: input.model ?? CUSTOM_PROVIDER_PLACEHOLDER_MODEL, protocol: input.protocol, endpoint: customProviderEndpoint(input.baseUrl, input.protocol) }; + const model: ProviderModel = { provider: id, model: input.model ?? CUSTOM_PROVIDER_PLACEHOLDER_MODEL, protocol: input.protocol, endpoint: customProviderEndpoint(input.baseUrl, input.protocol), ...(input.headers && Object.keys(input.headers).length ? { headers: input.headers } : {}) }; const definition: ProviderDefinition = { id, name: input.name, apiKeyEnv: envKeyFor(id), models: [model], custom: true }; const index = providerRegistry.findIndex((entry) => entry.id === id); if (index >= 0) providerRegistry[index] = definition; diff --git a/src/providers/types.ts b/src/providers/types.ts index 510678e..b8fe499 100644 --- a/src/providers/types.ts +++ b/src/providers/types.ts @@ -10,6 +10,13 @@ export interface ProviderModel { maxOutputTokens?: number; /** Input modalities the upstream accepts, restricted to what clients understand ("text", "image"). */ modalities?: string[]; + /** + * Extra HTTP headers sent with every upstream request for this model — + * attribution headers (OpenRouter's HTTP-Referer/X-Title) or the auth header + * shape a private gateway expects. Merged over the defaults, so a gateway + * that wants something other than `Authorization: Bearer` can say so. + */ + headers?: Record; } export interface ProviderDefinition { diff --git a/src/runtime.ts b/src/runtime.ts index 422b09a..0e4da2f 100644 --- a/src/runtime.ts +++ b/src/runtime.ts @@ -22,6 +22,8 @@ export interface CustomProviderState { baseUrl: string; protocol: string; model?: string; + /** Extra request headers for this endpoint. Non-secret by contract: API keys never enter runtime state. */ + headers?: Record; } /** diff --git a/src/server.ts b/src/server.ts index 92d3279..9869b06 100644 --- a/src/server.ts +++ b/src/server.ts @@ -49,6 +49,7 @@ function body(request: IncomingMessage, limit = MAX_BODY_BYTES): Promise }); } function json(response: ServerResponse, status: number, value: unknown) { response.writeHead(status, { "content-type": "application/json" }); response.end(JSON.stringify(value)); } +function unauthorized(response: ServerResponse) { return json(response, 401, { error: { message: "Invalid API key", type: "authentication_error" } }); } function debug(config: Config, message: string) { if (config.logLevel === "debug") console.error(`[adapter] ${message}`); } function streamOptions(config: Config, provider: ProviderModel, sessionId: string, onUsage?: (usage: TokenUsage) => void): StreamUsageOptions { return { provider: provider.provider, model: provider.model, protocol: provider.protocol, sessionId, onUsage, onDiagnostic: (message) => debug(config, `stream diagnostic: ${message}`) }; @@ -64,11 +65,18 @@ function respondError(response: ServerResponse, error: unknown) { const type = status >= 500 ? "upstream_error" : "invalid_request_error"; json(response, status, { error: { message: error instanceof Error ? error.message : "Invalid request", type } }); } -/** Anthropic's Messages API authenticates via x-api-key + anthropic-version, not Authorization: Bearer. Both forms are sent for an "anthropic" upstream so a gateway that only accepts Bearer still works. */ +/** + * Anthropic's Messages API authenticates via x-api-key + anthropic-version, + * not Authorization: Bearer. Both forms are sent for an "anthropic" upstream + * so a gateway that only accepts Bearer still works. + * + * A provider's own `headers` are merged last and therefore win, which is the + * point: a private gateway may expect its key in a header of its own naming. + */ function authHeaders(provider: ProviderModel, apiKey: string): Record { const headers: Record = { authorization: `Bearer ${apiKey}`, "content-type": "application/json" }; if (provider.protocol === "anthropic") { headers["x-api-key"] = apiKey; headers["anthropic-version"] = "2023-06-01"; } - return headers; + return { ...headers, ...provider.headers }; } /** POST the converted payload upstream; network failures surface as HTTP 502. */ async function forward(config: Config, provider: ProviderModel, apiKey: string, payload: unknown, signal: AbortSignal): Promise { @@ -84,8 +92,15 @@ async function forward(config: Config, provider: ProviderModel, apiKey: string, /** Base delay, growth factor, and ceiling for the retry backoff below. */ const RETRY_BASE_MS = 300; const RETRY_MAX_DELAY_MS = 4_000; -/** Upstream statuses worth a retry: rate-limited or transiently unavailable. Other 4xx are client/config errors that retrying cannot fix. */ -const RETRYABLE_STATUS = new Set([429, 502, 503, 504]); +/** + * Upstream statuses worth a retry: rate-limited, timed out, or transiently + * unavailable. 408 is a request timeout — transient by definition — and joins + * the set for the same reason 504 is in it. 500 stays out: it is as often a + * deterministic upstream fault as a transient one, and retrying it only + * delays a failure the caller still has to handle. Other 4xx are + * client/config errors that retrying cannot fix. + */ +const RETRYABLE_STATUS = new Set([408, 429, 502, 503, 504]); /** Sleep that resolves early (no error) when `signal` aborts, so backoff never outlives a disconnected client or the idle watchdog. */ function abortableSleep(ms: number, signal: AbortSignal): Promise { @@ -139,6 +154,72 @@ export async function forwardWithRetry( throw new Error("forwardWithRetry: unreachable"); } +/** + * Flat per-image token allowance for the estimator below. Anthropic bills an + * image by its pixel area, which a request body does not carry, so a single + * mid-sized-screenshot figure stands in for it — far closer than counting the + * base64 payload as if it were prose. + */ +const IMAGE_TOKEN_ESTIMATE = 1_600; + +/** + * Approximate the input tokens of an Anthropic Messages request. Used by + * `/v1/messages/count_tokens` when the upstream protocol has no equivalent + * endpoint to ask. Four characters per token is the conventional English + * approximation; this is explicitly an estimate, but one made from the actual + * request, unlike the blind fallback a 404 used to force on the client. + */ +export function estimateInputTokens(input: any): number { + let chars = 0; + const walk = (value: unknown): void => { + if (typeof value === "string") { chars += value.length; return; } + if (Array.isArray(value)) { value.forEach(walk); return; } + if (!value || typeof value !== "object") return; + // An image block's base64 payload is not prose; charge it a flat rate + // instead of walking into the data string. + if ((value as { type?: string }).type === "image") { chars += IMAGE_TOKEN_ESTIMATE * 4; return; } + Object.values(value).forEach(walk); + }; + walk(input?.system); walk(input?.messages); walk(input?.tools); + return Math.max(1, Math.ceil(chars / 4)); +} + +/** + * Ask a native Anthropic upstream for an exact count. Any failure (endpoint + * absent, network, malformed body) resolves to `undefined` so the caller + * falls back to the local estimate — a token count must never be the reason a + * session fails. + */ +async function upstreamTokenCount(provider: ProviderModel, apiKey: string, payload: unknown): Promise { + try { + const upstream = await fetch(`${provider.endpoint}/count_tokens`, { method: "POST", headers: authHeaders(provider, apiKey), body: JSON.stringify(payload), signal: AbortSignal.timeout(10_000) }); + if (!upstream.ok) return undefined; + const value = await upstream.json(); + return Number.isFinite(value?.input_tokens) ? Number(value.input_tokens) : undefined; + } catch { return undefined; } +} + +/** + * Bounded concurrency for upstream requests. `limit <= 0` disables the gate + * entirely, which stays the default: a lone client rarely overlaps requests, + * and queueing locally only helps when several do. Acquisitions are served + * first-come-first-served. + */ +export function createConcurrencyGate(limit: number) { + if (limit <= 0) return { acquire: async () => () => {} }; + let active = 0; + const waiting: Array<() => void> = []; + const releaseOne = () => { const next = waiting.shift(); if (next) next(); else active--; }; + return { + acquire: async (): Promise<() => void> => { + if (active < limit) active++; + else await new Promise((resolve) => waiting.push(resolve)); + let released = false; + return () => { if (released) return; released = true; releaseOne(); }; + }, + }; +} + /** Abort when the upstream sends nothing (headers or body) for this long. */ const UPSTREAM_IDLE_MS = 180_000; @@ -190,7 +271,7 @@ async function proxyRequest( route: { fallbackModel: string; payload: (input: any, model: string, provider: ProviderModel) => unknown }, maxBodyBytes = MAX_BODY_BYTES, ): Promise<{ input: any; model: string; provider: ProviderModel; watched: Response } | undefined> { - if (!authorized(request, token)) { json(response, 401, { error: { message: "Invalid API key", type: "authentication_error" } }); return undefined; } + if (!authorized(request, token)) { unauthorized(response); return undefined; } try { const input = JSON.parse(await body(request, maxBodyBytes)); // The endpoint decides which model id counts as "configured" (/v1/responses @@ -225,6 +306,7 @@ export async function startAdapter(config: Config, options: AdapterOptions = {}) const maxBodyBytes = options.maxBodyBytes ?? MAX_BODY_BYTES; const closeGraceMs = options.closeGraceMs ?? 3_000; const sessionId = randomUUID(); + const gate = createConcurrencyGate(config.maxConcurrency ?? 0); const sessionFor = (input: any) => typeof input?.session_id === "string" ? input.session_id : sessionId; /** Record one usage row, swallowing persistence errors into debug logs. */ const safeRecord = (usage: TokenUsage) => { void collector.record(usage).catch((error) => debug(config, `usage record error=${error instanceof Error ? error.message : "unknown"}`)); }; @@ -237,9 +319,46 @@ export async function startAdapter(config: Config, options: AdapterOptions = {}) // Usage statistics are read through `agentx usage` (direct storage access); // the former unauthenticated /usage/* HTTP endpoints were removed — see // docs/remaining-simplification-todos.md item B. - if (pathname === "/v1/models" && request.method === "GET") return json(response, 200, { - data: providers.filter((item) => (!config.provider || item.provider === config.provider) && !isPlaceholderModel(item)).map((item) => ({ id: item.model, object: "model", owned_by: item.provider })) - }); + // The model catalog reveals which provider and models this machine is + // configured for, so it is authenticated like the proxy routes. /health + // stays open: it carries no configuration and is what a supervisor polls. + const visibleModels = () => providers.filter((item) => (!config.provider || item.provider === config.provider) && !isPlaceholderModel(item)); + const modelEntry = (item: ProviderModel) => ({ id: item.model, object: "model", owned_by: item.provider }); + if (pathname === "/v1/models" && request.method === "GET") { + if (!authorized(request, token)) return unauthorized(response); + return json(response, 200, { data: visibleModels().map(modelEntry) }); + } + if (pathname.startsWith("/v1/models/") && request.method === "GET") { + if (!authorized(request, token)) return unauthorized(response); + const id = decodeURIComponent(pathname.slice("/v1/models/".length)); + const match = visibleModels().find((item) => item.model === id); + if (!match) return json(response, 404, { error: { message: `Model "${id}" not found`, type: "not_found" } }); + return json(response, 200, modelEntry(match)); + } + if (pathname === "/v1/messages/count_tokens" && request.method === "POST") { + if (!authorized(request, token)) return unauthorized(response); + try { + const input = JSON.parse(await body(request, maxBodyBytes)); + const model = honorRequestedModel(input.model, config.model, config.provider); + const provider = providerFor(model, config.provider); + // Only a native Anthropic upstream can answer this exactly; the other + // protocols have no equivalent endpoint, so the estimate stands in. + const exact = provider.protocol === "anthropic" ? await upstreamTokenCount(provider, apiKeyFor(provider, config.apiKey), { ...input, model }) : undefined; + return json(response, 200, { input_tokens: exact ?? estimateInputTokens(input) }); + } catch (error) { return respondError(response, error); } + } + // Hold a concurrency slot for the whole proxied exchange, streaming + // included, and free it once the response is done either way. Acquiring + // after the client already left would otherwise strand the slot. + // + // Authentication comes first: a request that is going to be rejected must + // not sit in the queue ahead of a legitimate one. + if (request.method === "POST" && (pathname === "/v1/responses" || pathname === "/v1/chat/completions" || pathname === "/v1/messages")) { + if (!authorized(request, token)) return unauthorized(response); + const release = await gate.acquire(); + if (response.closed || response.writableEnded) release(); + else { response.on("close", release); response.on("finish", release); } + } if (pathname === "/v1/responses" && request.method === "POST") { // Honor auto routing here too: Codex echoes OPENAI_MODEL=auto back. const proxied = await proxyRequest(config, request, response, token, { @@ -292,7 +411,7 @@ export async function startAdapter(config: Config, options: AdapterOptions = {}) fallbackModel: config.model, payload: (input, model, provider) => { if (provider.protocol === "anthropic") return { ...input, model }; // already Anthropic-shaped; zero conversion - return provider.protocol === "responses" ? toResponsesRequest(input, model) : toChatRequest(input, model, provider.provider); + return provider.protocol === "responses" ? toResponsesRequest(input, model, provider) : toChatRequest(input, model, provider); }, }, maxBodyBytes); if (!proxied) return; diff --git a/src/usage/cli.ts b/src/usage/cli.ts index 2103d8f..3837c1d 100644 --- a/src/usage/cli.ts +++ b/src/usage/cli.ts @@ -1,5 +1,5 @@ import { defaultUsageStore } from "./storage.js"; -import type { ModelUsageStat, UsagePeriod, UsageTotals } from "./types.js"; +import type { ModelUsageStat, UsagePeriod, UsageStore, UsageTotals } from "./types.js"; export function formatPeriod(period: UsagePeriod): string { return { today: "Today", week: "This week", month: "This month", all: "All time" }[period]; @@ -32,10 +32,39 @@ export function renderUsageStats(stats: { models: ModelUsageStat[]; totals: Usag return lines.join("\n"); } -export async function runUsageStats(period: UsagePeriod = "all"): Promise { - const store = await defaultUsageStore(); +/** Human-readable totals for one client session. */ +export function renderSessionTotals(sessionId: string, totals: UsageTotals): string { + return [ + `Token Usage (session ${sessionId})`, + "", + `Input: ${compact(totals.inputTokens)}`, + `Output: ${compact(totals.outputTokens)}`, + `Total: ${compact(totals.totalTokens)}`, + ].join("\n"); +} + +export interface UsageQuery { + period?: UsagePeriod; + /** Emit machine-readable JSON instead of the table, for scripts and status lines. */ + json?: boolean; + /** Restrict the report to one client session id. */ + sessionId?: string; +} + +/** `injected` is a seam so tests never read or write the real usage file. */ +export async function runUsageStats(query: UsageQuery = {}, injected?: UsageStore): Promise { + const period = query.period ?? "all"; + const store = injected ?? await defaultUsageStore(); try { + if (query.sessionId) { + const totals = await store.sessionTotals(query.sessionId); + return query.json + ? JSON.stringify({ sessionId: query.sessionId, totals }, null, 2) + : renderSessionTotals(query.sessionId, totals); + } const [models, totals] = await Promise.all([store.modelStats(period), store.totals(period)]); - return renderUsageStats({ models, totals, period }); + return query.json + ? JSON.stringify({ period, totals, models }, null, 2) + : renderUsageStats({ models, totals, period }); } finally { await store.close(); } } diff --git a/test/config.test.ts b/test/config.test.ts index 326902f..d6616ac 100644 --- a/test/config.test.ts +++ b/test/config.test.ts @@ -3,7 +3,7 @@ import { join } from "node:path"; import { tmpdir } from "node:os"; import test from "node:test"; import assert from "node:assert/strict"; -import { loadConfig, parseCliOptions } from "../src/config.js"; +import { loadConfig, parseCliOptions, parseHeaderFlags } from "../src/config.js"; import { defaultModelFor } from "../src/selection.js"; import { loadLastSelection, saveLastModel } from "../src/runtime.js"; @@ -88,3 +88,18 @@ test("--verbose enables debug logging", () => { assert.equal(loadConfig({ verbose: "true" }).logLevel, "debug"); assert.equal(loadConfig({}).logLevel, "info"); }); + +test("--header is repeatable and keeps separators inside the value", () => { + const headers = parseHeaderFlags(["--header", "X-Title=agentx", "--header=HTTP-Referer:https://example.com/a=b", "--header", "X-Empty="]); + assert.deepEqual(headers, { "X-Title": "agentx", "HTTP-Referer": "https://example.com/a=b", "X-Empty": "" }); +}); + +test("--header ignores a malformed pair and a flag with no value", () => { + assert.deepEqual(parseHeaderFlags(["--header", "nonsense", "--header", "--verbose"]), {}); +}); + +test("max concurrency defaults to unlimited and rejects a negative value", () => { + assert.equal(loadConfig({}).maxConcurrency, 0); + assert.equal(loadConfig({ "max-concurrency": "4" }).maxConcurrency, 4); + assert.throws(() => loadConfig({ "max-concurrency": "-1" }), /Invalid max concurrency/); +}); diff --git a/test/doctor.test.ts b/test/doctor.test.ts index 1c5a1de..9e7f966 100644 --- a/test/doctor.test.ts +++ b/test/doctor.test.ts @@ -1,6 +1,7 @@ import test from "node:test"; import assert from "node:assert/strict"; -import { renderDoctor, type DoctorResult } from "../src/doctor.js"; +import { checkUpstream, modelListUrl, renderDoctor, type DoctorResult } from "../src/doctor.js"; +import type { ProviderModel } from "../src/providers/types.js"; const base: DoctorResult = { nodeVersion: "v24.0.0", @@ -16,6 +17,7 @@ const base: DoctorResult = { codexFound: false, portAvailable: true, networkChecksSkipped: false, + upstream: { checked: true, reachable: true, authorized: true }, issues: [], }; @@ -48,3 +50,55 @@ test("flags when network checks are skipped", () => { const output = renderDoctor({ ...base, networkChecksSkipped: true }); assert.match(output, /network checks skipped/); }); + +// --- upstream probe --------------------------------------------------------- +// `--offline` used to gate nothing at all: doctor made no network call, so it +// could report a key as "found" but never as accepted. + +const chatModel: ProviderModel = { provider: "opencode", model: "glm-5", protocol: "chat-completions", endpoint: "https://upstream.invalid/v1/chat/completions" }; + +test("checkUpstream derives the model list URL from the request endpoint", async () => { + let requested = ""; + const result = await checkUpstream(chatModel, "key", (async (url: any) => { requested = String(url); return new Response("{}", { status: 200 }); }) as typeof fetch); + assert.equal(requested, "https://upstream.invalid/v1/models"); + assert.deepEqual(result, { checked: true, reachable: true, authorized: true }); +}); + +test("checkUpstream reports a rejected key as reachable but unauthorized", async () => { + const result = await checkUpstream(chatModel, "bad", (async () => new Response("{}", { status: 401 })) as typeof fetch); + assert.equal(result.reachable, true); + assert.equal(result.authorized, false); + assert.match(result.message ?? "", /rejected the API key/); +}); + +test("checkUpstream reports an unreachable endpoint without throwing", async () => { + const result = await checkUpstream(chatModel, "key", (async () => { throw new Error("getaddrinfo ENOTFOUND"); }) as typeof fetch); + assert.deepEqual(result, { checked: true, reachable: false, authorized: false, message: "getaddrinfo ENOTFOUND" }); +}); + +test("checkUpstream sends the provider's own headers and Anthropic's auth pair", async () => { + const anthropicModel: ProviderModel = { provider: "custom", model: "claude-x", protocol: "anthropic", endpoint: "https://gw.invalid/v1/messages", headers: { "X-Gateway": "on" } }; + let seen: Headers | undefined; + await checkUpstream(anthropicModel, "key", (async (_url: any, init?: any) => { seen = new Headers(init?.headers); return new Response("{}", { status: 200 }); }) as typeof fetch); + assert.equal(seen?.get("x-gateway"), "on"); + assert.equal(seen?.get("x-api-key"), "key"); + assert.equal(seen?.get("anthropic-version"), "2023-06-01"); +}); + +test("checkUpstream skips the probe when there is no key to test", async () => { + const result = await checkUpstream(chatModel, "", (async () => { throw new Error("must not be called"); }) as typeof fetch); + assert.match(result.message ?? "", /no API key/); +}); + +test("renders the upstream line only once the probe actually ran", () => { + assert.doesNotMatch(renderDoctor({ ...base, upstream: { checked: false, reachable: false, authorized: false } }), /Upstream/); + assert.match(renderDoctor(base), /✓ Upstream/); + assert.match(renderDoctor({ ...base, upstream: { checked: true, reachable: true, authorized: false, message: "rejected the API key (HTTP 401)" } }), /✗ Upstream {6}rejected the API key/); +}); + +test("modelListUrl handles every protocol's endpoint shape", () => { + assert.equal(modelListUrl("https://api.deepseek.com/v1/chat/completions"), "https://api.deepseek.com/v1/models"); + assert.equal(modelListUrl("https://opencode.ai/zen/go/v1/responses"), "https://opencode.ai/zen/go/v1/models"); + assert.equal(modelListUrl("https://gw.invalid/v1/messages"), "https://gw.invalid/v1/models"); + assert.equal(modelListUrl("https://gw.invalid/custom/path"), "https://gw.invalid/custom/models"); +}); diff --git a/test/server-endpoints.test.ts b/test/server-endpoints.test.ts new file mode 100644 index 0000000..f87ca8e --- /dev/null +++ b/test/server-endpoints.test.ts @@ -0,0 +1,203 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { createConcurrencyGate, estimateInputTokens, startAdapter } from "../src/server.js"; +import { createUsageStore } from "../src/usage/storage.js"; +import { registerCustomProvider, unregisterCustomProvider } from "../src/providers/registry.js"; +import type { Config } from "../src/config.js"; + +// One custom provider per protocol, registered for the whole file so these +// tests never reach the network or leak registry state. +let chatProvider: string; +let anthropicProvider: string; +test.before(() => { + chatProvider = registerCustomProvider({ name: "Endpoint Test Chat", baseUrl: "http://chat-upstream.invalid", protocol: "chat-completions", model: "chat-x" }).id; + anthropicProvider = registerCustomProvider({ name: "Endpoint Test Anthropic", baseUrl: "http://anthropic-upstream.invalid", protocol: "anthropic", model: "claude-y" }).id; +}); +test.after(() => { unregisterCustomProvider(chatProvider); unregisterCustomProvider(anthropicProvider); }); + +function config(overrides: Partial): Config { + return { host: "127.0.0.1", port: 0, model: "chat-x", provider: chatProvider, apiKey: "test-key", logLevel: "info", retry: 0, ...overrides }; +} + +async function adapterFor(overrides: Partial = {}) { + const store = await createUsageStore({ backend: "memory" }); + return startAdapter(config(overrides), { store }); +} + +test("estimateInputTokens charges an image a flat rate instead of counting its base64 as prose", () => { + // 400 content characters plus the "user" role string, at 4 chars per token. + const prose = estimateInputTokens({ messages: [{ role: "user", content: "a".repeat(400) }] }); + assert.equal(prose, 101); + const withImage = estimateInputTokens({ messages: [{ role: "user", content: [{ type: "image", source: { type: "base64", media_type: "image/png", data: "a".repeat(2_000_000) } }] }] }); + // A 2MB base64 payload would be 500K tokens if counted as text. + assert.ok(withImage < 5_000, `expected a flat image allowance, got ${withImage}`); +}); + +test("/v1/messages/count_tokens estimates locally when the upstream protocol has no such endpoint", async () => { + const adapter = await adapterFor(); + try { + const response = await fetch(`http://127.0.0.1:${adapter.port}/v1/messages/count_tokens`, { + method: "POST", + headers: { "content-type": "application/json", authorization: `Bearer ${adapter.token}` }, + body: JSON.stringify({ model: "chat-x", messages: [{ role: "user", content: "a".repeat(40) }] }), + }); + assert.equal(response.status, 200); + assert.deepEqual(await response.json(), { input_tokens: 11 }); + } finally { await adapter.close(); } +}); + +test("/v1/messages/count_tokens asks a native Anthropic upstream for the exact count", async () => { + const originalFetch = globalThis.fetch; + let requestedUrl = ""; + globalThis.fetch = (async (input: any, init?: any) => { + const url = typeof input === "string" ? input : input.url; + if (url.includes("127.0.0.1")) return originalFetch(input, init); + requestedUrl = url; + return new Response(JSON.stringify({ input_tokens: 4242 }), { status: 200, headers: { "content-type": "application/json" } }); + }) as typeof fetch; + const adapter = await adapterFor({ model: "claude-y", provider: anthropicProvider }); + try { + const response = await fetch(`http://127.0.0.1:${adapter.port}/v1/messages/count_tokens`, { + method: "POST", + headers: { "content-type": "application/json", authorization: `Bearer ${adapter.token}` }, + body: JSON.stringify({ model: "claude-y", messages: [{ role: "user", content: "hi" }] }), + }); + assert.deepEqual(await response.json(), { input_tokens: 4242 }); + assert.equal(requestedUrl, "http://anthropic-upstream.invalid/v1/messages/count_tokens"); + } finally { globalThis.fetch = originalFetch; await adapter.close(); } +}); + +test("count_tokens falls back to the estimate when the upstream refuses the request", async () => { + const originalFetch = globalThis.fetch; + globalThis.fetch = (async (input: any, init?: any) => { + const url = typeof input === "string" ? input : input.url; + if (url.includes("127.0.0.1")) return originalFetch(input, init); + return new Response("nope", { status: 404 }); + }) as typeof fetch; + const adapter = await adapterFor({ model: "claude-y", provider: anthropicProvider }); + try { + const response = await fetch(`http://127.0.0.1:${adapter.port}/v1/messages/count_tokens`, { + method: "POST", + headers: { "content-type": "application/json", authorization: `Bearer ${adapter.token}` }, + body: JSON.stringify({ model: "claude-y", messages: [{ role: "user", content: "a".repeat(40) }] }), + }); + assert.equal(response.status, 200); + assert.deepEqual(await response.json(), { input_tokens: 11 }); + } finally { globalThis.fetch = originalFetch; await adapter.close(); } +}); + +test("the model catalog and count_tokens reject an unauthenticated caller", async () => { + const adapter = await adapterFor(); + try { + for (const [method, path] of [["GET", "/v1/models"], ["GET", "/v1/models/chat-x"], ["POST", "/v1/messages/count_tokens"]] as const) { + const response = await fetch(`http://127.0.0.1:${adapter.port}${path}`, { method, ...(method === "POST" ? { body: "{}" } : {}) }); + assert.equal(response.status, 401, `${method} ${path}`); + } + // /health stays open: it carries no configuration and is what a supervisor polls. + assert.equal((await fetch(`http://127.0.0.1:${adapter.port}/health`)).status, 200); + } finally { await adapter.close(); } +}); + +test("/v1/models/{id} serves one model and 404s an unknown id", async () => { + const adapter = await adapterFor(); + const headers = { authorization: `Bearer ${adapter.token}` }; + try { + const found = await fetch(`http://127.0.0.1:${adapter.port}/v1/models/chat-x`, { headers }); + assert.equal(found.status, 200); + assert.deepEqual(await found.json(), { id: "chat-x", object: "model", owned_by: chatProvider }); + const missing = await fetch(`http://127.0.0.1:${adapter.port}/v1/models/nope`, { headers }); + assert.equal(missing.status, 404); + } finally { await adapter.close(); } +}); + +test("the concurrency gate admits up to the limit and queues the rest in order", async () => { + const gate = createConcurrencyGate(2); + const first = await gate.acquire(); + const second = await gate.acquire(); + let thirdAdmitted = false; + const third = gate.acquire().then((release) => { thirdAdmitted = true; return release; }); + await new Promise((resolve) => setImmediate(resolve)); + assert.equal(thirdAdmitted, false, "the third request must wait while two are in flight"); + first(); + (await third)(); + assert.equal(thirdAdmitted, true); + second(); +}); + +test("a zero limit disables the gate entirely", async () => { + const gate = createConcurrencyGate(0); + const releases = await Promise.all([gate.acquire(), gate.acquire(), gate.acquire()]); + assert.equal(releases.length, 3); + releases.forEach((release) => release()); +}); + +test("releasing a slot twice does not hand out an extra permit", async () => { + const gate = createConcurrencyGate(1); + const release = await gate.acquire(); + release(); + release(); + let secondAdmitted = false; + const second = gate.acquire().then((r) => { secondAdmitted = true; return r; }); + await new Promise((resolve) => setImmediate(resolve)); + assert.equal(secondAdmitted, true); + let thirdAdmitted = false; + void gate.acquire().then(() => { thirdAdmitted = true; }); + await new Promise((resolve) => setImmediate(resolve)); + assert.equal(thirdAdmitted, false, "the double release must not have freed a second slot"); + (await second)(); +}); + +test("a custom provider's own headers are sent upstream and override the defaults", async () => { + const id = registerCustomProvider({ + name: "Header Test Gateway", baseUrl: "http://gateway.invalid", protocol: "chat-completions", model: "gw-x", + headers: { "HTTP-Referer": "https://agentx.example", authorization: "Token gateway-scheme" }, + }).id; + const originalFetch = globalThis.fetch; + let seen: Headers | undefined; + globalThis.fetch = (async (input: any, init?: any) => { + const url = typeof input === "string" ? input : input.url; + if (url.includes("127.0.0.1")) return originalFetch(input, init); + seen = new Headers(init?.headers); + return new Response(JSON.stringify({ id: "c1", choices: [{ message: { content: "ok" }, finish_reason: "stop" }], usage: { prompt_tokens: 1, completion_tokens: 1 } }), { status: 200, headers: { "content-type": "application/json" } }); + }) as typeof fetch; + const adapter = await adapterFor({ model: "gw-x", provider: id }); + try { + await fetch(`http://127.0.0.1:${adapter.port}/v1/messages`, { + method: "POST", + headers: { "content-type": "application/json", authorization: `Bearer ${adapter.token}` }, + body: JSON.stringify({ model: "gw-x", messages: [{ role: "user", content: "hi" }] }), + }); + assert.equal(seen?.get("http-referer"), "https://agentx.example"); + // A gateway may want its key in a scheme of its own; the provider's header wins. + assert.equal(seen?.get("authorization"), "Token gateway-scheme"); + } finally { globalThis.fetch = originalFetch; await adapter.close(); unregisterCustomProvider(id); } +}); + +test("an unauthenticated proxy request is rejected before it can take a concurrency slot", async () => { + // With one slot, a request that never authenticates must not be able to sit + // in the queue ahead of a legitimate one. + const adapter = await adapterFor({ maxConcurrency: 1 }); + const originalFetch = globalThis.fetch; + let upstreamCalls = 0; + globalThis.fetch = (async (input: any, init?: any) => { + const url = typeof input === "string" ? input : input.url; + if (url.includes("127.0.0.1")) return originalFetch(input, init); + upstreamCalls++; + return new Response(JSON.stringify({ id: "c1", choices: [{ message: { content: "ok" }, finish_reason: "stop" }] }), { status: 200, headers: { "content-type": "application/json" } }); + }) as typeof fetch; + try { + const rejected = await Promise.all([0, 1, 2].map(() => fetch(`http://127.0.0.1:${adapter.port}/v1/messages`, { + method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ messages: [] }), + }))); + assert.deepEqual(rejected.map((r) => r.status), [401, 401, 401]); + assert.equal(upstreamCalls, 0); + // The slot is still free for a caller that does authenticate. + const allowed = await fetch(`http://127.0.0.1:${adapter.port}/v1/messages`, { + method: "POST", + headers: { "content-type": "application/json", authorization: `Bearer ${adapter.token}` }, + body: JSON.stringify({ model: "chat-x", messages: [{ role: "user", content: "hi" }] }), + }); + assert.equal(allowed.status, 200); + assert.equal(upstreamCalls, 1); + } finally { globalThis.fetch = originalFetch; await adapter.close(); } +}); diff --git a/test/server-retry.test.ts b/test/server-retry.test.ts index 9653c33..c98b223 100644 --- a/test/server-retry.test.ts +++ b/test/server-retry.test.ts @@ -35,7 +35,18 @@ test("gives up after exhausting retries on a persistent 503", async () => { } finally { globalThis.fetch = originalFetch; } }); -test("does not retry a 500: only 429/502/503/504 are treated as transient", async () => { +test("retries a 408 request timeout like the other transient statuses", async () => { + let calls = 0; + const originalFetch = globalThis.fetch; + globalThis.fetch = (async () => { calls++; return upstreamResponse(408); }) as typeof fetch; + try { + const result = await forwardWithRetry(config, provider, "key", {}, new AbortController().signal, 3, instantSleep); + assert.equal(result.status, 408); + assert.equal(calls, 4); + } finally { globalThis.fetch = originalFetch; } +}); + +test("does not retry a 500: only 408/429/502/503/504 are treated as transient", async () => { let calls = 0; const originalFetch = globalThis.fetch; globalThis.fetch = (async () => { calls++; return upstreamResponse(500); }) as typeof fetch; diff --git a/test/tools.test.ts b/test/tools.test.ts index 9bda699..4695c26 100644 --- a/test/tools.test.ts +++ b/test/tools.test.ts @@ -1,6 +1,6 @@ import test from "node:test"; import assert from "node:assert/strict"; -import { fromResponsesResponse, toResponsesRequest } from "../src/convert/index.js"; +import { fromChatResponse, fromResponsesResponse, toChatRequest, toResponsesRequest } from "../src/convert/index.js"; test("converts tools and tool results", () => { const result = toResponsesRequest({ tools: [{ name: "bash", description: "Run a command", input_schema: { type: "object" } }], messages: [{ role: "assistant", content: [{ type: "tool_use", id: "call-1", name: "bash", input: { command: "pwd" } }] }, { role: "user", content: [{ type: "tool_result", tool_use_id: "call-1", content: "/tmp" }] }] }, "gpt-5.6-luna"); @@ -12,3 +12,96 @@ test("converts function calls to tool use", () => { const result = fromResponsesResponse({ id: "r2", output: [{ type: "function_call", call_id: "call-1", name: "bash", arguments: '{"command":"pwd"}' }] }, "gpt-5.6-luna") as any; assert.deepEqual(result.content, [{ type: "tool_use", id: "call-1", name: "bash", input: { command: "pwd" } }]); assert.equal(result.stop_reason, "tool_use"); }); + +// --- tool results carrying media ------------------------------------------- +// Anthropic allows a tool_result's content to be a block array (Claude Code's +// Read tool returns one for an image). Serializing that array as JSON used to +// send the image's base64 upstream as literal text. + +const imageResult = { + messages: [ + { role: "assistant", content: [{ type: "tool_use", id: "call-1", name: "read", input: { path: "/a.png" } }] }, + { + role: "user", + content: [{ + type: "tool_result", tool_use_id: "call-1", + content: [{ type: "text", text: "read ok" }, { type: "image", source: { type: "base64", media_type: "image/png", data: "AAA" } }], + }], + }, + ], +}; + +test("inlines tool-result media as Responses input_image parts instead of a JSON string", () => { + const input = toResponsesRequest(imageResult as any, "gpt-5.6-luna").input as any[]; + assert.deepEqual(input[1], { + type: "function_call_output", + call_id: "call-1", + output: [{ type: "input_text", text: "read ok" }, { type: "input_image", image_url: "data:image/png;base64,AAA" }], + }); +}); + +test("drops tool-result media for a Responses model known not to accept images", () => { + const provider = { provider: "opencode", model: "text-only", protocol: "responses", endpoint: "https://upstream.invalid/responses", modalities: ["text"] } as const; + const input = toResponsesRequest(imageResult as any, "text-only", provider as any).input as any[]; + assert.deepEqual(input[1], { + type: "function_call_output", + call_id: "call-1", + output: "read ok\n[1 image returned by the tool; omitted because this model does not accept image input]", + }); +}); + +test("lifts tool-result media into a following user message for chat completions", () => { + const messages = (toChatRequest(imageResult as any, "glm-5") as any).messages; + // The tool message itself must stay a plain string: chat completions has no + // other shape for it. + assert.deepEqual(messages[1], { role: "tool", tool_call_id: "call-1", content: "read ok" }); + assert.deepEqual(messages[2], { + role: "user", + content: [{ type: "text", text: "Attached media from tool result:" }, { type: "image_url", image_url: { url: "data:image/png;base64,AAA" } }], + }); +}); + +test("keeps a media-only tool result non-empty and says so when the image cannot be forwarded", () => { + const mediaOnly = { + messages: [ + { role: "assistant", content: [{ type: "tool_use", id: "call-1", name: "read", input: {} }] }, + { role: "user", content: [{ type: "tool_result", tool_use_id: "call-1", content: [{ type: "image", source: { type: "base64", media_type: "image/png", data: "AAA" } }] }] }, + ], + }; + const forwarded = (toChatRequest(mediaOnly as any, "glm-5") as any).messages; + assert.equal(forwarded[1].content, "[1 image returned by the tool]"); + const provider = { provider: "opencode", model: "text-only", protocol: "chat-completions", endpoint: "https://upstream.invalid/chat/completions", modalities: ["text"] }; + const dropped = (toChatRequest(mediaOnly as any, "text-only", provider as any) as any).messages; + assert.match(dropped[1].content, /omitted because this model does not accept image input/); + assert.equal(dropped.length, 2); // no synthetic user message +}); + +test("announces a dropped image even when the tool result had text of its own", () => { + const provider = { provider: "opencode", model: "text-only", protocol: "chat-completions", endpoint: "https://upstream.invalid/chat/completions", modalities: ["text"] }; + const messages = (toChatRequest(imageResult as any, "text-only", provider as any) as any).messages; + // The text alone would leave the model silently missing content it was + // never told about. + assert.equal(messages[1].content, "read ok\n[1 image returned by the tool; omitted because this model does not accept image input]"); + assert.equal(messages.length, 2); + const responses = toResponsesRequest(imageResult as any, "text-only", provider as any).input as any[]; + assert.match(String(responses[1].output), /omitted because this model does not accept image input/); +}); + +test("flattens a text-only tool_result block array instead of forwarding its JSON", () => { + const blocks = { + messages: [ + { role: "assistant", content: [{ type: "tool_use", id: "call-1", name: "bash", input: {} }] }, + { role: "user", content: [{ type: "tool_result", tool_use_id: "call-1", content: [{ type: "text", text: "line one" }, { type: "text", text: "line two" }] }] }, + ], + }; + assert.equal((toChatRequest(blocks as any, "glm-5") as any).messages[1].content, "line one\nline two"); +}); + +test("maps an unrecognized finish_reason to end_turn rather than max_tokens", () => { + const ended = fromChatResponse({ choices: [{ message: { content: "hi" }, finish_reason: "something_new" }] }, "glm-5") as any; + assert.equal(ended.stop_reason, "end_turn"); + const filtered = fromChatResponse({ choices: [{ message: { content: "" }, finish_reason: "content_filter" }] }, "glm-5") as any; + assert.equal(filtered.stop_reason, "refusal"); + const truncated = fromChatResponse({ choices: [{ message: { content: "hi" }, finish_reason: "length" }] }, "glm-5") as any; + assert.equal(truncated.stop_reason, "max_tokens"); +}); diff --git a/test/usage-storage.test.ts b/test/usage-storage.test.ts index 7ea0d59..708a0ac 100644 --- a/test/usage-storage.test.ts +++ b/test/usage-storage.test.ts @@ -6,6 +6,7 @@ import { join } from "node:path"; import { createUsageStore, defaultUsageStore, defaultUsageLocation, periodStart, sqliteAvailable } from "../src/usage/storage.js"; import { TokenUsageCollector, normalizeUsage } from "../src/usage/collector.js"; import type { TokenUsage } from "../src/usage/types.js"; +import { runUsageStats } from "../src/usage/cli.js"; const now = Date.now(); @@ -143,3 +144,23 @@ test("cache write tokens survive normalization, storage, and stats", async () => assert.equal(models[0].cacheWriteTokens, 200); await store.close(); }); + +test("usage reports session totals and machine-readable JSON", async () => { + const store = await createUsageStore({ backend: "memory" }); + const collector = new TokenUsageCollector(store); + await collector.record(sample({ provider: "opencode", model: "glm-5", sessionId: "s-json" })); + await collector.record(sample({ provider: "opencode", model: "glm-5", sessionId: "other" })); + + const session = await runUsageStats({ sessionId: "s-json" }, store); + assert.match(session, /Token Usage \(session s-json\)/); + assert.match(session, /Total:\s+150/); + + const asJson = JSON.parse(await runUsageStats({ json: true }, store)); + assert.equal(asJson.period, "all"); + assert.equal(asJson.totals.totalTokens, 300); + assert.equal(asJson.models[0].model, "glm-5"); + + const sessionJson = JSON.parse(await runUsageStats({ sessionId: "s-json", json: true }, store)); + assert.deepEqual(sessionJson, { sessionId: "s-json", totals: { inputTokens: 100, outputTokens: 50, totalTokens: 150 } }); + await store.close(); +});