From a80fe68d578492bc0f500725df7724df62c5b029 Mon Sep 17 00:00:00 2001 From: zenord Date: Sun, 16 Aug 2026 10:08:03 +0800 Subject: [PATCH] Add ACP-backed role sessions --- README.md | 277 ++++++++---------------- config.example.json | 63 ++++-- gori-agent.sh | 20 +- package-lock.json | 16 +- package.json | 9 +- src/acp/backend-registry.ts | 23 ++ src/acp/backends/kimi.ts | 8 + src/acp/client.ts | 125 +++++++++++ src/acp/discovery.ts | 34 +++ src/acp/probe.ts | 10 + src/acp/session-manager.ts | 150 +++++++++++++ src/acp/types.ts | 44 ++++ src/acp/worker.ts | 128 +++++++++++ src/cli.ts | 152 ++++--------- src/cli/doctor.ts | 83 +++---- src/cli/setup.ts | 333 +++++------------------------ src/config.ts | 187 +++++++++++----- src/core/command-router.ts | 69 ++---- src/core/durable-session-store.ts | 150 +++++++++++++ src/core/gateway.ts | 137 ++++++------ src/platforms/qq/adapter.ts | 17 +- src/platforms/qq/gateway-client.ts | 2 +- src/platforms/qq/types.ts | 2 + src/roles/role-registry.ts | 51 +++++ src/roles/skill-loader.ts | 24 +++ src/server.ts | 114 +++++----- test/acp-client.test.ts | 23 ++ test/acp-session-manager.test.ts | 58 +++++ test/acp-worker.test.ts | 37 ++++ test/config.test.ts | 31 +++ test/durable-session-store.test.ts | 41 ++++ test/fixtures/fake-acp-agent.mjs | 73 +++++++ test/gateway.test.ts | 45 ++++ test/qq-adapter.test.ts | 34 +++ 34 files changed, 1683 insertions(+), 887 deletions(-) create mode 100644 src/acp/backend-registry.ts create mode 100644 src/acp/backends/kimi.ts create mode 100644 src/acp/client.ts create mode 100644 src/acp/discovery.ts create mode 100644 src/acp/probe.ts create mode 100644 src/acp/session-manager.ts create mode 100644 src/acp/types.ts create mode 100644 src/acp/worker.ts create mode 100644 src/core/durable-session-store.ts create mode 100644 src/roles/role-registry.ts create mode 100644 src/roles/skill-loader.ts create mode 100644 test/acp-client.test.ts create mode 100644 test/acp-session-manager.test.ts create mode 100644 test/acp-worker.test.ts create mode 100644 test/config.test.ts create mode 100644 test/durable-session-store.test.ts create mode 100644 test/fixtures/fake-acp-agent.mjs create mode 100644 test/gateway.test.ts create mode 100644 test/qq-adapter.test.ts diff --git a/README.md b/README.md index 8f4988c..37fe327 100644 --- a/README.md +++ b/README.md @@ -1,122 +1,118 @@ -# gori-agent-gateway +# gori-agent -Hermes-gateway-style multi-IM Agent Gateway for routing instant-message webhooks to configurable CLI agents. The primary product flow is `gori-agent setup`: discover local agents first, pick one, then choose the IM platform to connect. +Multi-IM gateway that drives coding agents through the official Agent Client Protocol (ACP): -## Quick start +```text +IM Adapter -> Gateway -> Role/Session Manager -> ACP Client -> kimi acp +``` -Requires Node.js 20+. +The MVP backend is Kimi Code ACP. Gateway code contains no Kimi CLI prompt fallback and no Pi/Codex private protocol. Future agents must enter through an ACP adapter such as `codex-acp` or a Pi ACP adapter. + +## Requirements and quick start + +Requires Node.js 20+ and an authenticated Kimi Code installation. ```bash npm install npm run build ./install.sh gori-agent setup -``` - -`install.sh` installs the command wrapper under `~/.gori-agent/bin`, stores runtime PID/log files under `~/.gori-agent/state` and `~/.gori-agent/logs`, and adds the command to `~/.bashrc`. Open a new terminal or run `source ~/.bashrc` before using `gori-agent` directly. - -If the package bin has not been linked, use the no-link fallback: - -```bash -npm run gori-agent -- setup -``` - -The setup wizard writes `config.json` only after confirmation. `config.json` is ignored by git and should hold local secrets. - -Start the gateway after setup: - -```bash +gori-agent doctor --config ./config.json gori-agent start --config ./config.json ``` -Fallback: - -```bash -npm run gori-agent -- start --config ./config.json -``` - -Convenience executable: - -```bash -gori-agent setup -gori-agent start -gori-agent start-daemon -gori-agent status -gori-agent logs -gori-agent stop -``` - -If the command is not linked into PATH, run the project-local wrapper: +`setup` probes ACP backends, creates an assistant role, optionally creates the `ops` role, preserves all existing platform settings, prints a migration summary, and only writes after confirmation. A v1 config is accepted at runtime only when its default agent is Kimi; other CLI agents require an explicit ACP backend. + +Convenience wrapper commands: ```bash +./gori-agent.sh start +./gori-agent.sh start-daemon ./gori-agent.sh status +./gori-agent.sh logs +./gori-agent.sh stop ``` -Use `--config path` after the command to override the config path. +`stop` sends SIGTERM and waits up to 10 seconds for graceful shutdown. For production, use systemd with `Restart=always` and `KillMode=control-group` so an unexpected gateway crash also cleans up ACP children. -## CLI commands +## CLI ```bash -gori-agent setup -gori-agent discover-agents [--json] +gori-agent setup [--config path] +gori-agent discover-backends [--json] gori-agent start [--config path] gori-agent status [--config path] gori-agent doctor [--config path] gori-agent print feishu [--config path] -gori-agent --help ``` -Package scripts mirror common commands: +`discover-agents` remains a deprecated alias for `discover-backends`. Discovery does not send a prompt. Kimi is ready when `kimi acp` is available; Codex and Pi report `needs-adapter` unless their ACP adapter executable exists. -```bash -npm run setup -npm run discover-agents -npm run doctor -``` +## Roles, sessions, and commands -## Agent-first setup flow +A conversation binding is keyed by `platform + chatId + role`. Each binding stores the ACP backend, native session ID, workspace, and role fingerprint. The selected role and bindings survive gateway restart. -`gori-agent setup` does the following: +- A new session receives a hidden bootstrap prompt containing persona, workspace, skill content, and policy. The binding is saved only after bootstrap succeeds. +- A role persona, workspace, policy, or skill-content change changes its fingerprint and causes a new native session. +- Idle workers are stopped and later cold-resumed with `session/resume` (or `session/load` when resume is unavailable). +- Failed prompts are never replayed automatically. +- `/new` cancels the current turn, unbinds the current chat/role, and creates a native session on the next prompt. It does not delete Kimi history. +- `/cancel` bypasses the per-chat lock and sends ACP `session/cancel`; an unresponsive worker is force-recycled after the grace period. -1. Discovers local CLI agents on `PATH` and always includes the built-in `echo` fallback. -2. Lets you pick one agent. -3. Lets you choose an IM platform: Feishu/Lark, WeChat external webhook, WeCom scaffold, QQ Bot webhook, or Generic webhook. -4. Builds a focused `config.json` containing the chosen default agent, the chosen agent plus `echo`, and the chosen platform enabled. -5. Asks before writing. +Chat commands: -Detected ready agents with known safe prompt modes: +- `/help` +- `/roles` +- `/role ` +- `/status` +- `/cancel` +- `/new` +- `/agents` and `/agent ` (deprecated aliases) -| Agent | Command | Args | Input mode | Permission mode | -| --- | --- | --- | --- | --- | -| Kimi | `kimi` | `-p` | final argv | `kimi -p` already runs non-interactively under Kimi Code's auto permission policy. `--yolo` cannot be combined with `-p`. | -| OpenCode | `opencode` | `run` | final argv | setup can append `--auto`. | -| Codex | `codex` | `exec` | final argv | default; add manual args if needed. | -| Claude | `claude` | `-p` | final argv | default; add manual args if needed. | -| Built-in echo | `node` | `scripts/echo-agent.js` | stdin | n/a | +Messages in one chat are serialized; different chats run concurrently. -When you pick Kimi, setup asks which Kimi Code model to use. It reads your local Kimi model aliases at setup time with `kimi provider list --json`, so custom providers from your actual Kimi config appear in the menu. Choosing the default keeps Kimi Code's configured `default_model`; choosing a model writes `-m -p` into that agent's args. You can also enter a custom model alias. +## Configuration v2 -Gemini, Qwen/Qwen Code, Copilot, Pi, and Hermes are probed with `--version` then `--help` only. If found, they are shown as `needs-config` unless a safe prompt mode is known; the wizard asks you to confirm command, args, input mode, extra auto/yolo permission args, and working directory before enabling them. +See `config.example.json`. Important sections: -Discovery never sends an actual prompt to an agent. Probes use `child_process.spawn(..., { shell: false })`, a 2s timeout, and a 4KB output cap. +- `backends[]`: ACP spawn command and arguments. The default is `/home/ubuntu/.kimi-code/bin/kimi acp`. +- `roles[]`: backend, absolute workspace, persona, skills, and permission policy. +- `skills[]`: readable skill files with a byte limit. Skill contents participate in the role fingerprint. +- `acp.stateFile`: defaults to `$GORI_AGENT_HOME/state/acp-sessions.json`. +- `acp.initializeTimeoutMs`, `promptTimeoutMs`, `cancelGraceMs`: request lifecycle limits. +- `acp.idleTimeoutMs`, `sweepIntervalMs`, `maxProcesses`: worker pool limits. +- `policy.allowedUsers`, `allowedChats`, `requireMentionInGroup`: inbound IM policy. +- `platforms`: Feishu, WeCom, QQ, generic webhook, and Weixin settings. Migration preserves this object, including credentials. -## Platform status +State writes use a serial promise queue, temporary file, fsync, and atomic rename. A single-instance lock protects the state file. Corrupt state is preserved and startup fails explicitly instead of overwriting it. + +## Permission policy + +ACP permission requests are decided by role policy, not by persona: + +- `deny`: reject every permission request. +- `allowlist`: require both an allowed tool name and, for bash/terminal, a matching raw command pattern. If the ACP request lacks enough command detail, it is denied. +- `auto`: approve everything; `doctor` prints a high-risk warning. + +The example `ops` role loads the existing read-only skill file at `/home/ubuntu/gori-space/gori-deploy/.kimi-code/skills/gori-update/SKILL.md`. Its allowlist covers selected `git`, build/deploy entry scripts, and read-only Docker status/log commands. `sudo`, force push, and arbitrary deletion are not allowed. Script-internal operations cannot be inspected by ACP once an explicitly allowed deployment script starts, so script paths must remain trusted. + +Kimi 0.36.1 does not advertise a generic model config option during ACP initialize; this MVP uses the model configured as Kimi's default and `doctor` reports that limitation. + +## Platform behavior | Platform | Inbound | Outbound | Notes | | --- | --- | --- | --- | -| Feishu/Lark | Implemented | Implemented | Replies to `im.message.receive_v1` via `/im/v1/messages/{message_id}/reply`. | -| WeCom | 501 scaffold | Implemented | Uses `gettoken` and `message/send`; inbound callback verification/encryption is not in v1. | -| personal WeChat | External webhook scaffold | Synchronous webhook response | Native iLink/personal WeChat integration is not included in v1. | -| QQ | Implemented | Implemented | Supports QQ official Bot WebSocket gateway mode and optional HTTP callback mode. | -| Generic webhook | Implemented | Synchronous JSON | HMAC-SHA256 signed JSON endpoint for local bridges and tests. | +| Feishu/Lark | Implemented | Implemented | `im.message.receive_v1` and message reply API. | +| WeCom | 501 scaffold | Implemented | Inbound verification/encryption is not implemented. | +| Weixin | External webhook scaffold | Synchronous | Native personal WeChat is not included. | +| QQ | Implemented | Implemented | WebSocket gateway and HTTP callback modes. | +| Generic webhook | Implemented | Synchronous JSON | HMAC-SHA256 signed JSON. | -## Hermes reference +QQ WebSocket mode uses intent `33554432` for `GROUP_AT_MESSAGE_CREATE` and `C2C_MESSAGE_CREATE`. The adapter supports top-level and nested `author.user_openid` fields, preserves callback ACK `{ "op": 12 }` even if ACP work later fails, and logs receive/send routing. No external QQ message is sent by the automated tests. -This project follows the Hermes gateway pattern: platform adapters normalize inbound messages, a central gateway applies policy/session/concurrency handling, then replies are sent through the originating adapter. The Feishu adapter ports the key Hermes behavior for tenant token caching, URL challenge handling, token verification, `im.message.receive_v1` parsing, mention cleanup, and message replies. +Endpoints: -## Endpoints - -- `GET /health` +- `GET /health` — aggregate ACP counters only; native session IDs are not exposed. - `GET /platforms` - `POST /webhook/feishu` - `POST /webhook/wecom` @@ -124,131 +120,36 @@ This project follows the Hermes gateway pattern: platform adapters normalize inb - `POST /webhook/generic` - `POST /webhook/weixin` -## Configuration - -Server execution still supports the existing environment variable: - -```bash -GORI_GATEWAY_CONFIG=./config.json npm start -``` - -The CLI resolves config in this order: - -1. `--config path` -2. `GORI_GATEWAY_CONFIG` -3. `./config.json` -4. seed from `config.example.json` and write to `./config.json` if confirmed - -Important sections: - -- `server.host`: bind address, default `0.0.0.0`. -- `server.port`: gateway port, default `3000`. -- `server.publicBaseUrl`: public HTTPS base URL used by print/setup hints, for example `https://agent.example.com`. -- `policy.allowedUsers`: allow only listed normalized user IDs when non-empty. -- `policy.allowedChats`: allow only listed normalized chat IDs when non-empty. -- `policy.requireMentionInGroup`: if true, group messages are ignored unless the adapter reports a bot mention. -- `defaultAgent`: selected when a chat has not chosen an agent. -- `agents[]`: named CLI agents using `child_process.spawn(command, args)` with `shell: false`. - -CLI agent options: - -- `inputMode: "stdin"`: sends the message text to stdin. -- `inputMode: "arg"`: appends the message text as the final argv item. -- `timeoutMs`: kills slow processes. -- `outputMaxBytes`: caps captured stdout/stderr. -- `cwd`: optional working directory for the agent process. - -## Chat commands - -The gateway handles these commands per chat before invoking an agent: - -- `/help` -- `/status` -- `/agents` -- `/agent ` -- `/new` - -## Feishu setup - -The wizard can collect Feishu values and `gori-agent print feishu --config ./config.json` prints the webhook URL and checklist. - -Manual steps: - -1. Create a Feishu/Lark custom app and enable bot messaging. -2. Configure event subscription for `im.message.receive_v1`. -3. Set the request URL to `https:///webhook/feishu`. -4. Put `appId`, `appSecret`, and `verificationToken` into your config. -5. Add bot display names to `platforms.feishu.botNames` so group mention detection can also work from flattened text. - -The adapter accepts Feishu URL verification challenges and returns `{ "challenge": "..." }`. - -## QQ setup - -QQ has two connection modes: - -- `websocket` (default/recommended): the gateway actively connects to QQ with OAuth access token and WebSocket. No public domain or callback URL is needed. -- `webhook`: QQ posts events to a public HTTPS callback URL. - -Manual WebSocket steps: - -1. Create a QQ official Bot. -2. Put `appId` and `clientSecret` into your config. -3. Keep `platforms.qq.connectionMode` as `"websocket"`. -4. Keep the default `intents` value `33554432` (`1 << 25`, `GROUP_AND_C2C_EVENT`) to receive `GROUP_AT_MESSAGE_CREATE` and `C2C_MESSAGE_CREATE`. -5. Run `gori-agent start --config ./config.json`; the process should log `QQ websocket ready` after successful authentication. - -Manual HTTP callback steps: - -1. Set `platforms.qq.connectionMode` to `"webhook"`. -2. Configure the callback URL to `https:///webhook/qq`. QQ callback URLs must use an allowed public HTTPS port such as 443, 8443, 8080, or 80. -3. Put `botSecret` into your config and keep `verifySignature` enabled so `X-Signature-Ed25519` callbacks are verified. - -The webhook adapter handles QQ `op: 13` callback URL validation and returns `{ "plain_token": "...", "signature": "..." }`. `GROUP_AT_MESSAGE_CREATE` and `C2C_MESSAGE_CREATE` callbacks return QQ HTTP callback ACK `{ "op": 12 }` immediately, then reply through the QQ Bot group/C2C message APIs. - -## Generic webhook - -Payload: +Generic webhook payload: ```json { "chat_id": "demo-chat", "user_id": "demo-user", - "text": "hello", - "message_id": "optional-message-id", + "text": "/status", + "message_id": "optional", "is_group": false, "mentions_bot": true } ``` -Signature header: +When `platforms.webhook.secret` is set, provide `X-Gori-Signature: sha256=` over the exact JSON body. -- Header name: `X-Gori-Signature` -- Format: `sha256=` or just `` -- Input: exact JSON request body bytes as sent by the client -- Algorithm: HMAC-SHA256 using `platforms.webhook.secret` +## Lifecycle and diagnostics -Example: +Shutdown order is QQ WebSocket stop, active ACP cancellation, worker termination, state flush/unlock. SIGINT and SIGTERM share the same idempotent shutdown path. Health statistics include active workers, in-flight turns, worker crashes, persisted bindings, and locked chats. + +`doctor` checks role/workspace/skill validity, ACP initialize and restore capabilities, state directory access, permission warnings, configured platforms, and QQ requirements. It performs no prompt and no external QQ test. + +## Development ```bash -body='{"chat_id":"demo-chat","user_id":"demo-user","text":"/status"}' -sig=$(printf '%s' "$body" | openssl dgst -sha256 -hmac 'replace-me' -hex | awk '{print $2}') -curl -sS http://localhost:3000/webhook/generic \ - -H 'Content-Type: application/json' \ - -H "X-Gori-Signature: sha256=$sig" \ - -d "$body" +npm run typecheck +npm test +npm run build +git diff --check ``` -Response shape: +Tests use Node `node:test` through `tsx`. The fake ACP subprocess covers initialize/new/resume/prompt/cancel, persistence, idle cold resume, Gateway commands, and QQ normalization/ACK behavior. -```json -{ "ok": true, "output": "...agent or command output..." } -``` - -## Security notes - -- Do not expose the gateway publicly without HTTPS and upstream authentication/rate limits. -- Keep platform secrets outside source control; use a copied config file or secret manager. -- CLI agents are spawned without a shell, but they can still execute arbitrary local code. Only configure trusted commands. -- Use `allowedUsers` and `allowedChats` in production. -- Keep `requireMentionInGroup` enabled for group chats to avoid accidental agent invocation. -- Set conservative `timeoutMs` and `outputMaxBytes` values for external-facing agents. +Legacy `src/agents/*` source remains only for compatibility/reference and is not imported by Gateway or Server. There is no runtime `CliAgent` fallback. diff --git a/config.example.json b/config.example.json index c365bce..606b594 100644 --- a/config.example.json +++ b/config.example.json @@ -1,4 +1,5 @@ { + "configVersion": 2, "server": { "host": "0.0.0.0", "port": 3000, @@ -9,23 +10,57 @@ "allowedChats": [], "requireMentionInGroup": true }, - "defaultAgent": "echo", - "agents": [ + "acp": { + "stateFile": "", + "initializeTimeoutMs": 10000, + "promptTimeoutMs": 600000, + "cancelGraceMs": 5000, + "idleTimeoutMs": 1800000, + "sweepIntervalMs": 60000, + "maxProcesses": 8 + }, + "backends": [ { - "name": "echo", - "command": "node", - "args": ["scripts/echo-agent.js"], - "inputMode": "stdin", - "timeoutMs": 30000, - "outputMaxBytes": 64000 + "id": "kimi", + "command": "/home/ubuntu/.kimi-code/bin/kimi", + "args": ["acp"], + "env": {} + } + ], + "skills": [ + { + "id": "gori-update", + "file": "/home/ubuntu/gori-space/gori-deploy/.kimi-code/skills/gori-update/SKILL.md", + "maxBytes": 256000 + } + ], + "defaultRole": "assistant", + "roles": [ + { + "id": "assistant", + "backend": "kimi", + "workspace": "/home/ubuntu/gori-space/gori-agent", + "persona": "", + "skills": [], + "policy": { + "permissionMode": "deny", + "allowedTools": [], + "allowedCommandPatterns": [] + } }, { - "name": "kimi", - "command": "kimi", - "args": ["-p"], - "inputMode": "arg", - "timeoutMs": 120000, - "outputMaxBytes": 64000 + "id": "ops", + "backend": "kimi", + "workspace": "/home/ubuntu/gori-space", + "persona": "你是 Gori 团队运维角色。严格遵循 gori-update skill;有风险或需要外部确认时停止并报告。", + "skills": ["gori-update"], + "policy": { + "permissionMode": "allowlist", + "allowedTools": ["read", "grep", "glob", "bash"], + "allowedCommandPatterns": [ + "^(?:.*\\\"command\\\":\\\")?(?:git (?:status|log|diff|pull --ff-only)|bash gori-deploy/(?:build\\.sh|dist/deploy-[a-z-]+\\.sh)|docker (?:ps|logs))" + ] + } } ], "platforms": { diff --git a/gori-agent.sh b/gori-agent.sh index feae508..3773f36 100755 --- a/gori-agent.sh +++ b/gori-agent.sh @@ -20,9 +20,10 @@ Commands: stop Stop background gateway restart Restart background gateway status Show gateway status - doctor Check configuration and local agents + doctor Check configuration, roles, and ACP backends logs Follow gateway log - discover-agents List detected local CLI agents + discover-backends List detected ACP backends + discover-agents Deprecated alias for discover-backends Environment: GORI_GATEWAY_CONFIG Config path, default: $CONFIG_FILE @@ -61,11 +62,20 @@ stop_daemon() { if ! is_running; then rm -f "$PID_FILE" echo "gori-agent is not running" - exit 0 + return 0 fi local pid pid="$(cat "$PID_FILE")" kill "$pid" + local waited=0 + while kill -0 "$pid" 2>/dev/null && [[ $waited -lt 100 ]]; do + sleep 0.1 + waited=$((waited + 1)) + done + if kill -0 "$pid" 2>/dev/null; then + echo "gori-agent did not stop gracefully within 10 seconds, pid $pid" >&2 + return 1 + fi rm -f "$PID_FILE" echo "gori-agent stopped, pid $pid" } @@ -128,9 +138,9 @@ case "$COMMAND" in ensure_build exec node "$ROOT_DIR/dist/cli.js" doctor --config "$CONFIG_FILE" ;; - discover-agents) + discover-backends|discover-agents) ensure_build - exec node "$ROOT_DIR/dist/cli.js" discover-agents --config "$CONFIG_FILE" + exec node "$ROOT_DIR/dist/cli.js" "$COMMAND" --config "$CONFIG_FILE" ;; logs) ensure_state diff --git a/package-lock.json b/package-lock.json index e3c2af3..67e6dcd 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,19 +1,20 @@ { - "name": "gori-agent-gateway", + "name": "gori-agent", "version": "0.1.0", "lockfileVersion": 3, "requires": true, "packages": { "": { - "name": "gori-agent-gateway", + "name": "gori-agent", "version": "0.1.0", "dependencies": { + "@agentclientprotocol/sdk": "1.3.0", "express": "^4.19.2", "ws": "^8.21.3", "zod": "^3.23.8" }, "bin": { - "gori-agent": "dist/cli.js" + "gori-agent": "bin/gori-agent" }, "devDependencies": { "@types/express": "^4.17.21", @@ -26,6 +27,15 @@ "node": ">=20" } }, + "node_modules/@agentclientprotocol/sdk": { + "version": "1.3.0", + "resolved": "https://registry.npmjs.org/@agentclientprotocol/sdk/-/sdk-1.3.0.tgz", + "integrity": "sha512-i3h/efaeuMUFAO1HSfo97QZQnnvMd7wWBYtBsdL6UMZg3a78sk3Ffya5Xu7C7tYsXomXoDXJBAzQF2PcFKAhIQ==", + "license": "Apache-2.0", + "peerDependencies": { + "zod": "^3.25.0 || ^4.0.0" + } + }, "node_modules/@esbuild/aix-ppc64": { "version": "0.28.2", "resolved": "https://registry.npmjs.org/@esbuild/aix-ppc64/-/aix-ppc64-0.28.2.tgz", diff --git a/package.json b/package.json index 29dc67d..856b237 100644 --- a/package.json +++ b/package.json @@ -1,7 +1,7 @@ { - "name": "gori-agent-gateway", + "name": "gori-agent", "version": "0.1.0", - "description": "Hermes-gateway-style multi-IM agent gateway for CLI agents.", + "description": "Hermes-style multi-IM gateway driving coding agents through ACP.", "type": "module", "main": "dist/server.js", "bin": { @@ -14,13 +14,16 @@ "dev": "tsx src/server.ts", "gori-agent": "node dist/cli.js", "setup": "node dist/cli.js setup", + "discover-backends": "node dist/cli.js discover-backends", "discover-agents": "node dist/cli.js discover-agents", - "doctor": "node dist/cli.js doctor" + "doctor": "node dist/cli.js doctor", + "test": "tsx --test test/*.test.ts" }, "engines": { "node": ">=20" }, "dependencies": { + "@agentclientprotocol/sdk": "1.3.0", "express": "^4.19.2", "ws": "^8.21.3", "zod": "^3.23.8" diff --git a/src/acp/backend-registry.ts b/src/acp/backend-registry.ts new file mode 100644 index 0000000..cd019e4 --- /dev/null +++ b/src/acp/backend-registry.ts @@ -0,0 +1,23 @@ +import type { AcpBackendConfig } from "../config.js"; +import type { AcpBackendSpec } from "./types.js"; + +export class AcpBackendRegistry { + private readonly backends = new Map(); + + constructor(configs: AcpBackendConfig[]) { + for (const config of configs) this.register({ ...config }); + } + + register(backend: AcpBackendSpec): void { + if (this.backends.has(backend.id)) throw new Error(`Duplicate ACP backend: ${backend.id}`); + this.backends.set(backend.id, backend); + } + + get(id: string): AcpBackendSpec { + const backend = this.backends.get(id); + if (!backend) throw new Error(`Unknown ACP backend: ${id}`); + return backend; + } + + list(): string[] { return [...this.backends.keys()].sort(); } +} diff --git a/src/acp/backends/kimi.ts b/src/acp/backends/kimi.ts new file mode 100644 index 0000000..d7f4281 --- /dev/null +++ b/src/acp/backends/kimi.ts @@ -0,0 +1,8 @@ +import type { AcpBackendConfig } from "../../config.js"; +import type { AcpBackendSpec } from "../types.js"; + +export function kimiBackend(config: AcpBackendConfig): AcpBackendSpec { + if (config.id !== "kimi") throw new Error(`Expected kimi backend, got '${config.id}'`); + if (!config.args.includes("acp")) throw new Error("Kimi ACP backend args must include 'acp'"); + return { ...config }; +} diff --git a/src/acp/client.ts b/src/acp/client.ts new file mode 100644 index 0000000..3c3ac3b --- /dev/null +++ b/src/acp/client.ts @@ -0,0 +1,125 @@ +import type { ChildProcessWithoutNullStreams } from "node:child_process"; +import { Readable, Writable } from "node:stream"; +import * as acp from "@agentclientprotocol/sdk"; +import type { AgentCapabilities, InitializeResponse, RequestPermissionRequest, RequestPermissionResponse, SessionNotification } from "@agentclientprotocol/sdk"; +import type { RolePolicy } from "../config.js"; + +export interface AcpClientOptions { + initializeTimeoutMs: number; + policy: RolePolicy; +} + +export class AcpClient { + private readonly connection: acp.ClientConnection; + private capabilities: AgentCapabilities = {}; + private activeSessionId?: string; + private collecting = false; + private chunks: string[] = []; + + constructor(private readonly child: ChildProcessWithoutNullStreams, private readonly options: AcpClientOptions) { + const app = acp.client({ name: "gori-agent" }) + .onRequest(acp.methods.client.session.requestPermission, ({ params }) => decidePermission(params, options.policy)) + .onNotification(acp.methods.client.session.update, ({ params }) => this.handleUpdate(params)); + const stream = acp.ndJsonStream( + Writable.toWeb(child.stdin) as WritableStream, + Readable.toWeb(child.stdout) as ReadableStream + ); + this.connection = app.connect(stream); + } + + async initialize(): Promise { + const response = await withTimeout(this.connection.agent.request(acp.methods.agent.initialize, { + protocolVersion: acp.PROTOCOL_VERSION, + clientCapabilities: {}, + clientInfo: { name: "gori-agent", version: "0.1.0" } + }), this.options.initializeTimeoutMs, "ACP initialize timed out"); + if (response.protocolVersion !== acp.PROTOCOL_VERSION) throw new Error(`Unsupported ACP protocol version ${response.protocolVersion}`); + this.capabilities = response.agentCapabilities || {}; + return response; + } + + async newSession(cwd: string): Promise { + const response = await this.connection.agent.request(acp.methods.agent.session.new, { cwd, mcpServers: [] }); + this.activeSessionId = response.sessionId; + return response.sessionId; + } + + async resumeSession(sessionId: string, cwd: string): Promise { + this.collecting = false; + this.chunks = []; + if (this.capabilities.sessionCapabilities?.resume) { + await this.connection.agent.request(acp.methods.agent.session.resume, { sessionId, cwd, mcpServers: [] }); + } else if (this.capabilities.loadSession) { + await this.connection.agent.request(acp.methods.agent.session.load, { sessionId, cwd, mcpServers: [] }); + } else { + throw new Error("ACP backend cannot resume or load sessions"); + } + this.activeSessionId = sessionId; + this.chunks = []; + } + + async prompt(text: string, cancellationSignal?: AbortSignal): Promise { + if (!this.activeSessionId) throw new Error("ACP session is not active"); + this.chunks = []; + this.collecting = true; + try { + await this.connection.agent.request(acp.methods.agent.session.prompt, { + sessionId: this.activeSessionId, + prompt: [{ type: "text", text }] + }, cancellationSignal ? { cancellationSignal } : undefined); + return this.chunks.join("").trim(); + } finally { + this.collecting = false; + } + } + + async cancel(): Promise { + if (this.activeSessionId) await this.connection.agent.notify(acp.methods.agent.session.cancel, { sessionId: this.activeSessionId }); + } + + async closeSession(): Promise { + if (this.activeSessionId && this.capabilities.sessionCapabilities?.close) { + await this.connection.agent.request(acp.methods.agent.session.close, { sessionId: this.activeSessionId }).catch(() => undefined); + } + } + + close(error?: unknown): void { this.connection.close(error); } + + private handleUpdate(notification: SessionNotification): void { + if (!this.collecting || notification.sessionId !== this.activeSessionId) return; + const update = notification.update; + if (update.sessionUpdate === "agent_message_chunk" && update.content.type === "text") this.chunks.push(update.content.text); + } +} + +export function decidePermission(request: RequestPermissionRequest, policy: RolePolicy): RequestPermissionResponse { + if (policy.permissionMode === "deny") return reject(request); + const allowOption = request.options.find((option) => option.kind === "allow_once") || request.options.find((option) => option.kind === "allow_always"); + if (!allowOption) return reject(request); + if (policy.permissionMode === "auto") return { outcome: { outcome: "selected", optionId: allowOption.optionId } }; + + const name = String(request.toolCall.name || request.toolCall.kind || "").toLowerCase(); + const title = String(request.toolCall.title || "").toLowerCase(); + const allowedTool = policy.allowedTools.some((tool) => name === tool.toLowerCase() || title.startsWith(tool.toLowerCase())); + if (!allowedTool) return reject(request); + if (name === "bash" || name === "terminal" || title.startsWith("bash") || title.startsWith("terminal")) { + if (request.toolCall.rawInput === undefined || policy.allowedCommandPatterns.length === 0) return reject(request); + const input = typeof request.toolCall.rawInput === "string" ? request.toolCall.rawInput : JSON.stringify(request.toolCall.rawInput); + if (!policy.allowedCommandPatterns.some((pattern) => new RegExp(pattern).test(input))) return reject(request); + } + return { outcome: { outcome: "selected", optionId: allowOption.optionId } }; +} + +function reject(request: RequestPermissionRequest): RequestPermissionResponse { + const option = request.options.find((item) => item.kind === "reject_once") || request.options.find((item) => item.kind === "reject_always"); + return option ? { outcome: { outcome: "selected", optionId: option.optionId } } : { outcome: { outcome: "cancelled" } }; +} + +async function withTimeout(promise: Promise, timeoutMs: number, message: string): Promise { + let timer: NodeJS.Timeout | undefined; + try { + return await Promise.race([promise, new Promise((_resolve, reject) => { timer = setTimeout(() => reject(new Error(message)), timeoutMs); })]); + } finally { + if (timer) clearTimeout(timer); + } +} diff --git a/src/acp/discovery.ts b/src/acp/discovery.ts new file mode 100644 index 0000000..99272aa --- /dev/null +++ b/src/acp/discovery.ts @@ -0,0 +1,34 @@ +import { spawn } from "node:child_process"; +import fs from "node:fs"; +import path from "node:path"; + +export type BackendStatus = "ready" | "not-found" | "needs-adapter"; +export interface DiscoveredBackend { id: string; command: string; args: string[]; status: BackendStatus; version?: string; reason?: string } + +export async function discoverBackends(): Promise { + const kimi = findExecutable("kimi") || (fs.existsSync("/home/ubuntu/.kimi-code/bin/kimi") ? "/home/ubuntu/.kimi-code/bin/kimi" : undefined); + return [ + kimi ? { id: "kimi", command: kimi, args: ["acp"], status: "ready", version: await version(kimi) } : { id: "kimi", command: "kimi", args: ["acp"], status: "not-found", reason: "Kimi executable not found" }, + { id: "codex", command: "codex-acp", args: [], status: findExecutable("codex-acp") ? "ready" : "needs-adapter", reason: "Requires codex-acp adapter" }, + { id: "pi", command: "pi-acp", args: [], status: findExecutable("pi-acp") ? "ready" : "needs-adapter", reason: "Requires a Pi ACP adapter" } + ]; +} + +function findExecutable(command: string): string | undefined { + for (const entry of (process.env.PATH || "").split(path.delimiter)) { + const candidate = path.join(entry, command); + try { fs.accessSync(candidate, fs.constants.X_OK); return candidate; } catch { /* continue */ } + } + return undefined; +} + +function version(command: string): Promise { + return new Promise((resolve) => { + const child = spawn(command, ["--version"], { stdio: ["ignore", "pipe", "pipe"] }); + let output = ""; + const timer = setTimeout(() => { child.kill(); resolve(undefined); }, 2_000); + child.stdout.on("data", (chunk) => { output += chunk; }); + child.on("close", () => { clearTimeout(timer); resolve(output.trim().split(/\r?\n/)[0] || undefined); }); + child.on("error", () => { clearTimeout(timer); resolve(undefined); }); + }); +} diff --git a/src/acp/probe.ts b/src/acp/probe.ts new file mode 100644 index 0000000..c7a5aac --- /dev/null +++ b/src/acp/probe.ts @@ -0,0 +1,10 @@ +import { spawn } from "node:child_process"; +import type { AcpBackendConfig } from "../config.js"; +import type { InitializeResponse } from "@agentclientprotocol/sdk"; +import { AcpClient } from "./client.js"; + +export async function probeAcpBackend(backend: AcpBackendConfig, timeoutMs = 10_000): Promise { + const child = spawn(backend.command, backend.args, { stdio: ["pipe", "pipe", "pipe"], env: { ...process.env, ...backend.env } }); + const client = new AcpClient(child, { initializeTimeoutMs: timeoutMs, policy: { permissionMode: "deny", allowedTools: [], allowedCommandPatterns: [] } }); + try { return await client.initialize(); } finally { client.close(); child.kill("SIGTERM"); } +} diff --git a/src/acp/session-manager.ts b/src/acp/session-manager.ts new file mode 100644 index 0000000..36d1436 --- /dev/null +++ b/src/acp/session-manager.ts @@ -0,0 +1,150 @@ +import type { AcpConfig } from "../config.js"; +import { chatKeyFor, type DurableSessionStore, type SessionBinding } from "../core/durable-session-store.js"; +import type { RoleRegistry } from "../roles/role-registry.js"; +import type { AcpBackendRegistry } from "./backend-registry.js"; +import type { ConversationRequest, ConversationResponse, ConversationRuntime, RuntimeStats } from "./types.js"; +import { AcpWorker } from "./worker.js"; + +export class AcpSessionManager implements ConversationRuntime { + private readonly workers = new Map(); + private readonly inFlight = new Map(); + private readonly sweeper: NodeJS.Timeout; + private crashes = 0; + private shuttingDown = false; + + constructor( + private readonly config: AcpConfig, + private readonly backends: AcpBackendRegistry, + private readonly roles: RoleRegistry, + private readonly store: DurableSessionStore + ) { + this.sweeper = setInterval(() => void this.sweep(), config.sweepIntervalMs); + this.sweeper.unref(); + } + + async prompt(request: ConversationRequest): Promise { + if (this.shuttingDown) throw new Error("ACP runtime is shutting down"); + const chatKey = chatKeyFor(request.platform, request.chatId); + const role = this.roles.get(this.store.getSelectedRole(chatKey, this.roles.defaultId())); + let binding = this.store.getBinding(chatKey, role.id); + if (binding && (binding.roleFingerprint !== role.fingerprint || binding.backendId !== role.backend || binding.workspace !== role.workspace)) { + await this.dropBinding(binding); + binding = undefined; + } + + const worker = await this.acquireWorker(role.id, binding); + this.inFlight.set(chatKey, worker); + try { + if (!binding) { + const now = Date.now(); + await worker.prompt(role.bootstrap); + binding = { + chatKey, roleId: role.id, backendId: role.backend, nativeSessionId: worker.nativeSessionId!, + workspace: role.workspace, roleFingerprint: role.fingerprint, createdAt: now, updatedAt: now + }; + await this.store.setBinding(binding); + } + const text = await worker.prompt(request.text); + await this.store.touchBinding(chatKey, role.id); + return { text: text || "(ACP agent returned no text)", roleId: role.id, backendId: role.backend }; + } catch (error) { + if (worker.nativeSessionId) this.workers.delete(workerKey(worker.backend.id, worker.nativeSessionId)); + await worker.terminate(); + throw error; + } finally { + if (this.inFlight.get(chatKey) === worker) this.inFlight.delete(chatKey); + } + } + + async cancel(platform: string, chatId: string): Promise { + const worker = this.inFlight.get(chatKeyFor(platform, chatId)); + return worker ? worker.cancel() : false; + } + + async reset(platform: string, chatId: string): Promise { + const chatKey = chatKeyFor(platform, chatId); + await this.cancel(platform, chatId); + const roleId = this.store.getSelectedRole(chatKey, this.roles.defaultId()); + const binding = await this.store.deleteBinding(chatKey, roleId); + if (binding) await this.stopWorker(binding.backendId, binding.nativeSessionId); + } + + async selectRole(platform: string, chatId: string, roleId: string): Promise { + this.roles.get(roleId); + await this.store.setSelectedRole(chatKeyFor(platform, chatId), roleId); + } + + selectedRole(platform: string, chatId: string): string { + return this.store.getSelectedRole(chatKeyFor(platform, chatId), this.roles.defaultId()); + } + + status(platform: string, chatId: string): Record { + const chatKey = chatKeyFor(platform, chatId); + const roleId = this.store.getSelectedRole(chatKey, this.roles.defaultId()); + const binding = this.store.getBinding(chatKey, roleId); + return { role: roleId, backend: this.roles.get(roleId).backend, persisted: Boolean(binding), running: this.inFlight.has(chatKey) }; + } + + stats(): RuntimeStats { + return { activeWorkers: this.workers.size, inFlight: this.inFlight.size, crashes: this.crashes, persistedBindings: this.store.stats().bindings }; + } + + async shutdown(): Promise { + if (this.shuttingDown) return; + this.shuttingDown = true; + clearInterval(this.sweeper); + await Promise.all([...this.inFlight.values()].map((worker) => worker.cancel().catch(() => false))); + await Promise.all([...this.workers.values()].map((worker) => worker.terminate())); + this.workers.clear(); + this.inFlight.clear(); + } + + private async acquireWorker(roleId: string, binding?: SessionBinding): Promise { + const role = this.roles.get(roleId); + if (binding) { + const existing = this.workers.get(workerKey(binding.backendId, binding.nativeSessionId)); + if (existing) return existing; + } + await this.ensureCapacity(); + const worker = new AcpWorker(this.backends.get(role.backend), role, this.config, (crashed, error) => { + this.crashes++; + if (crashed.nativeSessionId) this.workers.delete(workerKey(crashed.backend.id, crashed.nativeSessionId)); + console.error(`ACP worker crash: ${error.message}`); + }); + const nativeSessionId = await worker.start(binding?.nativeSessionId); + this.workers.set(workerKey(role.backend, nativeSessionId), worker); + return worker; + } + + private async ensureCapacity(): Promise { + if (this.workers.size < this.config.maxProcesses) return; + const candidate = [...this.workers.entries()].filter(([, worker]) => !worker.inFlight).sort((a, b) => a[1].lastUsedAt - b[1].lastUsedAt)[0]; + if (!candidate) throw new Error(`ACP worker limit reached (${this.config.maxProcesses})`); + this.workers.delete(candidate[0]); + await candidate[1].terminate(); + } + + private async sweep(): Promise { + const cutoff = Date.now() - this.config.idleTimeoutMs; + const expired = [...this.workers.entries()].filter(([, worker]) => !worker.inFlight && worker.lastUsedAt < cutoff); + for (const [key, worker] of expired) { + this.workers.delete(key); + await worker.terminate(); + } + } + + private async dropBinding(binding: SessionBinding): Promise { + await this.store.deleteBinding(binding.chatKey, binding.roleId); + await this.stopWorker(binding.backendId, binding.nativeSessionId); + } + + private async stopWorker(backendId: string, nativeSessionId: string): Promise { + const key = workerKey(backendId, nativeSessionId); + const worker = this.workers.get(key); + if (!worker) return; + this.workers.delete(key); + await worker.terminate(); + } +} + +function workerKey(backendId: string, nativeSessionId: string): string { return `${backendId}:${nativeSessionId}`; } diff --git a/src/acp/types.ts b/src/acp/types.ts new file mode 100644 index 0000000..41f0359 --- /dev/null +++ b/src/acp/types.ts @@ -0,0 +1,44 @@ +import type { RolePolicy } from "../config.js"; + +export interface AcpBackendSpec { + id: string; + command: string; + args: string[]; + env: Record; +} + +export interface ConversationRequest { + platform: string; + chatId: string; + userId: string; + text: string; + messageId?: string; +} + +export interface ConversationResponse { + text: string; + roleId: string; + backendId: string; +} + +export interface RuntimeStats { + activeWorkers: number; + inFlight: number; + crashes: number; + persistedBindings: number; +} + +export interface ConversationRuntime { + prompt(request: ConversationRequest): Promise; + cancel(platform: string, chatId: string): Promise; + reset(platform: string, chatId: string): Promise; + selectRole(platform: string, chatId: string, roleId: string): Promise; + selectedRole(platform: string, chatId: string): string; + status(platform: string, chatId: string): Record; + stats(): RuntimeStats; + shutdown(): Promise; +} + +export interface PermissionContext { + policy: RolePolicy; +} diff --git a/src/acp/worker.ts b/src/acp/worker.ts new file mode 100644 index 0000000..335fbc4 --- /dev/null +++ b/src/acp/worker.ts @@ -0,0 +1,128 @@ +import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process"; +import type { AcpConfig } from "../config.js"; +import type { ResolvedRole } from "../roles/role-registry.js"; +import { AcpClient } from "./client.js"; +import type { AcpBackendSpec } from "./types.js"; + +export class AcpWorker { + private child?: ChildProcessWithoutNullStreams; + private client?: AcpClient; + private abort?: AbortController; + private exited = false; + private stopping = false; + private stderrBytes = 0; + nativeSessionId?: string; + lastUsedAt = Date.now(); + inFlight = false; + + constructor( + readonly backend: AcpBackendSpec, + readonly role: ResolvedRole, + private readonly config: AcpConfig, + private readonly onCrash: (worker: AcpWorker, error: Error) => void + ) {} + + async start(nativeSessionId?: string): Promise { + this.child = spawn(this.backend.command, this.backend.args, { + cwd: this.role.workspace, + env: { ...process.env, ...this.backend.env }, + shell: false, + stdio: ["pipe", "pipe", "pipe"] + }); + this.child.stderr.on("data", (chunk: Buffer) => this.logStderr(chunk)); + this.child.once("error", (error) => this.crashed(error)); + this.child.once("close", (code, signal) => { + this.exited = true; + if (!this.stopping) this.crashed(new Error(`ACP worker exited (code=${code ?? "null"}, signal=${signal ?? "null"})`)); + }); + this.client = new AcpClient(this.child, { initializeTimeoutMs: this.config.initializeTimeoutMs, policy: this.role.policy }); + try { + await this.client.initialize(); + if (nativeSessionId) { + await this.client.resumeSession(nativeSessionId, this.role.workspace); + this.nativeSessionId = nativeSessionId; + } else { + this.nativeSessionId = await this.client.newSession(this.role.workspace); + } + return this.nativeSessionId; + } catch (error) { + await this.terminate(); + throw error; + } + } + + async prompt(text: string): Promise { + if (!this.client || !this.nativeSessionId || this.exited) throw new Error("ACP worker is not available"); + if (this.inFlight) throw new Error("ACP worker already has an in-flight turn"); + this.inFlight = true; + this.lastUsedAt = Date.now(); + this.abort = new AbortController(); + let timeout: NodeJS.Timeout | undefined; + const timeoutPromise = new Promise((_resolve, reject) => { + timeout = setTimeout(() => { + void this.client?.cancel(); + this.abort?.abort(); + reject(new Error(`ACP prompt timed out after ${this.config.promptTimeoutMs}ms`)); + setTimeout(() => { if (!this.exited) void this.terminate(); }, this.config.cancelGraceMs).unref(); + }, this.config.promptTimeoutMs); + }); + try { + return await Promise.race([this.client.prompt(text, this.abort.signal), timeoutPromise]); + } finally { + if (timeout) clearTimeout(timeout); + this.inFlight = false; + this.abort = undefined; + this.lastUsedAt = Date.now(); + } + } + + async cancel(): Promise { + if (!this.inFlight || !this.client) return false; + await this.client.cancel(); + this.abort?.abort(); + setTimeout(() => { if (this.inFlight) void this.terminate(); }, this.config.cancelGraceMs).unref(); + return true; + } + + async terminate(): Promise { + if (this.stopping) return; + this.stopping = true; + if (this.client && !this.exited) { + await Promise.race([ + this.client.closeSession().catch(() => undefined), + new Promise((resolve) => setTimeout(resolve, this.config.cancelGraceMs)) + ]); + } + this.client?.close(); + if (this.child && !this.exited) { + this.child.kill("SIGTERM"); + await waitForExit(this.child, this.config.cancelGraceMs); + if (!this.exited) { + this.child.kill("SIGKILL"); + await waitForExit(this.child, this.config.cancelGraceMs); + } + } + } + + private logStderr(chunk: Buffer): void { + const remaining = Math.max(0, 16_384 - this.stderrBytes); + if (!remaining) return; + const text = chunk.subarray(0, remaining).toString("utf8").trimEnd(); + this.stderrBytes += Buffer.byteLength(text); + if (text) console.error(`[acp:${this.backend.id}] ${text}`); + } + + private crashed(error: Error): void { + if (this.stopping) return; + this.stopping = true; + this.onCrash(this, error); + } +} + +function waitForExit(child: ChildProcessWithoutNullStreams, timeoutMs: number): Promise { + if (child.exitCode !== null || child.signalCode !== null) return Promise.resolve(); + return new Promise((resolve) => { + const timer = setTimeout(resolve, timeoutMs); + child.once("close", () => { clearTimeout(timer); resolve(); }); + }); +} diff --git a/src/cli.ts b/src/cli.ts index bde7469..bf5b1e9 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -1,145 +1,85 @@ #!/usr/bin/env node import process from "node:process"; -import { discoverAgents } from "./agents/discovery.js"; +import { discoverBackends } from "./acp/discovery.js"; import { loadConfigFile } from "./cli/config-file.js"; import { runDoctor } from "./cli/doctor.js"; +import { localBaseUrl } from "./cli/net.js"; import { printFeishu } from "./cli/print.js"; import { runSetup } from "./cli/setup.js"; -import { localBaseUrl } from "./cli/net.js"; +import { defaultStateFile } from "./config.js"; import { startServer } from "./server.js"; -interface ParsedArgs { - command?: string; - rest: string[]; - configPath?: string; - json: boolean; - help: boolean; -} +interface ParsedArgs { command?: string; rest: string[]; configPath?: string; json: boolean; help: boolean } async function main(argv: string[]): Promise { const parsed = parseArgs(argv); - if (parsed.help || !parsed.command) { - printHelp(); - return 0; - } - - if (parsed.command === "setup") { - await runSetup(parsed.configPath); - return 0; - } - - if (parsed.command === "discover-agents") { - const agents = await discoverAgents(); - if (parsed.json) { - console.log(JSON.stringify(agents, null, 2)); - } else { - for (const agent of agents) { - const version = agent.version ? ` (${agent.version})` : ""; - const reason = agent.reason ? ` - ${agent.reason}` : ""; - console.log(`${agent.name}\t${agent.status}\t${agent.command} ${agent.args.join(" ")}${version}${reason}`.trim()); - } + if (parsed.help || !parsed.command) { printHelp(); return 0; } + if (parsed.command === "setup") { await runSetup(parsed.configPath); return 0; } + if (parsed.command === "discover-backends" || parsed.command === "discover-agents") { + const backends = await discoverBackends(); + if (parsed.json) console.log(JSON.stringify(backends, null, 2)); + else for (const backend of backends) { + const version = backend.version ? ` (${backend.version})` : ""; + const reason = backend.reason ? ` - ${backend.reason}` : ""; + console.log(`${backend.id}\t${backend.status}\t${backend.command} ${backend.args.join(" ")}${version}${reason}`.trim()); } return 0; } - if (parsed.command === "start") { const loaded = loadConfigFile(parsed.configPath); - startServer(loaded.config); - return await new Promise(() => undefined); + const running = await startServer(loaded.config); + return await new Promise((resolve) => { + let stopping = false; + const stop = (signal: string): void => { + if (stopping) return; + stopping = true; + console.log(`Received ${signal}; shutting down...`); + void running.shutdown().then(() => resolve(0), (error) => { console.error(error); resolve(1); }); + }; + process.once("SIGINT", () => stop("SIGINT")); + process.once("SIGTERM", () => stop("SIGTERM")); + }); } - - if (parsed.command === "status") { - await printStatus(parsed.configPath); - return 0; - } - - if (parsed.command === "doctor") { - const loaded = loadConfigFile(parsed.configPath); - return runDoctor(loaded.config, loaded.path); - } - + if (parsed.command === "status") { await printStatus(parsed.configPath); return 0; } + if (parsed.command === "doctor") { const loaded = loadConfigFile(parsed.configPath); return runDoctor(loaded.config, loaded.path); } if (parsed.command === "print") { - const topic = parsed.rest[0]; const loaded = loadConfigFile(parsed.configPath); - if (topic === "feishu") { - await printFeishu(loaded.config); - return 0; - } - console.error(`Unknown print topic: ${topic || "(missing)"}`); - console.error("Available: feishu"); - return 1; + if (parsed.rest[0] === "feishu") { await printFeishu(loaded.config); return 0; } + console.error(`Unknown print topic: ${parsed.rest[0] || "(missing)"}`); return 1; } - - console.error(`Unknown command: ${parsed.command}`); - printHelp(); - return 1; + console.error(`Unknown command: ${parsed.command}`); printHelp(); return 1; } function parseArgs(argv: string[]): ParsedArgs { - const rest: string[] = []; - let command: string | undefined; - let configPath: string | undefined; - let json = false; - let help = false; - + const rest: string[] = []; let command: string | undefined; let configPath: string | undefined; let json = false; let help = false; for (let index = 0; index < argv.length; index++) { const arg = argv[index]; - if (arg === "--help" || arg === "-h") { - help = true; - } else if (arg === "--json") { - json = true; - } else if (arg === "--config") { - configPath = argv[++index]; - if (!configPath) throw new Error("--config requires a path"); - } else if (!command) { - command = arg; - } else { - rest.push(arg); - } + if (arg === "--help" || arg === "-h") help = true; + else if (arg === "--json") json = true; + else if (arg === "--config") { configPath = argv[++index]; if (!configPath) throw new Error("--config requires a path"); } + else if (!command) command = arg; else rest.push(arg); } - return { command, rest, configPath, json, help }; } async function printStatus(configPath?: string): Promise { - const loaded = loadConfigFile(configPath); - const baseUrl = localBaseUrl(loaded.config); + const loaded = loadConfigFile(configPath); const baseUrl = localBaseUrl(loaded.config); console.log(`Config: ${loaded.path}${loaded.exists ? "" : " (seeded from config.example.json)"}`); console.log(`Server: ${loaded.config.server.host}:${loaded.config.server.port}`); - console.log(`Default agent: ${loaded.config.defaultAgent}`); - console.log(`Agents: ${loaded.config.agents.map((agent) => agent.name).join(", ")}`); - console.log(`Enabled platforms: ${Object.entries(loaded.config.platforms).filter(([, value]) => value.enabled).map(([name]) => name).join(", ") || "none"}`); - + console.log(`Default role: ${loaded.config.defaultRole}`); + console.log(`Roles: ${loaded.config.roles.map(({ id }) => id).join(", ")}`); + console.log(`Backends: ${loaded.config.backends.map(({ id }) => id).join(", ")}`); + console.log(`State file: ${defaultStateFile(loaded.config)}`); try { - const controller = new AbortController(); - const timer = setTimeout(() => controller.abort(), 1_000); - const response = await fetch(`${baseUrl}/health`, { signal: controller.signal }); - clearTimeout(timer); - console.log(`Health probe: HTTP ${response.status}`); - } catch { - console.log("Health probe: not reachable on local URL"); - } + const response = await fetch(`${baseUrl}/health`, { signal: AbortSignal.timeout(1_000) }); + console.log(`Health probe: HTTP ${response.status} ${await response.text()}`); + } catch { console.log("Health probe: not reachable on local URL"); } } function printHelp(): void { - console.log(`gori-agent - local CLI-agent gateway setup and operations - -Usage: - gori-agent setup [--config path] - gori-agent discover-agents [--json] - gori-agent start [--config path] - gori-agent status [--config path] - gori-agent doctor [--config path] - gori-agent print feishu [--config path] - gori-agent --help - -No-link fallback: - npm run gori-agent -- setup`); + console.log(`gori-agent - multi-IM ACP agent gateway\n\nUsage:\n gori-agent setup [--config path]\n gori-agent discover-backends [--json]\n gori-agent start [--config path]\n gori-agent status [--config path]\n gori-agent doctor [--config path]\n gori-agent print feishu [--config path]\n\nDeprecated alias: discover-agents`); } -main(process.argv.slice(2)).then((code) => { - if (Number.isInteger(code)) process.exitCode = code; -}).catch((error) => { - console.error(error instanceof Error ? error.message : String(error)); - process.exitCode = 1; +main(process.argv.slice(2)).then((code) => { if (Number.isInteger(code)) process.exitCode = code; }).catch((error) => { + console.error(error instanceof Error ? error.message : String(error)); process.exitCode = 1; }); diff --git a/src/cli/doctor.ts b/src/cli/doctor.ts index ad1956b..24ac68b 100644 --- a/src/cli/doctor.ts +++ b/src/cli/doctor.ts @@ -1,63 +1,64 @@ import fs from "node:fs"; -import type { AppConfig } from "../config.js"; -import { discoverAgents } from "../agents/discovery.js"; +import path from "node:path"; +import { probeAcpBackend } from "../acp/probe.js"; +import { defaultStateFile, type AppConfig } from "../config.js"; +import { RoleRegistry } from "../roles/role-registry.js"; import { printUrlHints } from "./net.js"; export async function runDoctor(config: AppConfig, configPath: string): Promise { let problems = 0; console.log(`Config: ${configPath}`); console.log(`Server: ${config.server.host}:${config.server.port}`); - - if (!config.agents.some((agent) => agent.name === config.defaultAgent)) { - console.log(`ERROR: defaultAgent '${config.defaultAgent}' is not in agents[].`); - problems++; - } else { - console.log(`Default agent: ${config.defaultAgent}`); + try { + new RoleRegistry(config); + console.log(`Default role: ${config.defaultRole}`); + } catch (error) { + problems += reportError(error instanceof Error ? error.message : String(error)); } - for (const agent of config.agents) { - if (agent.cwd && !fs.existsSync(agent.cwd)) { - console.log(`WARN: agent '${agent.name}' cwd does not exist: ${agent.cwd}`); + for (const role of config.roles) { + if (!fs.existsSync(role.workspace)) problems += reportError(`role '${role.id}' workspace does not exist: ${role.workspace}`); + if (role.policy.permissionMode === "auto") console.log(`WARN: role '${role.id}' auto-approves every permission request.`); + if (role.policy.permissionMode === "allowlist" && role.policy.allowedTools.includes("bash") && role.policy.allowedCommandPatterns.length === 0) { + console.log(`WARN: role '${role.id}' allows bash by name but has no command patterns; bash requests will be denied.`); } } + for (const backend of config.backends) { + try { + const initialized = await probeAcpBackend(backend, config.acp.initializeTimeoutMs); + const capabilities = initialized.agentCapabilities; + console.log(`ACP '${backend.id}': ready (${initialized.agentInfo?.name || "unknown"} ${initialized.agentInfo?.version || ""})`); + console.log(` load=${Boolean(capabilities?.loadSession)} resume=${Boolean(capabilities?.sessionCapabilities?.resume)} list=${Boolean(capabilities?.sessionCapabilities?.list)} close=${Boolean(capabilities?.sessionCapabilities?.close)}`); + if (!capabilities?.loadSession && !capabilities?.sessionCapabilities?.resume) problems += reportError(`backend '${backend.id}' cannot restore sessions`); + console.log(" model config: not advertised by initialize; using the Kimi default model"); + } catch (error) { + problems += reportError(`ACP '${backend.id}' initialize failed: ${error instanceof Error ? error.message : String(error)}`); + } + } + + const stateFile = defaultStateFile(config); + try { + const stateDirectory = path.dirname(stateFile); + fs.mkdirSync(stateDirectory, { recursive: true }); + fs.accessSync(stateDirectory, fs.constants.R_OK | fs.constants.W_OK); + console.log(`State file: ${stateFile}`); + } catch (error) { + problems += reportError(`state directory is unavailable: ${error instanceof Error ? error.message : String(error)}`); + } + const enabledPlatforms = Object.entries(config.platforms).filter(([, value]) => value.enabled).map(([name]) => name); console.log(`Enabled platforms: ${enabledPlatforms.join(", ") || "none"}`); - if (enabledPlatforms.length === 0) { - console.log("WARN: no IM platform is enabled."); - } - if (config.platforms.feishu.enabled) { - if (!config.platforms.feishu.appId) problems += error("Feishu appId is empty."); - if (!config.platforms.feishu.appSecret || config.platforms.feishu.appSecret === "replace-me") warn("Feishu appSecret is missing or placeholder."); - if (!config.platforms.feishu.verificationToken || config.platforms.feishu.verificationToken === "replace-me") warn("Feishu verificationToken is missing or placeholder."); - } if (config.platforms.qq.enabled) { console.log(`QQ connection mode: ${config.platforms.qq.connectionMode}`); - if (!config.platforms.qq.appId) problems += error("QQ appId is empty."); - if (!config.platforms.qq.clientSecret || config.platforms.qq.clientSecret === "replace-me") warn("QQ clientSecret is missing or placeholder."); - if (config.platforms.qq.connectionMode === "webhook" && config.platforms.qq.verifySignature) { - const callbackSecret = config.platforms.qq.botSecret || config.platforms.qq.clientSecret; - if (!callbackSecret || callbackSecret === "replace-me") problems += error("QQ botSecret or clientSecret is required when verifySignature is true."); + if (!config.platforms.qq.appId) problems += reportError("QQ appId is empty."); + if (!config.platforms.qq.clientSecret || config.platforms.qq.clientSecret === "replace-me") console.log("WARN: QQ clientSecret is missing or placeholder."); + if (config.platforms.qq.connectionMode === "webhook" && config.platforms.qq.verifySignature && !(config.platforms.qq.botSecret || config.platforms.qq.clientSecret)) { + problems += reportError("QQ botSecret or clientSecret is required for webhook signature verification."); } } - - const discovered = await discoverAgents(); - console.log("Discovered local agents:"); - for (const agent of discovered) { - const version = agent.version ? ` (${agent.version})` : ""; - const reason = agent.reason ? ` - ${agent.reason}` : ""; - console.log(`- ${agent.name}: ${agent.status}${version}${reason}`); - } - await printUrlHints(config); return problems > 0 ? 1 : 0; } -function error(message: string): 1 { - console.log(`ERROR: ${message}`); - return 1; -} - -function warn(message: string): void { - console.log(`WARN: ${message}`); -} +function reportError(message: string): 1 { console.log(`ERROR: ${message}`); return 1; } diff --git a/src/cli/setup.ts b/src/cli/setup.ts index 3a0adaa..c2bd8f2 100644 --- a/src/cli/setup.ts +++ b/src/cli/setup.ts @@ -1,297 +1,66 @@ import path from "node:path"; -import type { AppConfig, CliAgentConfig } from "../config.js"; -import type { DiscoveredAgent } from "../agents/discovery.js"; -import { discoverAgents, echoAgent } from "../agents/discovery.js"; -import { listConfiguredKimiModels } from "../agents/kimi-models.js"; -import { createPromptSession, type Choice } from "./prompt.js"; +import { discoverBackends } from "../acp/discovery.js"; +import type { AppConfig } from "../config.js"; +import { createPromptSession } from "./prompt.js"; import { loadConfigFile, projectRoot, writeConfigFile } from "./config-file.js"; -import { configureFeishu } from "./setup-feishu.js"; -import { configureGenericWebhook } from "./setup-webhook.js"; -import { configureWeixin } from "./setup-weixin.js"; -import { printUrlHints } from "./net.js"; - -type PlatformName = keyof AppConfig["platforms"]; - -const PLATFORM_CHOICES: Choice[] = [ - { label: "Feishu/Lark", value: "feishu", hint: "full inbound/outbound adapter" }, - { label: "WeChat external webhook", value: "weixin", hint: "personal WeChat bridge scaffold" }, - { label: "WeCom scaffold", value: "wecom", hint: "inbound is not implemented in v1" }, - { label: "QQ Bot webhook", value: "qq", hint: "official QQ Bot HTTP callback" }, - { label: "Generic webhook", value: "webhook", hint: "signed JSON webhook" } -]; +const GORI_SKILL = "/home/ubuntu/gori-space/gori-deploy/.kimi-code/skills/gori-update/SKILL.md"; export async function runSetup(configPath?: string): Promise { const loaded = loadConfigFile(configPath); const prompt = createPromptSession(); try { - console.log("gori-agent setup"); - console.log(`Config target: ${loaded.path}${loaded.exists ? "" : " (will create from config.example.json)"}`); - await printUrlHints(loaded.config); + console.log("gori-agent ACP setup"); + console.log(`Config target: ${loaded.path}`); + const discovered = await discoverBackends(); + for (const backend of discovered) console.log(`- ${backend.id}: ${backend.status}${backend.version ? ` (${backend.version})` : ""}${backend.reason ? ` - ${backend.reason}` : ""}`); + const kimi = discovered.find((backend) => backend.id === "kimi" && backend.status === "ready"); + if (!kimi) throw new Error("Kimi ACP backend is not available"); - const discovered = await discoverAgents(projectRoot()); - console.log("\nDetected local agents:"); - for (const agent of discovered) { - const version = agent.version ? ` (${agent.version})` : ""; - const reason = agent.reason ? ` - ${agent.reason}` : ""; - console.log(`- ${agent.label} [${agent.name}]: ${agent.status}${version}${reason}`); + const existingRole = loaded.config.roles.find((role) => role.id === loaded.config.defaultRole); + const workspace = path.resolve(await prompt.ask("Assistant workspace", existingRole?.workspace || projectRoot())); + const includeOps = await prompt.askBoolean("Include ops role with the existing gori-update skill", loaded.config.roles.some((role) => role.id === "ops")); + const roles: AppConfig["roles"] = [{ + id: "assistant", backend: "kimi", workspace, persona: existingRole?.persona || "", skills: [], + policy: { permissionMode: "deny", allowedTools: [], allowedCommandPatterns: [] } + }]; + const skills: AppConfig["skills"] = []; + if (includeOps) { + skills.push({ id: "gori-update", file: GORI_SKILL, maxBytes: 256_000 }); + roles.push(opsRole()); } - - const selectedAgent = await prompt.choose( - "\nChoose an agent", - discovered.map((agent) => ({ - label: `${agent.label} (${agent.name})`, - value: agent, - hint: agent.status === "ready" ? agent.command : agent.reason - })), - Math.max(0, discovered.findIndex((agent) => agent.name === loaded.config.defaultAgent)) - ); - const agentConfig = await configureAgent(prompt, selectedAgent, loaded.config); - - const selectedPlatform = await prompt.choose("\nChoose an IM platform", PLATFORM_CHOICES, 0); - const keepOtherPlatforms = await prompt.askBoolean("Preserve existing enabled settings for other platforms", false); - const nextConfig = await buildNextConfig(prompt, loaded.config, agentConfig, selectedPlatform, keepOtherPlatforms); - - console.log("\nPlanned config summary:"); - console.log(`- defaultAgent: ${nextConfig.defaultAgent}`); - console.log(`- agents: ${nextConfig.agents.map((agent) => agent.name).join(", ")}`); - console.log(`- enabled platforms: ${Object.entries(nextConfig.platforms).filter(([, value]) => value.enabled).map(([name]) => name).join(", ") || "none"}`); - console.log("- secrets are not printed"); - - if (await prompt.askBoolean(`Write config to ${loaded.path}`, false)) { - writeConfigFile(loaded.path, nextConfig); - console.log(`Wrote ${loaded.path}`); - } else { - console.log("No changes written."); - } - } finally { - prompt.close(); - } -} - -async function configureAgent(prompt: ReturnType, discovered: DiscoveredAgent, config: AppConfig): Promise { - const existing = config.agents.find((agent) => agent.name === discovered.name); - const defaultCwd = existing?.cwd || discovered.cwd || projectRoot(); - - if (discovered.status === "needs-config") { - console.log(`\n${discovered.label} needs manual command details before it is enabled.`); - const command = await prompt.ask("Command", existing?.command || discovered.command); - const argsText = await prompt.ask("Args, separated by spaces", (existing?.args || discovered.args).join(" ")); - const inputMode = await prompt.choose("Input mode", [ - { label: "Append prompt as final argv", value: "arg" as const }, - { label: "Send prompt to stdin", value: "stdin" as const } - ], (existing?.inputMode || discovered.inputMode) === "stdin" ? 1 : 0); - const extraPermissionArgs = await prompt.ask("Extra auto/yolo permission args, if this agent needs them", ""); - const cwd = await prompt.ask("Working directory", defaultCwd); - return agentFromParts(discovered.name, command, [...splitArgs(argsText), ...splitArgs(extraPermissionArgs)], inputMode, cwd); - } - - const cwd = await prompt.ask("Working directory", defaultCwd); - const command = existing?.command || discovered.command; - const baseArgs = existing?.args || discovered.args; - const modelArgs = discovered.name === "kimi" ? await configureKimiArgs(prompt, baseArgs, command) : baseArgs; - const args = await configurePermissionArgs(prompt, discovered, modelArgs); - return agentFromParts(discovered.name, command, args, existing?.inputMode || discovered.inputMode, cwd); -} - -async function configurePermissionArgs(prompt: ReturnType, discovered: DiscoveredAgent, args: string[]): Promise { - if (!discovered.permissionModes || discovered.permissionModes.length === 0) return args; - - const selected = await prompt.choose( - `Permission mode for ${discovered.label}`, - discovered.permissionModes.map((mode) => ({ - label: mode.label, - value: mode, - hint: mode.warning - })), - 0 - ); - - if (selected.warning) console.log(`Note: ${selected.warning}`); - return mergePermissionArgs(args, selected.args); -} - -function mergePermissionArgs(args: string[], permissionArgs: string[]): string[] { - const merged = [...args]; - for (const arg of permissionArgs) { - if (!merged.includes(arg)) merged.push(arg); - } - return merged; -} - -async function configureKimiArgs(prompt: ReturnType, existingArgs: string[], kimiCommand: string): Promise { - console.log("\nKimi Code supports temporary model selection with --model / -m."); - console.log("Setup reads your local Kimi provider/model list via `kimi provider list --json`."); - console.log("Leave it on default to use default_model from your Kimi Code config."); - - const configuredModels = await listConfiguredKimiModels(kimiCommand); - const modelChoices: Choice[] = [ - { label: "Use Kimi default model", value: "", hint: "respect default_model in Kimi Code config" }, - ...configuredModels.map((model) => ({ - label: model.id, - value: model.id, - hint: model.source === "configured" ? "from local Kimi provider config" : "built-in fallback" - })), - { label: "Custom model alias", value: "__custom__", hint: "type any provider/model alias" } - ]; - - const existingModel = findKimiModel(existingArgs); - const defaultIndex = existingModel - ? modelChoices.findIndex((choice) => choice.value === existingModel) - : 0; - const modelChoice = await prompt.choose("Choose Kimi Code model", modelChoices, defaultIndex >= 0 ? defaultIndex : modelChoices.length - 1); - const model = modelChoice === "__custom__" - ? await prompt.ask("Custom Kimi model alias", existingModel || configuredModels[0]?.id || "kimi-for-coding") - : modelChoice; - - return withKimiModel(existingArgs, model); -} - -function findKimiModel(args: string[]): string | undefined { - for (let index = 0; index < args.length; index += 1) { - const arg = args[index]; - if ((arg === "-m" || arg === "--model") && args[index + 1]) return args[index + 1]; - if (arg.startsWith("--model=")) return arg.slice("--model=".length); - } - return undefined; -} - -function withKimiModel(args: string[], model: string): string[] { - const cleaned: string[] = []; - for (let index = 0; index < args.length; index += 1) { - const arg = args[index]; - if (arg === "-m" || arg === "--model") { - index += 1; - continue; - } - if (arg.startsWith("--model=")) continue; - cleaned.push(arg); - } - - if (!model) return ensureKimiPromptArg(cleaned); - return ["-m", model, ...ensureKimiPromptArg(cleaned)]; -} - -function ensureKimiPromptArg(args: string[]): string[] { - return args.includes("-p") || args.includes("--prompt") ? args : [...args, "-p"]; -} - -async function buildNextConfig( - prompt: ReturnType, - existing: AppConfig, - selectedAgent: CliAgentConfig, - selectedPlatform: PlatformName, - keepOtherPlatforms: boolean -): Promise { - const platforms: AppConfig["platforms"] = keepOtherPlatforms - ? structuredClone(existing.platforms) - : disableAllPlatforms(existing.platforms); - - if (selectedPlatform === "feishu") { - platforms.feishu = await configureFeishu(prompt, existing.platforms.feishu); - } else if (selectedPlatform === "webhook") { - platforms.webhook = await configureGenericWebhook(prompt, existing.platforms.webhook); - } else if (selectedPlatform === "weixin") { - platforms.weixin = await configureWeixin(prompt, existing.platforms.weixin); - } else if (selectedPlatform === "wecom") { - console.log("\nWarning: WeCom inbound webhook is a scaffold and returns 501 in v1."); - platforms.wecom = { - enabled: true, - corpId: await prompt.ask("WeCom corpId", existing.platforms.wecom.corpId), - agentId: await prompt.ask("WeCom agentId", existing.platforms.wecom.agentId), - secret: await prompt.ask(existing.platforms.wecom.secret ? "WeCom secret (leave blank to keep existing)" : "WeCom secret") || existing.platforms.wecom.secret + const next: AppConfig = { + ...loaded.config, + configVersion: 2, + backends: [{ id: "kimi", command: kimi.command, args: ["acp"], env: {} }], + skills, + defaultRole: "assistant", + roles, + platforms: structuredClone(loaded.config.platforms) }; - } else if (selectedPlatform === "qq") { - const connectionMode = await prompt.choose("QQ connection mode", [ - { label: "WebSocket gateway", value: "websocket" as const, hint: "recommended; no public callback URL needed" }, - { label: "HTTP callback webhook", value: "webhook" as const, hint: "requires public HTTPS callback URL" } - ], existing.platforms.qq.connectionMode === "webhook" ? 1 : 0); + console.log("\nMigration summary:"); + console.log(`- default role: ${next.defaultRole}`); + console.log(`- roles: ${next.roles.map(({ id }) => id).join(", ")}`); + console.log(`- backend: ${kimi.command} acp`); + console.log("- all existing platform settings and credentials are preserved"); + if (await prompt.askBoolean(`Write config to ${loaded.path}`, false)) { + writeConfigFile(loaded.path, next); + console.log(`Wrote ${loaded.path}`); + } else console.log("No changes written."); + } finally { prompt.close(); } +} - if (connectionMode === "websocket") { - console.log("\nQQ WebSocket gateway will actively connect to QQ; no public domain is needed."); - const intentsText = await prompt.ask("QQ gateway intents", String(existing.platforms.qq.intents)); - platforms.qq = { - enabled: true, - connectionMode, - appId: await prompt.ask("QQ appId", existing.platforms.qq.appId), - clientSecret: await prompt.ask(existing.platforms.qq.clientSecret ? "QQ clientSecret (leave blank to keep existing)" : "QQ clientSecret") || existing.platforms.qq.clientSecret, - botSecret: existing.platforms.qq.botSecret, - verifySignature: existing.platforms.qq.verifySignature, - botNames: splitCommaList(await prompt.ask("QQ bot names, comma separated", existing.platforms.qq.botNames.join(", "))), - intents: parsePositiveInt(intentsText, existing.platforms.qq.intents), - shard: existing.platforms.qq.shard - }; - } else { - console.log("\nQQ Bot HTTP callback endpoint: /webhook/qq"); - platforms.qq = { - enabled: true, - connectionMode, - appId: await prompt.ask("QQ appId", existing.platforms.qq.appId), - clientSecret: await prompt.ask(existing.platforms.qq.clientSecret ? "QQ clientSecret (leave blank to keep existing)" : "QQ clientSecret") || existing.platforms.qq.clientSecret, - botSecret: await prompt.ask(existing.platforms.qq.botSecret ? "QQ botSecret for callback signing (leave blank to keep existing)" : "QQ botSecret for callback signing, blank to reuse clientSecret", "") || existing.platforms.qq.botSecret, - verifySignature: await prompt.askBoolean("Verify QQ callback signatures", existing.platforms.qq.verifySignature), - botNames: splitCommaList(await prompt.ask("QQ bot names, comma separated", existing.platforms.qq.botNames.join(", "))), - intents: existing.platforms.qq.intents, - shard: existing.platforms.qq.shard - }; +function opsRole(): AppConfig["roles"][number] { + return { + id: "ops", backend: "kimi", workspace: "/home/ubuntu/gori-space", + persona: "你是 Gori 团队运维角色。严格遵循 gori-update skill;有风险或需要外部确认时停止并报告。", + skills: ["gori-update"], + policy: { + permissionMode: "allowlist", + allowedTools: ["read", "grep", "glob", "bash"], + allowedCommandPatterns: [ + "^(?:.*\\\"command\\\":\\\")?(?:git (?:status|log|diff|pull --ff-only)|bash gori-deploy/(?:build\\.sh|dist/deploy-[a-z-]+\\.sh)|docker (?:ps|logs))" + ] } - } - - const echo = echoConfig(); - const agents = selectedAgent.name === "echo" ? [echo] : [selectedAgent, echo]; - return { - server: existing.server, - policy: existing.policy, - defaultAgent: selectedAgent.name, - agents, - platforms }; } - -function disableAllPlatforms(platforms: AppConfig["platforms"]): AppConfig["platforms"] { - return { - feishu: { ...platforms.feishu, enabled: false }, - wecom: { ...platforms.wecom, enabled: false }, - qq: { ...platforms.qq, enabled: false }, - webhook: { ...platforms.webhook, enabled: false }, - weixin: { ...platforms.weixin, enabled: false } - }; -} - -function echoConfig(): CliAgentConfig { - const echo = echoAgent(projectRoot()); - return { - name: echo.name, - command: echo.command, - args: echo.args, - inputMode: echo.inputMode, - cwd: echo.cwd, - timeoutMs: 30_000, - outputMaxBytes: 64_000 - }; -} - -function agentFromParts(name: string, command: string, args: string[], inputMode: "stdin" | "arg", cwd: string): CliAgentConfig { - return { - name, - command, - args, - inputMode, - cwd: path.resolve(cwd), - timeoutMs: 120_000, - outputMaxBytes: 64_000 - }; -} - -function splitArgs(value: string): string[] { - return value.split(" ").map((part) => part.trim()).filter(Boolean); -} - -function parsePositiveInt(value: string, fallback: number): number { - const parsed = Number.parseInt(value, 10); - return Number.isInteger(parsed) && parsed > 0 ? parsed : fallback; -} - -function splitCommaList(value: string): string[] { - return value.split(",").map((part) => part.trim()).filter(Boolean); -} diff --git a/src/config.ts b/src/config.ts index 0b74f70..f6f3023 100644 --- a/src/config.ts +++ b/src/config.ts @@ -3,102 +3,179 @@ import path from "node:path"; import process from "node:process"; import { z } from "zod"; -const cliAgentSchema = z.object({ - name: z.string().min(1), - command: z.string().min(1), - args: z.array(z.string()).default([]), - inputMode: z.enum(["stdin", "arg"]).default("stdin"), - timeoutMs: z.number().int().positive().default(120_000), - outputMaxBytes: z.number().int().positive().default(64_000), - cwd: z.string().optional() -}); - const policySchema = z.object({ allowedUsers: z.array(z.string()).default([]), allowedChats: z.array(z.string()).default([]), requireMentionInGroup: z.boolean().default(false) }); +const rolePolicySchema = z.object({ + permissionMode: z.enum(["deny", "allowlist", "auto"]).default("deny"), + allowedTools: z.array(z.string()).default([]), + allowedCommandPatterns: z.array(z.string()).default([]) +}); + +const backendSchema = z.object({ + id: z.string().min(1), + command: z.string().min(1), + args: z.array(z.string()).default([]), + env: z.record(z.string()).default({}) +}); + +const skillSchema = z.object({ + id: z.string().min(1), + file: z.string().min(1), + maxBytes: z.number().int().positive().default(256_000) +}); + +const roleSchema = z.object({ + id: z.string().min(1), + backend: z.string().min(1), + workspace: z.string().min(1), + persona: z.string().default(""), + skills: z.array(z.string()).default([]), + policy: rolePolicySchema.default({}) +}); + +const acpSchema = z.object({ + stateFile: z.string().default(""), + initializeTimeoutMs: z.number().int().positive().default(10_000), + promptTimeoutMs: z.number().int().positive().default(600_000), + cancelGraceMs: z.number().int().positive().default(5_000), + idleTimeoutMs: z.number().int().positive().default(1_800_000), + sweepIntervalMs: z.number().int().positive().default(60_000), + maxProcesses: z.number().int().positive().default(8) +}); + const feishuSchema = z.object({ - enabled: z.boolean().default(false), - appId: z.string().default(""), - appSecret: z.string().default(""), - verificationToken: z.string().default(""), - botNames: z.array(z.string()).default([]) + enabled: z.boolean().default(false), appId: z.string().default(""), appSecret: z.string().default(""), + verificationToken: z.string().default(""), botNames: z.array(z.string()).default([]) }); - const wecomSchema = z.object({ - enabled: z.boolean().default(false), - corpId: z.string().default(""), - agentId: z.string().default(""), - secret: z.string().default("") + enabled: z.boolean().default(false), corpId: z.string().default(""), agentId: z.string().default(""), secret: z.string().default("") }); - const qqSchema = z.object({ - enabled: z.boolean().default(false), - connectionMode: z.enum(["websocket", "webhook"]).default("websocket"), - appId: z.string().default(""), - clientSecret: z.string().default(""), - botSecret: z.string().default(""), - verifySignature: z.boolean().default(true), - botNames: z.array(z.string()).default([]), + enabled: z.boolean().default(false), connectionMode: z.enum(["websocket", "webhook"]).default("websocket"), + appId: z.string().default(""), clientSecret: z.string().default(""), botSecret: z.string().default(""), + verifySignature: z.boolean().default(true), botNames: z.array(z.string()).default([]), intents: z.number().int().positive().default(1 << 25), shard: z.tuple([z.number().int().nonnegative(), z.number().int().positive()]).default([0, 1]) }); - -const webhookSchema = z.object({ - enabled: z.boolean().default(true), - secret: z.string().default("") -}); - +const webhookSchema = z.object({ enabled: z.boolean().default(true), secret: z.string().default("") }); const weixinSchema = z.object({ - enabled: z.boolean().default(false), - mode: z.enum(["external-webhook", "not-implemented"]).default("external-webhook"), - secret: z.string().default("") + enabled: z.boolean().default(false), mode: z.enum(["external-webhook", "not-implemented"]).default("external-webhook"), secret: z.string().default("") }); export const configSchema = z.object({ + configVersion: z.literal(2), server: z.object({ - host: z.string().default("0.0.0.0"), - port: z.number().int().positive().max(65_535).default(3000), - publicBaseUrl: z.string().default("") + host: z.string().default("0.0.0.0"), port: z.number().int().positive().max(65_535).default(3000), publicBaseUrl: z.string().default("") }).default({}), policy: policySchema.default({}), - defaultAgent: z.string().min(1).default("echo"), - agents: z.array(cliAgentSchema).min(1), + acp: acpSchema.default({}), + backends: z.array(backendSchema).min(1), + skills: z.array(skillSchema).default([]), + defaultRole: z.string().min(1), + roles: z.array(roleSchema).min(1), platforms: z.object({ - feishu: feishuSchema.default({}), - wecom: wecomSchema.default({}), - qq: qqSchema.default({}), - webhook: webhookSchema.default({}), - weixin: weixinSchema.default({}) + feishu: feishuSchema.default({}), wecom: wecomSchema.default({}), qq: qqSchema.default({}), + webhook: webhookSchema.default({}), weixin: weixinSchema.default({}) }).default({}) }); export type AppConfig = z.infer; -export type CliAgentConfig = z.infer; export type GatewayPolicy = z.infer; +export type AcpConfig = AppConfig["acp"]; +export type AcpBackendConfig = AppConfig["backends"][number]; +export type RoleConfig = AppConfig["roles"][number]; +export type RolePolicy = RoleConfig["policy"]; +export type SkillConfig = AppConfig["skills"][number]; + +// Kept only so legacy CliAgent source remains type-checkable; it is not used by the runtime. +export interface CliAgentConfig { + name: string; + command: string; + args: string[]; + inputMode: "stdin" | "arg"; + timeoutMs: number; + outputMaxBytes: number; + cwd?: string; +} + +interface V1AgentConfig { + name?: string; + command?: string; + args?: string[]; + cwd?: string; +} export function resolveConfigPath(configPath = process.env.GORI_GATEWAY_CONFIG): string { return path.resolve(configPath || "config.example.json"); } export function parseConfig(rawConfig: unknown): AppConfig { - const config = configSchema.parse(rawConfig); + const candidate = isRecord(rawConfig) && rawConfig.configVersion === 2 ? rawConfig : migrateV1Config(rawConfig); + const config = configSchema.parse(candidate); + validateUnique(config.backends.map((item) => item.id), "backend"); + validateUnique(config.roles.map((item) => item.id), "role"); + validateUnique(config.skills.map((item) => item.id), "skill"); - if (!config.agents.some((agent) => agent.name === config.defaultAgent)) { - throw new Error(`defaultAgent '${config.defaultAgent}' is not present in agents`); + const backendIds = new Set(config.backends.map((item) => item.id)); + const skillIds = new Set(config.skills.map((item) => item.id)); + if (!config.roles.some((role) => role.id === config.defaultRole)) throw new Error(`defaultRole '${config.defaultRole}' is not present in roles`); + for (const role of config.roles) { + if (!path.isAbsolute(role.workspace)) throw new Error(`role '${role.id}' workspace must be absolute`); + if (!backendIds.has(role.backend)) throw new Error(`role '${role.id}' references unknown backend '${role.backend}'`); + for (const skill of role.skills) if (!skillIds.has(skill)) throw new Error(`role '${role.id}' references unknown skill '${skill}'`); + for (const pattern of role.policy.allowedCommandPatterns) { + try { new RegExp(pattern); } catch { throw new Error(`role '${role.id}' has invalid command pattern '${pattern}'`); } + } } - return config; } +export function migrateV1Config(rawConfig: unknown): unknown { + if (!isRecord(rawConfig)) throw new Error("configuration must be an object"); + const agents = Array.isArray(rawConfig.agents) ? rawConfig.agents.filter(isRecord) as V1AgentConfig[] : []; + const defaultAgent = typeof rawConfig.defaultAgent === "string" ? rawConfig.defaultAgent : ""; + const selected = agents.find((agent) => agent.name === defaultAgent); + if (!selected || selected.name !== "kimi" || !selected.command || path.basename(selected.command) !== "kimi") { + throw new Error("v1 migration only supports a Kimi default agent; configure an ACP backend explicitly for other agents"); + } + return { + configVersion: 2, + server: rawConfig.server, + policy: rawConfig.policy, + acp: {}, + backends: [{ id: "kimi", command: selected.command, args: ["acp"] }], + skills: [], + defaultRole: "assistant", + roles: [{ + id: "assistant", backend: "kimi", workspace: path.resolve(selected.cwd || process.cwd()), persona: "", skills: [], + policy: { permissionMode: "deny", allowedTools: [], allowedCommandPatterns: [] } + }], + platforms: rawConfig.platforms + }; +} + +export function defaultStateFile(config: AppConfig): string { + if (config.acp.stateFile) return path.resolve(config.acp.stateFile); + const home = process.env.GORI_AGENT_HOME || process.cwd(); + return path.join(path.resolve(home), "state", "acp-sessions.json"); +} + export function loadConfig(configPath = process.env.GORI_GATEWAY_CONFIG): AppConfig { return loadConfigFromPath(resolveConfigPath(configPath)); } export function loadConfigFromPath(configPath: string): AppConfig { - const raw = fs.readFileSync(configPath, "utf8"); - const parsed = JSON.parse(raw) as unknown; - return parseConfig(parsed); + return parseConfig(JSON.parse(fs.readFileSync(configPath, "utf8")) as unknown); +} + +function validateUnique(ids: string[], label: string): void { + if (new Set(ids).size !== ids.length) throw new Error(`${label} IDs must be unique`); +} + +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); } diff --git a/src/core/command-router.ts b/src/core/command-router.ts index 5c46013..afe344f 100644 --- a/src/core/command-router.ts +++ b/src/core/command-router.ts @@ -1,57 +1,26 @@ -import { AgentRegistry } from "../agents/agent-registry.js"; -import { SessionStore } from "./session-store.js"; -import type { IncomingMessage } from "./types.js"; +export type CommandKind = "help" | "roles" | "role" | "status" | "cancel" | "new"; -export interface CommandResult { - handled: boolean; - text?: string; +export interface ParsedCommand { + kind: CommandKind; + argument?: string; + deprecatedAlias?: boolean; } export class CommandRouter { - constructor( - private readonly agents: AgentRegistry, - private readonly sessions: SessionStore - ) {} - - route(message: IncomingMessage, sessionId: string): CommandResult { - const text = message.text.trim(); - if (!text.startsWith("/")) return { handled: false }; - - const [command, ...args] = text.split(/\s+/); - switch (command) { - case "/help": - return { - handled: true, - text: [ - "Commands:", - "/help - show this help", - "/status - show gateway/session status", - "/agents - list available agents", - "/agent - select an agent for this chat", - "/new - reset this chat session" - ].join("\n") - }; - case "/status": { - const session = this.sessions.getSession(sessionId); - return { - handled: true, - text: `OK\nplatform=${message.platform}\nchat=${message.chatId}\nagent=${session.selectedAgent || this.agents.defaultName()}` - }; - } - case "/agents": - return { handled: true, text: `Available agents: ${this.agents.list().join(", ")}` }; - case "/agent": { - const agentName = args[0]; - if (!agentName) return { handled: true, text: "Usage: /agent " }; - if (!this.agents.has(agentName)) return { handled: true, text: `Unknown agent: ${agentName}` }; - this.sessions.setSelectedAgent(sessionId, agentName); - return { handled: true, text: `Selected agent: ${agentName}` }; - } - case "/new": - this.sessions.reset(sessionId); - return { handled: true, text: "Started a new session for this chat." }; - default: - return { handled: false }; + parse(text: string): ParsedCommand | undefined { + const trimmed = text.trim(); + if (!trimmed.startsWith("/")) return undefined; + const [command, ...args] = trimmed.split(/\s+/); + switch (command.toLowerCase()) { + case "/help": return { kind: "help" }; + case "/roles": return { kind: "roles" }; + case "/role": return { kind: "role", argument: args[0] }; + case "/status": return { kind: "status" }; + case "/cancel": return { kind: "cancel" }; + case "/new": return { kind: "new" }; + case "/agents": return { kind: "roles", deprecatedAlias: true }; + case "/agent": return { kind: "role", argument: args[0], deprecatedAlias: true }; + default: return undefined; } } } diff --git a/src/core/durable-session-store.ts b/src/core/durable-session-store.ts new file mode 100644 index 0000000..c29c5de --- /dev/null +++ b/src/core/durable-session-store.ts @@ -0,0 +1,150 @@ +import fs from "node:fs"; +import path from "node:path"; + +export interface ChatState { + selectedRole: string; + createdAt: number; + updatedAt: number; +} + +export interface SessionBinding { + chatKey: string; + roleId: string; + backendId: string; + nativeSessionId: string; + workspace: string; + roleFingerprint: string; + createdAt: number; + updatedAt: number; +} + +interface StoreData { + version: 1; + chats: Record; + bindings: Record; +} + +const EMPTY: StoreData = { version: 1, chats: {}, bindings: {} }; + +export class DurableSessionStore { + private data: StoreData = structuredClone(EMPTY); + private queue: Promise = Promise.resolve(); + private lockFd?: fs.promises.FileHandle; + private closed = false; + + constructor(readonly file: string) {} + + async open(): Promise { + await fs.promises.mkdir(path.dirname(this.file), { recursive: true }); + const lockFile = `${this.file}.lock`; + try { + this.lockFd = await acquireLock(lockFile); + await this.lockFd.writeFile(`${process.pid}\n`); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "EEXIST") throw new Error(`Session state is locked by another gori-agent instance: ${lockFile}`); + throw error; + } + try { + const raw = await fs.promises.readFile(this.file, "utf8"); + const parsed = JSON.parse(raw) as StoreData; + if (parsed.version !== 1 || !parsed.chats || !parsed.bindings) throw new Error("unsupported or malformed state data"); + this.data = parsed; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") { + await this.releaseLock(); + throw new Error(`Cannot read session state '${this.file}'; original file was preserved: ${error instanceof Error ? error.message : String(error)}`); + } + } + } + + getSelectedRole(chatKey: string, defaultRole: string): string { + return this.data.chats[chatKey]?.selectedRole || defaultRole; + } + + async setSelectedRole(chatKey: string, roleId: string): Promise { + const now = Date.now(); + const current = this.data.chats[chatKey]; + this.data.chats[chatKey] = { selectedRole: roleId, createdAt: current?.createdAt || now, updatedAt: now }; + await this.persist(); + } + + getBinding(chatKey: string, roleId: string): SessionBinding | undefined { + const value = this.data.bindings[bindingKey(chatKey, roleId)]; + return value ? structuredClone(value) : undefined; + } + + async setBinding(binding: SessionBinding): Promise { + this.data.bindings[bindingKey(binding.chatKey, binding.roleId)] = structuredClone(binding); + await this.persist(); + } + + async touchBinding(chatKey: string, roleId: string): Promise { + const binding = this.data.bindings[bindingKey(chatKey, roleId)]; + if (!binding) return; + binding.updatedAt = Date.now(); + await this.persist(); + } + + async deleteBinding(chatKey: string, roleId: string): Promise { + const key = bindingKey(chatKey, roleId); + const existing = this.data.bindings[key]; + delete this.data.bindings[key]; + if (existing) await this.persist(); + return existing; + } + + stats(): { chats: number; bindings: number } { + return { chats: Object.keys(this.data.chats).length, bindings: Object.keys(this.data.bindings).length }; + } + + async flush(): Promise { await this.queue; } + + async close(): Promise { + if (this.closed) return; + this.closed = true; + await this.flush(); + await this.releaseLock(); + } + + private persist(): Promise { + if (this.closed) return Promise.reject(new Error("Session store is closed")); + const snapshot = JSON.stringify(this.data, null, 2) + "\n"; + this.queue = this.queue.then(() => atomicWrite(this.file, snapshot)); + return this.queue; + } + + private async releaseLock(): Promise { + if (!this.lockFd) return; + await this.lockFd.close(); + this.lockFd = undefined; + await fs.promises.unlink(`${this.file}.lock`).catch((error: NodeJS.ErrnoException) => { + if (error.code !== "ENOENT") throw error; + }); + } +} + +export function chatKeyFor(platform: string, chatId: string): string { return `${platform}:${chatId}`; } +export function bindingKey(chatKey: string, roleId: string): string { return `${chatKey}\u0000${roleId}`; } + +async function acquireLock(lockFile: string): Promise { + return fs.promises.open(lockFile, "wx", 0o600); +} + +async function atomicWrite(file: string, content: string): Promise { + const temp = `${file}.${process.pid}.${Date.now()}.tmp`; + const handle = await fs.promises.open(temp, "wx", 0o600); + try { + await handle.writeFile(content, "utf8"); + await handle.sync(); + } finally { + await handle.close(); + } + try { + await fs.promises.rename(temp, file); + const dir = await fs.promises.open(path.dirname(file), "r"); + try { await dir.sync(); } finally { await dir.close(); } + } catch (error) { + await fs.promises.unlink(temp).catch(() => undefined); + throw error; + } +} diff --git a/src/core/gateway.ts b/src/core/gateway.ts index 6e76075..f6a7c5e 100644 --- a/src/core/gateway.ts +++ b/src/core/gateway.ts @@ -1,112 +1,99 @@ +import type { ConversationRuntime } from "../acp/types.js"; import type { GatewayPolicy } from "../config.js"; -import { AgentRegistry } from "../agents/agent-registry.js"; +import type { RoleRegistry } from "../roles/role-registry.js"; import type { PlatformAdapter } from "./adapter.js"; -import { CommandRouter } from "./command-router.js"; -import { SessionStore, sessionIdFor } from "./session-store.js"; +import { CommandRouter, type ParsedCommand } from "./command-router.js"; +import { chatKeyFor } from "./durable-session-store.js"; import type { IncomingMessage } from "./types.js"; -export interface GatewayResult { - ok: boolean; - reply?: string; - ignored?: boolean; - error?: string; -} +export interface GatewayResult { ok: boolean; reply?: string; ignored?: boolean; error?: string } export class Gateway { private readonly locks = new Map>(); - readonly commandRouter: CommandRouter; + readonly commandRouter = new CommandRouter(); - constructor( - private readonly policy: GatewayPolicy, - private readonly agents: AgentRegistry, - private readonly sessions: SessionStore - ) { - this.commandRouter = new CommandRouter(agents, sessions); - } + constructor(private readonly policy: GatewayPolicy, private readonly runtime: ConversationRuntime, private readonly roles: RoleRegistry) {} async receive(message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean } = {}): Promise { const policyError = this.checkPolicy(message); - if (policyError) return { ok: true, ignored: true, error: policyError }; + if (policyError) { + console.log(`Message ignored by policy: ${policyError} (${message.platform} ${message.chatId} ${message.userId})`); + return { ok: true, ignored: true, error: policyError }; + } - const sessionId = sessionIdFor(message.platform, message.chatId); - return this.withChatLock(sessionId, async () => { - const command = this.commandRouter.route(message, sessionId); - if (command.handled) { - const reply = command.text || ""; - await adapter.sendMessage({ - target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw }, - text: reply, - replyTo: message.messageId - }); - return { ok: true, reply }; - } - - const session = this.sessions.getSession(sessionId); - const agent = this.agents.get(session.selectedAgent); + const command = this.commandRouter.parse(message.text); + if (command?.kind === "cancel") return this.reply(message, adapter, await this.cancelText(message), options); + if (command?.kind === "new") await this.runtime.cancel(message.platform, message.chatId); + const chatKey = chatKeyFor(message.platform, message.chatId); + return this.withChatLock(chatKey, async () => { try { - const response = await agent.run({ - input: message.text, - sessionId, - platform: message.platform, - chatId: message.chatId, - userId: message.userId, - messageId: message.messageId - }); - await adapter.sendMessage({ - target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw }, - text: response.text, - replyTo: message.messageId - }); - return { ok: true, reply: response.text }; + const reply = command ? await this.executeCommand(command, message) : (await this.runtime.prompt({ + platform: message.platform, chatId: message.chatId, userId: message.userId, text: message.text, messageId: message.messageId + })).text; + return this.reply(message, adapter, reply, options); } catch (error) { const errorText = error instanceof Error ? error.message : String(error); const reply = `Agent error: ${errorText}`; - if (!options.synchronous) { - await adapter.sendMessage({ - target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw }, - text: reply, - replyTo: message.messageId - }); - } + if (!options.synchronous) await this.send(message, adapter, reply); return { ok: false, error: errorText, reply }; } }); } - stats(): { sessions: number; lockedChats: number } { - return { ...this.sessions.stats(), lockedChats: this.locks.size }; + stats(): ReturnType & { lockedChats: number } { return { ...this.runtime.stats(), lockedChats: this.locks.size }; } + + private async executeCommand(command: ParsedCommand, message: IncomingMessage): Promise { + const prefix = command.deprecatedAlias ? "Deprecated alias; use /role or /roles.\n" : ""; + switch (command.kind) { + case "help": return ["Commands:", "/roles", "/role ", "/status", "/cancel", "/new", "/help"].join("\n"); + case "roles": return `${prefix}Available roles: ${this.roles.list().join(", ")}`; + case "role": + if (!command.argument) return `${prefix}Usage: /role `; + if (!this.roles.has(command.argument)) return `${prefix}Unknown role: ${command.argument}`; + await this.runtime.cancel(message.platform, message.chatId); + await this.runtime.selectRole(message.platform, message.chatId, command.argument); + return `${prefix}Selected role: ${command.argument}`; + case "status": { + const status = this.runtime.status(message.platform, message.chatId); + return `OK\n${Object.entries(status).map(([key, value]) => `${key}=${value}`).join("\n")}`; + } + case "new": + await this.runtime.reset(message.platform, message.chatId); + return "Started a new native ACP session for this chat and role."; + case "cancel": return this.cancelText(message); + } + } + + private async cancelText(message: IncomingMessage): Promise { + return await this.runtime.cancel(message.platform, message.chatId) ? "Cancellation requested." : "No active turn to cancel."; + } + + private async reply(message: IncomingMessage, adapter: PlatformAdapter, reply: string, options: { synchronous?: boolean }): Promise { + if (!options.synchronous) await this.send(message, adapter, reply); + return { ok: true, reply }; + } + + private send(message: IncomingMessage, adapter: PlatformAdapter, text: string): Promise { + return adapter.sendMessage({ target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw }, text, replyTo: message.messageId }); } private checkPolicy(message: IncomingMessage): string | undefined { - if (this.policy.allowedUsers.length > 0 && !this.policy.allowedUsers.includes(message.userId)) { - return `User not allowed: ${message.userId}`; - } - if (this.policy.allowedChats.length > 0 && !this.policy.allowedChats.includes(message.chatId)) { - return `Chat not allowed: ${message.chatId}`; - } - if (this.policy.requireMentionInGroup && message.isGroup && !message.mentionsBot) { - return "Mention required in group chat"; - } + if (this.policy.allowedUsers.length > 0 && !this.policy.allowedUsers.includes(message.userId)) return `User not allowed: ${message.userId}`; + if (this.policy.allowedChats.length > 0 && !this.policy.allowedChats.includes(message.chatId)) return `Chat not allowed: ${message.chatId}`; + if (this.policy.requireMentionInGroup && message.isGroup && !message.mentionsBot) return "Mention required in group chat"; return undefined; } private async withChatLock(key: string, fn: () => Promise): Promise { const previous = this.locks.get(key) || Promise.resolve(); let release!: () => void; - const current = new Promise((resolve) => { - release = resolve; - }); + const current = new Promise((resolve) => { release = resolve; }); const queued = previous.then(() => current); this.locks.set(key, queued); - await previous; - try { - return await fn(); - } finally { + try { return await fn(); } finally { release(); - if (this.locks.get(key) === queued) { - this.locks.delete(key); - } + if (this.locks.get(key) === queued) this.locks.delete(key); } } } diff --git a/src/platforms/qq/adapter.ts b/src/platforms/qq/adapter.ts index ed43491..a804b2a 100644 --- a/src/platforms/qq/adapter.ts +++ b/src/platforms/qq/adapter.ts @@ -44,7 +44,11 @@ export class QqAdapter implements PlatformAdapter { handleDispatch(payload: QqWebhookPayload): boolean { const message = this.normalizeMessage(payload); - if (!message) return false; + if (!message) { + console.log(`QQ recv ${payload.t} ignored (empty text or missing ids)`); + return false; + } + console.log(`QQ recv ${payload.t} chat=${message.chatId} user=${message.userId} msgId=${message.messageId || ""} text=${JSON.stringify(message.text)}`); void this.gateway.receive(message, this).catch((error) => { console.error("QQ gateway error", error); }); @@ -54,11 +58,13 @@ export class QqAdapter implements PlatformAdapter { async sendMessage(message: OutgoingMessage): Promise { const token = await this.getAccessToken(); const raw = (message.target.raw || {}) as Record; - const groupOpenId = typeof raw.group_openid === "string" ? raw.group_openid : undefined; - const userOpenId = typeof raw.user_openid === "string" ? raw.user_openid : undefined; + const author = (raw.author || {}) as Record; + const groupOpenId = stringField(raw.group_openid) || stringField(raw.group_id); + const userOpenId = stringField(raw.user_openid) || stringField(author.user_openid); const isGroup = Boolean(groupOpenId) || message.target.chatId.startsWith("group:"); const targetId = groupOpenId || userOpenId || message.target.chatId.replace(/^group:/, "").replace(/^user:/, ""); const path = isGroup ? `/v2/groups/${encodeURIComponent(targetId)}/messages` : `/v2/users/${encodeURIComponent(targetId)}/messages`; + console.log(`QQ send -> ${path} msgId=${message.replyTo || ""} length=${message.text.length}`); const response = await fetch(`https://api.sgroup.qq.com${path}`, { method: "POST", @@ -141,10 +147,15 @@ function chatIdFor(data: QqWebhookEventData, isGroup: boolean): string | undefin if (data.channel_id) return `channel:${data.channel_id}`; if (data.guild_id) return `guild:${data.guild_id}`; if (!isGroup && data.user_openid) return `user:${data.user_openid}`; + if (!isGroup && data.author?.user_openid) return `user:${data.author.user_openid}`; if (!isGroup && data.author?.id) return `user:${data.author.id}`; return undefined; } +function stringField(value: unknown): string | undefined { + return typeof value === "string" && value ? value : undefined; +} + function stripConfiguredBotNames(text: string, botNames: string[]): string { let cleaned = text; for (const name of botNames) { diff --git a/src/platforms/qq/gateway-client.ts b/src/platforms/qq/gateway-client.ts index 93f0983..921ceef 100644 --- a/src/platforms/qq/gateway-client.ts +++ b/src/platforms/qq/gateway-client.ts @@ -131,7 +131,7 @@ export class QqGatewayClient { shard: this.config.shard, properties: { $os: process.platform, - $browser: "gori-agent-gateway", + $browser: "gori-agent", $device: os.hostname() } } diff --git a/src/platforms/qq/types.ts b/src/platforms/qq/types.ts index 33397a3..7cd0f48 100644 --- a/src/platforms/qq/types.ts +++ b/src/platforms/qq/types.ts @@ -50,6 +50,8 @@ export interface QqWebhookEventData { author?: { id?: string; user_openid?: string; + member_openid?: string; + union_openid?: string; username?: string; }; } diff --git a/src/roles/role-registry.ts b/src/roles/role-registry.ts new file mode 100644 index 0000000..3e5f957 --- /dev/null +++ b/src/roles/role-registry.ts @@ -0,0 +1,51 @@ +import crypto from "node:crypto"; +import type { AppConfig, RoleConfig } from "../config.js"; +import { SkillLoader, type LoadedSkill } from "./skill-loader.js"; + +export interface ResolvedRole extends RoleConfig { + loadedSkills: LoadedSkill[]; + fingerprint: string; + bootstrap: string; +} + +export class RoleRegistry { + private readonly roles = new Map(); + + constructor(private readonly config: AppConfig) { + const loader = new SkillLoader(config.skills); + for (const role of config.roles) { + const loadedSkills = role.skills.map((id) => loader.load(id)); + const fingerprint = crypto.createHash("sha256").update(JSON.stringify({ + id: role.id, + backend: role.backend, + workspace: role.workspace, + persona: role.persona, + policy: role.policy, + skills: loadedSkills.map(({ id, file, hash }) => ({ id, file, hash })) + })).digest("hex"); + this.roles.set(role.id, { ...role, loadedSkills, fingerprint, bootstrap: buildBootstrap(role, loadedSkills) }); + } + } + + get(id?: string): ResolvedRole { + const roleId = id || this.config.defaultRole; + const role = this.roles.get(roleId); + if (!role) throw new Error(`Unknown role: ${roleId}`); + return role; + } + + has(id: string): boolean { return this.roles.has(id); } + list(): string[] { return [...this.roles.keys()].sort(); } + defaultId(): string { return this.config.defaultRole; } +} + +function buildBootstrap(role: RoleConfig, skills: LoadedSkill[]): string { + return [ + "Initialize this ACP session with the following role. Treat these instructions as persistent context. Reply only with READY.", + `Role: ${role.id}`, + `Workspace: ${role.workspace}`, + role.persona ? `Persona:\n${role.persona}` : "Persona: general coding assistant", + `Permission policy enforced by the ACP client: ${JSON.stringify(role.policy)}`, + ...skills.map((skill) => `Skill ${skill.id} (${skill.file}):\n${skill.content}`) + ].join("\n\n"); +} diff --git a/src/roles/skill-loader.ts b/src/roles/skill-loader.ts new file mode 100644 index 0000000..122ee12 --- /dev/null +++ b/src/roles/skill-loader.ts @@ -0,0 +1,24 @@ +import crypto from "node:crypto"; +import fs from "node:fs"; +import type { SkillConfig } from "../config.js"; + +export interface LoadedSkill { + id: string; + file: string; + content: string; + hash: string; +} + +export class SkillLoader { + constructor(private readonly skills: SkillConfig[]) {} + + load(id: string): LoadedSkill { + const config = this.skills.find((skill) => skill.id === id); + if (!config) throw new Error(`Unknown skill: ${id}`); + const stat = fs.statSync(config.file); + if (!stat.isFile()) throw new Error(`Skill is not a file: ${config.file}`); + if (stat.size > config.maxBytes) throw new Error(`Skill '${id}' exceeds ${config.maxBytes} bytes`); + const content = fs.readFileSync(config.file, "utf8"); + return { id, file: config.file, content, hash: crypto.createHash("sha256").update(content).digest("hex") }; + } +} diff --git a/src/server.ts b/src/server.ts index 392e116..da44dcd 100644 --- a/src/server.ts +++ b/src/server.ts @@ -1,73 +1,59 @@ import express from "express"; import type { Server } from "node:http"; import { fileURLToPath } from "node:url"; -import { loadConfig, type AppConfig } from "./config.js"; -import { AgentRegistry } from "./agents/agent-registry.js"; -import { CliAgent } from "./agents/cli-agent.js"; +import { AcpBackendRegistry } from "./acp/backend-registry.js"; +import { AcpSessionManager } from "./acp/session-manager.js"; +import { defaultStateFile, loadConfig, type AppConfig } from "./config.js"; +import type { PlatformAdapter } from "./core/adapter.js"; +import { DurableSessionStore } from "./core/durable-session-store.js"; import { Gateway } from "./core/gateway.js"; import { PlatformRegistry } from "./core/platform-registry.js"; -import { SessionStore } from "./core/session-store.js"; -import type { PlatformAdapter } from "./core/adapter.js"; import { FeishuAdapter } from "./platforms/feishu/adapter.js"; -import { WeComAdapter } from "./platforms/wecom/adapter.js"; import { QqAdapter } from "./platforms/qq/adapter.js"; import { QqGatewayClient } from "./platforms/qq/gateway-client.js"; +import { WeComAdapter } from "./platforms/wecom/adapter.js"; import { GenericWebhookAdapter } from "./platforms/webhook/adapter.js"; import { WeixinAdapter } from "./platforms/weixin/adapter.js"; +import { RoleRegistry } from "./roles/role-registry.js"; export interface GatewayRuntime { app: express.Express; + gateway: Gateway; + sessionManager: AcpSessionManager; + store: DurableSessionStore; qqGatewayClient?: QqGatewayClient; + shutdown(): Promise; } -export function createGatewayRuntime(config: AppConfig): GatewayRuntime { - const sessions = new SessionStore(); - const agents = new AgentRegistry(config.defaultAgent); - for (const agentConfig of config.agents) { - agents.register(new CliAgent(agentConfig)); - } - - const gateway = new Gateway(config.policy, agents, sessions); +export async function createGatewayRuntime(config: AppConfig): Promise { + const store = new DurableSessionStore(defaultStateFile(config)); + await store.open(); + const roles = new RoleRegistry(config); + const sessionManager = new AcpSessionManager(config.acp, new AcpBackendRegistry(config.backends), roles, store); + const gateway = new Gateway(config.policy, sessionManager, roles); const platforms = new PlatformRegistry(); const qqAdapter = new QqAdapter(config.platforms.qq, gateway); const adapters: PlatformAdapter[] = [ - new FeishuAdapter(config.platforms.feishu, gateway), - new WeComAdapter(config.platforms.wecom), - qqAdapter, - new GenericWebhookAdapter(config.platforms.webhook, gateway), - new WeixinAdapter(config.platforms.weixin, gateway) + new FeishuAdapter(config.platforms.feishu, gateway), new WeComAdapter(config.platforms.wecom), qqAdapter, + new GenericWebhookAdapter(config.platforms.webhook, gateway), new WeixinAdapter(config.platforms.weixin, gateway) ]; for (const adapter of adapters) platforms.register(adapter); const app = express(); - app.use(express.json({ - limit: "1mb", - verify: (req, _res, buf) => { - (req as express.Request & { rawBody?: Buffer }).rawBody = Buffer.from(buf); - } - })); - - app.get("/health", (_req, res) => { - res.json({ ok: true, gateway: gateway.stats() }); - }); - - app.get("/platforms", (_req, res) => { - res.json({ ok: true, platforms: platforms.list(), agents: agents.list() }); - }); + app.use(express.json({ limit: "1mb", verify: (req, _res, buf) => { + (req as express.Request & { rawBody?: Buffer }).rawBody = Buffer.from(buf); + } })); + app.get("/health", (_req, res) => res.json({ ok: true, acp: gateway.stats() })); + app.get("/platforms", (_req, res) => res.json({ ok: true, platforms: platforms.list(), roles: roles.list(), backends: config.backends.map(({ id }) => id) })); function mountWebhook(routeName: string, adapterName = routeName): void { app.post(`/webhook/${routeName}`, async (req, res) => { try { const response = await platforms.get(adapterName).handleWebhook({ - req, - body: req.body, - headers: req.headers, - query: req.query, + req, body: req.body, headers: req.headers, query: req.query, rawBody: (req as express.Request & { rawBody?: Buffer }).rawBody }); - if (response.headers) { - for (const [key, value] of Object.entries(response.headers)) res.setHeader(key, value); - } + if (response.headers) for (const [key, value] of Object.entries(response.headers)) res.setHeader(key, value); res.status(response.status || 200).json(response.body ?? { ok: true }); } catch (error) { const message = error instanceof Error ? error.message : String(error); @@ -76,36 +62,42 @@ export function createGatewayRuntime(config: AppConfig): GatewayRuntime { } }); } - - for (const platformName of ["feishu", "wecom", "qq", "weixin"]) { - mountWebhook(platformName); - } + for (const name of ["feishu", "wecom", "qq", "weixin"]) mountWebhook(name); mountWebhook("generic", "webhook"); const qqGatewayClient = config.platforms.qq.enabled && config.platforms.qq.connectionMode === "websocket" - ? new QqGatewayClient(config.platforms.qq, qqAdapter) - : undefined; - - return { app, qqGatewayClient }; + ? new QqGatewayClient(config.platforms.qq, qqAdapter) : undefined; + let closing: Promise | undefined; + const runtime: GatewayRuntime = { + app, gateway, sessionManager, store, qqGatewayClient, + shutdown: () => closing ||= (async () => { + qqGatewayClient?.stop(); + await sessionManager.shutdown(); + await store.close(); + })() + }; + return runtime; } -export function createApp(config: AppConfig): express.Express { - return createGatewayRuntime(config).app; -} +export async function createApp(config: AppConfig): Promise { return (await createGatewayRuntime(config)).app; } -export function startServer(config: AppConfig): Server { - const runtime = createGatewayRuntime(config); +export interface RunningServer { server: Server; runtime: GatewayRuntime; shutdown(): Promise } + +export async function startServer(config: AppConfig): Promise { + const runtime = await createGatewayRuntime(config); const server = runtime.app.listen(config.server.port, config.server.host, () => { - console.log(`gori-agent-gateway listening on ${config.server.host}:${config.server.port}`); - if (runtime.qqGatewayClient) { - console.log("QQ websocket gateway enabled; connecting to QQ..."); - runtime.qqGatewayClient.start(); - } + console.log(`gori-agent listening on ${config.server.host}:${config.server.port}`); + if (runtime.qqGatewayClient) { console.log("QQ websocket gateway enabled; connecting to QQ..."); runtime.qqGatewayClient.start(); } }); - server.on("close", () => runtime.qqGatewayClient?.stop()); - return server; + let closing: Promise | undefined; + return { server, runtime, shutdown: () => closing ||= (async () => { + runtime.qqGatewayClient?.stop(); + const runtimeShutdown = runtime.shutdown(); + await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); + await runtimeShutdown; + })() }; } if (process.argv[1] && fileURLToPath(import.meta.url) === process.argv[1]) { - startServer(loadConfig()); + void startServer(loadConfig()).catch((error) => { console.error(error); process.exitCode = 1; }); } diff --git a/test/acp-client.test.ts b/test/acp-client.test.ts new file mode 100644 index 0000000..442bf04 --- /dev/null +++ b/test/acp-client.test.ts @@ -0,0 +1,23 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import type { RequestPermissionRequest } from "@agentclientprotocol/sdk"; +import { decidePermission } from "../src/acp/client.js"; + +const request = (rawInput: unknown): RequestPermissionRequest => ({ + sessionId: "session", + toolCall: { toolCallId: "call", title: "bash git status", kind: "execute", name: "bash", rawInput }, + options: [ + { optionId: "allow", name: "Allow", kind: "allow_once" }, + { optionId: "reject", name: "Reject", kind: "reject_once" } + ] +}); + +test("permission deny and allowlist fail closed", () => { + assert.equal(decidePermission(request("git status"), { permissionMode: "deny", allowedTools: [], allowedCommandPatterns: [] }).outcome.outcome, "selected"); + const allowed = decidePermission(request("git status"), { permissionMode: "allowlist", allowedTools: ["bash"], allowedCommandPatterns: ["^git status$"] }); + assert.deepEqual(allowed.outcome, { outcome: "selected", optionId: "allow" }); + const missingDetail = decidePermission(request(undefined), { permissionMode: "allowlist", allowedTools: ["bash"], allowedCommandPatterns: ["^git status$"] }); + assert.deepEqual(missingDetail.outcome, { outcome: "selected", optionId: "reject" }); + const destructive = decidePermission(request("rm -rf /"), { permissionMode: "allowlist", allowedTools: ["bash"], allowedCommandPatterns: ["^git status$"] }); + assert.deepEqual(destructive.outcome, { outcome: "selected", optionId: "reject" }); +}); diff --git a/test/acp-session-manager.test.ts b/test/acp-session-manager.test.ts new file mode 100644 index 0000000..1d226bb --- /dev/null +++ b/test/acp-session-manager.test.ts @@ -0,0 +1,58 @@ +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import test from "node:test"; +import { AcpBackendRegistry } from "../src/acp/backend-registry.js"; +import { AcpSessionManager } from "../src/acp/session-manager.js"; +import { parseConfig } from "../src/config.js"; +import { DurableSessionStore } from "../src/core/durable-session-store.js"; +import { RoleRegistry } from "../src/roles/role-registry.js"; + +const fixture = path.resolve("test/fixtures/fake-acp-agent.mjs"); + +function config(stateFile: string, logFile: string) { + return parseConfig({ + configVersion: 2, acp: { stateFile, promptTimeoutMs: 2_000, cancelGraceMs: 100, idleTimeoutMs: 100, sweepIntervalMs: 20, maxProcesses: 2 }, + backends: [{ id: "kimi", command: process.execPath, args: [fixture], env: { FAKE_ACP_LOG: logFile } }], + skills: [], defaultRole: "assistant", + roles: [{ id: "assistant", backend: "kimi", workspace: path.resolve("."), persona: "test", skills: [], policy: { permissionMode: "deny" } }], + platforms: {} + }); +} + +async function runtime(stateFile: string, logFile: string) { + const cfg = config(stateFile, logFile); const store = new DurableSessionStore(stateFile); await store.open(); + return { cfg, store, manager: new AcpSessionManager(cfg.acp, new AcpBackendRegistry(cfg.backends), new RoleRegistry(cfg), store) }; +} + +test("creates, persists, idles, and resumes the same native session", async () => { + const dir = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-acp-")); const state = path.join(dir, "state.json"); const log = path.join(dir, "fake.log"); + const first = await runtime(state, log); + const response = await first.manager.prompt({ platform: "qq", chatId: "chat", userId: "user", text: "one" }); + const sessionId = response.text.split(":")[1]; + assert.match(response.text, /reply:fake-/); + await new Promise((resolve) => setTimeout(resolve, 180)); + await first.manager.prompt({ platform: "qq", chatId: "chat", userId: "user", text: "two" }); + await first.manager.shutdown(); await first.store.close(); + + const second = await runtime(state, log); + const resumed = await second.manager.prompt({ platform: "qq", chatId: "chat", userId: "user", text: "three" }); + assert.equal(resumed.text, `reply:${sessionId}:three`); + const entries = (await fs.promises.readFile(log, "utf8")).trim().split("\n").map(JSON.parse); + assert.ok(entries.some((entry) => entry.method === "session/resume" && entry.sessionId === sessionId)); + await second.manager.shutdown(); await second.store.close(); +}); + +test("cancel reaches a hanging ACP prompt and new unbinds", async () => { + const dir = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-acp-cancel-")); const state = path.join(dir, "state.json"); const log = path.join(dir, "fake.log"); + const current = await runtime(state, log); + await current.manager.prompt({ platform: "qq", chatId: "chat", userId: "user", text: "ready" }); + const hanging = current.manager.prompt({ platform: "qq", chatId: "chat", userId: "user", text: "hang" }); + await new Promise((resolve) => setTimeout(resolve, 50)); + assert.equal(await current.manager.cancel("qq", "chat"), true); + await hanging; + await current.manager.reset("qq", "chat"); + assert.equal(current.store.stats().bindings, 0); + await current.manager.shutdown(); await current.store.close(); +}); diff --git a/test/acp-worker.test.ts b/test/acp-worker.test.ts new file mode 100644 index 0000000..dd8a62f --- /dev/null +++ b/test/acp-worker.test.ts @@ -0,0 +1,37 @@ +import assert from "node:assert/strict"; +import path from "node:path"; +import test from "node:test"; +import { AcpWorker } from "../src/acp/worker.js"; +import { parseConfig } from "../src/config.js"; +import { RoleRegistry } from "../src/roles/role-registry.js"; + +const fixture = path.resolve("test/fixtures/fake-acp-agent.mjs"); + +function makeWorker(promptTimeoutMs = 80) { + const config = parseConfig({ + configVersion: 2, + acp: { promptTimeoutMs, cancelGraceMs: 30 }, + backends: [{ id: "kimi", command: process.execPath, args: [fixture] }], + defaultRole: "assistant", + roles: [{ id: "assistant", backend: "kimi", workspace: path.resolve("."), policy: { permissionMode: "deny" } }], + platforms: {} + }); + let crashes = 0; + const worker = new AcpWorker(config.backends[0], new RoleRegistry(config).get(), config.acp, () => { crashes++; }); + return { worker, crashes: () => crashes }; +} + +test("worker times out a hung prompt and can be terminated", async () => { + const { worker } = makeWorker(); + await worker.start(); + await assert.rejects(worker.prompt("hang"), /timed out/); + await worker.terminate(); +}); + +test("worker reports an unexpected ACP subprocess exit", async () => { + const current = makeWorker(500); + await current.worker.start(); + await assert.rejects(current.worker.prompt("crash")); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(current.crashes(), 1); +}); diff --git a/test/config.test.ts b/test/config.test.ts new file mode 100644 index 0000000..d4ffd43 --- /dev/null +++ b/test/config.test.ts @@ -0,0 +1,31 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { parseConfig } from "../src/config.js"; + +const platforms = { + qq: { enabled: true, appId: "id", clientSecret: "secret", connectionMode: "websocket" as const }, + feishu: {}, wecom: {}, webhook: {}, weixin: {} +}; + +test("migrates a v1 Kimi config and preserves QQ fields", () => { + const config = parseConfig({ + server: { port: 8787 }, policy: {}, defaultAgent: "kimi", + agents: [{ name: "kimi", command: "/home/ubuntu/.kimi-code/bin/kimi", args: ["-p"], cwd: "/tmp" }], platforms + }); + assert.equal(config.configVersion, 2); + assert.deepEqual(config.backends[0].args, ["acp"]); + assert.equal(config.roles[0].workspace, "/tmp"); + assert.equal(config.platforms.qq.clientSecret, "secret"); + assert.equal(config.platforms.qq.connectionMode, "websocket"); +}); + +test("does not silently migrate a non-Kimi CLI agent", () => { + assert.throws(() => parseConfig({ defaultAgent: "echo", agents: [{ name: "echo", command: "node" }], platforms }), /only supports a Kimi/); +}); + +test("validates role references and absolute workspace", () => { + assert.throws(() => parseConfig({ + configVersion: 2, backends: [{ id: "kimi", command: "kimi", args: ["acp"] }], defaultRole: "a", + roles: [{ id: "a", backend: "missing", workspace: "relative" }], platforms + }), /workspace must be absolute/); +}); diff --git a/test/durable-session-store.test.ts b/test/durable-session-store.test.ts new file mode 100644 index 0000000..31eaeae --- /dev/null +++ b/test/durable-session-store.test.ts @@ -0,0 +1,41 @@ +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import test from "node:test"; +import { DurableSessionStore } from "../src/core/durable-session-store.js"; + +test("persists chat role and native session across reopen", async () => { + const dir = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-store-")); + const file = path.join(dir, "state.json"); + const first = new DurableSessionStore(file); + await first.open(); + await first.setSelectedRole("qq:chat", "ops"); + await first.setBinding({ chatKey: "qq:chat", roleId: "ops", backendId: "kimi", nativeSessionId: "native-1", workspace: "/tmp", roleFingerprint: "fp", createdAt: 1, updatedAt: 1 }); + await first.close(); + assert.equal((await fs.promises.readdir(dir)).some((name) => name.endsWith(".tmp")), false); + + const second = new DurableSessionStore(file); + await second.open(); + assert.equal(second.getSelectedRole("qq:chat", "assistant"), "ops"); + assert.equal(second.getBinding("qq:chat", "ops")?.nativeSessionId, "native-1"); + await second.close(); +}); + +test("preserves corrupt state and fails explicitly", async () => { + const dir = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-store-bad-")); + const file = path.join(dir, "state.json"); + await fs.promises.writeFile(file, "not-json"); + const store = new DurableSessionStore(file); + await assert.rejects(store.open(), /original file was preserved/); + assert.equal(await fs.promises.readFile(file, "utf8"), "not-json"); +}); + +test("refuses a second writer lock", async () => { + const dir = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-store-lock-")); + const file = path.join(dir, "state.json"); + const first = new DurableSessionStore(file); await first.open(); + const second = new DurableSessionStore(file); + await assert.rejects(second.open(), /locked by another/); + await first.close(); +}); diff --git a/test/fixtures/fake-acp-agent.mjs b/test/fixtures/fake-acp-agent.mjs new file mode 100644 index 0000000..8ead3fc --- /dev/null +++ b/test/fixtures/fake-acp-agent.mjs @@ -0,0 +1,73 @@ +#!/usr/bin/env node +import fs from "node:fs"; +import { Readable, Writable } from "node:stream"; +import * as acp from "@agentclientprotocol/sdk"; + +const pending = new Map(); +const logFile = process.env.FAKE_ACP_LOG; +const log = (entry) => { if (logFile) fs.appendFileSync(logFile, `${JSON.stringify(entry)}\n`); }; + +const app = acp.agent({ name: "fake-acp-agent" }) + .onRequest(acp.methods.agent.initialize, ({ params }) => { + log({ method: "initialize" }); + return { + protocolVersion: params.protocolVersion, + agentCapabilities: { loadSession: true, sessionCapabilities: { resume: {}, close: {}, list: {} } }, + agentInfo: { name: "fake-acp-agent", version: "1" } + }; + }) + .onRequest(acp.methods.agent.session.new, ({ params }) => { + const sessionId = `fake-${Date.now()}-${Math.random().toString(16).slice(2)}`; + log({ method: "session/new", sessionId, cwd: params.cwd }); + return { sessionId }; + }) + .onRequest(acp.methods.agent.session.load, ({ params }) => { log({ method: "session/load", sessionId: params.sessionId }); return {}; }) + .onRequest(acp.methods.agent.session.resume, ({ params }) => { log({ method: "session/resume", sessionId: params.sessionId }); return {}; }) + .onRequest(acp.methods.agent.session.close, ({ params }) => { log({ method: "session/close", sessionId: params.sessionId }); return {}; }) + .onRequest(acp.methods.agent.session.list, () => ({ sessions: [] })) + .onRequest(acp.methods.agent.session.prompt, async ({ params, client, signal }) => { + const text = params.prompt.filter((item) => item.type === "text").map((item) => item.text).join(""); + log({ method: "session/prompt", sessionId: params.sessionId, text }); + if (text === "crash") process.exit(9); + if (text === "permission") { + const response = await client.request(acp.methods.client.session.requestPermission, { + sessionId: params.sessionId, + toolCall: { toolCallId: "bash-1", title: "bash git status", kind: "execute", name: "bash", rawInput: "git status" }, + options: [ + { optionId: "allow", name: "Allow", kind: "allow_once" }, + { optionId: "reject", name: "Reject", kind: "reject_once" } + ] + }); + await update(client, params.sessionId, response.outcome.outcome === "selected" ? response.outcome.optionId : "cancelled"); + return { stopReason: "end_turn" }; + } + if (text === "hang") { + await new Promise((resolve) => { + const done = () => resolve(undefined); + pending.set(params.sessionId, done); + signal.addEventListener("abort", done, { once: true }); + }); + pending.delete(params.sessionId); + return { stopReason: "cancelled" }; + } + await update(client, params.sessionId, text.includes("Initialize this ACP session") ? "READY" : `reply:${params.sessionId}:${text}`); + return { stopReason: "end_turn" }; + }) + .onNotification(acp.methods.agent.session.cancel, ({ params }) => { + log({ method: "session/cancel", sessionId: params.sessionId }); + pending.get(params.sessionId)?.(); + }); + +function update(client, sessionId, text) { + return client.notify(acp.methods.client.session.update, { + sessionId, + update: { sessionUpdate: "agent_message_chunk", content: { type: "text", text } } + }); +} + +const stream = acp.ndJsonStream( + Writable.toWeb(process.stdout), + Readable.toWeb(process.stdin) +); +const connection = app.connect(stream); +await connection.closed; diff --git a/test/gateway.test.ts b/test/gateway.test.ts new file mode 100644 index 0000000..bb6fcb8 --- /dev/null +++ b/test/gateway.test.ts @@ -0,0 +1,45 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import type { ConversationRuntime } from "../src/acp/types.js"; +import { parseConfig } from "../src/config.js"; +import type { PlatformAdapter } from "../src/core/adapter.js"; +import { Gateway } from "../src/core/gateway.js"; +import type { IncomingMessage } from "../src/core/types.js"; +import { RoleRegistry } from "../src/roles/role-registry.js"; + +class FakeRuntime implements ConversationRuntime { + role = "assistant"; prompts = 0; cancelled = 0; resets = 0; release?: () => void; + async prompt() { this.prompts++; await new Promise((resolve) => { this.release = resolve; }); return { text: "done", roleId: this.role, backendId: "kimi" }; } + async cancel() { this.cancelled++; this.release?.(); return true; } + async reset() { this.resets++; } + async selectRole(_p: string, _c: string, role: string) { this.role = role; } + selectedRole() { return this.role; } + status() { return { role: this.role, running: Boolean(this.release) }; } + stats() { return { activeWorkers: 0, inFlight: 0, crashes: 0, persistedBindings: 0 }; } + async shutdown() {} +} + +const cfg = parseConfig({ configVersion: 2, backends: [{ id: "kimi", command: "kimi", args: ["acp"] }], defaultRole: "assistant", roles: [ + { id: "assistant", backend: "kimi", workspace: "/tmp" }, { id: "ops", backend: "kimi", workspace: "/tmp" } +], platforms: {} }); +const adapter: PlatformAdapter = { name: "test", async handleWebhook() { return {}; }, async sendMessage() {} }; +const message = (text: string): IncomingMessage => ({ platform: "qq", chatId: "chat", userId: "user", text }); + +test("cancel bypasses the chat lock and new resets", async () => { + const runtime = new FakeRuntime(); const gateway = new Gateway(cfg.policy, runtime, new RoleRegistry(cfg)); + const turn = gateway.receive(message("work"), adapter, { synchronous: true }); + await new Promise((resolve) => setTimeout(resolve, 10)); + const cancelled = await gateway.receive(message("/cancel"), adapter, { synchronous: true }); + assert.equal(cancelled.reply, "Cancellation requested."); + await turn; + const reset = await gateway.receive(message("/new"), adapter, { synchronous: true }); + assert.match(reset.reply || "", /new native ACP session/); + assert.equal(runtime.resets, 1); +}); + +test("role aliases, role selection, and status are routed", async () => { + const runtime = new FakeRuntime(); const gateway = new Gateway(cfg.policy, runtime, new RoleRegistry(cfg)); + assert.match((await gateway.receive(message("/agents"), adapter, { synchronous: true })).reply || "", /Deprecated alias/); + assert.equal((await gateway.receive(message("/role ops"), adapter, { synchronous: true })).reply, "Selected role: ops"); + assert.match((await gateway.receive(message("/status"), adapter, { synchronous: true })).reply || "", /role=ops/); +}); diff --git a/test/qq-adapter.test.ts b/test/qq-adapter.test.ts new file mode 100644 index 0000000..8030388 --- /dev/null +++ b/test/qq-adapter.test.ts @@ -0,0 +1,34 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { parseConfig } from "../src/config.js"; +import { QqAdapter } from "../src/platforms/qq/adapter.js"; + +const config = parseConfig({ configVersion: 2, backends: [{ id: "kimi", command: "kimi", args: ["acp"] }], defaultRole: "assistant", roles: [{ id: "assistant", backend: "kimi", workspace: "/tmp" }], platforms: { qq: { + enabled: true, connectionMode: "webhook", appId: "id", clientSecret: "secret", verifySignature: false, botNames: ["Bot"] +} } }); + +test("normalizes GROUP and C2C author openid while ACK remains immediate", async () => { + const received: any[] = []; + const gateway = { receive: async (message: unknown) => { received.push(message); throw new Error("ACP failed"); } }; + const adapter = new QqAdapter(config.platforms.qq, gateway as never); + const group = await adapter.handleWebhook({ body: { op: 0, t: "GROUP_AT_MESSAGE_CREATE", d: { id: "m1", group_openid: "g1", author: { user_openid: "u1" }, content: "@Bot hi" } }, headers: {}, query: {}, req: {} as never }); + const c2c = await adapter.handleWebhook({ body: { op: 0, t: "C2C_MESSAGE_CREATE", d: { id: "m2", author: { user_openid: "u2" }, content: "hello" } }, headers: {}, query: {}, req: {} as never }); + assert.deepEqual(group.body, { op: 12 }); assert.deepEqual(c2c.body, { op: 12 }); + await new Promise((resolve) => setImmediate(resolve)); + assert.equal(received[0].chatId, "group:g1"); assert.equal(received[0].userId, "u1"); + assert.equal(received[1].chatId, "user:u2"); assert.equal(received[1].userId, "u2"); +}); + +test("sendMessage uses nested author.user_openid for C2C endpoint", async () => { + const original = globalThis.fetch; const urls: string[] = []; + globalThis.fetch = (async (input: string | URL | Request) => { + urls.push(String(input)); + if (String(input).includes("getAppAccessToken")) return new Response(JSON.stringify({ access_token: "token", expires_in: 7200 }), { status: 200 }); + return new Response(JSON.stringify({ id: "sent" }), { status: 200 }); + }) as typeof fetch; + try { + const adapter = new QqAdapter(config.platforms.qq, { receive: async () => ({ ok: true }) } as never); + await adapter.sendMessage({ target: { platform: "qq", chatId: "user:u2", raw: { author: { user_openid: "u2" } } }, text: "reply", replyTo: "m2" }); + assert.ok(urls.some((url) => url.endsWith("/v2/users/u2/messages"))); + } finally { globalThis.fetch = original; } +});