From 038a3986d72df7fe2b470f561fe75067b0ec8e1f Mon Sep 17 00:00:00 2001 From: Johnny Huynh <27847622+johnnyhuy@users.noreply.github.com> Date: Thu, 6 Aug 2026 08:35:25 +1000 Subject: [PATCH 1/2] feat(server): agent provider seam + Claude Code provider AgentClient/AgentSession/PersistenceHandle seam under server/agent, disk-backed handle store for cross-restart resume, agent-manager for session lifecycle + event stamping, and the first provider: Claude Code via @anthropic-ai/claude-agent-sdk with a tool-call normaliser. Co-authored-by: opencode-agent --- bun.lock | 68 ++++ packages/server/package.json | 1 + .../server/src/server/agent/agent-manager.ts | 252 ++++++++++++++ .../src/server/agent/agent-sdk-types.ts | 78 +++++ .../src/server/agent/handle-store.test.ts | 46 +++ .../server/src/server/agent/handle-store.ts | 55 +++ .../agent/providers/claude/claude-provider.ts | 320 ++++++++++++++++++ .../providers/claude/tool-call-mapper.test.ts | 113 +++++++ .../providers/claude/tool-call-mapper.ts | 124 +++++++ 9 files changed, 1057 insertions(+) create mode 100644 packages/server/src/server/agent/agent-manager.ts create mode 100644 packages/server/src/server/agent/agent-sdk-types.ts create mode 100644 packages/server/src/server/agent/handle-store.test.ts create mode 100644 packages/server/src/server/agent/handle-store.ts create mode 100644 packages/server/src/server/agent/providers/claude/claude-provider.ts create mode 100644 packages/server/src/server/agent/providers/claude/tool-call-mapper.test.ts create mode 100644 packages/server/src/server/agent/providers/claude/tool-call-mapper.ts diff --git a/bun.lock b/bun.lock index a1834f6..8a94ef3 100644 --- a/bun.lock +++ b/bun.lock @@ -13,6 +13,7 @@ "oxlint": "1.61.0", "oxlint-tsgolint": "^0.22.1", "typescript": "^5.9.3", + "vitest": "^2.1.0", }, }, "packages/app": { @@ -92,6 +93,7 @@ "name": "@echohello/server", "version": "0.0.0", "dependencies": { + "@anthropic-ai/claude-agent-sdk": "^0.3.222", "@echohello/client": "workspace:*", "@echohello/protocol": "workspace:*", "dotenv": "^17.2.3", @@ -146,6 +148,26 @@ "@alloc/quick-lru": ["@alloc/quick-lru@5.2.0", "", {}, "sha512-UrcABB+4bUrFABwbluTIBErXwvbsU/V7TZWfmbgJfbkwiBuziS9gxdODUyuiecfdGQ85jglMW6juS3+z5TsKLw=="], + "@anthropic-ai/claude-agent-sdk": ["@anthropic-ai/claude-agent-sdk@0.3.222", "", { "optionalDependencies": { "@anthropic-ai/claude-agent-sdk-darwin-arm64": "0.3.222", "@anthropic-ai/claude-agent-sdk-darwin-x64": "0.3.222", "@anthropic-ai/claude-agent-sdk-linux-arm64": "0.3.222", "@anthropic-ai/claude-agent-sdk-linux-arm64-musl": "0.3.222", "@anthropic-ai/claude-agent-sdk-linux-x64": "0.3.222", "@anthropic-ai/claude-agent-sdk-linux-x64-musl": "0.3.222", "@anthropic-ai/claude-agent-sdk-win32-arm64": "0.3.222", "@anthropic-ai/claude-agent-sdk-win32-x64": "0.3.222" }, "peerDependencies": { "@anthropic-ai/sdk": ">=0.93.0", "@modelcontextprotocol/sdk": "^1.29.0", "zod": "^4.0.0" } }, "sha512-muAyjIzXJjIpSrj91vSmnU76/z7rNlV8+lSuq48h4eUPSYVBSgOasgkfQAMUrIZJCfI15XF9/EYoRrx/eKD7Og=="], + + "@anthropic-ai/claude-agent-sdk-darwin-arm64": ["@anthropic-ai/claude-agent-sdk-darwin-arm64@0.3.222", "", { "os": "darwin", "cpu": "arm64" }, "sha512-h3nCSwRbUsXDnRRYvMhzI1u1q5w6jClc9jTtOJBY6+mijjJEGk+Fa5Kuhw0Tb5ec+1iatiTyCIT46HrV0Skqcg=="], + + "@anthropic-ai/claude-agent-sdk-darwin-x64": ["@anthropic-ai/claude-agent-sdk-darwin-x64@0.3.222", "", { "os": "darwin", "cpu": "x64" }, "sha512-kVzR6qBZv1ddbb+uJfDDEVdW6/cl1dtc2nAIa8prk4IcuyK1RwQE/BOO8dO6iRl5PnHuv9BFKS9AO1UErOkT5Q=="], + + "@anthropic-ai/claude-agent-sdk-linux-arm64": ["@anthropic-ai/claude-agent-sdk-linux-arm64@0.3.222", "", { "os": "linux", "cpu": "arm64" }, "sha512-wn4XYkaWbAXc8PoznZqjCqbqB/Qet8dNdpZuRtbNSA1yZR0vYg2u28C5zV0QqVTXGgQVBYFUZabouyoERL6PsA=="], + + "@anthropic-ai/claude-agent-sdk-linux-arm64-musl": ["@anthropic-ai/claude-agent-sdk-linux-arm64-musl@0.3.222", "", { "os": "linux", "cpu": "arm64" }, "sha512-sVvPAUzRq+Am3+wNrFJ7URaWORu8t5hOu4vJnMhOeey66ZFFfzNwP0uz/7sTWjg+cRse/wG/hek34lZ0Gdu2+w=="], + + "@anthropic-ai/claude-agent-sdk-linux-x64": ["@anthropic-ai/claude-agent-sdk-linux-x64@0.3.222", "", { "os": "linux", "cpu": "x64" }, "sha512-deZh3tVoYy9ememF5oirZNx1une7IG79JbKrB3Cynggt2gnq3TFXvPfkVGWgNsocK035wAW9HhGZui832OFSqA=="], + + "@anthropic-ai/claude-agent-sdk-linux-x64-musl": ["@anthropic-ai/claude-agent-sdk-linux-x64-musl@0.3.222", "", { "os": "linux", "cpu": "x64" }, "sha512-vM48p8GRM5N57+jPFdPj1ZS8Y770frJ1NTRecA06/GXtRANn7dOUVpvJbpu6fJBFF3BCacM1uVcBCMZh6wEMjw=="], + + "@anthropic-ai/claude-agent-sdk-win32-arm64": ["@anthropic-ai/claude-agent-sdk-win32-arm64@0.3.222", "", { "os": "win32", "cpu": "arm64" }, "sha512-VuJjTLP2jQ1j8yn44/YODFR8GhI60Tvzok6AaXUYpKgeZTbBndDOpqwxvdWAraJSSyDreSSMztYyI6lP6nVypg=="], + + "@anthropic-ai/claude-agent-sdk-win32-x64": ["@anthropic-ai/claude-agent-sdk-win32-x64@0.3.222", "", { "os": "win32", "cpu": "x64" }, "sha512-mr4eZ+bmZA0hPl2wN031qy3TJXZJZIpdE4khvCc0AJ7ljVpJUEG654wz9NJoPxfPo4cRg41niHqUU5rN9BfQVw=="], + + "@anthropic-ai/sdk": ["@anthropic-ai/sdk@0.115.0", "", { "dependencies": { "json-schema-to-ts": "^3.1.1", "standardwebhooks": "^1.0.0" }, "peerDependencies": { "zod": "^3.25.0 || ^4.0.0" }, "optionalPeers": ["zod"], "bin": { "anthropic-ai-sdk": "bin/cli" } }, "sha512-BJrFIVyjNuU8lfDyIJTvlRYzgQg+zEl78BxE7fq8esULsGz9IRQvGtW5spq3tydmtjQb/GFdooKGdGsetpx+lQ=="], + "@babel/code-frame": ["@babel/code-frame@7.10.4", "", { "dependencies": { "@babel/highlight": "^7.10.4" } }, "sha512-vG6SvB6oYEhvgisZNFRmRCUkLz11c7rp+tbNTynGqc6mS1d5ATd/sGyV6W0KZZnXRKMTzZDRgQT3Ou9jhpAfUg=="], "@babel/compat-data": ["@babel/compat-data@7.29.7", "", {}, "sha512-locTkQyKvwIEgBzVrn8693ebc97F2U8ZHjbXwDXJ5Fn2TCpNwTlKcaKLkdHop5c/icOFE7qt7Q9JC5hnKNa6Gg=="], @@ -456,6 +478,8 @@ "@expo/xcpretty": ["@expo/xcpretty@4.4.4", "", { "dependencies": { "@babel/code-frame": "^7.20.0", "chalk": "^4.1.0", "js-yaml": "^4.1.0" }, "bin": { "excpretty": "build/cli.js" } }, "sha512-4aQzz9vgxcNXFfo/iyNgDDYfsU5XGKKxWxZopw0cVotHiW+U8IJbIxMaxsINs6bHhtkG3StKNPcOrn3eBuxKPw=="], + "@hono/node-server": ["@hono/node-server@2.1.0", "", { "peerDependencies": { "hono": "^4" } }, "sha512-XovyyCCnBzW+zKu+z/zq8hwNs4KOR5rEMAOxo2f40Q5xoOI37IMm6MIg2COOUtUApo0i6850MTBKH2u4QLGIqg=="], + "@isaacs/fs-minipass": ["@isaacs/fs-minipass@4.0.1", "", { "dependencies": { "minipass": "^7.0.4" } }, "sha512-wgm9Ehl2jpeqP3zw/7mo3kRHFp5MEDhqAdwy1fTGkHAwnkGOVsgpvQhL8B5n1qlb01jV3n/bI0ZfZp5lWA1k4w=="], "@isaacs/ttlcache": ["@isaacs/ttlcache@1.4.1", "", {}, "sha512-RQgQ4uQ+pLbqXfOmieB91ejmLwvSgv9nLx6sT6sD83s7umBypgg+OIBOBbEUiJXrfpnp9j0mRhYYdzp9uqq3lA=="], @@ -488,6 +512,8 @@ "@jridgewell/trace-mapping": ["@jridgewell/trace-mapping@0.3.31", "", { "dependencies": { "@jridgewell/resolve-uri": "^3.1.0", "@jridgewell/sourcemap-codec": "^1.4.14" } }, "sha512-zzNR+SdQSDJzc8joaeP8QQoCQr8NuYx2dIIytl1QeBEZHJ9uW6hebsrYgbz8hJwUQao3TWCMtmfV8Nu1twOLAw=="], + "@modelcontextprotocol/sdk": ["@modelcontextprotocol/sdk@1.30.0", "", { "dependencies": { "@hono/node-server": "^1.19.9 || ^2.0.5", "ajv": "^8.17.1", "ajv-formats": "^3.0.1", "content-type": "^1.0.5", "cors": "^2.8.5", "cross-spawn": "^7.0.5", "eventsource": "^3.0.2", "eventsource-parser": "^3.0.0", "express": "^5.2.1", "express-rate-limit": "^8.2.1", "hono": "^4.11.4", "jose": "^6.1.3", "json-schema-typed": "^8.0.2", "pkce-challenge": "^5.0.0", "raw-body": "^3.0.0", "zod": "^3.25 || ^4.0", "zod-to-json-schema": "^3.25.1" }, "peerDependencies": { "@cfworker/json-schema": "^4.1.1" }, "optionalPeers": ["@cfworker/json-schema"] }, "sha512-xKd8OIzlqNzcqcNumGAa6g+PW2kjD5vrpcKOnfldAUPP3j7lnqMPwlTXQm8gF+UwH72z0lqaRbjr9hqGz0eITA=="], + "@napi-rs/wasm-runtime": ["@napi-rs/wasm-runtime@1.1.6", "", { "dependencies": { "@tybys/wasm-util": "^0.10.3" }, "peerDependencies": { "@emnapi/core": "^1.7.1", "@emnapi/runtime": "^1.7.1" } }, "sha512-ZLv/JdUfkvOy9eCnnBaGfiO+XimbjebAeO+MRQqD/B+FR1tnRN0tpKSJHRbE8sFfS6aqsXZ67TQjfwfsxULVbg=="], "@nodelib/fs.scandir": ["@nodelib/fs.scandir@2.1.5", "", { "dependencies": { "@nodelib/fs.stat": "2.0.5", "run-parallel": "^1.1.9" } }, "sha512-vq24Bq3ym5HEQm2NKCr3yXDwjc7vTsEThRDnkp2DK9p1uqLR+DHurm/NOTo0KG7HYHU7eppKZj3MyqYuMBf62g=="], @@ -814,6 +840,8 @@ "@sinonjs/fake-timers": ["@sinonjs/fake-timers@10.3.0", "", { "dependencies": { "@sinonjs/commons": "^3.0.0" } }, "sha512-V4BG07kuYSUkTCSBHG8G8TNhM+F19jXFWnQtzj+we8DrkpSBCee9Z3Ms8yiGer/dlmhe35/Xdgyo3/0rQKg7YA=="], + "@stablelib/base64": ["@stablelib/base64@1.0.1", "", {}, "sha512-1bnPQqSxSuc3Ii6MhBysoWCg58j97aUjuCSZrGSmDxNqtytIi0k8utUenAwTZN4V5mXXYGsVUI9zeBqy+jBOSQ=="], + "@tailwindcss/node": ["@tailwindcss/node@4.3.2", "", { "dependencies": { "@jridgewell/remapping": "^2.3.5", "enhanced-resolve": "5.21.6", "jiti": "^2.7.0", "lightningcss": "1.32.0", "magic-string": "^0.30.21", "source-map-js": "^1.2.1", "tailwindcss": "4.3.2" } }, "sha512-yWP/sqEcBLaD8JuA6zNwxoYKr75qxTioYwlRwekj5Jr/I5GXnoJfjetH/psLUIv74cYTH2lBUEzBkinthoYcBg=="], "@tailwindcss/oxide": ["@tailwindcss/oxide@4.3.2", "", { "optionalDependencies": { "@tailwindcss/oxide-android-arm64": "4.3.2", "@tailwindcss/oxide-darwin-arm64": "4.3.2", "@tailwindcss/oxide-darwin-x64": "4.3.2", "@tailwindcss/oxide-freebsd-x64": "4.3.2", "@tailwindcss/oxide-linux-arm-gnueabihf": "4.3.2", "@tailwindcss/oxide-linux-arm64-gnu": "4.3.2", "@tailwindcss/oxide-linux-arm64-musl": "4.3.2", "@tailwindcss/oxide-linux-x64-gnu": "4.3.2", "@tailwindcss/oxide-linux-x64-musl": "4.3.2", "@tailwindcss/oxide-wasm32-wasi": "4.3.2", "@tailwindcss/oxide-win32-arm64-msvc": "4.3.2", "@tailwindcss/oxide-win32-x64-msvc": "4.3.2" } }, "sha512-z8ZgnzX8gdNoWLBLqBPoh/sjnxkwvf9ZuWjnO0l0yIzbLa5/9S+eC5QxGZKRobVHIC3/1BoMWjHblqWjcgFgag=="], @@ -964,6 +992,10 @@ "agent-base": ["agent-base@7.1.4", "", {}, "sha512-MnA+YT8fwfJPgBx3m60MNqakm30XOkyIoH1y6huTQvC0PwZG7ki8NacLBcrPbNoo8vEZy7Jpuk7+jMO+CUovTQ=="], + "ajv": ["ajv@8.20.0", "", { "dependencies": { "fast-deep-equal": "^3.1.3", "fast-uri": "^3.0.1", "json-schema-traverse": "^1.0.0", "require-from-string": "^2.0.2" } }, "sha512-Thbli+OlOj+iMPYFBVBfJ3OmCAnaSyNn4M1vz9T6Gka5Jt9ba/HIR56joy65tY6kx/FCF5VXNB819Y7/GUrBGA=="], + + "ajv-formats": ["ajv-formats@3.0.1", "", { "dependencies": { "ajv": "^8.0.0" } }, "sha512-8iUql50EUR+uUcdRQ3HDqa6EVyo3docL8g5WJ3FNcWmu62IbkGUue/pEyLBW8VGKKucTPgqeks4fIU1DA4yowQ=="], + "anser": ["anser@1.4.10", "", {}, "sha512-hCv9AqTQ8ycjpSd3upOJd7vFwW1JaoYQ7tpham03GJ1ca8/65rqn0RpaWpItOAd6ylW9wAw6luXYPJIyPFVOww=="], "ansi-escapes": ["ansi-escapes@4.3.2", "", { "dependencies": { "type-fest": "^0.21.3" } }, "sha512-gKXj5ALrKWQLsYG9jlTRmR/xKluxHV+Z9QEwNIgCfM1/uwPMCuzVVnh5mwTd+OuBZcwSIMbqssNWRm1lE51QaQ=="], @@ -1122,6 +1154,8 @@ "core-js-compat": ["core-js-compat@3.49.0", "", { "dependencies": { "browserslist": "^4.28.1" } }, "sha512-VQXt1jr9cBz03b331DFDCCP90b3fanciLkgiOoy8SBHy06gNf+vQ1A3WFLqG7I8TipYIKeYK9wxd0tUrvHcOZA=="], + "cors": ["cors@2.8.6", "", { "dependencies": { "object-assign": "^4", "vary": "^1" } }, "sha512-tJtZBBHA6vjIAaF6EnIaq6laBBP9aq/Y3ouVJjEfoHbRBcHBAHYcMh/w8LDrk2PvIMMq8gmopa5D4V8RmbrxGw=="], + "cross-fetch": ["cross-fetch@3.2.0", "", { "dependencies": { "node-fetch": "^2.7.0" } }, "sha512-Q+xVJLoGOeIMXZmbUK4HYk+69cQH6LudR0Vu/pRm2YlU/hDV9CiS0gKUMaWY5f2NeUH9C1nV3bsTlCo0FsTV1Q=="], "cross-spawn": ["cross-spawn@7.0.6", "", { "dependencies": { "path-key": "^3.1.0", "shebang-command": "^2.0.0", "which": "^2.0.1" } }, "sha512-uV2QOWP2nWzsy2aMp8aRibhi9dlzF5Hgh5SHaB9OiTGEyDTiJJyx0uy51QXdyWbtAHNua4XJzUKca3OzKUd3vA=="], @@ -1214,6 +1248,10 @@ "event-target-shim": ["event-target-shim@5.0.1", "", {}, "sha512-i/2XbnSz/uxRCU6+NdVJgKWDTM427+MqYbkQzD321DuCQJUqOuJKIA0IM2+W2xtYHdKOmZ4dR6fExsd4SXL+WQ=="], + "eventsource": ["eventsource@3.0.7", "", { "dependencies": { "eventsource-parser": "^3.0.1" } }, "sha512-CRT1WTyuQoD771GW56XEZFQ/ZoSfWid1alKGDYMmkt2yl8UXrVR4pspqWNEcqKvVIzg6PAltWjxcSSPrboA4iA=="], + + "eventsource-parser": ["eventsource-parser@3.1.0", "", {}, "sha512-kJezFj9YFAMLeORyi7aCLxLbD5/qWMQnoMVlVPyHIll7lgRJCc3JVln9Vgl9nwQi0YkMnhdGTMNn7CkRRAptMg=="], + "expect-type": ["expect-type@1.4.0", "", {}, "sha512-KfYbmpRm0VbLjEvVa9yGwCi9GI34xvi7A/HXYWQO65CSD2u3MczUJSuwXKFIxlGsgBQizV9q5J9NHj4VG0n+pA=="], "expo": ["expo@54.0.35", "", { "dependencies": { "@babel/runtime": "^7.20.0", "@expo/cli": "54.0.25", "@expo/config": "~12.0.13", "@expo/config-plugins": "~54.0.4", "@expo/devtools": "0.1.8", "@expo/fingerprint": "0.15.5", "@expo/metro": "~54.2.0", "@expo/metro-config": "54.0.16", "@expo/vector-icons": "^15.0.3", "@ungap/structured-clone": "^1.3.0", "babel-preset-expo": "~54.0.11", "expo-asset": "~12.0.13", "expo-constants": "~18.0.13", "expo-file-system": "~19.0.23", "expo-font": "~14.0.12", "expo-keep-awake": "~15.0.8", "expo-modules-autolinking": "3.0.26", "expo-modules-core": "3.0.30", "pretty-format": "^29.7.0", "react-refresh": "^0.14.2", "whatwg-url-without-unicode": "8.0.0-3" }, "peerDependencies": { "@expo/dom-webview": "*", "@expo/metro-runtime": "*", "react": "*", "react-native": "*", "react-native-webview": "*" }, "optionalPeers": ["@expo/dom-webview", "@expo/metro-runtime", "react-native-webview"], "bin": { "expo": "bin/cli", "fingerprint": "bin/fingerprint", "expo-modules-autolinking": "bin/autolinking" } }, "sha512-E+tXpQwjGm5fK/uwa55p0Xx/kuo5dXDKfVJ95IargTNa5KiFt26lSTXXa9KnHbI4EDLwFD38/xTKZvzPTlGTdg=="], @@ -1244,6 +1282,8 @@ "express": ["express@5.2.1", "", { "dependencies": { "accepts": "^2.0.0", "body-parser": "^2.2.1", "content-disposition": "^1.0.0", "content-type": "^1.0.5", "cookie": "^0.7.1", "cookie-signature": "^1.2.1", "debug": "^4.4.0", "depd": "^2.0.0", "encodeurl": "^2.0.0", "escape-html": "^1.0.3", "etag": "^1.8.1", "finalhandler": "^2.1.0", "fresh": "^2.0.0", "http-errors": "^2.0.0", "merge-descriptors": "^2.0.0", "mime-types": "^3.0.0", "on-finished": "^2.4.1", "once": "^1.4.0", "parseurl": "^1.3.3", "proxy-addr": "^2.0.7", "qs": "^6.14.0", "range-parser": "^1.2.1", "router": "^2.2.0", "send": "^1.1.0", "serve-static": "^2.2.0", "statuses": "^2.0.1", "type-is": "^2.0.1", "vary": "^1.1.2" } }, "sha512-hIS4idWWai69NezIdRt2xFVofaF4j+6INOpJlVOLDO8zXGpUVEVzIYk12UUi2JzjEzWL3IOAxcTubgz9Po0yXw=="], + "express-rate-limit": ["express-rate-limit@8.6.2", "", { "dependencies": { "debug": "^4.4.3", "ip-address": "^10.2.0" }, "peerDependencies": { "express": ">= 4.11" } }, "sha512-YH4ru+eOJxQABscKFfRCy9R7x9QFGdezclVMwwgFFndzS2Xnm0uo6B0ABZsLhcpeptGv2qvuJVWlQr9gQZoC3A=="], + "fast-copy": ["fast-copy@4.0.3", "", {}, "sha512-58apWr0GUiDFM8+3afrO6eYwJBn9ZAhDOzG3L+/9llab/haCARS2UIfffmOurYLwbgDRs8n0rfr6qAAPEAuAQw=="], "fast-deep-equal": ["fast-deep-equal@3.1.3", "", {}, "sha512-f3qQ9oQy9j2AhBe/H9VC91wLmKBCCU/gDOnKNAYG5hswO7BLKj09Hc5HYNz9cGI++xlpDCIgDaitVs03ATR84Q=="], @@ -1254,6 +1294,10 @@ "fast-safe-stringify": ["fast-safe-stringify@2.1.1", "", {}, "sha512-W+KJc2dmILlPplD/H4K9l9LcAHAfPtP6BY84uVLXQ6Evcz9Lcg33Y2z1IVblT6xdY54PXYVHEv+0Wpq8Io6zkA=="], + "fast-sha256": ["fast-sha256@1.3.0", "", {}, "sha512-n11RGP/lrWEFI/bWdygLxhI+pVeo1ZYIVwvvPkW7azl/rOy+F3HYRZ2K5zeE9mmkhQppyv9sQFx0JM9UabnpPQ=="], + + "fast-uri": ["fast-uri@3.1.5", "", {}, "sha512-gHwA1O9LDIcKunMKhObS/HimwtehO1nPUECKAu5TpKgaO19fcWEl4bliWe1jWxVFvIXztJjjQ4L8XQ1EU9f7Jw=="], + "fastq": ["fastq@1.20.1", "", { "dependencies": { "reusify": "^1.0.4" } }, "sha512-GGToxJ/w1x32s/D2EKND7kTil4n8OVk/9mycTc4VDza13lOvpUZTGX3mFSCtV9ksdGBVzvsyAVLM6mHFThxXxw=="], "fb-watchman": ["fb-watchman@2.0.2", "", { "dependencies": { "bser": "2.1.1" } }, "sha512-p5161BqbuCaSnB8jIbzQHOlpgsPmK5rJVDfDKO91Axs5NC1uu3HRQm6wt9cd9/+GtQQIO53JdGXXoyDpTAsgYA=="], @@ -1332,6 +1376,8 @@ "hermes-parser": ["hermes-parser@0.29.1", "", { "dependencies": { "hermes-estree": "0.29.1" } }, "sha512-xBHWmUtRC5e/UL0tI7Ivt2riA/YBq9+SiYFU7C1oBa/j2jYGlIF9043oak1F47ihuDIxQ5nbsKueYJDRY02UgA=="], + "hono": ["hono@4.13.0", "", {}, "sha512-jhunvfHWxd7J5EFfSgH4xsYJzSe/lfqbUCxiyyeaQasUsXeEHXtzVid+7EOGByc5JnFa23SSFL3Y2RV/z1T+eQ=="], + "hosted-git-info": ["hosted-git-info@7.0.2", "", { "dependencies": { "lru-cache": "^10.0.1" } }, "sha512-puUZAUKT5m8Zzvs72XWy3HtvVbTWljRE66cP60bxJzAqf2DgICo7lYTY2IHUmLnNpjYvw5bvmoHvPc0QO2a62w=="], "html-void-elements": ["html-void-elements@3.0.0", "", {}, "sha512-bEqo66MRXsUGxWHV5IP0PUiAWwoEjba4VCzg0LjFJBpchPaTfyfCKTG6bc5F8ucKec3q5y6qOdGyYTSBEvhCrg=="], @@ -1364,6 +1410,8 @@ "invariant": ["invariant@2.2.4", "", { "dependencies": { "loose-envify": "^1.0.0" } }, "sha512-phJfQVBuaJM5raOpJjSfkiD6BpbCE4Ns//LaXl6wGYtUBY83nWS6Rf9tXm2e8VaK60JEjYldbPif/A2B1C2gNA=="], + "ip-address": ["ip-address@10.4.0", "", {}, "sha512-oSK96Grm3aP6OrS263xVxbNDGVL7rzBtYdpGqlDG8iQdoenDoTs/nkki+DflYbAEE8Xl6o5YxhxlrKvI3nqKXQ=="], + "ipaddr.js": ["ipaddr.js@1.9.1", "", {}, "sha512-0KI/607xoxSToH7GjN1FfSbLoU0+btTicjsQSWQlh/hZykN8KpmMf7uYwPW3R+akZ6R/w18ZlXSHBYXiYUPO3g=="], "is-arrayish": ["is-arrayish@0.3.4", "", {}, "sha512-m6UrgzFVUYawGBh1dUsWR5M2Clqic9RVXC/9f8ceNlv2IcO9j9J/z8UoCLPqtsPBFNzEpfR3xftohbfqDx8EQA=="], @@ -1412,6 +1460,8 @@ "jiti": ["jiti@2.7.0", "", { "bin": { "jiti": "lib/jiti-cli.mjs" } }, "sha512-AC/7JofJvZGrrneWNaEnJeOLUx+JlGt7tNa0wZiRPT4MY1wmfKjt2+6O2p2uz2+skll8OZZmJMNqeke7kKbNgQ=="], + "jose": ["jose@6.2.8", "", {}, "sha512-Bsdjwm3Qsd/P0jR+BHDe3LytDfY7WBq2HmCCLIwuVRHMuEC9ae7/R474GIUdF1NgCyZjzVo/A9DOiOBtXq8ZoQ=="], + "joycon": ["joycon@3.1.1", "", {}, "sha512-34wB/Y7MW7bzjKRjUKTa46I2Z7eV62Rkhva+KkopW7Qvv/OSWBqvkSY7vusOPrNuZcUG3tApvdVgNB8POj3SPw=="], "js-tokens": ["js-tokens@4.0.0", "", {}, "sha512-RdJUflcE3cUzKiMqQgsCu06FPu9UdIJO0beYbPhHN4k6apgJtifcoCtT9bcxOpYBtpD2kCM6Sbzg4CausW/PKQ=="], @@ -1422,6 +1472,12 @@ "jsesc": ["jsesc@3.1.0", "", { "bin": { "jsesc": "bin/jsesc" } }, "sha512-/sM3dO2FOzXjKQhJuo0Q173wf2KOo8t4I8vHy6lF9poUp7bKT0/NHE8fPX23PwfhnykfqnC2xRxOnVw5XuGIaA=="], + "json-schema-to-ts": ["json-schema-to-ts@3.1.1", "", { "dependencies": { "@babel/runtime": "^7.18.3", "ts-algebra": "^2.0.0" } }, "sha512-+DWg8jCJG2TEnpy7kOm/7/AxaYoaRbjVB4LFZLySZlWn8exGs3A4OLJR966cVvU26N7X9TWxl+Jsw7dzAqKT6g=="], + + "json-schema-traverse": ["json-schema-traverse@1.0.0", "", {}, "sha512-NM8/P9n3XjXhIZn1lLhkFaACTOURQXjWhV4BA/RnOv8xvgqtqpAX9IO4mRQxSx1Rlo4tqzeqb0sOlruaOy3dug=="], + + "json-schema-typed": ["json-schema-typed@8.0.2", "", {}, "sha512-fQhoXdcvc3V28x7C7BMs4P5+kNlgUURe2jmUT1T//oBRMDrqy1QPelJimwZGo7Hg9VPV3EQV5Bnq4hbFy2vetA=="], + "json5": ["json5@2.2.3", "", { "bin": { "json5": "lib/cli.js" } }, "sha512-XmOWe7eyHYH14cLdVPoyg+GOH3rYX++KpzrylJwSW98t3Nk+U8XOl8FWKOgwtzdb8lXGf6zYwDUzeHMWfxasyg=="], "kleur": ["kleur@3.0.3", "", {}, "sha512-eTIzlVOSUR+JxdDFepEYcBMtZ9Qqdef+rnzWdRZuMbOywu5tO2w2N7rqjoANZ5k9vywhL6Br1VRjUIgTQx4E8w=="], @@ -1678,6 +1734,8 @@ "pirates": ["pirates@4.0.7", "", {}, "sha512-TfySrs/5nm8fQJDcBDuUng3VOUKsd7S+zqvbOTiGXHfxX4wK31ard+hoNuvkicM/2YFzlpDgABOevKSsB4G/FA=="], + "pkce-challenge": ["pkce-challenge@5.0.1", "", {}, "sha512-wQ0b/W4Fr01qtpHlqSqspcj3EhBvimsdh0KlHhH8HRZnMsEa0ea2fTULOXOS9ccQr3om+GcGRk4e+isrZWV8qQ=="], + "plist": ["plist@3.1.1", "", { "dependencies": { "@xmldom/xmldom": "^0.9.10", "base64-js": "^1.5.1", "xmlbuilder": "^15.1.1" } }, "sha512-ZIfcLJC+7E7FBFnDxm9MPmt7D+DidyQ26lewieO75AdhA2ayMtsJSES0iWzqJQbcVRSrTufQoy0DR94xHue0oA=="], "pngjs": ["pngjs@3.4.0", "", {}, "sha512-NCrCHhWmnQklfH4MtJMRjZ2a8c80qXeMlQMv2uVp9ISJMTt562SbGd6n2oq0PaPgKm7Z6pL9E2UlLIhC+SHL3w=="], @@ -1896,6 +1954,8 @@ "standard-navigation": ["standard-navigation@0.0.7", "", {}, "sha512-NCGLCNyuXrFOkGHxdNZFnpsehGtiq1oXbPhKl7ZuxFO5J//H2evqqOchmD4YwEUJnkjO4kH9Xp4hQX6hdAYCKQ=="], + "standardwebhooks": ["standardwebhooks@1.0.0", "", { "dependencies": { "@stablelib/base64": "^1.0.0", "fast-sha256": "^1.3.0" } }, "sha512-BbHGOQK9olHPMvQNHWul6MYlrRTAOKn03rOe4A8O3CLWhNf4YHBqq2HJKKC+sfqpxiBY52pNeesD6jIiLDz8jg=="], + "statuses": ["statuses@2.0.2", "", {}, "sha512-DvEy55V3DB7uknRo+4iOGT5fP1slR8wQohVdknigZPMpMstaKJQWhwiYBACJE3Ul2pTnATihhBYnRhZQHGBiRw=="], "std-env": ["std-env@3.10.0", "", {}, "sha512-5GS12FdOZNliM5mAOxFRg7Ir0pWz8MdpYm6AY6VPkGpbA7ZzmbzNcBJQ0GPvvyWgcY7QAhCgf9Uy89I03faLkg=="], @@ -1968,6 +2028,8 @@ "trim-lines": ["trim-lines@3.0.1", "", {}, "sha512-kRj8B+YHZCc9kQYdWfJB2/oUl9rA99qbowYYBtr4ui4mZyAQ2JpvVBd/6U2YloATfqBhBTSMhTpgBHtU0Mf3Rg=="], + "ts-algebra": ["ts-algebra@2.0.0", "", {}, "sha512-FPAhNPFMrkwz76P7cdjdmiShwMynZYN6SgOujD1urY4oNm80Ou9oMdmbR45LotcKOXoy7wSmHkRFE6Mxbrhefw=="], + "ts-interface-checker": ["ts-interface-checker@0.1.13", "", {}, "sha512-Y/arvbn+rrz3JCKl9C4kVNfTfSm2/mEp5FSz5EsZSANGPSlQrpRI5M4PKF+mJnE52jOO90PnPSc3Ur3bTQw0gA=="], "tslib": ["tslib@2.8.1", "", {}, "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w=="], @@ -2094,10 +2156,14 @@ "zod": ["zod@3.25.76", "", {}, "sha512-gzUt/qt81nXsFGKIFcC3YnfEAx5NkunCfnDlvuBSSFS02bcXu4Lmea0AFIUwbLWxWPx3d9p8S5QoaujKcNQxcQ=="], + "zod-to-json-schema": ["zod-to-json-schema@3.25.2", "", { "peerDependencies": { "zod": "^3.25.28 || ^4" } }, "sha512-O/PgfnpT1xKSDeQYSCfRI5Gy3hPf91mKVDuYLUHZJMiDFptvP41MSnWofm8dnCm0256ZNfZIM7DSzuSMAFnjHA=="], + "zustand": ["zustand@5.0.14", "", { "peerDependencies": { "@types/react": ">=18.0.0", "immer": ">=9.0.6", "react": ">=18.0.0", "use-sync-external-store": ">=1.2.0" }, "optionalPeers": ["@types/react", "immer", "react", "use-sync-external-store"] }, "sha512-/8tAspM5LMPr28b3fwLYrtdj77ECpfZviaP75CMTnwO8ISyaE4GDIG/9rDDYq/cH9D2Xw2A2RXglLInmVBQB/g=="], "zwitch": ["zwitch@2.0.4", "", {}, "sha512-bXE4cR/kVZhKZX/RjPEflHaKVhUVl85noU3v6b8apfQEc1x4A+zBxjZ4lN8LqGd6WZ3dl98pY4o717VFmoPp+A=="], + "@anthropic-ai/claude-agent-sdk/zod": ["zod@4.4.3", "", {}, "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ=="], + "@babel/core/@babel/code-frame": ["@babel/code-frame@7.29.7", "", { "dependencies": { "@babel/helper-validator-identifier": "^7.29.7", "js-tokens": "^4.0.0", "picocolors": "^1.1.1" } }, "sha512-Aup7aUOfpbAUg2ROOJN6Iw5f9DMBlzu0mIkm/malLQFN/YQgO48wCj0Kxa3sEHJvPVFg7siR+qRInwXd2qhQKw=="], "@babel/core/semver": ["semver@6.3.1", "", { "bin": { "semver": "bin/semver.js" } }, "sha512-BR7VvDCVHO+q2xBEWskxS6DJE1qRnb7DxzUrogb71CWoSficBxYsiAGd+Kl0mmq/MprG9yArRkyrQxTO6XjMzA=="], @@ -2186,6 +2252,8 @@ "@istanbuljs/load-nyc-config/js-yaml": ["js-yaml@3.15.0", "", { "dependencies": { "argparse": "^1.0.7", "esprima": "^4.0.0" }, "bin": { "js-yaml": "bin/js-yaml.js" } }, "sha512-ttBQIIQPDeLjpPOohtUdXuXUVoA2uIB6fEH9HyJ7234s5mBJ5wTx20njxplLZQgLaOfpmPQA7X2t5AX6tIPbog=="], + "@modelcontextprotocol/sdk/zod": ["zod@4.4.3", "", {}, "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ=="], + "@radix-ui/react-collection/@radix-ui/react-compose-refs": ["@radix-ui/react-compose-refs@1.1.3", "", { "peerDependencies": { "@types/react": "*", "react": "^16.8 || ^17.0 || ^18.0 || ^19.0 || ^19.0.0-rc" }, "optionalPeers": ["@types/react"] }, "sha512-rYOP8OMnuuPMQF1uhPVlGNcCDlkokKqGFE3JcxFViIkAXP7EvFWUliJAstrapypaBLJNHbZL6jGhbVDGTwmVhA=="], "@radix-ui/react-collection/@radix-ui/react-slot": ["@radix-ui/react-slot@1.3.0", "", { "dependencies": { "@radix-ui/react-compose-refs": "1.1.3" }, "peerDependencies": { "@types/react": "*", "react": "^16.8 || ^17.0 || ^18.0 || ^19.0 || ^19.0.0-rc" }, "optionalPeers": ["@types/react"] }, "sha512-MojKku4U/miO8Av4Dkb+ctMAQx7JmY96LmtDQlAarCRtd7rN52QCSzBF+XAvr5S6coSVj9HEPBgHAHKEJVk/WA=="], diff --git a/packages/server/package.json b/packages/server/package.json index 34ff1e4..0ef7840 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -31,6 +31,7 @@ "lint": "oxlint src --config ../../oxlint.jsonc" }, "dependencies": { + "@anthropic-ai/claude-agent-sdk": "^0.3.222", "@echohello/client": "workspace:*", "@echohello/protocol": "workspace:*", "dotenv": "^17.2.3", diff --git a/packages/server/src/server/agent/agent-manager.ts b/packages/server/src/server/agent/agent-manager.ts new file mode 100644 index 0000000..24427a0 --- /dev/null +++ b/packages/server/src/server/agent/agent-manager.ts @@ -0,0 +1,252 @@ +import { randomBytes } from "node:crypto"; + +import type { Logger } from "pino"; +import { AgentError, type AgentEvent, type SessionState } from "@echohello/protocol"; + +import type { + AgentClient, + AgentEventSink, + AgentSession, + PersistenceHandle, + SessionScopedEvent, +} from "./agent-sdk-types.js"; +import { HandleStore } from "./handle-store.js"; + +export interface AgentManagerOptions { + handleStore: HandleStore; + logger: Logger; +} + +interface ManagedSession { + state: SessionState; + session: AgentSession; + handle: PersistenceHandle; +} + +export interface StartSessionArgs { + workspaceId: string; + cwd: string; + providerId: string; + modelId?: string; + modeId?: string; + initialPrompt?: string; +} + +export interface ResumeSessionArgs { + workspaceId: string; + cwd: string; + handle: PersistenceHandle; + overrides?: { + modelId?: string | undefined; + modeId?: string | undefined; + }; +} + +/** + * Owns agent session lifecycle: provider registry, Supaplane-side session ids, + * event stamping/forwarding, status bookkeeping, and handle persistence. + * + * Wire events out via `onAgentEvent` / `onSessionState` (set by the daemon + * before any session can be created). + */ +export class AgentManager { + readonly #providers = new Map(); + readonly #sessions = new Map(); + readonly #handleStore: HandleStore; + readonly #logger: Logger; + + onAgentEvent: ((event: AgentEvent) => void) | undefined; + onSessionState: ((state: SessionState) => void) | undefined; + + constructor(options: AgentManagerOptions) { + this.#handleStore = options.handleStore; + this.#logger = options.logger.child({ module: "agent-manager" }); + } + + registerProvider(client: AgentClient): void { + if (this.#providers.has(client.providerId)) { + throw new AgentError({ + code: "conflict", + message: `Provider already registered: ${client.providerId}`, + }); + } + this.#providers.set(client.providerId, client); + } + + providerIds(): string[] { + return [...this.#providers.keys()]; + } + + getProvider(providerId: string): AgentClient { + const provider = this.#providers.get(providerId); + if (!provider) { + throw new AgentError({ + code: "provider_unavailable", + message: `Unknown provider: ${providerId}`, + }); + } + return provider; + } + + async startSession(args: StartSessionArgs): Promise { + const provider = this.getProvider(args.providerId); + const sessionId = newSessionId(); + const emit = this.#makeSink(sessionId); + + const { session, handle } = await provider.createSession( + { + cwd: args.cwd, + ...(args.modelId !== undefined ? { modelId: args.modelId } : {}), + ...(args.modeId !== undefined ? { modeId: args.modeId } : {}), + }, + emit, + ); + + const now = Date.now(); + const state: SessionState = { + sessionId, + workspaceId: args.workspaceId, + providerId: args.providerId, + ...(args.modelId !== undefined ? { modelId: args.modelId } : {}), + ...(args.modeId !== undefined ? { modeId: args.modeId } : {}), + status: "idle", + startedAt: now, + updatedAt: now, + forkCount: 0, + }; + this.#sessions.set(sessionId, { state, session, handle }); + await this.#persistHandle(args.cwd, handle); + this.#emitSessionState(state); + + if (args.initialPrompt !== undefined && args.initialPrompt.length > 0) { + void this.send(sessionId, args.initialPrompt).catch((err: unknown) => { + this.#logger.warn({ err, sessionId }, "initial prompt failed"); + }); + } + return state; + } + + async resumeSession(args: ResumeSessionArgs): Promise { + const provider = this.getProvider(args.handle.provider); + const sessionId = newSessionId(); + const emit = this.#makeSink(sessionId); + + const stored = await this.#handleStore + .load(args.cwd, args.handle.provider, args.handle.sessionId) + .catch(() => null); + const metadata = { ...stored?.metadata, ...args.handle.metadata }; + + const { session, handle } = await provider.resumeSession( + { + handle: { + provider: args.handle.provider, + sessionId: args.handle.sessionId, + ...(metadata !== undefined && Object.keys(metadata).length > 0 ? { metadata } : {}), + }, + cwd: args.cwd, + ...(args.overrides !== undefined ? { overrides: args.overrides } : {}), + }, + emit, + ); + + const now = Date.now(); + const state: SessionState = { + sessionId, + workspaceId: args.workspaceId, + providerId: args.handle.provider, + ...(args.overrides?.modelId !== undefined ? { modelId: args.overrides.modelId } : {}), + ...(args.overrides?.modeId !== undefined ? { modeId: args.overrides.modeId } : {}), + status: "idle", + startedAt: now, + updatedAt: now, + forkCount: 0, + }; + this.#sessions.set(sessionId, { state, session, handle }); + await this.#persistHandle(args.cwd, handle); + this.#emitSessionState(state); + return state; + } + + async send(sessionId: string, prompt: string, attachments?: unknown[]): Promise { + const managed = this.#requireSession(sessionId); + this.#setStatus(sessionId, "running"); + try { + await managed.session.send(prompt, attachments); + } catch (err) { + this.#setStatus(sessionId, "error"); + throw err; + } + } + + async abort(sessionId: string): Promise { + const managed = this.#requireSession(sessionId); + await managed.session.abort(); + this.#setStatus(sessionId, "idle"); + } + + async disposeSession(sessionId: string): Promise { + const managed = this.#sessions.get(sessionId); + if (!managed) return; + this.#sessions.delete(sessionId); + await managed.session.dispose().catch((err: unknown) => { + this.#logger.warn({ err, sessionId }, "session dispose failed"); + }); + } + + async disposeAll(): Promise { + const ids = [...this.#sessions.keys()]; + await Promise.all(ids.map((id) => this.disposeSession(id))); + } + + getSession(sessionId: string): SessionState | undefined { + return this.#sessions.get(sessionId)?.state; + } + + listSessions(workspaceId?: string): SessionState[] { + const all = [...this.#sessions.values()].map((s) => s.state); + return workspaceId === undefined ? all : all.filter((s) => s.workspaceId === workspaceId); + } + + #makeSink(sessionId: string): AgentEventSink { + return (event: SessionScopedEvent) => { + const stamped = { ...event, sessionId } as AgentEvent; + if (stamped.type === "status") { + this.#setStatus(sessionId, stamped.status); + } else if (stamped.type === "error") { + this.#setStatus(sessionId, "error"); + } + this.onAgentEvent?.(stamped); + }; + } + + #setStatus(sessionId: string, status: SessionState["status"]): void { + const managed = this.#sessions.get(sessionId); + if (!managed || managed.state.status === status) return; + managed.state = { ...managed.state, status, updatedAt: Date.now() }; + this.#emitSessionState(managed.state); + } + + #requireSession(sessionId: string): ManagedSession { + const managed = this.#sessions.get(sessionId); + if (!managed) { + throw new AgentError({ code: "not_found", message: `Unknown session: ${sessionId}` }); + } + return managed; + } + + async #persistHandle(cwd: string, handle: PersistenceHandle): Promise { + try { + await this.#handleStore.save(cwd, handle); + } catch (err) { + this.#logger.warn({ err, cwd }, "failed to persist session handle"); + } + } + + #emitSessionState(state: SessionState): void { + this.onSessionState?.(state); + } +} + +function newSessionId(): string { + return `ses_${randomBytes(9).toString("base64url")}`; +} diff --git a/packages/server/src/server/agent/agent-sdk-types.ts b/packages/server/src/server/agent/agent-sdk-types.ts new file mode 100644 index 0000000..e04efc0 --- /dev/null +++ b/packages/server/src/server/agent/agent-sdk-types.ts @@ -0,0 +1,78 @@ +import type { AgentEvent, ProviderModel, ProviderMode } from "@echohello/protocol"; + +/** + * The seam between the daemon and agent runtimes (see docs/providers.md). + * Each provider (claude, opencode, cursor-acp, ...) implements `AgentClient`; + * the agent-manager owns session lifecycle on top of it. + */ + +/** + * Everything needed to resume an upstream session after a daemon restart. + * Persisted per-session at `/agents//.json`. + */ +export interface PersistenceHandle { + /** Provider id, e.g. "claude". */ + provider: string; + /** The upstream provider's own session id. */ + sessionId: string; + /** cwd, model, modeId, thinkingOption, systemPrompt, ... */ + metadata?: Record | undefined; +} + +type DistributiveOmit = T extends unknown ? Omit : never; + +/** + * An `AgentEvent` without `sessionId`. Providers emit these; the agent-manager + * stamps the Supaplane-side session id before broadcasting. + */ +export type SessionScopedEvent = DistributiveOmit; + +export type AgentEventSink = (event: SessionScopedEvent) => void; + +/** A live agent session. Implementations wrap one upstream conversation. */ +export interface AgentSession { + /** Send a user prompt into the session. Resolves once accepted, not when the turn ends. */ + send(prompt: string, attachments?: unknown[]): Promise; + /** Interrupt the in-flight turn, if any. */ + abort(): Promise; + /** Tear down the underlying provider process/resources. */ + dispose(): Promise; +} + +export interface AgentSessionHandle { + session: AgentSession; + handle: PersistenceHandle; +} + +export interface CreateSessionArgs { + cwd: string; + modelId?: string; + modeId?: string; + signal?: AbortSignal; +} + +export interface ResumeSessionArgs { + handle: PersistenceHandle; + cwd: string; + overrides?: { + modelId?: string | undefined; + modeId?: string | undefined; + }; +} + +export interface ImportableSession { + sessionId: string; + title?: string; + cwd?: string; + updatedAt: number; +} + +export interface AgentClient { + readonly providerId: string; + createSession(args: CreateSessionArgs, emit: AgentEventSink): Promise; + resumeSession(args: ResumeSessionArgs, emit: AgentEventSink): Promise; + listModels(): Promise; + listModes(): Promise; + listImportableSessions(): Promise; + getDiagnostic?(): Promise<{ diagnostic: string }>; +} diff --git a/packages/server/src/server/agent/handle-store.test.ts b/packages/server/src/server/agent/handle-store.test.ts new file mode 100644 index 0000000..938f5bc --- /dev/null +++ b/packages/server/src/server/agent/handle-store.test.ts @@ -0,0 +1,46 @@ +import { mkdtemp, readFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +import { describe, expect, it } from "vitest"; + +import { HandleStore, sanitizeCwd } from "./handle-store.js"; + +describe("sanitizeCwd", () => { + it("replaces path separators and collapses dashes", () => { + expect(sanitizeCwd("/Users/foo/my-project")).toBe("Users-foo-my-project"); + }); + + it("returns 'root' for an all-symbol cwd", () => { + expect(sanitizeCwd("///")).toBe("root"); + }); +}); + +describe("HandleStore", () => { + it("round-trips a persistence handle", async () => { + const home = await mkdtemp(join(tmpdir(), "supaplane-handles-")); + const store = new HandleStore(home); + const handle = { + provider: "claude", + sessionId: "abc-123", + metadata: { cwd: "/tmp/proj", modelId: "sonnet" }, + }; + await store.save("/tmp/proj", handle); + + const loaded = await store.load("/tmp/proj", "claude", "abc-123"); + expect(loaded).toEqual(handle); + + const raw = await readFile(join(home, "agents", "tmp-proj", "abc-123.json"), "utf8"); + expect(JSON.parse(raw)).toEqual(handle); + }); + + it("returns null for missing or mismatched handles", async () => { + const home = await mkdtemp(join(tmpdir(), "supaplane-handles-")); + const store = new HandleStore(home); + await store.save("/tmp/proj", { provider: "claude", sessionId: "abc-123" }); + + expect(await store.load("/tmp/proj", "claude", "nope")).toBeNull(); + expect(await store.load("/tmp/other", "claude", "abc-123")).toBeNull(); + expect(await store.load("/tmp/proj", "opencode", "abc-123")).toBeNull(); + }); +}); diff --git a/packages/server/src/server/agent/handle-store.ts b/packages/server/src/server/agent/handle-store.ts new file mode 100644 index 0000000..0cfc448 --- /dev/null +++ b/packages/server/src/server/agent/handle-store.ts @@ -0,0 +1,55 @@ +import { mkdir, readFile, writeFile } from "node:fs/promises"; +import { join } from "node:path"; + +import type { PersistenceHandle } from "./agent-sdk-types.js"; + +/** + * Persists `PersistenceHandle`s to disk so sessions can be resumed across + * daemon restarts. Layout: `/agents//.json`. + */ +export class HandleStore { + readonly #home: string; + + constructor(supaplaneHome: string) { + this.#home = supaplaneHome; + } + + async save(cwd: string, handle: PersistenceHandle): Promise { + const dir = this.#dirFor(cwd); + await mkdir(dir, { recursive: true }); + const path = join(dir, `${handle.sessionId}.json`); + await writeFile(path, JSON.stringify(handle, null, 2), "utf8"); + } + + async load(cwd: string, provider: string, sessionId: string): Promise { + const path = join(this.#dirFor(cwd), `${sessionId}.json`); + let raw: string; + try { + raw = await readFile(path, "utf8"); + } catch { + return null; + } + const parsed: unknown = JSON.parse(raw); + if (typeof parsed !== "object" || parsed === null) return null; + const record = parsed as Record; + if (record.provider !== provider || record.sessionId !== sessionId) return null; + const handle: PersistenceHandle = { provider, sessionId }; + if (typeof record.metadata === "object" && record.metadata !== null) { + handle.metadata = record.metadata as Record; + } + return handle; + } + + #dirFor(cwd: string): string { + return join(this.#home, "agents", sanitizeCwd(cwd)); + } +} + +export function sanitizeCwd(cwd: string): string { + return ( + cwd + .replace(/[^a-zA-Z0-9]/g, "-") + .replace(/-+/g, "-") + .replace(/^-|-$/g, "") || "root" + ); +} diff --git a/packages/server/src/server/agent/providers/claude/claude-provider.ts b/packages/server/src/server/agent/providers/claude/claude-provider.ts new file mode 100644 index 0000000..69446c5 --- /dev/null +++ b/packages/server/src/server/agent/providers/claude/claude-provider.ts @@ -0,0 +1,320 @@ +import { execFile } from "node:child_process"; + +import { + listSessions, + query, + type Options, + type PermissionMode, + type Query, + type SDKUserMessage, +} from "@anthropic-ai/claude-agent-sdk"; +import { ProviderError, type ProviderModel, type ProviderMode } from "@echohello/protocol"; + +import type { + AgentClient, + AgentEventSink, + AgentSession, + AgentSessionHandle, + CreateSessionArgs, + ImportableSession, + ResumeSessionArgs, +} from "../../agent-sdk-types.js"; +import { ClaudeToolCallMapper } from "./tool-call-mapper.js"; + +const CLAUDE_MODELS: ProviderModel[] = [ + { id: "sonnet", label: "Claude Sonnet", reasoning: true, vision: true }, + { id: "opus", label: "Claude Opus", reasoning: true, vision: true }, + { id: "haiku", label: "Claude Haiku", reasoning: false, vision: true }, +]; + +const CLAUDE_MODES: ProviderMode[] = [ + { id: "default", label: "Default", isUnattended: false, features: [] }, + { id: "plan", label: "Plan", isUnattended: false, features: ["read-only"] }, + { id: "accept-edits", label: "Accept edits", isUnattended: false, features: ["auto-edit"] }, + { + id: "bypass", + label: "Bypass permissions", + description: "Auto-accept all tool calls. Unattended runs only.", + isUnattended: true, + features: ["unattended"], + }, +]; + +const MODE_TO_PERMISSION: Record = { + default: "default", + plan: "plan", + "accept-edits": "acceptEdits", + bypass: "bypassPermissions", +}; + +export interface ClaudeAgentClientOptions { + /** Provider binary + args, e.g. `["claude"]`. */ + command?: string[]; +} + +/** + * Claude Code provider (`@anthropic-ai/claude-agent-sdk`). Spawns one + * streaming-input `query()` per session; user prompts are pushed through an + * async queue after the `system/init` message yields the upstream session id. + */ +export class ClaudeAgentClient implements AgentClient { + readonly providerId = "claude"; + readonly #command: string[]; + + constructor(options: ClaudeAgentClientOptions = {}) { + this.#command = options.command ?? ["claude"]; + } + + createSession(args: CreateSessionArgs, emit: AgentEventSink): Promise { + return this.#spawn(args, emit); + } + + resumeSession(args: ResumeSessionArgs, emit: AgentEventSink): Promise { + const modelId = + args.overrides?.modelId ?? + (typeof args.handle.metadata?.modelId === "string" + ? args.handle.metadata.modelId + : undefined); + const modeId = + args.overrides?.modeId ?? + (typeof args.handle.metadata?.modeId === "string" ? args.handle.metadata.modeId : undefined); + return this.#spawn( + { + cwd: args.cwd, + ...(modelId !== undefined ? { modelId } : {}), + ...(modeId !== undefined ? { modeId } : {}), + }, + emit, + args.handle.sessionId, + ); + } + + listModels(): Promise { + return Promise.resolve(CLAUDE_MODELS); + } + + listModes(): Promise { + return Promise.resolve(CLAUDE_MODES); + } + + async listImportableSessions(): Promise { + try { + const sessions = await listSessions(); + return sessions.map((s) => { + const title = s.customTitle ?? s.summary; + return { + sessionId: s.sessionId, + ...(title.length > 0 ? { title } : {}), + ...(s.cwd !== undefined ? { cwd: s.cwd } : {}), + updatedAt: s.lastModified, + }; + }); + } catch { + return []; + } + } + + async getDiagnostic(): Promise<{ diagnostic: string }> { + const [bin, ...baseArgs] = this.#command; + if (bin === undefined) { + throw new ProviderError({ message: "claude provider has no command configured" }); + } + return new Promise((resolvePromise, rejectPromise) => { + execFile(bin, [...baseArgs, "--version"], { timeout: 5000 }, (err, stdout, stderr) => { + if (err) { + rejectPromise( + new ProviderError({ + code: "provider_unavailable", + message: `claude CLI unavailable: ${stderr || err.message}`, + cause: err, + }), + ); + return; + } + resolvePromise({ diagnostic: `claude ${stdout.trim()}` }); + }); + }); + } + + async #spawn( + args: CreateSessionArgs, + emit: AgentEventSink, + resumeSessionId?: string, + ): Promise { + const queue = new MessageQueue(); + const abortController = new AbortController(); + const options: Options = { + cwd: args.cwd, + abortController, + permissionMode: resolvePermissionMode(args.modeId), + ...(args.modelId !== undefined ? { model: args.modelId } : {}), + ...(resumeSessionId !== undefined ? { resume: resumeSessionId } : {}), + }; + + let upstreamSessionId: string | undefined; + let upstreamModel: string | undefined; + const mapper = new ClaudeToolCallMapper(); + + let resolveInit!: () => void; + let rejectInit!: (err: unknown) => void; + const initSeen = new Promise((resolvePromise, rejectPromise) => { + resolveInit = resolvePromise; + rejectInit = rejectPromise; + }); + + let q: Query; + try { + q = query({ prompt: queue.iterate(), options }); + } catch (err) { + throw new ProviderError({ message: "failed to start claude session", cause: err }); + } + + const consume = async (): Promise => { + try { + for await (const message of q) { + if (message.type === "system" && message.subtype === "init") { + upstreamSessionId = message.session_id; + upstreamModel = message.model; + resolveInit(); + } + for (const event of mapper.map(message)) { + emit(event); + } + } + } catch (err) { + const error = + err instanceof Error ? err : new ProviderError({ message: String(err), cause: err }); + if (abortController.signal.aborted) return; + rejectInit(error); + emit({ + type: "error", + code: "provider_error", + message: error.message, + ts: Date.now(), + }); + } + }; + void consume(); + + const timeout = setTimeout(() => { + rejectInit(new ProviderError({ code: "timeout", message: "claude init timed out" })); + }, 30_000); + try { + await initSeen; + } catch (err) { + queue.close(); + q.close(); + throw err instanceof Error ? new ProviderError({ message: err.message, cause: err }) : err; + } finally { + clearTimeout(timeout); + } + + if (upstreamSessionId === undefined) { + throw new ProviderError({ message: "claude init did not yield a session id" }); + } + const sessionId = upstreamSessionId; + const session = new ClaudeAgentSession(q, queue, abortController, () => sessionId); + + return { + session, + handle: { + provider: this.providerId, + sessionId, + metadata: { + cwd: args.cwd, + ...(upstreamModel !== undefined ? { modelId: upstreamModel } : {}), + ...(args.modeId !== undefined ? { modeId: args.modeId } : {}), + }, + }, + }; + } +} + +class ClaudeAgentSession implements AgentSession { + readonly #query: Query; + readonly #queue: MessageQueue; + readonly #abortController: AbortController; + readonly #sessionId: () => string; + + constructor( + query_: Query, + queue: MessageQueue, + abortController: AbortController, + sessionId: () => string, + ) { + this.#query = query_; + this.#queue = queue; + this.#abortController = abortController; + this.#sessionId = sessionId; + } + + send(prompt: string, attachments?: unknown[]): Promise { + void attachments; + const message: SDKUserMessage = { + type: "user", + session_id: this.#sessionId(), + parent_tool_use_id: null, + message: { role: "user", content: prompt }, + }; + this.#queue.push(message); + return Promise.resolve(); + } + + async abort(): Promise { + try { + await this.#query.interrupt(); + } catch { + this.#abortController.abort(); + } + } + + dispose(): Promise { + this.#queue.close(); + this.#query.close(); + this.#abortController.abort(); + return Promise.resolve(); + } +} + +function resolvePermissionMode(modeId: string | undefined): PermissionMode { + if (modeId === undefined) return "default"; + return MODE_TO_PERMISSION[modeId] ?? "default"; +} + +class MessageQueue { + readonly #pending: SDKUserMessage[] = []; + #waiter: (() => void) | null = null; + #closed = false; + + push(message: SDKUserMessage): void { + if (this.#closed) return; + this.#pending.push(message); + this.#wake(); + } + + close(): void { + this.#closed = true; + this.#wake(); + } + + async *iterate(): AsyncIterable { + for (;;) { + const next = this.#pending.shift(); + if (next !== undefined) { + yield next; + continue; + } + if (this.#closed) return; + await new Promise((resolvePromise) => { + this.#waiter = resolvePromise; + }); + this.#waiter = null; + } + } + + #wake(): void { + const waiter = this.#waiter; + this.#waiter = null; + waiter?.(); + } +} diff --git a/packages/server/src/server/agent/providers/claude/tool-call-mapper.test.ts b/packages/server/src/server/agent/providers/claude/tool-call-mapper.test.ts new file mode 100644 index 0000000..c402217 --- /dev/null +++ b/packages/server/src/server/agent/providers/claude/tool-call-mapper.test.ts @@ -0,0 +1,113 @@ +import type { SDKMessage } from "@anthropic-ai/claude-agent-sdk"; +import { describe, expect, it } from "vitest"; + +import { ClaudeToolCallMapper } from "./tool-call-mapper.js"; + +function assistantMessage(content: unknown[]): SDKMessage { + return { + type: "assistant", + uuid: "u-1", + session_id: "s-1", + parent_tool_use_id: null, + message: { role: "assistant", content }, + } as unknown as SDKMessage; +} + +describe("ClaudeToolCallMapper", () => { + it("maps text blocks to message.final", () => { + const mapper = new ClaudeToolCallMapper(); + const events = mapper.map(assistantMessage([{ type: "text", text: "hello" }])); + expect(events).toEqual([ + { type: "message.final", partId: "u-1:0", text: "hello", ts: expect.any(Number) }, + ]); + }); + + it("maps thinking blocks to reasoning deltas", () => { + const mapper = new ClaudeToolCallMapper(); + const events = mapper.map(assistantMessage([{ type: "thinking", thinking: "hmm" }])); + expect(events).toEqual([ + { + type: "message.delta", + partId: "u-1:0", + text: "hmm", + reasoning: true, + ts: expect.any(Number), + }, + ]); + }); + + it("maps tool_use to tool.start and tool_result to tool.result with duration", () => { + const mapper = new ClaudeToolCallMapper(); + const [start] = mapper.map( + assistantMessage([{ type: "tool_use", id: "tool-1", name: "Read", input: { path: "/x" } }]), + ); + expect(start).toMatchObject({ type: "tool.start", toolCallId: "tool-1", name: "Read" }); + + const userMessage = { + type: "user", + uuid: "u-2", + session_id: "s-1", + parent_tool_use_id: null, + message: { + role: "user", + content: [{ type: "tool_result", tool_use_id: "tool-1", content: "file contents" }], + }, + } as unknown as SDKMessage; + const [result] = mapper.map(userMessage); + expect(result).toMatchObject({ + type: "tool.result", + toolCallId: "tool-1", + output: "file contents", + }); + expect(result?.type === "tool.result" && result.durationMs >= 0).toBe(true); + }); + + it("maps result success to idle status and errors to error events", () => { + const mapper = new ClaudeToolCallMapper(); + const success = { + type: "result", + subtype: "success", + uuid: "u-3", + session_id: "s-1", + } as unknown as SDKMessage; + expect(mapper.map(success)).toEqual([ + { type: "status", status: "idle", ts: expect.any(Number) }, + ]); + + const failure = { + type: "result", + subtype: "error_during_execution", + errors: ["boom"], + uuid: "u-4", + session_id: "s-1", + } as unknown as SDKMessage; + expect(mapper.map(failure)).toEqual([ + { + type: "error", + code: "provider_error", + message: "error_during_execution: boom", + ts: expect.any(Number), + }, + ]); + }); + + it("maps session_state_changed to status events", () => { + const mapper = new ClaudeToolCallMapper(); + const changed = { + type: "system", + subtype: "session_state_changed", + state: "requires_action", + uuid: "u-5", + session_id: "s-1", + } as unknown as SDKMessage; + expect(mapper.map(changed)).toEqual([ + { type: "status", status: "waiting", ts: expect.any(Number) }, + ]); + }); + + it("ignores unrelated message types", () => { + const mapper = new ClaudeToolCallMapper(); + const other = { type: "system", subtype: "init", session_id: "s-1" } as unknown as SDKMessage; + expect(mapper.map(other)).toEqual([]); + }); +}); diff --git a/packages/server/src/server/agent/providers/claude/tool-call-mapper.ts b/packages/server/src/server/agent/providers/claude/tool-call-mapper.ts new file mode 100644 index 0000000..998a1a9 --- /dev/null +++ b/packages/server/src/server/agent/providers/claude/tool-call-mapper.ts @@ -0,0 +1,124 @@ +import type { SDKMessage } from "@anthropic-ai/claude-agent-sdk"; + +import type { SessionScopedEvent } from "../../agent-sdk-types.js"; + +interface TextBlock { + type: "text"; + text: string; +} +interface ThinkingBlock { + type: "thinking"; + thinking: string; +} +interface ToolUseBlock { + type: "tool_use"; + id: string; + name: string; + input: unknown; +} +interface ToolResultBlock { + type: "tool_result"; + tool_use_id: string; + content: unknown; + is_error?: boolean; +} +type ContentBlock = TextBlock | ThinkingBlock | ToolUseBlock | ToolResultBlock; + +/** + * Maps provider-native Claude Agent SDK messages onto Supaplane + * `SessionScopedEvent`s. Stateful: tracks tool_use start times so + * `tool.result` events can carry `durationMs`. + */ +export class ClaudeToolCallMapper { + readonly #toolStarts = new Map(); + + map(message: SDKMessage): SessionScopedEvent[] { + const ts = Date.now(); + switch (message.type) { + case "assistant": { + if (message.error !== undefined) { + return [ + { + type: "error", + code: "provider_error", + message: `assistant message error: ${message.error}`, + ts, + }, + ]; + } + const blocks = message.message.content as unknown as ContentBlock[]; + const events: SessionScopedEvent[] = []; + blocks.forEach((block, index) => { + if (block.type === "text" && block.text.length > 0) { + events.push({ + type: "message.final", + partId: `${message.uuid}:${index}`, + text: block.text, + ts, + }); + } else if (block.type === "thinking" && block.thinking.length > 0) { + events.push({ + type: "message.delta", + partId: `${message.uuid}:${index}`, + text: block.thinking, + reasoning: true, + ts, + }); + } else if (block.type === "tool_use") { + this.#toolStarts.set(block.id, ts); + events.push({ + type: "tool.start", + toolCallId: block.id, + name: block.name, + input: block.input, + ts, + }); + } + }); + return events; + } + case "user": { + const content = message.message.content; + if (typeof content === "string") return []; + const blocks = content as unknown as ContentBlock[]; + const events: SessionScopedEvent[] = []; + for (const block of blocks) { + if (block.type !== "tool_result") continue; + const startedAt = this.#toolStarts.get(block.tool_use_id); + this.#toolStarts.delete(block.tool_use_id); + events.push({ + type: "tool.result", + toolCallId: block.tool_use_id, + output: block.content, + durationMs: startedAt === undefined ? 0 : Math.max(0, ts - startedAt), + ts, + }); + } + return events; + } + case "result": { + if (message.subtype === "success") { + return [{ type: "status", status: "idle", ts }]; + } + return [ + { + type: "error", + code: "provider_error", + message: `${message.subtype}: ${message.errors.join("; ")}`, + ts, + }, + ]; + } + case "system": { + if (message.subtype === "session_state_changed") { + if (message.state === "running") return [{ type: "status", status: "running", ts }]; + if (message.state === "idle") return [{ type: "status", status: "idle", ts }]; + return [{ type: "status", status: "waiting", ts }]; + } + return []; + } + default: + return []; + } + } +} From 3c7f3ce6e550fced5894d68db9aeb1257c2e1550 Mon Sep 17 00:00:00 2001 From: Johnny Huynh <27847622+johnnyhuy@users.noreply.github.com> Date: Thu, 6 Aug 2026 08:35:35 +1000 Subject: [PATCH 2/2] feat(server): wire command dispatcher + workspace registry into daemon session.start/resume/send/abort and workspace.open/refresh now run end-to-end; unimplemented commands return explicit error frames; hello_ack advertises actually-registered providers. Adds daemon E2E tests with a fake provider covering the full session lifecycle. Co-authored-by: opencode-agent --- packages/server/src/daemon.ts | 38 +++ packages/server/src/handshake.ts | 4 +- .../server/src/server/command-dispatcher.ts | 159 ++++++++++++ .../daemon-e2e/session-lifecycle.test.ts | 243 ++++++++++++++++++ packages/server/src/server/exports.ts | 17 ++ .../server/src/server/workspace-registry.ts | 51 ++++ packages/server/src/websocket-server.ts | 10 + 7 files changed, 521 insertions(+), 1 deletion(-) create mode 100644 packages/server/src/server/command-dispatcher.ts create mode 100644 packages/server/src/server/daemon-e2e/session-lifecycle.test.ts create mode 100644 packages/server/src/server/workspace-registry.ts diff --git a/packages/server/src/daemon.ts b/packages/server/src/daemon.ts index 763b10d..95a4e73 100644 --- a/packages/server/src/daemon.ts +++ b/packages/server/src/daemon.ts @@ -6,6 +6,11 @@ import { loadOrCreateIdentity } from "./handshake.js"; import { createHttpApp } from "./http-app.js"; import { createLogger } from "./logger.js"; import { resolveSupaplaneHome, SUPAPLANE_VERSION } from "./paths.js"; +import { AgentManager } from "./server/agent/agent-manager.js"; +import { HandleStore } from "./server/agent/handle-store.js"; +import { ClaudeAgentClient } from "./server/agent/providers/claude/claude-provider.js"; +import { CommandDispatcher } from "./server/command-dispatcher.js"; +import { WorkspaceRegistry } from "./server/workspace-registry.js"; import { SupaplaneWebsocketServer } from "./websocket-server.js"; export interface DaemonHandle { @@ -15,6 +20,8 @@ export interface DaemonHandle { stop: () => Promise; httpServer: ReturnType; wsServer: SupaplaneWebsocketServer; + agentManager: AgentManager; + workspaces: WorkspaceRegistry; } /** @@ -54,14 +61,42 @@ export async function startDaemon(args?: { }); const httpServer = createServer(httpApp); + + const handleStore = new HandleStore(supaplaneHome); + const agentManager = new AgentManager({ handleStore, logger }); + agentManager.registerProvider(new ClaudeAgentClient()); + const workspaces = new WorkspaceRegistry(); + const wsServer = new SupaplaneWebsocketServer({ httpServer, logger, identity, ...(config.daemonAuthToken ? { authToken: config.daemonAuthToken } : {}), serverVersion: SUPAPLANE_VERSION, + providers: agentManager.providerIds(), }); + agentManager.onAgentEvent = (event) => wsServer.broadcast({ kind: "event", event }); + agentManager.onSessionState = (session) => wsServer.broadcast({ kind: "session_state", session }); + + const dispatcher = new CommandDispatcher({ + workspaces, + agents: agentManager, + broadcast: (event) => wsServer.broadcast(event), + logger, + }); + wsServer.setCommandHandler((cmd, session) => + dispatcher.handle(cmd, { + clientId: session.clientId, + sendError: (error) => + wsServer.sendTo(session.socket, { + type: "error", + code: error.code, + message: error.message, + }), + }), + ); + await new Promise((resolve, reject) => { const onError = (err: Error) => { httpServer.off("listening", onListening); @@ -84,8 +119,11 @@ export async function startDaemon(args?: { supaplaneHome, httpServer, wsServer, + agentManager, + workspaces, async stop(): Promise { logger.info("stopping daemon"); + await agentManager.disposeAll(); await new Promise((resolve, reject) => { httpServer.close((err) => (err ? reject(err) : resolve())); }); diff --git a/packages/server/src/handshake.ts b/packages/server/src/handshake.ts index 862c393..192cd0d 100644 --- a/packages/server/src/handshake.ts +++ b/packages/server/src/handshake.ts @@ -15,6 +15,8 @@ import { getOrCreateServerId } from "./server-id.js"; export interface WebsocketSessionContext { serverId: string; serverVersion: string; + /** Provider ids advertised in the `hello_ack` capabilities. */ + providers: readonly string[]; daemonLabel?: string; logger: Logger; } @@ -51,7 +53,7 @@ export function handleHello(args: { serverVersion: args.ctx.serverVersion, protocolVersion: SUPAPLANE_PROTOCOL_VERSION, capabilities: { - providers: ["opencode", "claude", "cursor"], + providers: [...args.ctx.providers], relay: false, worktrees: false, scheduling: false, diff --git a/packages/server/src/server/command-dispatcher.ts b/packages/server/src/server/command-dispatcher.ts new file mode 100644 index 0000000..db33e6d --- /dev/null +++ b/packages/server/src/server/command-dispatcher.ts @@ -0,0 +1,159 @@ +import type { Logger } from "pino"; +import { + AgentError, + SupaplaneError, + type ClientCommand, + type ServerEvent, +} from "@echohello/protocol"; + +import type { AgentManager } from "./agent/agent-manager.js"; +import type { WorkspaceRegistry } from "./workspace-registry.js"; + +export interface CommandContext { + clientId: string; + /** Report an error back to the originating socket (no ack frames exist for one-way commands). */ + sendError: (error: { code: string; message: string }) => void; +} + +export interface CommandDispatcherOptions { + workspaces: WorkspaceRegistry; + agents: AgentManager; + broadcast: (event: ServerEvent) => void; + logger: Logger; +} + +/** + * Routes validated `ClientCommand`s from the WebSocket server to the + * workspace registry / agent manager, and broadcasts resulting state. + * + * Commands that are valid per the protocol but not yet implemented at the + * daemon layer are answered with a `bad_request` error frame, never dropped. + */ +export class CommandDispatcher { + readonly #workspaces: WorkspaceRegistry; + readonly #agents: AgentManager; + readonly #broadcast: (event: ServerEvent) => void; + readonly #logger: Logger; + + constructor(options: CommandDispatcherOptions) { + this.#workspaces = options.workspaces; + this.#agents = options.agents; + this.#broadcast = options.broadcast; + this.#logger = options.logger.child({ module: "command-dispatcher" }); + } + + handle(cmd: ClientCommand, ctx: CommandContext): void { + switch (cmd.type) { + case "workspace.open": { + const workspace = this.#workspaces.open(cmd.cwd); + this.#broadcast({ kind: "workspace_state", workspace }); + return; + } + case "workspace.refresh": { + const workspace = this.#workspaces.refresh(cmd.workspaceId); + if (workspace === undefined) { + this.#fail( + ctx, + new AgentError({ + code: "not_found", + message: `Unknown workspace: ${cmd.workspaceId}`, + }), + ); + return; + } + this.#broadcast({ kind: "workspace_state", workspace }); + return; + } + case "session.start": { + void this.#guard(ctx, async () => { + const workspace = this.#requireWorkspace(cmd.workspaceId); + const state = await this.#agents.startSession({ + workspaceId: workspace.workspaceId, + cwd: workspace.cwd, + providerId: cmd.providerId, + ...(cmd.modelId !== undefined ? { modelId: cmd.modelId } : {}), + ...(cmd.modeId !== undefined ? { modeId: cmd.modeId } : {}), + ...(cmd.initialPrompt !== undefined ? { initialPrompt: cmd.initialPrompt } : {}), + }); + this.#broadcast({ kind: "session_state", session: state }); + }); + return; + } + case "session.resume": { + void this.#guard(ctx, async () => { + const cwd = + typeof cmd.handle.metadata?.cwd === "string" ? cmd.handle.metadata.cwd : process.cwd(); + const workspace = this.#workspaces.open(cwd); + const state = await this.#agents.resumeSession({ + workspaceId: workspace.workspaceId, + cwd: workspace.cwd, + handle: cmd.handle, + ...(cmd.overrides !== undefined ? { overrides: cmd.overrides } : {}), + }); + this.#broadcast({ kind: "session_state", session: state }); + }); + return; + } + case "session.send": { + void this.#guard(ctx, async () => { + await this.#agents.send(cmd.sessionId, cmd.prompt, cmd.attachments); + }); + return; + } + case "session.abort": { + void this.#guard(ctx, async () => { + await this.#agents.abort(cmd.sessionId); + }); + return; + } + case "session.fork": + case "diff.open": + case "file.open": + case "git.checkout": + case "permission.resolve": { + this.#fail( + ctx, + new AgentError({ code: "bad_request", message: `${cmd.type} is not implemented yet` }), + ); + return; + } + case "ping": + case "pong": + case "subscribe": + case "unsubscribe": { + // Handled by the WebSocket server transport layer. + return; + } + } + } + + #requireWorkspace(workspaceId: string) { + const workspace = this.#workspaces.get(workspaceId); + if (workspace === undefined) { + throw new AgentError({ code: "not_found", message: `Unknown workspace: ${workspaceId}` }); + } + return workspace; + } + + async #guard(ctx: CommandContext, run: () => Promise): Promise { + try { + await run(); + } catch (err) { + this.#fail(ctx, err); + } + } + + #fail(ctx: CommandContext, err: unknown): void { + if (err instanceof SupaplaneError) { + this.#logger.warn( + { clientId: ctx.clientId, code: err.code, err: err.message }, + "command failed", + ); + ctx.sendError({ code: err.code, message: err.message }); + return; + } + const message = err instanceof Error ? err.message : String(err); + this.#logger.warn({ clientId: ctx.clientId, err: message }, "command failed"); + ctx.sendError({ code: "internal", message }); + } +} diff --git a/packages/server/src/server/daemon-e2e/session-lifecycle.test.ts b/packages/server/src/server/daemon-e2e/session-lifecycle.test.ts new file mode 100644 index 0000000..8de83a5 --- /dev/null +++ b/packages/server/src/server/daemon-e2e/session-lifecycle.test.ts @@ -0,0 +1,243 @@ +import { mkdtemp } from "node:fs/promises"; +import type { AddressInfo } from "node:net"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +import { SupaplaneClient } from "@echohello/client"; +import type { ServerEvent, SessionState, WorkspaceState } from "@echohello/protocol"; +import { afterEach, beforeEach, describe, expect, it } from "vitest"; + +import { startDaemon, type DaemonHandle } from "../../daemon.js"; +import type { + AgentClient, + AgentEventSink, + AgentSession, + AgentSessionHandle, + CreateSessionArgs, + ImportableSession, + ResumeSessionArgs, +} from "../agent/agent-sdk-types.js"; + +class FakeAgentSession implements AgentSession { + readonly #emit: AgentEventSink; + + constructor(emit: AgentEventSink) { + this.#emit = emit; + } + + send(prompt: string): Promise { + const ts = Date.now(); + this.#emit({ type: "status", status: "running", ts }); + this.#emit({ type: "message.delta", partId: "p1", text: `echo: ${prompt}`, ts }); + this.#emit({ type: "message.final", partId: "p1", text: `echo: ${prompt}`, ts: Date.now() }); + this.#emit({ type: "status", status: "idle", ts: Date.now() }); + return Promise.resolve(); + } + + abort(): Promise { + this.#emit({ type: "status", status: "idle", ts: Date.now() }); + return Promise.resolve(); + } + + dispose(): Promise { + return Promise.resolve(); + } +} + +class FakeProvider implements AgentClient { + readonly providerId = "fake"; + readonly createdWith: CreateSessionArgs[] = []; + + createSession(args: CreateSessionArgs, emit: AgentEventSink): Promise { + this.createdWith.push(args); + return Promise.resolve({ + session: new FakeAgentSession(emit), + handle: { + provider: this.providerId, + sessionId: "upstream-1", + metadata: { cwd: args.cwd }, + }, + }); + } + + resumeSession(args: ResumeSessionArgs, emit: AgentEventSink): Promise { + return Promise.resolve({ + session: new FakeAgentSession(emit), + handle: args.handle, + }); + } + + listModels(): Promise { + return Promise.resolve([]); + } + + listModes(): Promise { + return Promise.resolve([]); + } + + listImportableSessions(): Promise { + return Promise.resolve([]); + } +} + +async function waitFor( + collected: T[], + predicate: (item: T) => boolean, + timeoutMs = 5000, +): Promise { + const deadline = Date.now() + timeoutMs; + for (;;) { + const found = collected.find(predicate); + if (found !== undefined) return found; + if (Date.now() > deadline) throw new Error("waitFor timed out"); + await new Promise((resolvePromise) => setTimeout(resolvePromise, 25)); + } +} + +describe("daemon e2e: session lifecycle", () => { + let daemon: DaemonHandle; + let client: SupaplaneClient; + let events: ServerEvent[]; + let provider: FakeProvider; + + beforeEach(async () => { + const supaplaneHome = await mkdtemp(join(tmpdir(), "supaplane-e2e-")); + daemon = await startDaemon({ + config: { listenPort: 0, logLevel: "error" }, + supaplaneHome, + }); + provider = new FakeProvider(); + daemon.agentManager.registerProvider(provider); + + const { port } = daemon.httpServer.address() as AddressInfo; + client = new SupaplaneClient({ + endpoint: `ws://127.0.0.1:${port}`, + clientId: "e2e-client", + clientType: "cli", + reconnect: false, + }); + events = []; + client.onServerEvent((event) => events.push(event)); + await client.connect(); + }); + + afterEach(async () => { + client.close(); + await daemon.stop(); + }); + + it("advertises registered providers in hello_ack", () => { + expect(client.helloAck?.capabilities.providers).toContain("claude"); + }); + + it("workspace.open → workspace_state", async () => { + client.sendCommand({ type: "workspace.open", cwd: "/tmp/supaplane-e2e-ws" }); + const event = await waitFor(events, (e) => e.kind === "workspace_state"); + if (event.kind !== "workspace_state") throw new Error("unreachable"); + const workspace: WorkspaceState = event.workspace; + expect(workspace.cwd).toBe("/tmp/supaplane-e2e-ws"); + expect(workspace.workspaceId).toMatch(/^ws_/); + }); + + it("session.start with initialPrompt streams agent events end-to-end", async () => { + client.sendCommand({ type: "workspace.open", cwd: "/tmp/supaplane-e2e-ws" }); + const wsEvent = await waitFor(events, (e) => e.kind === "workspace_state"); + if (wsEvent.kind !== "workspace_state") throw new Error("unreachable"); + + client.sendCommand({ + type: "session.start", + workspaceId: wsEvent.workspace.workspaceId, + providerId: "fake", + modelId: "test-model", + initialPrompt: "hello agent", + }); + + const sessionEvent = await waitFor( + events, + (e) => e.kind === "session_state" && e.session.providerId === "fake", + ); + if (sessionEvent.kind !== "session_state") throw new Error("unreachable"); + const session: SessionState = sessionEvent.session; + expect(session.sessionId).toMatch(/^ses_/); + expect(session.workspaceId).toBe(wsEvent.workspace.workspaceId); + expect(session.modelId).toBe("test-model"); + + expect(provider.createdWith).toHaveLength(1); + expect(provider.createdWith[0]?.cwd).toBe("/tmp/supaplane-e2e-ws"); + expect(provider.createdWith[0]?.modelId).toBe("test-model"); + + const final = await waitFor( + events, + (e) => + e.kind === "event" && + e.event.type === "message.final" && + e.event.sessionId === session.sessionId, + ); + if (final.kind !== "event" || final.event.type !== "message.final") { + throw new Error("unreachable"); + } + expect(final.event.text).toBe("echo: hello agent"); + + await waitFor( + events, + (e) => + e.kind === "session_state" && + e.session.sessionId === session.sessionId && + e.session.status === "idle", + ); + }); + + it("session.send + session.abort on a running session", async () => { + client.sendCommand({ type: "workspace.open", cwd: "/tmp/supaplane-e2e-ws" }); + const wsEvent = await waitFor(events, (e) => e.kind === "workspace_state"); + if (wsEvent.kind !== "workspace_state") throw new Error("unreachable"); + + client.sendCommand({ + type: "session.start", + workspaceId: wsEvent.workspace.workspaceId, + providerId: "fake", + }); + const sessionEvent = await waitFor(events, (e) => e.kind === "session_state"); + if (sessionEvent.kind !== "session_state") throw new Error("unreachable"); + const sessionId = sessionEvent.session.sessionId; + + client.sendCommand({ type: "session.send", sessionId, prompt: "second turn", attachments: [] }); + const final = await waitFor( + events, + (e) => + e.kind === "event" && + e.event.type === "message.final" && + e.event.sessionId === sessionId && + e.event.text === "echo: second turn", + ); + expect(final.kind).toBe("event"); + + client.sendCommand({ type: "session.abort", sessionId }); + const idle = await waitFor( + events, + (e) => + e.kind === "session_state" && + e.session.sessionId === sessionId && + e.session.status === "idle", + ); + if (idle.kind !== "session_state") throw new Error("unreachable"); + expect(idle.session.status).toBe("idle"); + }); + + it("unknown provider surfaces an error frame, not silence", async () => { + client.sendCommand({ type: "workspace.open", cwd: "/tmp/supaplane-e2e-ws" }); + const wsEvent = await waitFor(events, (e) => e.kind === "workspace_state"); + if (wsEvent.kind !== "workspace_state") throw new Error("unreachable"); + + client.sendCommand({ + type: "session.start", + workspaceId: wsEvent.workspace.workspaceId, + providerId: "does-not-exist", + }); + + await new Promise((resolvePromise) => setTimeout(resolvePromise, 300)); + // Error frames travel outside ServerEvent; assert no session was created. + expect(provider.createdWith).toHaveLength(0); + expect(daemon.agentManager.listSessions()).toHaveLength(0); + }); +}); diff --git a/packages/server/src/server/exports.ts b/packages/server/src/server/exports.ts index 436c064..c8230c7 100644 --- a/packages/server/src/server/exports.ts +++ b/packages/server/src/server/exports.ts @@ -7,3 +7,20 @@ export { loadDaemonConfig, DaemonConfigSchema, type DaemonConfig } from "../conf export { resolveSupaplaneHome, SUPAPLANE_VERSION } from "../paths.js"; export { getOrCreateServerId } from "../server-id.js"; export { loadOrCreateDaemonKeyPair, type DaemonKeyPair } from "../daemon-keypair.js"; +export { AgentManager } from "./agent/agent-manager.js"; +export { HandleStore, sanitizeCwd } from "./agent/handle-store.js"; +export { ClaudeAgentClient } from "./agent/providers/claude/claude-provider.js"; +export { ClaudeToolCallMapper } from "./agent/providers/claude/tool-call-mapper.js"; +export { CommandDispatcher, type CommandContext } from "./command-dispatcher.js"; +export { WorkspaceRegistry } from "./workspace-registry.js"; +export type { + AgentClient, + AgentEventSink, + AgentSession, + AgentSessionHandle, + CreateSessionArgs, + ImportableSession, + PersistenceHandle, + ResumeSessionArgs, + SessionScopedEvent, +} from "./agent/agent-sdk-types.js"; diff --git a/packages/server/src/server/workspace-registry.ts b/packages/server/src/server/workspace-registry.ts new file mode 100644 index 0000000..76d217f --- /dev/null +++ b/packages/server/src/server/workspace-registry.ts @@ -0,0 +1,51 @@ +import { basename, resolve } from "node:path"; + +import type { WorkspaceState } from "@echohello/protocol"; + +/** + * In-memory registry of open workspaces. Keyed by workspace id; `open()` + * dedupes on the resolved cwd so the same directory never yields two + * workspaces. + */ +export class WorkspaceRegistry { + readonly #byId = new Map(); + readonly #idByCwd = new Map(); + #counter = 0; + + open(cwd: string): WorkspaceState { + const resolved = resolve(cwd); + const existingId = this.#idByCwd.get(resolved); + if (existingId !== undefined) { + const existing = this.#byId.get(existingId); + if (existing !== undefined) return existing; + } + const now = Date.now(); + const workspace: WorkspaceState = { + workspaceId: `ws_${++this.#counter}`, + cwd: resolved, + repoName: basename(resolved), + dirty: false, + freshness: "active", + updatedAt: now, + }; + this.#byId.set(workspace.workspaceId, workspace); + this.#idByCwd.set(resolved, workspace.workspaceId); + return workspace; + } + + refresh(workspaceId: string): WorkspaceState | undefined { + const workspace = this.#byId.get(workspaceId); + if (workspace === undefined) return undefined; + const updated: WorkspaceState = { ...workspace, updatedAt: Date.now() }; + this.#byId.set(workspaceId, updated); + return updated; + } + + get(workspaceId: string): WorkspaceState | undefined { + return this.#byId.get(workspaceId); + } + + list(): WorkspaceState[] { + return [...this.#byId.values()]; + } +} diff --git a/packages/server/src/websocket-server.ts b/packages/server/src/websocket-server.ts index 93b8e6d..53c2948 100644 --- a/packages/server/src/websocket-server.ts +++ b/packages/server/src/websocket-server.ts @@ -29,6 +29,8 @@ export interface WebsocketServerOptions { daemonLabel?: string; serverVersion: string; onCommand?: (cmd: ClientCommand, session: SessionRecord) => void; + /** Provider ids advertised in the `hello_ack` capabilities. */ + providers?: readonly string[]; } /** @@ -49,6 +51,7 @@ export class SupaplaneWebsocketServer { #serverVersion: string; #daemonLabel: string | undefined; #onCommand?: (cmd: ClientCommand, session: SessionRecord) => void; + #providers: readonly string[]; constructor(options: WebsocketServerOptions) { this.#logger = options.logger.child({ module: "ws-server" }); @@ -56,6 +59,7 @@ export class SupaplaneWebsocketServer { this.#authToken = options.authToken; this.#serverVersion = options.serverVersion; this.#daemonLabel = options.daemonLabel; + this.#providers = options.providers ?? []; if (options.onCommand) { this.#onCommand = options.onCommand; } @@ -77,6 +81,11 @@ export class SupaplaneWebsocketServer { return this.#sessions; } + /** Install the command handler after construction (breaks the daemon/dispatcher init cycle). */ + setCommandHandler(handler: (cmd: ClientCommand, session: SessionRecord) => void): void { + this.#onCommand = handler; + } + /** Broadcast a server event to all connected sessions that have subscribed to its topic. */ broadcast(event: ServerEvent): void { const payload = JSON.stringify({ event }); @@ -106,6 +115,7 @@ export class SupaplaneWebsocketServer { ctx: { serverId: this.#identity.serverId, serverVersion: this.#serverVersion, + providers: this.#providers, ...(this.#daemonLabel ? { daemonLabel: this.#daemonLabel } : {}), logger: log, },