diff --git a/CHANGELOG.md b/CHANGELOG.md index dd66972..602c43f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,15 @@ ## Unreleased +- Gateway (#144): conversations queue on the shared `SessionManager` instead of + their own promise chains. Beyond `max_conversations`, the least recently used + conversation is evicted (it was the oldest inserted). Only Telegram's `/start` + command (`/start`, `/start@bot`, `/start `) is rewritten to "hello"; + `/started ...` and other platforms pass through. Shutdown waits for adapters + to stop (up to 5 s) before exiting. +- **Breaking (webhook):** `POST /message` returns the run's token `usage` + instead of `null`, and an agent failure is `502 {"error":"agent error: ..."}` + instead of `200` with the error in `reply`. - Profiles (#117): `~/.lich/profiles/.json` (plus an optional `.md` used as the system prompt) is merged between the global and project config. Select one with `--profile`, `LICH_PROFILE`, a project `profile` key, or the diff --git a/docs/getting-started.md b/docs/getting-started.md index 671e64e..3c28953 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -94,7 +94,7 @@ curl -s -X POST http://localhost:8089/message \ Observed response shape (the `reply` text is the model's answer; `usage` is `null` on this endpoint): ```json -{"reply":"Hello! How can I help you today? I can assist with coding, file management, running commands, web searches, and more — just let me know what you'd like to do.","usage":null} +{"reply":"Hello! How can I help you today? I can assist with coding, file management, running commands, web searches, and more — just let me know what you'd like to do.","usage":{"prompt_tokens":812,"completion_tokens":41,"total_tokens":853}} ``` Stop the gateway with Ctrl+C (SIGINT) or `kill` (SIGTERM); both shut down adapters cleanly. diff --git a/docs/user-guide/gateway.md b/docs/user-guide/gateway.md index 94f5829..458cbc5 100644 --- a/docs/user-guide/gateway.md +++ b/docs/user-guide/gateway.md @@ -4,7 +4,7 @@ ## How it works -`lich gateway ` runs a long-lived process that forwards inbound chat messages to **one shared agent** and routes replies back. Per-conversation memory is keyed `platform:chat_id` (Telegram/Discord chat ids, Twitch channel names, webhook `chat_id` field): each conversation keeps its own bounded history capped at 40 messages (oldest evicted; conversations beyond 200 are evicted oldest-first). Messages for the same conversation are serialized, so overlapping messages never interleave histories; different conversations can run concurrently. Failures become a safe one-line reply: `agent error: `. +`lich gateway ` runs a long-lived process that forwards inbound chat messages to **one shared agent** and routes replies back. Per-conversation memory is keyed `platform:chat_id` (Telegram/Discord chat ids, Twitch channel names, webhook `chat_id` field): each conversation keeps its own bounded history capped at 40 messages (oldest evicted; beyond 200 conversations, the least recently used one is evicted). Messages for the same conversation are serialized, so overlapping messages never interleave histories; different conversations can run concurrently. Failures become a safe one-line reply, `agent error: `; the webhook returns it as HTTP `502 {"error":"agent error: ..."}`. ```mermaid flowchart LR @@ -25,7 +25,7 @@ lich gateway telegram discord twitch # no webhook server lich gateway # defaults to webhook ``` -The gateway is silent after startup: Telegram/Discord/Twitch respond only in chats, channels, or servers the bot can see or has joined, and the webhook only serves HTTP. Telegram media messages arrive as the placeholder text `media not supported yet`; other non-text events are ignored. Telegram `/start` is answered like a plain "hello". +The gateway is silent after startup: Telegram/Discord/Twitch respond only in chats, channels, or servers the bot can see or has joined, and the webhook only serves HTTP. Telegram media messages arrive as the placeholder text `media not supported yet`; other non-text events are ignored. Telegram `/start` (also `/start@bot` and `/start `) is answered like a plain "hello"; other text, and other platforms, are passed through as written. **Security defaults:** the webhook binds loopback only; Telegram/Discord/Twitch default-deny until you configure allowlists; the gateway agent uses a read-only tool subset (no `terminal`, no file writes) unless you override `gateway.tools_enabled`. @@ -40,7 +40,7 @@ lich gateway webhook ```sh curl -s -X POST http://127.0.0.1:8089/message \ -H "content-type: application/json" -d '{"text": "hello"}' -# -> {"reply":"...","usage":null} +# -> {"reply":"...","usage":{"prompt_tokens":...,"completion_tokens":...,"total_tokens":...}} curl -s http://127.0.0.1:8089/health # -> {"status":"ok"} @@ -203,10 +203,10 @@ Request: Only `text` is required. Any `platform` field in the body is ignored; the conversation key is always `webhook:`. Success (`200`): ```json -{"reply":"Hello! How can I help you today? ...","usage":null} +{"reply":"Hello! How can I help you today? ...","usage":{"prompt_tokens":812,"completion_tokens":24,"total_tokens":836}} ``` -`reply` is the agent's final answer; `usage` is always `null` on this endpoint (webhook replies are formatted without usage stats, unlike the chat/TUI footers). Errors: +`reply` is the agent's final answer; `usage` is the run's token totals. Errors: | Status | When | | --- | --- | @@ -214,8 +214,7 @@ Only `text` is required. Any `platform` field in the body is ignored; the conver | `401` | `LICH_GATEWAY_TOKEN` is set and the `x-lich-token` header does not match. | | `404` | Anything other than `POST /message` or `GET /health`. | | `500` | Internal dispatch failure: `{"error":"internal error"}`. | - -Agent-level failures (e.g. every provider failed) return `200` with `reply` set to a sanitized one-line `agent error: ...` string, so callers always get a deliverable text. +| `502` | The agent run failed (e.g. every provider failed): `{"error":"agent error: ..."}`, a sanitized one-line message. | ### `GET /health` diff --git a/docs/user-guide/godot.md b/docs/user-guide/godot.md index 2e4fcc7..be54ee7 100644 --- a/docs/user-guide/godot.md +++ b/docs/user-guide/godot.md @@ -44,10 +44,10 @@ curl -s -X POST http://127.0.0.1:8089/message \ -d '{"text":"round 1: hero1 at full. goblin is the only living enemy.","chat_id":"run-1"}' ``` -Success is exactly one JSON object. This endpoint always sends `usage: null` — it does not forward provider token counts: +Success is exactly one JSON object, with the run's token totals in `usage`: ```json -{"reply":"...","usage":null} +{"reply":"...","usage":{"prompt_tokens":812,"completion_tokens":24,"total_tokens":836}} ``` | Status | Body | @@ -56,10 +56,11 @@ Success is exactly one JSON object. This endpoint always sends `usage: null` — | `401` | `{"error":"unauthorized"}` when `LICH_GATEWAY_TOKEN` is set and `x-lich-token` does not match | | `404` | `{"error":"not found"}` for any other method or path | | `500` | `{"error":"internal error"}` if the handler throws before a response is sent | +| `502` | `{"error":"agent error: ..."}` when the agent run fails (a sanitized one-line message) | -A failed run is still `200`. `reply` is then a sanitized `agent error: ...` line. There is no streaming, pagination, or cursor. Full platform notes: [webhook API](gateway.md#webhook-api-reference). +There is no streaming, pagination, or cursor. Full platform notes: [webhook API](gateway.md#webhook-api-reference). -Memory is keyed `platform:chat_id`, capped at 40 messages (oldest dropped) and 200 conversations (oldest dropped). For a roguelike, `chat_id` = the run id gives the commander that process's memory of the run. A new run id starts a fresh history. That history is in memory only — restarting the gateway clears it. Restate facts the digest still needs. Durable notes are a different file, below. +Memory is keyed `platform:chat_id`, capped at 40 messages (oldest dropped) and 200 conversations (least recently used dropped). For a roguelike, `chat_id` = the run id gives the commander that process's memory of the run. A new run id starts a fresh history. That history is in memory only — restarting the gateway clears it. Restate facts the digest still needs. Durable notes are a different file, below. ## Wire the example plugin @@ -157,6 +158,6 @@ Plugins run in-process with the agent's privileges (files, network, environment) ## Limits - A long run evicts gateway history past 40 messages. Put habits that still matter in the digest, or in `memory.jsonl` if they must survive a restart. -- More than 200 concurrent `chat_id`s on one process drops the oldest conversation. Fine for one developer machine; a host of many runs should know the cap. +- More than 200 concurrent `chat_id`s on one process drops the least recently used conversation. Fine for one developer machine; a host of many runs should know the cap. - Combat must finish if the gateway is down, the call times out, or `orders.jsonl` is empty or garbage. The game's fallback is the game's — this repo does not ship one. - Same-`chat_id` calls run one after another. Different `chat_id`s run concurrently. The plugin itself makes no concurrency guarantee; one bridge per `chat_id` is the intended pattern. diff --git a/examples/persona_orchestrator/README.md b/examples/persona_orchestrator/README.md index 34bc95b..b15bb41 100644 --- a/examples/persona_orchestrator/README.md +++ b/examples/persona_orchestrator/README.md @@ -49,7 +49,7 @@ curl -s -X POST http://127.0.0.1:8090/message \ -d '{"text":"round 1: hero1 at full. goblin is the only living enemy.","chat_id":"npc:commander:run-1"}' ``` -Success is `{reply, usage}`. `usage` is the run's `usage_total`. The CLI webhook still sends `usage: null`; a Godot client that only reads `reply` needs no change. Missing `text` is `400 {"error":"text is required"}`. A bad token is `401`. Wrong `Content-Type` is `415`. Oversized body is `413`. Unknown or missing `chat_id` is still `200` with `reply` starting `agent error: unknown persona` and `usage: null` — this example does not default `chat_id` to `"default"`, because a persona cannot be inferred. +Success is `{reply, usage}`. `usage` is the run's `usage_total`. The CLI webhook sends the same shape, and returns `502` on an agent failure; this example keeps `200` with an `agent error:` reply. Missing `text` is `400 {"error":"text is required"}`. A bad token is `401`. Wrong `Content-Type` is `415`. Oversized body is `413`. Unknown or missing `chat_id` is still `200` with `reply` starting `agent error: unknown persona` and `usage: null` — this example does not default `chat_id` to `"default"`, because a persona cannot be inferred. From a source checkout, with the working directory at the repo root: diff --git a/src/gateway/bus.ts b/src/gateway/bus.ts index b531fec..61e6573 100644 --- a/src/gateway/bus.ts +++ b/src/gateway/bus.ts @@ -1,17 +1,19 @@ /** * GatewayBus: conversation-keyed runner over one shared Agent. * - * Per-conversation history lives in a bounded Map; concurrent messages for - * the same conversation are serialized through a promise chain so history - * never interleaves. Agent failures become sanitized reply strings. + * Per-conversation history lives in a bounded Map evicted by least recent + * use; messages for the same conversation are serialized on the shared + * SessionManager so history never interleaves. Agent failures become + * sanitized reply strings. */ import type { Agent } from "../agent/agent.js"; import type { AgentConfig } from "../agent/config.js"; import { history_after_run_error } from "../agent/loop.js"; import type { Message } from "../providers/types.js"; +import { create_session_manager, type SessionManager } from "../session/manager.js"; import { logger } from "../util/log.js"; import { check_gateway_sender } from "./access.js"; -import { sanitize_agent_error } from "./types.js"; +import { sanitize_agent_error, type GatewayReply } from "./types.js"; export interface GatewayBusOptions { history_cap?: number; @@ -33,7 +35,7 @@ export class GatewayBus { private readonly agent_factory: () => Agent; private agent: Agent | undefined; private readonly histories: Map = new Map(); - private readonly chains: Map> = new Map(); + private readonly sessions: SessionManager = create_session_manager(); private readonly history_cap: number; private readonly max_conversations: number; private stop_logging: (() => void) | undefined; @@ -48,32 +50,19 @@ export class GatewayBus { } } - /** Serializes runs per conversation and resolves to the reply text. */ + /** Serializes runs per conversation and resolves to the reply text (undefined when empty or denied). */ async handle(platform: string, chat_id: string, user_id: string, text: string): Promise { + const reply = await this.reply(platform, chat_id, user_id, text); + return reply === undefined || reply.text.length === 0 ? undefined : reply.text; + } + + /** Like `handle`, with the run's usage and whether it failed; undefined when the sender is denied. */ + async reply(platform: string, chat_id: string, user_id: string, text: string): Promise { if (check_gateway_sender(this.config, platform, chat_id, user_id) === false) { return undefined; } const key = conversation_key(platform, chat_id); - const previous = this.chains.get(key) ?? Promise.resolve(); - const run = previous.then(() => this.run_once(key, platform, chat_id, user_id, text)); - let tracked: Promise = Promise.resolve(); - tracked = run.then( - () => { - this.release_chain(key, tracked); - }, - () => { - this.release_chain(key, tracked); - }, - ); - this.chains.set(key, tracked); - return run; - } - - /** Drops a settled chain entry unless a newer message re-queued the key. */ - private release_chain(key: string, tracked: Promise): void { - if (this.chains.get(key) === tracked) { - this.chains.delete(key); - } + return this.sessions.enqueue(key, () => this.run_once(key, platform, chat_id, user_id, text)); } /** Unsubscribes the debug tool logger (bus owns no other resources). */ @@ -88,21 +77,37 @@ export class GatewayBus { chat_id: string, user_id: string, text: string, - ): Promise { - const input = text.startsWith("/start") === true ? "hello" : text; + ): Promise { + const input = is_telegram_start(platform, text) === true ? "hello" : text; const history = this.history_for(key); const agent = this.ensure_agent(); try { const result = await agent.run({ input, history, label: `gw:${platform}:${chat_id}` }); - this.histories.set(key, cap_history(result.messages, this.history_cap)); - return final_reply_text(result.outcome.final?.content); + this.store_history(key, cap_history(result.messages, this.history_cap)); + return { text: result.outcome.final?.content ?? "", usage: result.usage_total }; } catch (error) { logger.error(`gateway bus run failed for ${key} (user ${user_id})`, error); const kept = history_after_run_error(error); if (kept !== undefined) { - this.histories.set(key, cap_history(kept, this.history_cap)); + this.store_history(key, cap_history(kept, this.history_cap)); } - return sanitize_agent_error(error); + return { text: sanitize_agent_error(error), failed: true }; + } + } + + /** + * Re-inserts the key so Map order tracks the most recent use, then trims to + * the cap: concurrent new chats can each pass `history_for` before any stores. + */ + private store_history(key: string, messages: Message[]): void { + this.histories.delete(key); + this.histories.set(key, messages); + while (this.histories.size > this.max_conversations) { + const oldest = this.histories.keys().next(); + if (oldest.done === true || oldest.value === key) { + break; + } + this.histories.delete(oldest.value); } } @@ -113,7 +118,7 @@ export class GatewayBus { return this.agent; } - /** Oldest-first eviction keeps the conversation map bounded. */ + /** Least-recently-used eviction keeps the conversation map bounded. */ private history_for(key: string): Message[] { while (this.histories.size >= this.max_conversations && this.histories.has(key) === false) { const oldest = this.histories.keys().next(); @@ -141,6 +146,11 @@ function conversation_key(platform: string, chat_id: string): string { return `${platform}:${chat_id}`; } +/** Telegram's bot-start command (`/start`, `/start@bot`, `/start `); not `/started` or other platforms. */ +function is_telegram_start(platform: string, text: string): boolean { + return platform === "telegram" && /^\/start(@\w+)?(\s|$)/.test(text); +} + /** * Newest `cap` messages, starting at the first user message in the window. * When the window has no user message, drop a leading tool-result fragment so @@ -195,7 +205,3 @@ function assistant_owning_tools(messages: readonly Message[], start: number): nu const calls = parent.tool_calls; return calls !== undefined && calls.length > 0 ? index : undefined; } - -function final_reply_text(content: string | undefined): string | undefined { - return content === undefined || content.length === 0 ? undefined : content; -} \ No newline at end of file diff --git a/src/gateway/runner.ts b/src/gateway/runner.ts index e3f0e47..298abe5 100644 --- a/src/gateway/runner.ts +++ b/src/gateway/runner.ts @@ -49,7 +49,7 @@ function is_known_platform(platform: string): boolean { function build_adapters(config: AgentConfig, bus: GatewayBus, platforms: readonly string[]): PlatformAdapter[] { const params: AdapterParams = { config, - handle_message: (platform, chat_id, user_id, text) => bus.handle(platform, chat_id, user_id, text), + handle_message: (platform, chat_id, user_id, text) => bus.reply(platform, chat_id, user_id, text), get_agent: () => { throw new Error("get_agent is reserved for future use"); }, @@ -109,16 +109,37 @@ async function start_all_adapters(adapters: readonly PlatformAdapter[]): Promise } } +/** Upper bound on waiting for adapters to stop, so a hung platform socket cannot block exit. */ +const SHUTDOWN_TIMEOUT_MS = 5000; + +/** Stops every adapter and waits for them (bounded), then releases the bus and agent. */ +export async function shutdown_gateway( + agent: Agent, + bus: GatewayBus, + adapters: readonly PlatformAdapter[], + timeout_ms = SHUTDOWN_TIMEOUT_MS, +): Promise { + const stops = Promise.allSettled( + adapters.map((adapter) => + adapter.stop().catch((error: unknown) => logger.warn(`gateway adapter stop failed: ${adapter.name}`, error)), + ), + ); + let timer: ReturnType | undefined; + const timeout = new Promise((resolve) => { + timer = setTimeout(() => { + logger.warn(`gateway adapters did not stop within ${timeout_ms}ms; exiting anyway`); + resolve(); + }, timeout_ms); + }); + await Promise.race([stops, timeout]); + clearTimeout(timer); + bus.stop(); + agent.close(); +} + function install_signal_handlers(agent: Agent, bus: GatewayBus, adapters: readonly PlatformAdapter[]): void { const shutdown = (): void => { - for (const adapter of adapters) { - void adapter - .stop() - .catch((error: unknown) => logger.warn(`gateway adapter stop failed: ${adapter.name}`, error)); - } - bus.stop(); - agent.close(); - process.exit(0); + void shutdown_gateway(agent, bus, adapters).finally(() => process.exit(0)); }; process.once("SIGINT", shutdown); process.once("SIGTERM", shutdown); diff --git a/src/gateway/types.ts b/src/gateway/types.ts index 94bd150..f2b20cd 100644 --- a/src/gateway/types.ts +++ b/src/gateway/types.ts @@ -5,6 +5,7 @@ */ import type { Agent } from "../agent/agent.js"; import type { AgentConfig } from "../agent/config.js"; +import type { Usage } from "../providers/types.js"; import { logger } from "../util/log.js"; export type PlatformName = "webhook" | "telegram" | "discord" | "twitch"; @@ -29,13 +30,28 @@ export interface PlatformAdapter { stop(): Promise; } -/** Handles one inbound message and resolves to the reply text. */ +/** One run's reply: text, token usage, and whether the agent failed (text is then the sanitized error). */ +export interface GatewayReply { + text: string; + usage?: Usage; + failed?: boolean; +} + +/** Handles one inbound message and resolves to the reply (text alone, or text plus usage/failure). */ export type InboundHandler = ( platform: string, chat_id: string, user_id: string, text: string, -) => Promise; +) => Promise; + +/** Reply text for chat platforms, whichever form the handler returned. */ +export function reply_text(reply: string | GatewayReply | undefined): string { + if (reply === undefined) { + return ""; + } + return typeof reply === "string" ? reply : reply.text; +} /** Dependencies handed to every adapter factory. */ export interface AdapterParams { @@ -97,7 +113,7 @@ export async function run_inbound_message( text: string, ): Promise { try { - return (await handle(platform, chat_id, user_id, text)) ?? ""; + return reply_text(await handle(platform, chat_id, user_id, text)); } catch (error) { logger.error(`gateway ${platform} message handling failed`, error); return sanitize_agent_error(error); diff --git a/src/gateway/webhook.ts b/src/gateway/webhook.ts index e66e747..fafb775 100644 --- a/src/gateway/webhook.ts +++ b/src/gateway/webhook.ts @@ -7,7 +7,7 @@ import { createServer, type IncomingMessage, type Server, type ServerResponse } import { logger } from "../util/log.js"; import { format_agent_reply } from "./format.js"; import { read_platform_token } from "./token_env.js"; -import type { AdapterParams, PlatformAdapter } from "./types.js"; +import { reply_text, type AdapterParams, type PlatformAdapter } from "./types.js"; export const DEFAULT_GATEWAY_PORT = 8089; export const DEFAULT_GATEWAY_HOST = "127.0.0.1"; @@ -105,7 +105,13 @@ async function handle_message_post( const chat_id = payload.chat_id ?? "default"; const user_id = payload.user_id ?? "anonymous"; const reply = await params.handle_message(platform, String(chat_id), String(user_id), String(text)); - respond_json_text(response, 200, format_agent_reply(reply ?? "", undefined, "webhook")); + if (typeof reply === "object" && reply.failed === true) { + // The text is already the sanitized agent error. + send_json(response, 502, { error: reply.text }); + return; + } + const usage = typeof reply === "object" ? reply.usage : undefined; + respond_json_text(response, 200, format_agent_reply(reply_text(reply), usage, "webhook")); } function content_length_exceeds(request: IncomingMessage, max_bytes: number): boolean { diff --git a/test/gateway.test.ts b/test/gateway.test.ts index f600656..c56c261 100644 --- a/test/gateway.test.ts +++ b/test/gateway.test.ts @@ -21,6 +21,7 @@ import { note_partial_messages } from "../src/agent/loop.js"; import { GatewayBus } from "../src/gateway/bus.js"; import type { Agent, AgentRunResult } from "../src/agent/agent.js"; import type { Message, Usage } from "../src/providers/types.js"; +import type { GatewayReply } from "../src/gateway/types.js"; import { TMP_BASE } from "./helpers/tmp_base.js"; const usage_zero: Usage = { prompt_tokens: 0, completion_tokens: 0, total_tokens: 0 }; @@ -383,7 +384,7 @@ describe("gateway bus", () => { } }); - it("releases settled promise chains so the chains map stays bounded (G-6)", async () => { + it("releases settled session queues so they stay bounded (G-6)", async () => { const work_dir = temp_work_dir(); try { const bus = new GatewayBus({ @@ -392,8 +393,8 @@ describe("gateway bus", () => { }); await bus.handle("webhook", "c1", "u1", "one"); await bus.handle("webhook", "c2", "u1", "two"); - const chains = (bus as unknown as { chains: Map> }).chains; - expect(chains.size).toBe(0); + const sessions = (bus as unknown as { sessions: { pending_count(): number } }).sessions; + expect(sessions.pending_count()).toBe(0); } finally { rmSync(work_dir, { recursive: true, force: true }); } @@ -421,6 +422,73 @@ describe("gateway bus", () => { } }); + it("evicts the least recently used conversation, not the oldest inserted (G-10)", async () => { + const work_dir = temp_work_dir(); + try { + const records: RunRecord[] = []; + const bus = new GatewayBus( + { config: config_for(work_dir), agent_factory: () => recording_agent(records) }, + { max_conversations: 2 }, + ); + await bus.handle("webhook", "a", "u1", "a1"); + await bus.handle("webhook", "b", "u1", "b1"); + await bus.handle("webhook", "a", "u1", "a2"); + await bus.handle("webhook", "c", "u1", "c1"); + await bus.handle("webhook", "a", "u1", "a3"); + await bus.handle("webhook", "b", "u1", "b2"); + const history_of = (input: string): readonly Message[] => records.find((run) => run.input === input)?.history ?? []; + // a was used after b, so c evicted b: a keeps its history and b starts over. + expect(history_of("a3").length).toBeGreaterThan(0); + expect(history_of("b2")).toEqual([]); + } finally { + rmSync(work_dir, { recursive: true, force: true }); + } + }); + + it("rewrites only telegram's /start command to hello (G-10)", async () => { + const work_dir = temp_work_dir(); + try { + const records: RunRecord[] = []; + const bus = new GatewayBus({ + config: config_for(work_dir, { allowed_users: { telegram: ["u1"] } }), + agent_factory: () => recording_agent(records), + }); + await bus.handle("telegram", "t1", "u1", "/start"); + await bus.handle("telegram", "t2", "u1", "/start@lich_bot deep-link"); + await bus.handle("telegram", "t3", "u1", "/started a thing"); + await bus.handle("webhook", "w1", "u1", "/start"); + expect(records.map((run) => run.input)).toEqual(["hello", "hello", "/started a thing", "/start"]); + } finally { + rmSync(work_dir, { recursive: true, force: true }); + } + }); + + it("reply() carries the run's usage, and marks agent failures", async () => { + const work_dir = temp_work_dir(); + try { + let fail = false; + const bus = new GatewayBus({ + config: config_for(work_dir), + agent_factory: () => + ({ + run: async (options: { input: string; history?: readonly Message[] }): Promise => { + if (fail === true) { + throw new Error("provider down"); + } + return { ...reply_result("ok", options.history ?? [], options.input), usage_total: usage_small }; + }, + }) as unknown as Agent, + }); + expect(await bus.reply("webhook", "c1", "u1", "hi")).toEqual({ text: "ok", usage: usage_small }); + fail = true; + const failed = await bus.reply("webhook", "c1", "u1", "again"); + expect(failed?.failed).toBe(true); + expect(failed?.text.startsWith("agent error: provider down")).toBe(true); + } finally { + rmSync(work_dir, { recursive: true, force: true }); + } + }); + it("keeps completed tool turns when a later model call throws", async () => { const work_dir = temp_work_dir(); try { @@ -567,7 +635,7 @@ describe("webhook adapter", () => { chat_id: string, user_id: string, text: string, - ) => Promise, + ) => Promise, ) { return create_webhook_adapter({ config: config_for(work_dir), @@ -584,6 +652,68 @@ describe("webhook adapter", () => { }); } + it("stops promptly while a keep-alive client is still connected", async () => { + const work_dir = temp_work_dir(); + let port: number | undefined; + const adapter = make_adapter(work_dir, "echo", (seen) => { + port = seen; + }); + try { + await adapter.start(); + const { request } = await import("node:http"); + const agent = new (await import("node:http")).Agent({ keepAlive: true }); + await new Promise((resolve, reject) => { + const req = request( + { host: "127.0.0.1", port, path: "/health", agent }, + (response) => { + response.resume(); + response.on("end", () => resolve()); + }, + ); + req.on("error", reject); + req.end(); + }); + const started = Date.now(); + await adapter.stop(); + expect(Date.now() - started).toBeLessThan(1000); + agent.destroy(); + } finally { + rmSync(work_dir, { recursive: true, force: true }); + } + }); + + it("returns the run's usage, and HTTP 502 with the error when the agent fails", async () => { + const work_dir = temp_work_dir(); + let port: number | undefined; + const adapter = make_adapter( + work_dir, + "unused", + (seen) => { + port = seen; + }, + async (_platform, _chat, _user, text) => + text === "fail" ? { text: "agent error: provider down", failed: true } : { text: `ok:${text}`, usage: usage_small }, + ); + try { + await adapter.start(); + const post = (text: string) => + fetch(`http://127.0.0.1:${String(port)}/message`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ text }), + }); + const ok = await post("hi"); + expect(ok.status).toBe(200); + expect(await ok.json()).toEqual({ reply: "ok:hi", usage: usage_small }); + const failed = await post("fail"); + expect(failed.status).toBe(502); + expect(await failed.json()).toEqual({ error: "agent error: provider down" }); + } finally { + await adapter.stop(); + rmSync(work_dir, { recursive: true, force: true }); + } + }); + it("serves POST /message, GET /health, and 404 for other paths", async () => { const work_dir = temp_work_dir(); let port: number | undefined; diff --git a/test/gateway_conversation_bounds.test.ts b/test/gateway_conversation_bounds.test.ts index fe69f4a..0ca46b1 100644 --- a/test/gateway_conversation_bounds.test.ts +++ b/test/gateway_conversation_bounds.test.ts @@ -123,6 +123,29 @@ describe("gateway max_conversations", () => { expect(c2?.history_len).toBeGreaterThan(0); }); + it("stays within max_conversations when new chats run concurrently", async () => { + const work_dir = make_temp_dir(); + const started: HeldRun[] = []; + const bus = new GatewayBus( + { + config: parse_agent_config({ + providers: [{ kind: "openai_compat", name: "main", model: "mock-model" }], + work_dir, + log_level: "error", + }), + agent_factory: () => holding_agent(started), + }, + { max_conversations: 2, history_cap: 40 }, + ); + const replies = ["a", "b", "c", "d"].map((chat) => bus.handle("webhook", chat, "u", `${chat}1`)); + for (const chat of ["a", "b", "c", "d"]) { + (await wait_for_start(started, `${chat}1`)).release(); + } + await Promise.all(replies); + const histories = (bus as unknown as { histories: Map }).histories; + expect([...histories.keys()]).toEqual(["webhook:c", "webhook:d"]); + }); + it("keeps an in-flight chain after history eviction so later turns stay serialized", async () => { const work_dir = make_temp_dir(); const started: HeldRun[] = []; @@ -148,8 +171,9 @@ describe("gateway max_conversations", () => { const other = bus.handle("webhook", "b", "u", "b1"); const b1 = await wait_for_start(started, "b1"); - const chains = (bus as unknown as { chains: Map> }).chains; - expect(chains.has("webhook:a")).toBe(true); + // Both conversations still hold a queue: evicting a's history does not drop its run queue. + const sessions = (bus as unknown as { sessions: { pending_count(): number } }).sessions; + expect(sessions.pending_count()).toBe(2); const third = bus.handle("webhook", "a", "u", "a3"); await new Promise((resolve) => setImmediate(resolve)); diff --git a/test/gateway_shutdown.test.ts b/test/gateway_shutdown.test.ts new file mode 100644 index 0000000..2802e43 --- /dev/null +++ b/test/gateway_shutdown.test.ts @@ -0,0 +1,51 @@ +/** + * Gateway shutdown (G-10): adapters are stopped and awaited (bounded by a + * timeout) before the bus and agent are released. + */ +import { describe, expect, it, vi } from "vitest"; +import type { Agent } from "../src/agent/agent.js"; +import type { GatewayBus } from "../src/gateway/bus.js"; +import { shutdown_gateway } from "../src/gateway/runner.js"; +import type { PlatformAdapter } from "../src/gateway/types.js"; +import { logger } from "../src/util/log.js"; + +function tracked(order: string[]): { agent: Agent; bus: GatewayBus } { + return { + agent: { close: () => order.push("agent.close") } as unknown as Agent, + bus: { stop: () => order.push("bus.stop") } as unknown as GatewayBus, + }; +} + +describe("shutdown_gateway", () => { + it("waits for every adapter to stop before releasing the bus and agent", async () => { + const order: string[] = []; + const { agent, bus } = tracked(order); + const slow: PlatformAdapter = { + name: "slow", + start: async () => undefined, + stop: () => new Promise((resolve) => setTimeout(() => resolve(order.push("slow.stop") as unknown as void), 20)), + }; + const broken: PlatformAdapter = { + name: "broken", + start: async () => undefined, + stop: async () => { + throw new Error("socket gone"); + }, + }; + vi.spyOn(logger, "warn").mockImplementation(() => undefined); + await shutdown_gateway(agent, bus, [slow, broken]); + expect(order).toEqual(["slow.stop", "bus.stop", "agent.close"]); + vi.restoreAllMocks(); + }); + + it("stops waiting after the timeout when an adapter hangs", async () => { + const order: string[] = []; + const { agent, bus } = tracked(order); + const hung: PlatformAdapter = { name: "hung", start: async () => undefined, stop: () => new Promise(() => undefined) }; + const warn = vi.spyOn(logger, "warn").mockImplementation(() => undefined); + await shutdown_gateway(agent, bus, [hung], 20); + expect(order).toEqual(["bus.stop", "agent.close"]); + expect(String(warn.mock.calls[0]?.[0])).toContain("did not stop within 20ms"); + vi.restoreAllMocks(); + }); +});