diff --git a/electron-builder.mcp-resources.yml b/electron-builder.mcp-resources.yml index d6144202..5cb52e3c 100644 --- a/electron-builder.mcp-resources.yml +++ b/electron-builder.mcp-resources.yml @@ -22,9 +22,16 @@ extraResources: - '@agent-relay/cloud/**' - '@agent-relay/config/**' - '@agent-relay/fleet/**' + - '@agent-relay/fleet/node_modules/@relaycast/sdk/**' + - '@agent-relay/fleet/node_modules/@relaycast/sdk/node_modules/zod/**' + - '@agent-relay/fleet/node_modules/@relaycast/types/**' + - '@agent-relay/fleet/node_modules/@relaycast/types/node_modules/zod/**' - '@agent-relay/harness-driver/**' - '@agent-relay/harnesses/**' - '@agent-relay/sdk/**' + - '@agent-relay/sdk/node_modules/@relaycast/sdk/**' + - '@agent-relay/sdk/node_modules/@relaycast/types/**' + - '@agent-relay/sdk/node_modules/zod/**' - '@agent-relay/utils/**' - '@aws-crypto/sha1-browser/**' - '@aws-crypto/sha1-browser/node_modules/@smithy/util-utf8/**' @@ -74,10 +81,6 @@ extraResources: - '@modelcontextprotocol/sdk/**' - '@posthog/core/**' - '@posthog/types/**' - - '@relaycast/sdk/**' - - '@relaycast/sdk/node_modules/zod/**' - - '@relaycast/types/**' - - '@relaycast/types/node_modules/zod/**' - '@relayfile/client/**' - '@relayflows/browser-primitive/**' - '@relayflows/browser-primitive/node_modules/@agent-relay/sdk/**' @@ -158,6 +161,13 @@ extraResources: - '@xterm/headless/**' - 'accepts/**' - 'agent-relay/**' + - 'ai-hist-native/**' + - 'ai-hist-native-darwin-arm64/**' + - 'ai-hist-native-darwin-x64/**' + - 'ai-hist-native-linux-arm64-gnu/**' + - 'ai-hist-native-linux-arm64-musl/**' + - 'ai-hist-native-linux-x64-gnu/**' + - 'ai-hist-native-linux-x64-musl/**' - 'ajv/**' - 'ajv-formats/**' - 'ansi-escapes/**' diff --git a/package-lock.json b/package-lock.json index c672cb1f..b21e03a4 100644 --- a/package-lock.json +++ b/package-lock.json @@ -8,13 +8,13 @@ "name": "pear-by-agent-relay", "version": "1.0.0", "dependencies": { - "@agent-relay/cloud": "^9.2.1", - "@agent-relay/factory": "^0.1.13", - "@agent-relay/fleet": "^9.2.1", - "@agent-relay/harness-driver": "^9.2.1", - "@agent-relay/harnesses": "^9.2.1", - "@agent-relay/integration-prompts": "^9.2.1", - "@agent-relay/sdk": "^9.2.1", + "@agent-relay/cloud": "^10.6.3", + "@agent-relay/factory": "^0.1.19", + "@agent-relay/fleet": "^10.6.3", + "@agent-relay/harness-driver": "^10.6.3", + "@agent-relay/harnesses": "^10.6.3", + "@agent-relay/integration-prompts": "^10.6.3", + "@agent-relay/sdk": "^10.6.3", "@agentworkforce/deploy": "^4.1.16", "@relayburn/sdk": "^4.0.0", "@relaycast/sdk": "^5.0.7", @@ -24,7 +24,7 @@ "@xterm/addon-web-links": "^0.11.0", "@xterm/addon-webgl": "^0.18.0", "@xterm/xterm": "^5.5.0", - "agent-relay": "^9.2.1", + "agent-relay": "^10.6.3", "agentworkforce": "^4.1.16", "ai-hist": "^0.2.3", "allotment": "^1.0.9", @@ -44,7 +44,7 @@ "pear": "bin/pear.mjs" }, "devDependencies": { - "@agent-relay/evals": "^9.2.1", + "@agent-relay/evals": "^10.6.3", "@playwright/test": "^1.57.0", "@tailwindcss/vite": "^4.0.0", "@testing-library/react": "^16.3.2", @@ -104,9 +104,9 @@ "integrity": "sha512-h6YkyIl0DZSIeVTqgPZGgN+E2I0COcF7XjljVKGzZGLUKzmBKhLROeyiy9efz/YIpApat3xjLZ4Ac20+EsydTA==" }, "node_modules/@agent-relay/broker-darwin-arm64": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/broker-darwin-arm64/-/broker-darwin-arm64-9.2.1.tgz", - "integrity": "sha512-aPQHV+Pjj57ei6rF583BltYxvmQDEEWP0DZHeZGPPRzc5x/w1jHyehXPEUGIKyb4vDWdEVbqkj27tvbuaRB0fA==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/broker-darwin-arm64/-/broker-darwin-arm64-10.6.3.tgz", + "integrity": "sha512-b6e+Yl/72QVuQVq12Xuv2SPGQ4rfMrHwpLFOsNceOpopeKJjsIW1DHMGje+F5wpYU0D+N6IBA82HQLUGfGGgKQ==", "cpu": [ "arm64" ], @@ -117,9 +117,9 @@ ] }, "node_modules/@agent-relay/broker-darwin-x64": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/broker-darwin-x64/-/broker-darwin-x64-9.2.1.tgz", - "integrity": "sha512-IVwsuYGSV5BpfPepdfoIihPVskdy+V9QtMqvnA0ATefppA3Ti20pyHkrpXZqDpl7ZXLMI3y2XnZAhm/NXH5ySg==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/broker-darwin-x64/-/broker-darwin-x64-10.6.3.tgz", + "integrity": "sha512-tFMM1NQdhepoDhR7Wt3D7XlrNH6Nekvy5ji3x09jv3H2g8MvwwdSu8oUmlrpdAMWXW4clq6gX91/otl9o53b6A==", "cpu": [ "x64" ], @@ -130,9 +130,9 @@ ] }, "node_modules/@agent-relay/broker-linux-arm64": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/broker-linux-arm64/-/broker-linux-arm64-9.2.1.tgz", - "integrity": "sha512-ZmcgruSukc7+XTAS4se+Y5bug07ZdEpAA4rhRB7fF8MJsAkMMKjPoFST0MCtg3pc5kBs1C0czzhVPR6M4KZ7hg==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/broker-linux-arm64/-/broker-linux-arm64-10.6.3.tgz", + "integrity": "sha512-Zmj4VfX0rCNOY1ABa8JNLbtQdLD7VEWI40mWudKQAqv9Ec0emPngHtwgukmOtuggQTDWV4segEnzgx1U9BgsOQ==", "cpu": [ "arm64" ], @@ -143,9 +143,9 @@ ] }, "node_modules/@agent-relay/broker-linux-x64": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/broker-linux-x64/-/broker-linux-x64-9.2.1.tgz", - "integrity": "sha512-8w+v/xOkfhFkE54R0XurZEaSSctEJiZk3pg2uNCY8bTiKCGBTy6JKPY2bWw3AhzrGR763lXa2v8Y5lZNQYJGXQ==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/broker-linux-x64/-/broker-linux-x64-10.6.3.tgz", + "integrity": "sha512-2ZLZCpcSmf+wDPpffwlKd3qRH8WSthAXXS+bUzRHRYdLNRsHS0jioA8RyNbSlLbJeK0IUyk81VdjktYdHHvI7A==", "cpu": [ "x64" ], @@ -156,9 +156,9 @@ ] }, "node_modules/@agent-relay/broker-win32-x64": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/broker-win32-x64/-/broker-win32-x64-9.2.1.tgz", - "integrity": "sha512-OxLBRAqX/E4lBw7XfR4++EQ2kzQDFeU/wOzOYhuHLLLWkGo99V2Vs+nDavH98nAjyZzO48qXhCaXd2FU7G8SXQ==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/broker-win32-x64/-/broker-win32-x64-10.6.3.tgz", + "integrity": "sha512-YuKjsYYM8NysV2jO6aWWok+I5uO6vFuESpbwzkDRhL4qWOG+QoF44G2FIGU79GXKnnhgPFxeAHJ0VJsvq8oHlg==", "cpu": [ "x64" ], @@ -169,11 +169,11 @@ ] }, "node_modules/@agent-relay/cloud": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/cloud/-/cloud-9.2.1.tgz", - "integrity": "sha512-ZcZo6qrz1kiPoEPRRnAmOGJvJdrkYK7ocyYWS1Iwpa98YsUD0/vzB5VWWAwkQ6Ow/ohtCaZ+fMOQn9lhPpoRHA==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/cloud/-/cloud-10.6.3.tgz", + "integrity": "sha512-LS3QPqC7EdJlMQBw15fPF8CqfEauF+zdYpDYlaIg1s0erUBDTEXVaWnt0B+bL9kxxQd+eciZJj/ECZB0QuJfPA==", "dependencies": { - "@agent-relay/config": "9.2.1", + "@agent-relay/config": "10.6.3", "@aws-sdk/client-s3": "3.1020.0", "ignore": "^7.0.5", "tar": "^7.5.10" @@ -183,23 +183,23 @@ } }, "node_modules/@agent-relay/config": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/config/-/config-9.2.1.tgz", - "integrity": "sha512-HPqhkUgBAZbYBwWBppj1YscQO/EDKZf2ltb0gF9HgFE08FR66hmQAnpCawanTtuZYDgojp8qgzPCaz4XfdDGTQ==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/config/-/config-10.6.3.tgz", + "integrity": "sha512-7FavBqVIPrqkWs7CSJJX7b6Dt7iuLKgSKcHK3peo4EtIWqBMEzPnvIjrji6UJ5CNd5vPuQaG2/47xmXCSDvumw==", "dependencies": { "zod": "^3.23.8", "zod-to-json-schema": "^3.23.1" } }, "node_modules/@agent-relay/evals": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/evals/-/evals-9.2.1.tgz", - "integrity": "sha512-xSqq+KtJ0UepHvipfHDu4e4/zBu8l19V4bE9ieS4JLUwMbVj2/4JzDvdASo0vIfubZ9d57IiJDeN/yqV/PWwaQ==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/evals/-/evals-10.6.3.tgz", + "integrity": "sha512-7Hjqe9WhcQKcgHVVrbTc/3cCqht4XomkKfv/+sSOoyelSFZuPRhAkLaodt68P+nxuNoSD/TOKUM++jPLyQeUnQ==", "dev": true, "license": "Apache-2.0", "dependencies": { - "@agent-relay/harness-driver": "9.2.1", - "@agent-relay/integration-prompts": "9.2.1" + "@agent-relay/harness-driver": "10.6.3", + "@agent-relay/integration-prompts": "10.6.3" } }, "node_modules/@agent-relay/events": { @@ -242,86 +242,154 @@ } }, "node_modules/@agent-relay/factory": { - "version": "0.1.13", - "resolved": "https://registry.npmjs.org/@agent-relay/factory/-/factory-0.1.13.tgz", - "integrity": "sha512-mdoC/5mqWNsrtWUWpanp1FMo5EhShH8113bpZAl9lw0Xj2TEKHtYlI7Y6rn1I2TqWaV89vuAhHJ8S5lCe+Clig==", + "version": "0.1.19", + "resolved": "https://registry.npmjs.org/@agent-relay/factory/-/factory-0.1.19.tgz", + "integrity": "sha512-I6FQKmeiKNx/jOI1zn7YM0HWizjFsciRBfxq8fCqTKHfye9Rdt/GUdukxtE6+X5VjfURPblLm/tNEqWCLeS2rg==", "license": "Apache-2.0", "dependencies": { - "@agent-relay/cloud": "^9.0.2", - "@agent-relay/fleet": "^9.0.2", - "@agent-relay/harness-driver": "^9.0.2", - "@agent-relay/integration-prompts": "^9.0.2", + "@agent-relay/cloud": "^10.6.0", + "@agent-relay/fleet": "^10.6.0", + "@agent-relay/harness-driver": "^10.6.0", + "@agent-relay/integration-prompts": "^10.6.0", + "@agent-relay/sdk": "^10.6.0", + "@relayfile/relay-helpers": "^0.4.2", "@relayfile/sdk": "^0.10.9", - "agent-relay": "^9.0.1", + "@relayflows/core": "^1.0.3", + "agent-relay": "^10.6.0", + "proper-lockfile": "^4.1.2", "zod": "^3.25.76" }, "bin": { "factory": "bin/factory.mjs" }, "engines": { - "node": ">=20" + "node": ">=20.18.1" } }, "node_modules/@agent-relay/fleet": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/fleet/-/fleet-9.2.1.tgz", - "integrity": "sha512-r27w0ghPDVKcIIUkMAQgzjx3Ev2ubRMR9G9vQTSvdaC224t2Mo/wzoCF6HEySBeaow68KzAlmBGFzUZIjzHItw==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/fleet/-/fleet-10.6.3.tgz", + "integrity": "sha512-8dopkmu+JTR06y11DFLHus09jkxB4haa74i/53qRpWS9EBGIF7SFObkACBkzwmyfngGfj5si9W4Rup2fD9d/ag==", "license": "Apache-2.0", "dependencies": { - "@agent-relay/harness-driver": "9.2.1", - "@agent-relay/harnesses": "9.2.1", - "@agent-relay/sdk": "9.2.1", + "@agent-relay/harness-driver": "10.6.3", + "@agent-relay/harnesses": "10.6.3", + "@relaycast/sdk": "^6.0.0", + "ws": "^8.18.3", "zod": "^3.23.8" } }, + "node_modules/@agent-relay/fleet/node_modules/@relaycast/sdk": { + "version": "6.2.0", + "resolved": "https://registry.npmjs.org/@relaycast/sdk/-/sdk-6.2.0.tgz", + "integrity": "sha512-7j/QGtBf6WNRa8tnnCT1rfDmcIl5QyhSQZQ//R3PePNKrRwQum0K1i3lR/dmd8DkO37Wgr8rBAh3i1xIHiLFeg==", + "dependencies": { + "@relaycast/types": "6.2.0", + "zod": "^4.3.6" + } + }, + "node_modules/@agent-relay/fleet/node_modules/@relaycast/sdk/node_modules/zod": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/zod/-/zod-4.4.3.tgz", + "integrity": "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ==", + "license": "MIT", + "funding": { + "url": "https://github.com/sponsors/colinhacks" + } + }, + "node_modules/@agent-relay/fleet/node_modules/@relaycast/types": { + "version": "6.2.0", + "resolved": "https://registry.npmjs.org/@relaycast/types/-/types-6.2.0.tgz", + "integrity": "sha512-9e5ywyO0HM3yNfGE2JvvGWMeyv1xwshuLNkyEPGOCPyyw/UlFFQLESGShEy2GMwowZPlw8sm5mVTzmJqRZ3CLQ==", + "dependencies": { + "zod": "^4.3.6" + } + }, + "node_modules/@agent-relay/fleet/node_modules/@relaycast/types/node_modules/zod": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/zod/-/zod-4.4.3.tgz", + "integrity": "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ==", + "license": "MIT", + "funding": { + "url": "https://github.com/sponsors/colinhacks" + } + }, "node_modules/@agent-relay/harness-driver": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/harness-driver/-/harness-driver-9.2.1.tgz", - "integrity": "sha512-I1wL3FIF5r5QSGyfJsWVXi5l5FbuVI4fsle6KBXke52jdhl/FOFLarhefd4LR1uuQF47+xaxYKzK4hbclkeG4A==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/harness-driver/-/harness-driver-10.6.3.tgz", + "integrity": "sha512-oviHDovfDPktg3gPOWXRluZO+nIOQBe5SSOSsQkoBTQvKUB6X7rcCEfBrqwI3rapU1ptUv5hYkFnYE8ei1k07A==", "license": "Apache-2.0", "dependencies": { - "@agent-relay/sdk": "9.2.1", + "@agent-relay/sdk": "10.6.3", "ws": "^8.18.3", "zod": "^3.23.8" }, "optionalDependencies": { - "@agent-relay/broker-darwin-arm64": "9.2.1", - "@agent-relay/broker-darwin-x64": "9.2.1", - "@agent-relay/broker-linux-arm64": "9.2.1", - "@agent-relay/broker-linux-x64": "9.2.1", - "@agent-relay/broker-win32-x64": "9.2.1" + "@agent-relay/broker-darwin-arm64": "10.6.3", + "@agent-relay/broker-darwin-x64": "10.6.3", + "@agent-relay/broker-linux-arm64": "10.6.3", + "@agent-relay/broker-linux-x64": "10.6.3", + "@agent-relay/broker-win32-x64": "10.6.3" } }, "node_modules/@agent-relay/harnesses": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/harnesses/-/harnesses-9.2.1.tgz", - "integrity": "sha512-DhJLASqGxz0jikgCt2px2LezNE4IIqSqG4BOVDjKiYuKWEWpxeXEBOupcihnGrjgBRcAfHCOwyUf/jawo/B4jQ==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/harnesses/-/harnesses-10.6.3.tgz", + "integrity": "sha512-EX3mJhVjd3RPd/ESX8Hku5DCSmnJ2mdy9cSrC5+TEq7rv43F+OgnQAt8YKc6dQhtSbdNTn0V0VxXv61g4oc47A==", "license": "Apache-2.0", "dependencies": { - "@agent-relay/harness-driver": "9.2.1", - "@agent-relay/sdk": "9.2.1" + "@agent-relay/harness-driver": "10.6.3", + "@agent-relay/sdk": "10.6.3" } }, "node_modules/@agent-relay/integration-prompts": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/integration-prompts/-/integration-prompts-9.2.1.tgz", - "integrity": "sha512-oqBoCiXnfA+Z30gFxLPqPVhhhmWR3rbglevVRsp9Cje9RymUAFr6I9DtHenYH98Rs8GUbEL7EW6erQZhlFmwsw==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/integration-prompts/-/integration-prompts-10.6.3.tgz", + "integrity": "sha512-SXieePr9JeF5Fw8B3tTBqCBi4tkaTKh+Xa0xJ/z6+qGFNpguCGz1Y8VTY8S213wSt0X1HunSkUjZV6kv2+yCwg==", "license": "Apache-2.0" }, "node_modules/@agent-relay/sdk": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/sdk/-/sdk-9.2.1.tgz", - "integrity": "sha512-hLIMSFxkxSRj5/RIXJmwm8pKh6zbwbCbpNuscgg+RpU66ZFYZ37wB2kWUt8RPmgFMahnttuZHe/K8CF05+nsjw==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/sdk/-/sdk-10.6.3.tgz", + "integrity": "sha512-5wnFbZ8tTwAuERkvd/dEJ3dWU0chxfYn3BwqLhUxMDj9PhtyT/KFlT/WA9DBJDhjxqfv1d250jvI4DuVQ1PIrw==", + "dependencies": { + "@relaycast/sdk": "^6.0.0", + "@relaycast/types": "^6.0.0", + "zod": "^4.3.6" + } + }, + "node_modules/@agent-relay/sdk/node_modules/@relaycast/sdk": { + "version": "6.2.0", + "resolved": "https://registry.npmjs.org/@relaycast/sdk/-/sdk-6.2.0.tgz", + "integrity": "sha512-7j/QGtBf6WNRa8tnnCT1rfDmcIl5QyhSQZQ//R3PePNKrRwQum0K1i3lR/dmd8DkO37Wgr8rBAh3i1xIHiLFeg==", "dependencies": { - "@relaycast/sdk": "^5.0.5" + "@relaycast/types": "6.2.0", + "zod": "^4.3.6" + } + }, + "node_modules/@agent-relay/sdk/node_modules/@relaycast/types": { + "version": "6.2.0", + "resolved": "https://registry.npmjs.org/@relaycast/types/-/types-6.2.0.tgz", + "integrity": "sha512-9e5ywyO0HM3yNfGE2JvvGWMeyv1xwshuLNkyEPGOCPyyw/UlFFQLESGShEy2GMwowZPlw8sm5mVTzmJqRZ3CLQ==", + "dependencies": { + "zod": "^4.3.6" + } + }, + "node_modules/@agent-relay/sdk/node_modules/zod": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/zod/-/zod-4.4.3.tgz", + "integrity": "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ==", + "license": "MIT", + "funding": { + "url": "https://github.com/sponsors/colinhacks" } }, "node_modules/@agent-relay/utils": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/@agent-relay/utils/-/utils-9.2.1.tgz", - "integrity": "sha512-zBrQ5rbj/L7Xd5Kfavddh+993VwLVGVTVw7hbOIxxtrVEOnP3YGd4Hnmisyr0x8whbZ1vyFZTnohfCjkj6AlTg==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/@agent-relay/utils/-/utils-10.6.3.tgz", + "integrity": "sha512-oGXmmjjAY+YfnMT6o/+mk4GVizUzRwSMaWY+n30HL1E4xXquFnkxUCQTNTWKotmSKVrIrkxAnzTQiFRdv/KJWw==", "dependencies": { - "@agent-relay/config": "9.2.1", + "@agent-relay/config": "10.6.3", "compare-versions": "^6.1.1" } }, @@ -414,9 +482,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "optional": true, "os": [ "linux" @@ -432,9 +497,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "optional": true, "os": [ "linux" @@ -5210,9 +5272,6 @@ "cpu": [ "arm" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -5233,9 +5292,6 @@ "cpu": [ "arm" ], - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -5256,9 +5312,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -5279,9 +5332,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -5302,9 +5352,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -5325,9 +5372,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -5535,9 +5579,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "optional": true, "os": [ "linux" @@ -5553,9 +5594,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "optional": true, "os": [ "linux" @@ -5620,10 +5658,86 @@ "@relayfile/sdk": ">=0.6.0 <1" } }, + "node_modules/@relayfile/adapter-linear": { + "version": "0.2.11", + "resolved": "https://registry.npmjs.org/@relayfile/adapter-linear/-/adapter-linear-0.2.11.tgz", + "integrity": "sha512-P+J/m/V2n+X8TOQqNI3jpZwuZA286L8AbMzoFOikCA7koX/8OUz7NYaU/xujAe+6XdWizLdLQQ2P8C0Bzl2gHg==", + "license": "MIT", + "optional": true, + "peer": true, + "dependencies": { + "@relayfile/adapter-core": "^0.2.26" + }, + "engines": { + "node": ">=18" + }, + "peerDependencies": { + "@relayfile/sdk": ">=0.6.0 <1" + } + }, + "node_modules/@relayfile/adapter-linear/node_modules/@relayfile/adapter-core": { + "version": "0.2.26", + "resolved": "https://registry.npmjs.org/@relayfile/adapter-core/-/adapter-core-0.2.26.tgz", + "integrity": "sha512-+FIosWgo2O0bfBAt0oV6tgYIzU1TTSBoGrjw4IwedQFuDbxgTODnn3taGiCpdAxIDj7/xq5vZRkjnl9wEPMD9w==", + "license": "MIT", + "optional": true, + "peer": true, + "dependencies": { + "@scalar/postman-to-openapi": "^0.6.0", + "cheerio": "^1.2.0", + "minimatch": "^10.0.3", + "yaml": "^2.8.1" + }, + "bin": { + "adapter-core": "dist/cli.js" + }, + "engines": { + "node": ">=18" + }, + "peerDependencies": { + "@relayfile/sdk": ">=0.6.0 <1" + } + }, + "node_modules/@relayfile/adapter-reddit": { + "version": "0.2.6", + "resolved": "https://registry.npmjs.org/@relayfile/adapter-reddit/-/adapter-reddit-0.2.6.tgz", + "integrity": "sha512-LYDCs8WpbyRZIAUmSsbi2GgOXdOx0RDQzof2XF7/22tuGtz2T3+7GVc+Rk1/qusV73hQfSAFz+nhcM7Hwv2L2A==", + "license": "Apache-2.0", + "dependencies": { + "@relayfile/adapter-core": "^0.5.8" + }, + "engines": { + "node": ">=18" + }, + "peerDependencies": { + "@relayfile/sdk": ">=0.6.0 <1" + } + }, + "node_modules/@relayfile/adapter-reddit/node_modules/@relayfile/adapter-core": { + "version": "0.5.8", + "resolved": "https://registry.npmjs.org/@relayfile/adapter-core/-/adapter-core-0.5.8.tgz", + "integrity": "sha512-R/Zwxb63uqSs4BRAWbrNo/ZTKvTWfrA5kva0H7zRWWG0SzO+OGzym+DuDlJSRpksonU4Lp1piofYZezYzmfV9g==", + "license": "Apache-2.0", + "dependencies": { + "@scalar/postman-to-openapi": "^0.6.0", + "cheerio": "^1.2.0", + "minimatch": "^10.0.3", + "yaml": "^2.8.1" + }, + "bin": { + "adapter-core": "dist/src/cli.js" + }, + "engines": { + "node": ">=18" + }, + "peerDependencies": { + "@relayfile/sdk": ">=0.6.0 <1" + } + }, "node_modules/@relayfile/client": { - "version": "0.10.19", - "resolved": "https://registry.npmjs.org/@relayfile/client/-/client-0.10.19.tgz", - "integrity": "sha512-1hoaPaMc/jB2sMq/a5gbEZ7vGHPAx8WaN23Jr03h8DTowmCeoqwu+cliMWRASo90VmdatJyfjTzQc8csHIXL2g==", + "version": "0.10.27", + "resolved": "https://registry.npmjs.org/@relayfile/client/-/client-0.10.27.tgz", + "integrity": "sha512-1ASWmrDDIZlMhQuGodR+vkYuJy6Dmkc06DAwidYKVJzTVgvefrhQaNyP+diOd0HLvKv/VGcDe67cCsqnpOMUBA==", "license": "Apache-2.0", "engines": { "node": ">=18" @@ -5703,6 +5817,53 @@ "linux" ] }, + "node_modules/@relayfile/relay-helpers": { + "version": "0.4.8", + "resolved": "https://registry.npmjs.org/@relayfile/relay-helpers/-/relay-helpers-0.4.8.tgz", + "integrity": "sha512-KmTyz+hP+B8d+fhNzO65En5r22dhiKTOT8m5j1hSRoaF6wTgoK/MwNIWUsZ4SGoyeiJGDtbDKm9vyEWbuHQpQg==", + "license": "Apache-2.0", + "dependencies": { + "@relayfile/adapter-core": "^0.5.8", + "@relayfile/adapter-linear": "^0.4.7", + "@relayfile/adapter-reddit": "^0.2.6" + } + }, + "node_modules/@relayfile/relay-helpers/node_modules/@relayfile/adapter-core": { + "version": "0.5.8", + "resolved": "https://registry.npmjs.org/@relayfile/adapter-core/-/adapter-core-0.5.8.tgz", + "integrity": "sha512-R/Zwxb63uqSs4BRAWbrNo/ZTKvTWfrA5kva0H7zRWWG0SzO+OGzym+DuDlJSRpksonU4Lp1piofYZezYzmfV9g==", + "license": "Apache-2.0", + "dependencies": { + "@scalar/postman-to-openapi": "^0.6.0", + "cheerio": "^1.2.0", + "minimatch": "^10.0.3", + "yaml": "^2.8.1" + }, + "bin": { + "adapter-core": "dist/src/cli.js" + }, + "engines": { + "node": ">=18" + }, + "peerDependencies": { + "@relayfile/sdk": ">=0.6.0 <1" + } + }, + "node_modules/@relayfile/relay-helpers/node_modules/@relayfile/adapter-linear": { + "version": "0.4.7", + "resolved": "https://registry.npmjs.org/@relayfile/adapter-linear/-/adapter-linear-0.4.7.tgz", + "integrity": "sha512-SBmYzvk1gIiRXWqGnGowE8UqYYSKpjo15nQ0nTNC/7KT9YrOLZ3loHBS1VjWGBxtzai5PLdpjUTIfG0FMtRVwQ==", + "license": "Apache-2.0", + "dependencies": { + "@relayfile/adapter-core": "^0.5.8" + }, + "engines": { + "node": ">=18" + }, + "peerDependencies": { + "@relayfile/sdk": ">=0.6.0 <1" + } + }, "node_modules/@relayfile/sdk": { "version": "0.10.19", "resolved": "https://registry.npmjs.org/@relayfile/sdk/-/sdk-0.10.19.tgz", @@ -6201,9 +6362,6 @@ "arm" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6218,9 +6376,6 @@ "arm" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -6235,9 +6390,6 @@ "arm64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6252,9 +6404,6 @@ "arm64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -6269,9 +6418,6 @@ "loong64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6286,9 +6432,6 @@ "loong64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -6303,9 +6446,6 @@ "ppc64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6320,9 +6460,6 @@ "ppc64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -6337,9 +6474,6 @@ "riscv64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6354,9 +6488,6 @@ "riscv64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -6371,9 +6502,6 @@ "s390x" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6388,9 +6516,6 @@ "x64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6405,9 +6530,6 @@ "x64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -7332,9 +7454,6 @@ "arm64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -7352,9 +7471,6 @@ "arm64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -7372,9 +7488,6 @@ "x64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -7392,9 +7505,6 @@ "x64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -8343,20 +8453,19 @@ } }, "node_modules/agent-relay": { - "version": "9.2.1", - "resolved": "https://registry.npmjs.org/agent-relay/-/agent-relay-9.2.1.tgz", - "integrity": "sha512-Mx31Z0/lZlTzQfJaKptipIm/Cq6+xTNObaACOeSSIU6ukEcrYqtVZD1zYkS/Se06xLiiQHgTa5RIS0eQQRq2Mw==", + "version": "10.6.3", + "resolved": "https://registry.npmjs.org/agent-relay/-/agent-relay-10.6.3.tgz", + "integrity": "sha512-sAsLByzz0z6vwDrKCOkyL210lSPizDyJ7C76iVXdhxaWzkJQbPw7Ubl8pTAulWwIHYXLrnqmVC609eKfkO+JLQ==", "license": "Apache-2.0", "dependencies": { - "@agent-relay/cloud": "9.2.1", - "@agent-relay/config": "9.2.1", - "@agent-relay/fleet": "9.2.1", - "@agent-relay/harness-driver": "9.2.1", - "@agent-relay/sdk": "9.2.1", - "@agent-relay/utils": "9.2.1", + "@agent-relay/cloud": "10.6.3", + "@agent-relay/config": "10.6.3", + "@agent-relay/fleet": "10.6.3", + "@agent-relay/harness-driver": "10.6.3", + "@agent-relay/sdk": "10.6.3", + "@agent-relay/utils": "10.6.3", "@modelcontextprotocol/sdk": "^1.0.0", - "@relaycast/sdk": "^5.0.5", - "@relayfile/client": "^0.10.19", + "@relayfile/client": "^0.10.21", "@relayflows/cli": "^1.0.1", "@xterm/headless": "^6.0.0", "commander": "^12.1.0", @@ -8372,6 +8481,9 @@ }, "engines": { "node": ">=20.9.0" + }, + "optionalDependencies": { + "ai-hist-native": "^0.4.1" } }, "node_modules/agent-trajectories": { @@ -8417,6 +8529,120 @@ "node": ">=18" } }, + "node_modules/ai-hist-native": { + "version": "0.4.3", + "resolved": "https://registry.npmjs.org/ai-hist-native/-/ai-hist-native-0.4.3.tgz", + "integrity": "sha512-9qjvXlh7YijscpGRlRQJBxnmIAhizm08jlfw6j+t95aDTWfwIXxuE14bcowLuCknU6OEFmyXsk7C9IoGa4PAlg==", + "license": "MIT", + "optional": true, + "engines": { + "node": ">= 18" + }, + "optionalDependencies": { + "ai-hist-native-darwin-arm64": "0.4.3", + "ai-hist-native-darwin-x64": "0.4.3", + "ai-hist-native-linux-arm64-gnu": "0.4.3", + "ai-hist-native-linux-arm64-musl": "0.4.3", + "ai-hist-native-linux-x64-gnu": "0.4.3", + "ai-hist-native-linux-x64-musl": "0.4.3" + } + }, + "node_modules/ai-hist-native-darwin-arm64": { + "version": "0.4.3", + "resolved": "https://registry.npmjs.org/ai-hist-native-darwin-arm64/-/ai-hist-native-darwin-arm64-0.4.3.tgz", + "integrity": "sha512-06zNVOSh4SptrzOsq5GafhF04yWkepnO5L0MFgBgYLAoXgky02nxtBVDBFeo9dVlFOeNCj8cMxGkEn6Bc+dY9w==", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "darwin" + ], + "engines": { + "node": ">= 18" + } + }, + "node_modules/ai-hist-native-darwin-x64": { + "version": "0.4.3", + "resolved": "https://registry.npmjs.org/ai-hist-native-darwin-x64/-/ai-hist-native-darwin-x64-0.4.3.tgz", + "integrity": "sha512-NxLYJSdoNn18mMqnPjZuOmijjAKz5Q8towr5cVcpe6plZu/bjFQlsaU1Zeb+5EiTBUZoVdIexkKIqVy6XGEw+g==", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "darwin" + ], + "engines": { + "node": ">= 18" + } + }, + "node_modules/ai-hist-native-linux-arm64-gnu": { + "version": "0.4.3", + "resolved": "https://registry.npmjs.org/ai-hist-native-linux-arm64-gnu/-/ai-hist-native-linux-arm64-gnu-0.4.3.tgz", + "integrity": "sha512-bFsu2Y23VXaQrxYHl3R7e9WGzIDQEyvrkzlPqNHswyGpYUV1gtt2mgfSDjg8pw1PAOaRiF4cbht7n5nqeZ0hug==", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "engines": { + "node": ">= 18" + } + }, + "node_modules/ai-hist-native-linux-arm64-musl": { + "version": "0.4.3", + "resolved": "https://registry.npmjs.org/ai-hist-native-linux-arm64-musl/-/ai-hist-native-linux-arm64-musl-0.4.3.tgz", + "integrity": "sha512-8xTE9/CcmL2i7fM6NxE+ytbUSsmse4+Y5wTDlAAvlrmWQd0vukDM3ENfm0SSawYsShCSoLJqvRDktGNzI6Pfbg==", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "engines": { + "node": ">= 18" + } + }, + "node_modules/ai-hist-native-linux-x64-gnu": { + "version": "0.4.3", + "resolved": "https://registry.npmjs.org/ai-hist-native-linux-x64-gnu/-/ai-hist-native-linux-x64-gnu-0.4.3.tgz", + "integrity": "sha512-yQfQdgpHt6haZV37/TsUUh8H7EWGGvvx6bn7h7xgaABp4egOFbEC0JXtTZEgMbVMZXn+QGg17sF+5Aq/H9BLuw==", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "engines": { + "node": ">= 18" + } + }, + "node_modules/ai-hist-native-linux-x64-musl": { + "version": "0.4.3", + "resolved": "https://registry.npmjs.org/ai-hist-native-linux-x64-musl/-/ai-hist-native-linux-x64-musl-0.4.3.tgz", + "integrity": "sha512-GUOrR8tgXBACvKZSdzQQXDp2JdLM0T69ASqQbwHgMCKRGxZeXOub0eA4ep5beQsCHs58tYvWTYHY6bxiu/8TWA==", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "engines": { + "node": ">= 18" + } + }, "node_modules/ajv": { "version": "8.20.0", "resolved": "https://registry.npmjs.org/ajv/-/ajv-8.20.0.tgz", @@ -12376,9 +12602,6 @@ "arm64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MPL-2.0", "optional": true, "os": [ @@ -12400,9 +12623,6 @@ "arm64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MPL-2.0", "optional": true, "os": [ @@ -12424,9 +12644,6 @@ "x64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MPL-2.0", "optional": true, "os": [ @@ -12448,9 +12665,6 @@ "x64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MPL-2.0", "optional": true, "os": [ @@ -13957,7 +14171,6 @@ "version": "4.1.2", "resolved": "https://registry.npmjs.org/proper-lockfile/-/proper-lockfile-4.1.2.tgz", "integrity": "sha512-TjNPblN4BwAWMXU8s9AEz4JmQxnD1NNL7bNOY/AKUzyamc379FWASUhc/K1pL2noVb+XmZKLL68cjzLsiOAMaA==", - "dev": true, "license": "MIT", "dependencies": { "graceful-fs": "^4.2.4", @@ -13969,7 +14182,6 @@ "version": "0.12.0", "resolved": "https://registry.npmjs.org/retry/-/retry-0.12.0.tgz", "integrity": "sha512-9LkiTwjUh6rT555DtE9rTX+BKByPfrMzEAtnlEtdEwr3Nkffwiihqe2bWADg+OQRjt9gl6ICdmB/ZFDCGAtSow==", - "dev": true, "license": "MIT", "engines": { "node": ">= 4" @@ -14780,7 +14992,6 @@ "version": "3.0.7", "resolved": "https://registry.npmjs.org/signal-exit/-/signal-exit-3.0.7.tgz", "integrity": "sha512-wnD2ZE+l+SPC/uoS0vXeE9L1+0wuaMqKlfz9AMUo38JsyLSBWSFcHR1Rri62LZc12vLr1gb3jl7iwQhgwpAbGQ==", - "dev": true, "license": "ISC" }, "node_modules/simple-update-notifier": { diff --git a/package.json b/package.json index 53c7acb2..d64e7bab 100644 --- a/package.json +++ b/package.json @@ -35,13 +35,13 @@ "lint": "eslint ." }, "dependencies": { - "@agent-relay/cloud": "^9.2.1", - "@agent-relay/factory": "^0.1.13", - "@agent-relay/fleet": "^9.2.1", - "@agent-relay/harness-driver": "^9.2.1", - "@agent-relay/harnesses": "^9.2.1", - "@agent-relay/integration-prompts": "^9.2.1", - "@agent-relay/sdk": "^9.2.1", + "@agent-relay/cloud": "^10.6.3", + "@agent-relay/factory": "^0.1.19", + "@agent-relay/fleet": "^10.6.3", + "@agent-relay/harness-driver": "^10.6.3", + "@agent-relay/harnesses": "^10.6.3", + "@agent-relay/integration-prompts": "^10.6.3", + "@agent-relay/sdk": "^10.6.3", "@agentworkforce/deploy": "^4.1.16", "@relayburn/sdk": "^4.0.0", "@relaycast/sdk": "^5.0.7", @@ -51,7 +51,7 @@ "@xterm/addon-web-links": "^0.11.0", "@xterm/addon-webgl": "^0.18.0", "@xterm/xterm": "^5.5.0", - "agent-relay": "^9.2.1", + "agent-relay": "^10.6.3", "agentworkforce": "^4.1.16", "ai-hist": "^0.2.3", "allotment": "^1.0.9", @@ -71,7 +71,7 @@ "protobufjs": "8.5.0" }, "devDependencies": { - "@agent-relay/evals": "^9.2.1", + "@agent-relay/evals": "^10.6.3", "@playwright/test": "^1.57.0", "@tailwindcss/vite": "^4.0.0", "@testing-library/react": "^16.3.2", diff --git a/src/main/broker.test.ts b/src/main/broker.test.ts index 0eb8a7a5..43ccbaee 100644 --- a/src/main/broker.test.ts +++ b/src/main/broker.test.ts @@ -739,9 +739,7 @@ describe('BrokerManager local + cloud coexistence', () => { projectId: PROJECT_ID, cwd: '/tmp/project-1', brokerName: 'pear-project-1', - connection: { - url: 'http://127.0.0.1:4242' - } + readBrokerSession: expect.any(Function) })) await manager.shutdown() @@ -2021,14 +2019,22 @@ exit 2 kind: 'worker_stream', name: 'claude-1', chunk: 'pong\n', - seq: 22 + seq: 22, + offset: 101 } listener?.(chunkEvent) listener?.(chunkEvent) const ptyCalls = (win.webContents.send as ReturnType).mock.calls .filter(([channel]) => channel === 'broker:pty-chunk') - expect(ptyCalls).toEqual([['broker:pty-chunk', PROJECT_ID, 'claude-1', 'pong\n']]) + expect(ptyCalls).toEqual([[ + 'broker:pty-chunk', + PROJECT_ID, + 'claude-1', + 'pong\n', + 101, + expect.any(Number) + ]]) await manager.shutdown() }) diff --git a/src/main/broker.ts b/src/main/broker.ts index a85348b5..ea6eeaab 100644 --- a/src/main/broker.ts +++ b/src/main/broker.ts @@ -363,6 +363,8 @@ export interface AttachTerminalResult { cols: number cursor: [number, number] screen: string + offset?: number + generation?: number } } @@ -1840,20 +1842,11 @@ export class BrokerManager { await this.stopSessionFleetSidecar(session) } - const url = getClientBaseUrl(session.client) - if (!url) { - console.warn(`[broker] Local fleet node skipped for project ${session.projectId}: broker URL unavailable`) - return - } - const sidecar = startPearFleetSidecar({ projectId: session.projectId, cwd: session.cwd, brokerName: session.name, - connection: { - url, - ...(getClientApiKey(session.client) ? { apiKey: getClientApiKey(session.client) } : {}) - }, + readBrokerSession: () => session.client.getSession(), log: (message) => console.log(`[broker] ${message}`), warn: (message) => console.warn(`[broker] ${message}`) }) @@ -2522,7 +2515,14 @@ export class BrokerManager { } const targetWindow = this.windowForSession(sessionKey, win) if (targetWindow && !targetWindow.isDestroyed()) { - targetWindow.webContents.send('broker:pty-chunk', projectId, event.name, event.chunk) + const offset = brokerEventNumber(event, 'offset') + targetWindow.webContents.send( + 'broker:pty-chunk', + projectId, + event.name, + event.chunk, + ...(offset !== undefined ? [offset, eventStreamGeneration] : []) + ) } this.rememberAgentSession(event.name, sessionKey) if (this.sessions.get(sessionKey)?.cloudSandboxId) { @@ -3445,7 +3445,11 @@ export class BrokerManager { rows: snapshot.rows, cols: snapshot.cols, cursor: snapshot.cursor, - screen: Buffer.from(snapshot.screen, 'base64').toString('utf-8') + screen: Buffer.from(snapshot.screen, 'base64').toString('utf-8'), + ...(typeof snapshot.offset === 'number' && Number.isFinite(snapshot.offset) + ? { offset: snapshot.offset } + : {}), + generation: session.eventStreamGeneration } } } catch (err) { @@ -3682,7 +3686,8 @@ export class BrokerManager { screen: format === 'ansi' ? Buffer.from(snapshot.screen, 'base64').toString('utf-8') - : snapshot.screen + : snapshot.screen, + ...(snapshot.offset !== undefined ? { offset: snapshot.offset } : {}) } } catch (err) { if (isMissingAgentError(err)) return null diff --git a/src/main/pear-fleet-node.test.ts b/src/main/pear-fleet-node.test.ts index 20d823b8..d0a04ecd 100644 --- a/src/main/pear-fleet-node.test.ts +++ b/src/main/pear-fleet-node.test.ts @@ -1,6 +1,11 @@ import { describe, expect, it, vi } from 'vitest' -import { invokeNodeHandler, nodeInfo, nodeManifest } from '@agent-relay/fleet' -import { createPearFleetNodeDefinition, PEAR_LOCAL_SPAWN_HARNESSES, reconnectDelayMs } from './pear-fleet-node' +import { invokeNodeHandler, nodeInfo } from '@agent-relay/fleet' +import { + createPearFleetNodeDefinition, + PEAR_LOCAL_SPAWN_HARNESSES, + resolvePearFleetConnection, + startPearFleetSidecar +} from './pear-fleet-node' const expectedCapabilities = Object.keys(PEAR_LOCAL_SPAWN_HARNESSES).map((cli) => `spawn:${cli}`) @@ -12,16 +17,19 @@ describe('Pear local fleet node', () => { brokerName: 'pear-project-1' }) - const manifest = nodeManifest(definition) - expect(manifest.name).toBe('pear-project-1-local-fleet') - expect(manifest.capabilities.map((capability) => capability.name)).toEqual(expectedCapabilities) - expect(manifest.capabilities).toEqual(expectedCapabilities.map((name) => expect.objectContaining({ - name, - metadata: expect.objectContaining({ - pearLocalNode: true, - clonePaths: { 'project-1': '/tmp/project-1' } - }) - }))) + const info = nodeInfo(definition) + expect(info.name).toBe('pear-project-1-local-fleet') + expect(info.capabilities).toEqual(expectedCapabilities) + expect(expectedCapabilities.map((name) => definition.capabilities[name])).toEqual( + expectedCapabilities.map(() => + expect.objectContaining({ + metadata: expect.objectContaining({ + pearLocalNode: true, + clonePaths: { 'project-1': '/tmp/project-1' } + }) + }) + ) + ) }) it('spawns non-Claude/Codex harnesses through the broker', async () => { @@ -110,26 +118,84 @@ describe('Pear local fleet node', () => { }) }) -describe('reconnectDelayMs', () => { - it('caps transient reconnects (already registered) at the fast ceiling', () => { - expect(reconnectDelayMs(1, true)).toBe(500) - expect(reconnectDelayMs(2, true)).toBe(1_000) - expect(reconnectDelayMs(4, true)).toBe(4_000) - // 500 * 2**4 = 8000 -> clamped to the 5s ceiling - expect(reconnectDelayMs(5, true)).toBe(5_000) - expect(reconnectDelayMs(50, true)).toBe(5_000) +describe('resolvePearFleetConnection', () => { + it('waits for the broker-minted node token and attaches to the same v10 node', async () => { + let now = 0 + let reads = 0 + const connection = await resolvePearFleetConnection(async () => { + reads += 1 + if (reads === 1) { + return { node_id: 'node-1', node_name: 'pear-project-1' } + } + return { + node_id: 'node-1', + node_name: 'pear-project-1', + node_token: 'nt_live_test', + relay_base_url: 'https://cast.example' + } + }, new AbortController().signal, { + timeoutMs: 1_000, + pollIntervalMs: 250, + now: () => now, + sleep: async (ms) => { + now += ms + } + }) + + expect(reads).toBe(2) + expect(connection).toEqual({ + connection: { + nodeId: 'node-1', + nodeToken: 'nt_live_test', + baseUrl: 'https://cast.example' + }, + nodeName: 'pear-project-1' + }) + }) + + it('fails clearly instead of silently disabling the provider when the token never arrives', async () => { + let now = 0 + await expect(resolvePearFleetConnection( + async () => ({ node_id: 'node-1', node_name: 'pear-project-1' }), + new AbortController().signal, + { + timeoutMs: 500, + pollIntervalMs: 250, + now: () => now, + sleep: async (ms) => { + now += ms + } + } + )).rejects.toThrow('timed out waiting for a node token for node-1') }) - it('backs off much harder while the node has never registered', () => { - // Grows past the registered ceiling so a wedged broker is not stormed. - expect(reconnectDelayMs(5, false)).toBe(8_000) - expect(reconnectDelayMs(6, false)).toBe(16_000) - expect(reconnectDelayMs(8, false)).toBe(60_000) - expect(reconnectDelayMs(50, false)).toBe(60_000) + it('bounds a broker session read that never settles', async () => { + let now = 0 + await expect(resolvePearFleetConnection( + () => new Promise(() => {}), + new AbortController().signal, + { + timeoutMs: 500, + now: () => now, + sleep: async (ms) => { + now += ms + } + } + )).rejects.toThrow('timed out waiting for the broker node id') + expect(now).toBe(500) }) - it('never returns a negative or sub-base delay for the first attempt', () => { - expect(reconnectDelayMs(0, false)).toBe(500) - expect(reconnectDelayMs(1, false)).toBe(500) + it('stops promptly while a broker session read is hung', async () => { + const sidecar = startPearFleetSidecar({ + projectId: 'project-1', + cwd: '/tmp/project-1', + brokerName: 'pear-project-1', + readBrokerSession: () => new Promise(() => {}) + }) + const registered = sidecar.registered.catch((error: unknown) => error) + + await sidecar.stop() + + await expect(registered).resolves.toMatchObject({ name: 'AbortError' }) }) }) diff --git a/src/main/pear-fleet-node.ts b/src/main/pear-fleet-node.ts index 288f25f8..0f343b6e 100644 --- a/src/main/pear-fleet-node.ts +++ b/src/main/pear-fleet-node.ts @@ -2,11 +2,12 @@ import { basename, resolve } from 'node:path' import { action, defineNode, - invokeNodeHandler, - nodeInfo, - nodeManifest, + startServeNode, type FleetCapabilityValue, - type FleetNodeDefinition + type FleetNodeDefinition, + type FleetNodeInfo, + type NodeEngineConnection, + type RunningNode } from '@agent-relay/fleet' import { aider, @@ -21,8 +22,6 @@ import { type PtyHarness } from '@agent-relay/harnesses' import { resolveStaticHarnessConfig, type RestartPolicy } from '@agent-relay/harness-driver' -import { PROTOCOL_VERSION, type NodeManifest } from '@agent-relay/harness-driver/protocol' -import WebSocket from 'ws' import { z } from 'zod' export const PEAR_LOCAL_SPAWN_HARNESSES = { @@ -37,19 +36,8 @@ export const PEAR_LOCAL_SPAWN_HARNESSES = { droid } satisfies Record -const RECONNECT_BASE_DELAY_MS = 500 -const RECONNECT_MAX_DELAY_MS = 5_000 -// When the sidecar has never completed a registration handshake, the broker is -// either wedged (accepts the WS upgrade but never answers `hello`) or does not -// accept this fleet node at all. Reconnecting on the normal 5s ceiling in that -// state produces a tight storm — leaked sockets plus a flood of "hello request -// timed out; reconnecting" warnings — that never self-resolves. Back off much -// harder while unregistered so we keep probing for recovery without hammering -// the broker. A transient drop after a healthy registration keeps the fast -// ceiling so we recover quickly from ordinary reconnects. -const RECONNECT_UNREGISTERED_MAX_DELAY_MS = 60_000 -const REQUEST_TIMEOUT_MS = 5_000 -const DEREGISTER_TIMEOUT_MS = 1_000 +const NODE_TOKEN_WAIT_MS = 15_000 +const NODE_TOKEN_POLL_MS = 250 const restartPolicySchema = z.object({ enabled: z.boolean().optional(), @@ -92,34 +80,44 @@ export interface PearFleetNodeOptions { } export interface PearFleetSidecarOptions extends PearFleetNodeOptions { - connection: { - url: string - apiKey?: string - } + readBrokerSession: () => Promise log?: (message: string) => void warn?: (message: string) => void } -export interface RunningPearFleetSidecar { - readonly registered: Promise - readonly done: Promise - stop(): Promise +export interface PearFleetBrokerSession { + relay_base_url?: string + node_id?: string + node_name?: string + node_token?: string } -type BrokerFrame = { - type?: string - request_id?: string - payload?: unknown +export interface ResolvedPearFleetConnection { + connection: NodeEngineConnection + nodeName: string } -type PendingRequest = { - resolve: (value: unknown) => void - reject: (error: Error) => void - timer: ReturnType +export interface PearFleetConnectionWaitOptions { + timeoutMs?: number + pollIntervalMs?: number + now?: () => number + sleep?: (ms: number, signal: AbortSignal) => Promise +} + +type BrokerSessionReadResult = + | { status: 'value'; session: PearFleetBrokerSession | null | undefined } + | { status: 'error'; error: unknown } + | { status: 'deadline' } + | { status: 'aborted' } + +export interface RunningPearFleetSidecar { + readonly registered: Promise + readonly done: Promise + stop(): Promise } export function createPearFleetNodeDefinition(options: PearFleetNodeOptions): FleetNodeDefinition { - const nodeName = pearFleetNodeName(options) + const nodeName = pearFleetProviderName(options) const clonePathKey = basename(options.cwd) || options.projectId || 'project' const clonePaths = { [clonePathKey]: options.cwd } const capabilities: Record = {} @@ -207,15 +205,132 @@ function resolvePearSpawnCwd(projectCwd: string, input: SpawnCapabilityInput): s return requested } +/** + * Resolve the broker's v10 fleet identity. The broker publishes `node_id` + * before its background token mint can complete, so a single session read can + * strand Pear's capability provider for the lifetime of the broker. Poll for a + * bounded window and attach to the same engine node once the token appears. + */ +export async function resolvePearFleetConnection( + readBrokerSession: () => Promise, + signal: AbortSignal, + options: PearFleetConnectionWaitOptions = {} +): Promise { + const timeoutMs = options.timeoutMs ?? NODE_TOKEN_WAIT_MS + const pollIntervalMs = options.pollIntervalMs ?? NODE_TOKEN_POLL_MS + const now = options.now ?? Date.now + const sleep = options.sleep ?? delay + const deadline = now() + timeoutMs + let lastSession: PearFleetBrokerSession | undefined + let lastError: unknown + + for (;;) { + if (signal.aborted) throw abortError() + const remainingMs = Math.max(0, deadline - now()) + const read = await readBrokerSessionBeforeDeadline( + readBrokerSession, + signal, + remainingMs, + sleep + ) + if (read.status === 'aborted') throw abortError() + if (read.status === 'deadline') break + if (read.status === 'value') { + lastSession = read.session ?? undefined + lastError = undefined + if (lastSession?.node_id && lastSession.node_token) { + return { + connection: { + nodeId: lastSession.node_id, + nodeToken: lastSession.node_token, + ...(lastSession.relay_base_url ? { baseUrl: lastSession.relay_base_url } : {}) + }, + nodeName: lastSession.node_name ?? lastSession.node_id + } + } + } else { + lastError = read.error + } + + if (now() >= deadline) break + await sleep(Math.min(pollIntervalMs, Math.max(0, deadline - now())), signal) + } + + if (signal.aborted) throw abortError() + if (!lastSession?.node_id) { + const suffix = lastError ? `: ${toError(lastError).message}` : '' + throw new Error(`Pear fleet provider timed out waiting for the broker node id${suffix}`) + } + throw new Error(`Pear fleet provider timed out waiting for a node token for ${lastSession.node_id}`) +} + +// A broker request can wedge independently of the polling clock. Race every +// session read against both the remaining deadline and stop signal so a hung +// getSession() cannot strand sidecar shutdown. The read promise itself may not +// be cancellable, but both of its settlement paths remain observed after the +// race, avoiding a late unhandled rejection. +async function readBrokerSessionBeforeDeadline( + readBrokerSession: () => Promise, + signal: AbortSignal, + remainingMs: number, + sleep: (ms: number, signal: AbortSignal) => Promise +): Promise { + if (signal.aborted) return { status: 'aborted' } + if (remainingMs <= 0) return { status: 'deadline' } + + const deadlineController = new AbortController() + const onAbort = (): void => deadlineController.abort() + signal.addEventListener('abort', onAbort, { once: true }) + let sessionPromise: Promise + try { + sessionPromise = readBrokerSession() + } catch (error) { + sessionPromise = Promise.reject(error) + } + // Give already-settled client mocks/caches one microtask to win before + // starting the deadline timer. Real in-flight I/O falls through and is + // bounded below. + const pendingRead = Symbol('pending broker session read') + try { + const immediate = await Promise.race([sessionPromise, Promise.resolve(pendingRead)]) + if (immediate !== pendingRead) { + signal.removeEventListener('abort', onAbort) + deadlineController.abort() + return signal.aborted + ? { status: 'aborted' } + : { status: 'value', session: immediate } + } + } catch (error) { + signal.removeEventListener('abort', onAbort) + deadlineController.abort() + return signal.aborted ? { status: 'aborted' } : { status: 'error', error } + } + const read = sessionPromise.then( + (session) => ({ status: 'value', session }), + (error) => ({ status: 'error', error }) + ) + const timeout = sleep(remainingMs, deadlineController.signal).then(() => + signal.aborted ? { status: 'aborted' } : { status: 'deadline' } + ) + + try { + return await Promise.race([read, timeout]) + } finally { + signal.removeEventListener('abort', onAbort) + deadlineController.abort() + } +} + export function startPearFleetSidecar(options: PearFleetSidecarOptions): RunningPearFleetSidecar { const controller = new AbortController() - let resolveRegistered: (manifest: NodeManifest) => void - let rejectRegistered: (error: Error) => void + let running: RunningNode | undefined let registeredSettled = false - const registered = new Promise((resolve, reject) => { - resolveRegistered = (manifest) => { + let resolveRegistered!: (info: FleetNodeInfo) => void + let rejectRegistered!: (error: Error) => void + const registered = new Promise((resolve, reject) => { + resolveRegistered = (info) => { registeredSettled = true - resolve(manifest) + resolve(info) } rejectRegistered = (error) => { registeredSettled = true @@ -224,276 +339,53 @@ export function startPearFleetSidecar(options: PearFleetSidecarOptions): Running }) const definition = createPearFleetNodeDefinition(options) - const done = runPearFleetSidecarLoop({ - ...options, - definition, - signal: controller.signal, - onRegistered: (manifest) => { - if (!registeredSettled) resolveRegistered(manifest) - }, - onInitialFailure: (error) => { - if (!registeredSettled) rejectRegistered(error) + const done = (async () => { + try { + const target = await resolvePearFleetConnection(options.readBrokerSession, controller.signal) + if (controller.signal.aborted) throw abortError() + running = startServeNode({ + definition, + connection: target.connection, + nameOverride: target.nodeName, + providerName: definition.name, + reconnect: true, + signal: controller.signal, + log: options.log, + warn: options.warn, + onRegistered: (info) => { + if (!registeredSettled) resolveRegistered(info) + } + }) + await running.done + if (!registeredSettled) { + rejectRegistered(new Error('Pear fleet provider exited before registering')) + } + } catch (error) { + if (!registeredSettled) rejectRegistered(toError(error)) + if (!controller.signal.aborted) throw error } - }) + })() return { registered, done, stop: async () => { controller.abort() + await running?.stop().catch(() => undefined) await done } } } -function pearFleetNodeName(options: PearFleetNodeOptions): string { +function pearFleetProviderName(options: PearFleetNodeOptions): string { const rawName = `${options.brokerName || options.projectId || 'pear'}-local-fleet` return rawName.replace(/[^\w.-]+/gu, '-').replace(/^-+|-+$/gu, '') || 'pear-local-fleet' } -async function runPearFleetSidecarLoop(options: PearFleetSidecarOptions & { - definition: FleetNodeDefinition - signal: AbortSignal - onRegistered: (manifest: NodeManifest) => void - onInitialFailure: (error: Error) => void -}): Promise { - let attempt = 0 - let registeredOnce = false - - while (!options.signal.aborted) { - try { - await runPearFleetSidecarConnection({ - ...options, - onRegistered: (manifest) => { - registeredOnce = true - options.onRegistered(manifest) - } - }) - attempt = 0 - } catch (error) { - const err = toError(error) - if (!registeredOnce) { - options.onInitialFailure(err) - } - if (options.signal.aborted) return - options.warn?.(`Pear fleet sidecar disconnected: ${err.message}; reconnecting`) - } - - if (options.signal.aborted) return - attempt += 1 - await delay(reconnectDelayMs(attempt, registeredOnce), options.signal) - } -} - -/** - * Backoff between reconnect attempts. Once the node has registered at least - * once, transient drops recover quickly on the normal ceiling; while it has - * never registered (wedged/unsupported broker) the ceiling widens so we probe - * for recovery instead of storming the broker every few seconds. - */ -export function reconnectDelayMs(attempt: number, registeredOnce: boolean): number { - const ceiling = registeredOnce ? RECONNECT_MAX_DELAY_MS : RECONNECT_UNREGISTERED_MAX_DELAY_MS - const exponent = Math.max(0, attempt - 1) - return Math.min(ceiling, RECONNECT_BASE_DELAY_MS * 2 ** exponent) -} - -function runPearFleetSidecarConnection(options: PearFleetSidecarOptions & { - definition: FleetNodeDefinition - signal: AbortSignal - onRegistered: (manifest: NodeManifest) => void -}): Promise { - return new Promise((resolve, reject) => { - const ws = new WebSocket(fleetWsUrl(options.connection.url), { - headers: options.connection.apiKey ? { 'X-API-Key': options.connection.apiKey } : undefined - }) - const pending = new Map() - let requestSeq = 0 - let settled = false - let nodeRegistered = false - - const settle = (fn: () => void): void => { - if (settled) return - settled = true - for (const pendingRequest of pending.values()) { - pendingRequest.reject(new Error('Pear fleet sidecar connection closed')) - } - pending.clear() - options.signal.removeEventListener('abort', abort) - fn() - } - - const sendRequest = (type: string, payload: unknown): Promise => { - if (ws.readyState !== WebSocket.OPEN) { - return Promise.reject(new Error('Pear fleet sidecar websocket is not open')) - } - const requestId = `pear_fleet_${Date.now()}_${++requestSeq}` - const frame = { - v: PROTOCOL_VERSION, - type, - request_id: requestId, - payload - } - - return new Promise((requestResolve, requestReject) => { - const timer = setTimeout(() => { - pending.delete(requestId) - requestReject(new Error(`Pear fleet sidecar ${type} request timed out after ${REQUEST_TIMEOUT_MS}ms`)) - }, REQUEST_TIMEOUT_MS) - const pendingRequest: PendingRequest = { - timer, - resolve: (value) => { - clearTimeout(timer) - requestResolve(value) - }, - reject: (error) => { - clearTimeout(timer) - requestReject(error) - } - } - pending.set(requestId, pendingRequest) - ws.send(JSON.stringify(frame), (error) => { - if (!error) return - if (!pending.has(requestId)) return - pending.delete(requestId) - pendingRequest.reject(error) - }) - }) - } - - const close = async (): Promise => { - if (ws.readyState === WebSocket.OPEN && nodeRegistered) { - await withTimeout(sendRequest('deregister_node', {}), DEREGISTER_TIMEOUT_MS).catch(() => undefined) - } - if (ws.readyState === WebSocket.OPEN || ws.readyState === WebSocket.CONNECTING) { - ws.close() - } - } - - const ignoreCloseError = (error: unknown): void => { - options.warn?.(`Pear fleet sidecar close skipped: ${toError(error).message}`) - } - - const abort = (): void => { - void close().catch(ignoreCloseError).finally(() => settle(resolve)) - } - - const sendHandlerResult = async (invocationId: string, output: unknown, error?: unknown): Promise => { - const payload = error - ? { invocation_id: invocationId, error: toError(error).message } - : { invocation_id: invocationId, output: output ?? null } - await sendRequest('handler_result', payload) - } - - const handleInvoke = async (payload: unknown): Promise => { - const record = asRecord(payload) - const invocationId = typeof record.invocation_id === 'string' ? record.invocation_id : undefined - const name = typeof record.name === 'string' ? record.name : undefined - if (!invocationId || !name) return - - try { - const output = await invokeNodeHandler(options.definition, name, record.input, { - node: nodeInfo(options.definition), - invocationId, - relay: { - sendMessage: (input) => sendRequest('send_message', { - to: input.to, - text: input.text, - from: input.from ?? options.definition.name, - ...(input.threadId ? { thread_id: input.threadId } : {}), - ...(input.workspaceId ? { workspace_id: input.workspaceId } : {}), - ...(input.workspaceAlias ? { workspace_alias: input.workspaceAlias } : {}), - ...(input.mode ? { mode: input.mode } : {}), - ...(input.data ? { data: input.data } : {}) - }) - }, - spawnAgent: (input) => sendRequest('spawn_agent', { - agent: input.agent, - ...(input.initialTask !== undefined ? { initial_task: input.initialTask } : {}), - skip_relay_prompt: input.skipRelayPrompt ?? false, - ...((input.invocationId ?? invocationId) ? { invocation_id: input.invocationId ?? invocationId } : {}) - }) - }) - await sendHandlerResult(invocationId, output) - } catch (error) { - await sendHandlerResult(invocationId, undefined, error) - } - } - - options.signal.addEventListener('abort', abort, { once: true }) - - ws.on('open', () => { - void (async () => { - const manifest = nodeManifest(options.definition) - await sendRequest('hello', { - client_name: 'pear-local-fleet', - client_version: '1.0.0' - }) - await sendRequest('register_node', { manifest }) - nodeRegistered = true - await sendRequest('register_handlers', { names: Object.keys(options.definition.capabilities) }) - options.log?.(`Pear fleet node "${manifest.name}" registered with ${manifest.capabilities.length} capabilities.`) - options.onRegistered(manifest) - })().catch((error) => { - settle(() => reject(toError(error))) - void close().catch(ignoreCloseError) - }) - }) - - ws.on('message', (data) => { - const frame = parseBrokerFrame(data) - if (!frame) return - if (frame.request_id && pending.has(frame.request_id)) { - const pendingRequest = pending.get(frame.request_id) - pending.delete(frame.request_id) - if (!pendingRequest) return - if (frame.type === 'error') { - pendingRequest.reject(frameError(frame.payload)) - } else { - pendingRequest.resolve(readOkResult(frame.payload)) - } - return - } - if (frame.type === 'invoke_handler') { - void handleInvoke(frame.payload).catch((error) => options.warn?.(toError(error).message)) - } - }) - - ws.on('close', () => settle(resolve)) - ws.on('error', (error) => settle(() => reject(toError(error)))) - }) -} - -function parseBrokerFrame(data: WebSocket.RawData): BrokerFrame | null { - try { - const text = Array.isArray(data) ? Buffer.concat(data).toString('utf8') : data.toString() - return JSON.parse(text) as BrokerFrame - } catch { - return null - } -} - -function readOkResult(payload: unknown): unknown { - return payload && typeof payload === 'object' && 'result' in payload - ? (payload as { result: unknown }).result - : payload -} - -function frameError(payload: unknown): Error { - if (payload && typeof payload === 'object') { - const record = payload as Record - const error = new Error(typeof record.message === 'string' ? record.message : 'Pear fleet sidecar request failed') - error.name = typeof record.code === 'string' ? record.code : 'PearFleetSidecarError' - return error - } - return new Error('Pear fleet sidecar request failed') -} - -function fleetWsUrl(baseUrl: string): string { - return `${baseUrl.replace(/\/+$/u, '').replace(/^http/u, 'ws')}/api/fleet/ws` -} - -function asRecord(value: unknown): Record { - return value && typeof value === 'object' ? value as Record : {} +function abortError(): Error { + const error = new Error('Pear fleet provider stopped') + error.name = 'AbortError' + return error } function toError(error: unknown): Error { @@ -515,13 +407,3 @@ function delay(ms: number, signal: AbortSignal): Promise { signal.addEventListener('abort', onAbort, { once: true }) }) } - -function withTimeout(promise: Promise, timeoutMs: number): Promise { - let timer: ReturnType | undefined - const timeout = new Promise((_, reject) => { - timer = setTimeout(() => reject(new Error(`Timed out after ${timeoutMs}ms`)), timeoutMs) - }) - return Promise.race([promise, timeout]).finally(() => { - if (timer) clearTimeout(timer) - }) -} diff --git a/src/preload/index.ts b/src/preload/index.ts index 6a7f9475..20874d94 100644 --- a/src/preload/index.ts +++ b/src/preload/index.ts @@ -293,9 +293,23 @@ const api = { onEvent: (callback: (event: unknown) => void) => subscribe('broker:event', callback), onEventStreamDiagnostic: (callback: (event: BrokerEventStreamDiagnostic) => void) => subscribe('broker:event-stream-diagnostic', callback), - onPtyChunk: (callback: (projectId: string, name: string, chunk: string) => void) => { - const handler = (_: unknown, projectId: string, name: string, chunk: string): void => - callback(projectId, name, chunk) + onPtyChunk: ( + callback: ( + projectId: string, + name: string, + chunk: string, + offset?: number, + generation?: number + ) => void + ) => { + const handler = ( + _: unknown, + projectId: string, + name: string, + chunk: string, + offset?: number, + generation?: number + ): void => callback(projectId, name, chunk, offset, generation) ipcRenderer.on('broker:pty-chunk', handler) return () => ipcRenderer.removeListener('broker:pty-chunk', handler) }, diff --git a/src/renderer/src/hooks/use-broker-events.ts b/src/renderer/src/hooks/use-broker-events.ts index 9c5a9204..534fd368 100644 --- a/src/renderer/src/hooks/use-broker-events.ts +++ b/src/renderer/src/hooks/use-broker-events.ts @@ -77,9 +77,9 @@ export function useBrokerEvents(): void { // PTY chunks ride a dedicated lightweight channel so per-character typing // doesn't pay for the broker:event metadata spread / structured clone. - const unsubPtyChunk = pear.broker.onPtyChunk((projectId, name, chunk) => { + const unsubPtyChunk = pear.broker.onPtyChunk((projectId, name, chunk, offset, generation) => { const key = getAgentKey(projectId, name) - appendPtyChunk(key, chunk) + appendPtyChunk(key, chunk, offset, generation) useTypingStore.getState().noteActivity(key) useAgentStore.getState().markAgentActive(projectId, name) }) diff --git a/src/renderer/src/lib/echo-router.test.ts b/src/renderer/src/lib/echo-router.test.ts index cf458979..2f77a69f 100644 --- a/src/renderer/src/lib/echo-router.test.ts +++ b/src/renderer/src/lib/echo-router.test.ts @@ -110,7 +110,8 @@ interface Harness { engine: PredictiveEchoEngine model: HeadlessModel setSrtt(value: number | null): void - dispose(): void + disposeEngine(): void + dispose(): Promise } function makeHarness(): Harness { @@ -133,6 +134,13 @@ function makeHarness(): Harness { isViewportPinned: () => false, scrollToBottom: () => {} }) + let engineDisposed = false + let disposePromise: Promise | null = null + const disposeEngine = (): void => { + if (engineDisposed) return + engineDisposed = true + engine.reset() + } return { router, live, @@ -141,10 +149,15 @@ function makeHarness(): Harness { setSrtt: (value) => { srtt = value }, + disposeEngine, dispose: () => { - router.dispose() - engine.reset() - live.dispose() + if (!disposePromise) { + disposePromise = router.dispose().then(() => { + disposeEngine() + live.dispose() + }) + } + return disposePromise } } } @@ -179,8 +192,8 @@ describe('echo-router — real engine + real parser equivalence', () => { harness = makeHarness() }) - afterEach(() => { - harness.dispose() + afterEach(async () => { + await harness.dispose() }) it('starts on the direct route and passes bytes through byte-for-byte', async () => { @@ -342,6 +355,27 @@ describe('echo-router — real engine + real parser equivalence', () => { expect(readScreen(harness.live)).toBe(await reference(a + b + flood)) }) + it('drains accepted engine bytes before v10 reset can discard them', async () => { + harness.setSrtt(80) + harness.router.onUserInput('x') + await drain(harness.live) + await settle(harness) + expect(harness.router.route()).toBe('engine') + + // Back up the real v10 engine tail, enqueue authoritative bytes, then + // begin teardown immediately. v10 reset drops a chunk whose processing + // has not started, so this assertion fails unless router.dispose waits for + // its accepted ordering chain before the owner resets the engine. + harness.model.writeDelayMs = 30 + const chunk = '\x1b[2J\x1b[Hqueued-before-dispose' + harness.router.onServerOutput(chunk) + await harness.router.dispose() + harness.disposeEngine() + await drain(harness.live) + + expect(readScreen(harness.live)).toBe(await reference(chunk)) + }) + it('skips optimistic echo while a chunk is mid-flight in the engine tail (stale-cursor strand)', async () => { // Get on the engine route with a prompt at row 0. harness.router.onServerOutput('\x1b[2J\x1b[H$ ') @@ -417,7 +451,7 @@ describe('echo-router — reseed capture timeout', () => { vi.useRealTimers() }) - it('aborts a stalled capture, releases held chunks directly, and stays direct', () => { + it('aborts a stalled capture, releases held chunks directly, and stays direct', async () => { const written: string[] = [] const engine = { seed: vi.fn(() => Promise.resolve()), @@ -454,10 +488,10 @@ describe('echo-router — reseed capture timeout', () => { // Later output flows normally on the direct route. router.onServerOutput('after') expect(written).toEqual(['held-1held-2', 'after']) - router.dispose() + await router.dispose() }) - it('releases held chunks directly when the engine changes during the capture', () => { + it('releases held chunks directly when the engine changes during the capture', async () => { const written: string[] = [] let captureCallback: (() => void) | null = null const makeEngine = (): PredictiveEchoWithStatus => ({ @@ -497,10 +531,10 @@ describe('echo-router — reseed capture timeout', () => { expect(router.route()).toBe('direct') expect(firstEngine.seed).not.toHaveBeenCalled() expect(engine.seed).not.toHaveBeenCalled() - router.dispose() + await router.dispose() }) - it('dispose drops held chunks and cancels the capture timeout', () => { + it('dispose flushes held chunks and cancels the capture timeout', async () => { const written: string[] = [] const router = createEchoRouter({ write: (data) => { @@ -523,8 +557,8 @@ describe('echo-router — reseed capture timeout', () => { router.onUserInput('a') router.onServerOutput('held') - router.dispose() + await router.dispose() vi.advanceTimersByTime(RESEED_CAPTURE_TIMEOUT_MS * 2) - expect(written).toEqual([]) + expect(written).toEqual(['held']) }) }) diff --git a/src/renderer/src/lib/echo-router.ts b/src/renderer/src/lib/echo-router.ts index ac6a4ac4..0d92e1f9 100644 --- a/src/renderer/src/lib/echo-router.ts +++ b/src/renderer/src/lib/echo-router.ts @@ -97,7 +97,10 @@ export interface EchoRouter { onUserInput(data: string): void // Which sink server bytes currently take. Exposed for tests/diagnostics. route(): 'direct' | 'engine' - dispose(): void + // Stop accepting new work and resolve only after every already-accepted + // server byte has reached the live terminal. Callers must await this before + // resetting the v10 engine, whose reset discards queued server output. + dispose(): Promise } export function createEchoRouter(deps: EchoRouterDeps): EchoRouter { @@ -105,7 +108,9 @@ export function createEchoRouter(deps: EchoRouterDeps): EchoRouter { let reseedPending = false let heldChunks: string[] = [] let reseedTimeout: ReturnType | null = null + let closing = false let disposed = false + let disposePromise: Promise | null = null // Single ordering domain for everything that reaches the live terminal. // Engine writes happen inside the engine's ASYNC tail (its model parse @@ -222,7 +227,7 @@ export function createEchoRouter(deps: EchoRouterDeps): EchoRouter { return { onServerOutput(combined: string): void { - if (disposed || combined.length === 0) return + if (closing || disposed || combined.length === 0) return if (reseedPending) { heldChunks.push(combined) return @@ -257,7 +262,7 @@ export function createEchoRouter(deps: EchoRouterDeps): EchoRouter { } }, onUserInput(data: string): void { - if (disposed) return + if (closing || disposed) return // While a reseed capture is in flight predictions stay off — the // engine's model isn't authoritative yet. Only the optimistic echo is // skipped; the keystroke still reaches the PTY via the send path. @@ -291,14 +296,21 @@ export function createEchoRouter(deps: EchoRouterDeps): EchoRouter { route(): 'direct' | 'engine' { return route }, - dispose(): void { - disposed = true + dispose(): Promise { + if (disposePromise) return disposePromise + closing = true reseedPending = false - heldChunks = [] if (reseedTimeout) { clearTimeout(reseedTimeout) reseedTimeout = null } + // A direct→engine capture may be holding accepted server chunks outside + // opChain. Queue them on the direct sink before sealing the chain. + releaseHeldDirect() + disposePromise = opChain.then(() => { + disposed = true + }) + return disposePromise } } } diff --git a/src/renderer/src/lib/ipc-mock.ts b/src/renderer/src/lib/ipc-mock.ts index dc301db3..825d4b03 100644 --- a/src/renderer/src/lib/ipc-mock.ts +++ b/src/renderer/src/lib/ipc-mock.ts @@ -115,7 +115,13 @@ interface MockState { brokerEventListeners: Set> brokerStatusListeners: Set> brokerDiagnosticListeners: Set> - ptyChunkListeners: Set<(projectId: string, name: string, chunk: string) => void> + ptyChunkListeners: Set<( + projectId: string, + name: string, + chunk: string, + offset?: number, + generation?: number + ) => void> menuListeners: Map void>> cloudAgentListeners: Set> proactiveAgentListeners: Set> @@ -130,7 +136,13 @@ export interface PearMockHarness { reset: () => void injectBrokerEvent: (event: BrokerEventLike) => void injectBrokerEvents: (events: BrokerEventLike[]) => void - injectPtyChunk: (projectId: string, name: string, chunk: string) => void + injectPtyChunk: ( + projectId: string, + name: string, + chunk: string, + offset?: number, + generation?: number + ) => void // Rendering-harness knobs: force the input SRTT the renderer polls (so the // predictive-echo engine route engages) and echo typed bytes back through // the PTY stream after a delay (so predictions reconcile end-to-end). @@ -983,7 +995,13 @@ export const pearMock: PearAPI = { onEvent: (callback: (event: unknown) => void) => noopUnsubscribe(state.brokerEventListeners, callback), onEventStreamDiagnostic: (callback: (event: BrokerEventStreamDiagnostic) => void) => noopUnsubscribe(state.brokerDiagnosticListeners, callback), - onPtyChunk: (callback: (projectId: string, name: string, chunk: string) => void) => + onPtyChunk: (callback: ( + projectId: string, + name: string, + chunk: string, + offset?: number, + generation?: number + ) => void) => noopUnsubscribe(state.ptyChunkListeners, callback), onStatus: (callback: (status: BrokerStatusEvent) => void) => noopUnsubscribe(state.brokerStatusListeners, callback), checkCliAvailable: async (_cli: string) => true @@ -1307,10 +1325,12 @@ export const pearMockHarness: PearMockHarness = { injectBrokerEvents: (events: BrokerEventLike[]) => { for (const event of events) handleInjectedBrokerEvent(event) }, - injectPtyChunk: (projectId: string, name: string, chunk: string) => { + injectPtyChunk: (projectId, name, chunk, offset, generation) => { const ptyKey = key(projectId, name) state.ptyChunks[ptyKey] = [...(state.ptyChunks[ptyKey] || []), chunk] - for (const listener of [...state.ptyChunkListeners]) listener(projectId, name, chunk) + for (const listener of [...state.ptyChunkListeners]) { + listener(projectId, name, chunk, offset, generation) + } }, setInputSrtt: (ms: number | null) => { mockInputSrttMs = ms diff --git a/src/renderer/src/lib/terminal-runtime-registry.dom.test.ts b/src/renderer/src/lib/terminal-runtime-registry.dom.test.ts index 10cb5c6a..9b1596a0 100644 --- a/src/renderer/src/lib/terminal-runtime-registry.dom.test.ts +++ b/src/renderer/src/lib/terminal-runtime-registry.dom.test.ts @@ -215,6 +215,50 @@ describe('terminal-runtime-registry — pty drain coalescing', () => { }) }) +describe('terminal-runtime-registry — v10 snapshot offset seeding', () => { + it('writes post-snapshot chunks without replaying chunks covered by the snapshot', async () => { + const ipc = await import('@/lib/ipc') + vi.mocked(ipc.pear.broker.attachTerminal).mockImplementationOnce(async () => { + // Both chunks arrive while attach IPC is in flight. The snapshot covers + // only the first offset; the second must survive the seed baseline. + ptyBuffer.appendPtyChunk('p:a', 'covered', 10, 7) + ptyBuffer.appendPtyChunk('p:a', 'fresh', 20, 7) + return { + name: 'a', + mode: 'auto_inject' as const, + pending: 0, + snapshot: { + rows: 24, + cols: 80, + cursor: [0, 0] as [number, number], + screen: 'snapshot', + offset: 10, + generation: 7 + } + } + }) + + const runtime = registry.acquireTerminalRuntime({ + projectId: 'p', + agentName: 'a', + terminalMode: 'drive', + theme: 'dark', + getInputSrtt: () => null + }) + const term = createdTerminals[0] + + runtime.mount(makeLayoutContainer()) + await flushAsync() + await flushAsync() + + expect(term.__writes).toContain('snapshot') + expect(term.__writes).toContain('fresh') + expect(term.__writes).not.toContain('covered') + + registry.disposeTerminalRuntime(runtime.key) + }) +}) + describe('terminal-runtime-registry — clearOnDataIf identity check', () => { it('only clears the on-data handler when the caller still owns the slot', () => { const runtime = registry.acquireTerminalRuntime({ diff --git a/src/renderer/src/lib/terminal-runtime-registry.ts b/src/renderer/src/lib/terminal-runtime-registry.ts index c2d6176d..11adfdbb 100644 --- a/src/renderer/src/lib/terminal-runtime-registry.ts +++ b/src/renderer/src/lib/terminal-runtime-registry.ts @@ -27,6 +27,7 @@ import { diagPtyEnabled, flushPtyChunksNow, getPtyChunkTotal, + getPtyChunksAfterOffset, getPtyChunksSinceTotal, subscribePtyBuffer } from '@/stores/pty-buffer-store' @@ -340,10 +341,13 @@ function createRuntime( // real engine + parser). const echoRouter = createEchoRouter({ write: (data, callback) => { - if (disposed || !term) return + // Teardown marks the public runtime disposed before awaiting the router + // drain. Keep this private sink alive until that drain completes so + // Relay v10 reset cannot discard already-accepted server output. + if (!term) return term.write(data, callback) }, - getEngine: () => (disposed ? null : predictiveEcho), + getEngine: () => predictiveEcho, buildModelSeed: () => (term ? buildModelSeedFromTerminal(term) : '\x1bc'), getInputSrtt: () => currentSrttGetter(), isViewportPinned: () => (term ? isViewportPinnedToBottom(term) : false), @@ -472,35 +476,35 @@ function createRuntime( }) } + const writeChunks = (newChunks: string[]): void => { + if (disposed || !term) return + if (newChunks.length === 0) return + activitySerial += 1 + lastOutputAt = Date.now() + // Optional diagnostic, gated on localStorage.PEAR_DIAG_PTY === '1'. + // See pty-buffer-store.ts for the enable instructions. Flag is + // cached to avoid a per-batch localStorage read. + if (diagPtyEnabled()) { + console.log(`[diag:runtime:writeChunks] key=${key} count=${newChunks.length} firstPreview="${newChunks[0]?.slice(0, 80).replace(/\n/g, '\\n').replace(/\r/g, '\\r').replace(/\x1b/g, '\\e')}"`) + } + // Typing-trace accounting stays per-chunk: recordChunkEchoed consumes + // one pending keystroke per call, so coalescing it would under-count. + for (const chunk of newChunks) recordChunkEchoed(chunk) + // Coalesce the frame's chunks into a single write. xterm's VT parser is + // a streaming state machine, so write(a)+write(b) ≡ write(a+b) — for the + // live terminal AND the predictive-echo headless model (which parses the + // bytes a second time). One write collapses N parser passes + N model + // writes + N promise ticks into one. Under heavy TUI redraw streaming + // that per-chunk fan-out was the drain hot path: the renderer couldn't + // keep up and input lagged. Byte content and order are unchanged, so the + // one-write-per-byte invariant holds. + const combined = newChunks.length === 1 ? newChunks[0] : newChunks.join('') + echoRouter.onServerOutput(combined) + } + const seedBufferSubscription = (): void => { if (unsubBuffer || !term || disposed) return - const writeChunks = (newChunks: string[]): void => { - if (disposed || !term) return - if (newChunks.length === 0) return - activitySerial += 1 - lastOutputAt = Date.now() - // Optional diagnostic, gated on localStorage.PEAR_DIAG_PTY === '1'. - // See pty-buffer-store.ts for the enable instructions. Flag is - // cached to avoid a per-batch localStorage read. - if (diagPtyEnabled()) { - console.log(`[diag:runtime:writeChunks] key=${key} count=${newChunks.length} firstPreview="${newChunks[0]?.slice(0, 80).replace(/\n/g, '\\n').replace(/\r/g, '\\r').replace(/\x1b/g, '\\e')}"`) - } - // Typing-trace accounting stays per-chunk: recordChunkEchoed consumes - // one pending keystroke per call, so coalescing it would under-count. - for (const chunk of newChunks) recordChunkEchoed(chunk) - // Coalesce the frame's chunks into a single write. xterm's VT parser is - // a streaming state machine, so write(a)+write(b) ≡ write(a+b) — for the - // live terminal AND the predictive-echo headless model (which parses the - // bytes a second time). One write collapses N parser passes + N model - // writes + N promise ticks into one. Under heavy TUI redraw streaming - // that per-chunk fan-out was the drain hot path: the renderer couldn't - // keep up and input lagged. Byte content and order are unchanged, so the - // one-write-per-byte invariant holds. - const combined = newChunks.length === 1 ? newChunks[0] : newChunks.join('') - echoRouter.onServerOutput(combined) - } - unsubBuffer = subscribePtyBuffer(key, writeChunks) // Initial replay: pull whatever is already in the buffer past the // snapshot baseline (writtenTotal). The listener only receives tails @@ -517,6 +521,7 @@ function createRuntime( attachInFlight = true let shouldReplay = true + let postSnapshotChunks: string[] = [] try { const result = await pear.broker.attachTerminal({ projectId: opts.projectId, @@ -535,11 +540,19 @@ function createRuntime( ) { term.write(result.snapshot.screen) await predictiveEcho?.seed(result.snapshot.screen) - // Drain any chunks that arrived during the IPC roundtrip but are - // still staged in pending. Without this, the next rAF would push - // them into the buffer AFTER we capture writtenChunks, and the - // subsequent subscribe would replay them on top of the snapshot. + // Drain chunks staged during the IPC roundtrip before capturing the + // buffer total. Relay v10 snapshots and worker_stream events share a + // cumulative raw-PTY offset, so preserve chunks strictly after the + // snapshot instead of treating the entire roundtrip window as covered + // (which could drop output emitted after snapshot capture). flushPtyChunksNow(key) + if (typeof result.snapshot.offset === 'number') { + postSnapshotChunks = getPtyChunksAfterOffset( + key, + result.snapshot.offset, + result.snapshot.generation + ) + } writtenTotal = getPtyChunkTotal(key) shouldReplay = false } @@ -564,6 +577,10 @@ function createRuntime( await predictiveEcho?.seed('') } + // These chunks were buffered before `writtenTotal` but landed after the + // authoritative v10 snapshot. Write them in byte order before subscribing; + // the subscription then catches anything that arrived after the baseline. + writeChunks(postSnapshotChunks) seedBufferSubscription() // SIGWINCH bounce: 200ms after attach completes, send a one-pixel @@ -718,32 +735,37 @@ function createRuntime( cancelPendingInit() reconciler.dispose() sizeSync.dispose() - echoRouter.dispose() - clearPtyBuffer(key) disposed = true currentToken = null unsubBuffer?.() unsubBuffer = null - disposePredictiveEcho?.() - disposePredictiveEcho = null - predictiveEcho = null - try { - webglAddon?.dispose() - } catch { - // Teardown: an already-disposed or context-lost WebGL addon can throw - // on dispose; we null it out next regardless, so silence is correct. - } - webglAddon = null - try { - term?.dispose() - } catch { - // Teardown: disposing an xterm instance twice (or after its host was - // detached) can throw; we null it out next regardless, so silence is correct. - } - term = null - if (host.parentElement) { - host.parentElement.removeChild(host) - } + clearPtyBuffer(key) + // Relay v10 intentionally drops queued engine output after reset. Drain + // the router ordering chain first, then reset/dispose the engine, model, + // and live terminal. The registry record is already gone, so this old + // runtime cannot be reacquired while its accepted bytes finish parsing. + void echoRouter.dispose().finally(() => { + disposePredictiveEcho?.() + disposePredictiveEcho = null + predictiveEcho = null + try { + webglAddon?.dispose() + } catch { + // Teardown: an already-disposed or context-lost WebGL addon can throw + // on dispose; we null it out next regardless, so silence is correct. + } + webglAddon = null + try { + term?.dispose() + } catch { + // Teardown: disposing an xterm instance twice (or after its host was + // detached) can throw; we null it out next regardless, so silence is correct. + } + term = null + if (host.parentElement) { + host.parentElement.removeChild(host) + } + }) }, isMounted(): boolean { return currentToken !== null @@ -787,7 +809,7 @@ function createRuntime( currentSrttGetter = getter }, getPredictiveEcho(): PredictiveEcho | null { - return predictiveEcho + return disposed ? null : predictiveEcho }, noteUserInput(data: string): void { if (disposed || !term) return diff --git a/src/renderer/src/stores/pty-buffer-store.test.ts b/src/renderer/src/stores/pty-buffer-store.test.ts index fa397f9d..1a3686ec 100644 --- a/src/renderer/src/stores/pty-buffer-store.test.ts +++ b/src/renderer/src/stores/pty-buffer-store.test.ts @@ -5,6 +5,7 @@ import { flushPtyChunksNow, getPtyChunks, getPtyChunkTotal, + getPtyChunksAfterOffset, getPtyChunksSinceTotal, subscribePtyBuffer } from './pty-buffer-store' @@ -147,12 +148,7 @@ describe('pty-buffer-store', () => { expect(getPtyChunks('k1')).toEqual(['first', 'second']) }) - it('duplicate appendPtyChunk replay arrives once at the listener per chunk', () => { - // AGENTS.md: "Add regression tests when touching PTY buffering. Include - // duplicate/replay cases." Renderer-side guarantee is: every appendPtyChunk - // call adds exactly one chunk to the buffer; if the broker sends the same - // chunk twice the listener sees two distinct entries — dedup is the - // broker/main's responsibility, NOT the renderer buffer's. + it('delivers repeated identity-less chunks because identical terminal traffic is valid', () => { const listener = vi.fn() subscribePtyBuffer('k1', listener) @@ -164,6 +160,24 @@ describe('pty-buffer-store', () => { expect(listener).toHaveBeenCalledWith(['dup', 'dup']) expect(getPtyChunks('k1')).toEqual(['dup', 'dup']) }) + + it('suppresses a replay only when generation, offset, and content all match', () => { + const listener = vi.fn() + subscribePtyBuffer('k1', listener) + + appendPtyChunk('k1', 'same', 10, 1) + appendPtyChunk('k1', 'same', 10, 1) + // Same content at a new offset is ordinary repeated terminal traffic. + appendPtyChunk('k1', 'same', 20, 1) + // Offset counters reset across generations; this is fresh traffic. + appendPtyChunk('k1', 'same', 10, 2) + // Identity alone is insufficient when content differs. + appendPtyChunk('k1', 'different', 10, 1) + flushRaf() + + expect(listener).toHaveBeenCalledWith(['same', 'same', 'same', 'different']) + expect(getPtyChunks('k1')).toEqual(['same', 'same', 'same', 'different']) + }) }) // Invariants for the monotonic replay-baseline API (getPtyChunkTotal / @@ -227,6 +241,46 @@ describe('pty-buffer-store replay baselines', () => { expect(getPtyChunksSinceTotal(key, getPtyChunkTotal(key))).toEqual([]) }) + it('uses v10 offsets to retain chunks emitted after an attach snapshot', () => { + const key = trackedKey('offset') + appendPtyChunk(key, 'snapshot-covered-a', 10, 2) + appendPtyChunk(key, 'snapshot-covered-b', 20, 2) + appendPtyChunk(key, 'after-snapshot-a', 30, 2) + appendPtyChunk(key, 'after-snapshot-b', 40, 2) + + expect(getPtyChunksAfterOffset(key, 20, 2)).toEqual([ + 'after-snapshot-a', + 'after-snapshot-b' + ]) + }) + + it('delivers identity-less chunks when offset correlation is uncertain', () => { + const key = trackedKey('offset-missing') + appendPtyChunk(key, 'known-covered', 10, 2) + appendPtyChunk(key, 'legacy-or-reset') + appendPtyChunk(key, 'known-fresh', 30, 2) + + expect(getPtyChunksAfterOffset(key, 20, 2)).toEqual([ + 'legacy-or-reset', + 'known-fresh' + ]) + }) + + it('scopes offsets to the listener generation across broker reconnects', () => { + const key = trackedKey('offset-generation') + appendPtyChunk(key, 'stale-high-offset', 5_000, 4) + appendPtyChunk(key, 'snapshot-covered', 10, 5) + appendPtyChunk(key, 'current-fresh', 20, 5) + appendPtyChunk(key, 'future-generation', 1, 6) + appendPtyChunk(key, 'unknown-generation', 30) + + expect(getPtyChunksAfterOffset(key, 10, 5)).toEqual([ + 'current-fresh', + 'future-generation', + 'unknown-generation' + ]) + }) + it('never replays pre-baseline chunks after the buffer trims', () => { const key = trackedKey('trim') appendPtyChunk(key, 'snapshot-covered') diff --git a/src/renderer/src/stores/pty-buffer-store.ts b/src/renderer/src/stores/pty-buffer-store.ts index dc7fa950..dc9acfc7 100644 --- a/src/renderer/src/stores/pty-buffer-store.ts +++ b/src/renderer/src/stores/pty-buffer-store.ts @@ -21,6 +21,14 @@ const MAX_PTY_BUFFER_CHUNKS = 10_000 type Listener = (newChunks: string[]) => void const buffers = new Map() +// Relay v10 attaches a cumulative raw-PTY byte offset to each worker_stream +// chunk. Keep offset + Pear listener generation aligned with `buffers` so +// attach seeding can discard only chunks the broker snapshot proves it already +// contains. Undefined metadata is legacy/unknown and is delivered when +// correlation is requested — a possible duplicate repaint is safer than +// dropping terminal bytes. +const offsets = new Map>() +const generations = new Map>() const listeners = new Map>() // Monotonic count of chunks ever flushed into each key's buffer. Unlike // `buffers.get(key).length` this never moves backwards on trim, so consumers @@ -31,7 +39,18 @@ const listeners = new Map>() const totals = new Map() // Chunks staged for the next animation frame, keyed by agent key. -const pending = new Map() +type BufferedPtyChunk = { chunk: string; offset?: number; generation?: number } +const pending = new Map() + +// Renderer-side final guard for duplicate IPC/listener delivery. Identity is +// the broker-listener generation plus Relay raw-PTY offset; content is hashed +// separately. Never suppress on either component alone. +const PTY_REPLAY_WINDOW = 512 +interface RendererPtyDedupeState { + recentByIdentity: Map + suppressedReplays: number +} +const replayDedupe = new Map() // Scheduled rAF handles per key so we can cancel on clear/dispose. type FrameHandle = | { kind: 'raf'; handle: number } @@ -74,18 +93,33 @@ function flushPending(key: string): void { if (!queued || queued.length === 0) return const existing = buffers.get(key) ?? [] - const combined = existing.concat(queued) + const existingOffsets = offsets.get(key) ?? [] + const existingGenerations = generations.get(key) ?? [] + const queuedChunks = queued.map((entry) => entry.chunk) + const queuedOffsets = queued.map((entry) => entry.offset) + const queuedGenerations = queued.map((entry) => entry.generation) + const combined = existing.concat(queuedChunks) + const combinedOffsets = existingOffsets.concat(queuedOffsets) + const combinedGenerations = existingGenerations.concat(queuedGenerations) const trimmed = combined.length > MAX_PTY_BUFFER_CHUNKS ? combined.slice(combined.length - MAX_PTY_BUFFER_CHUNKS) : combined + const trimmedOffsets = combinedOffsets.length > MAX_PTY_BUFFER_CHUNKS + ? combinedOffsets.slice(combinedOffsets.length - MAX_PTY_BUFFER_CHUNKS) + : combinedOffsets + const trimmedGenerations = combinedGenerations.length > MAX_PTY_BUFFER_CHUNKS + ? combinedGenerations.slice(combinedGenerations.length - MAX_PTY_BUFFER_CHUNKS) + : combinedGenerations buffers.set(key, trimmed) + offsets.set(key, trimmedOffsets) + generations.set(key, trimmedGenerations) totals.set(key, (totals.get(key) ?? 0) + queued.length) const keyListeners = listeners.get(key) if (!keyListeners || keyListeners.size === 0) return for (const listener of [...keyListeners]) { try { - listener(queued) + listener(queuedChunks) } catch (err) { console.error('[pty-buffer-store] listener threw', err) } @@ -116,6 +150,31 @@ export function getPtyChunksSinceTotal(key: string, baselineTotal: number): stri return buffer.slice(buffer.length - missed) } +// Return every retained chunk that is not proven to be represented by a v10 +// attach snapshot. `offset > snapshotOffset` is authoritative post-snapshot +// output. Missing offsets are deliberately delivered (legacy broker or mixed +// generation): never drop a PTY chunk on uncertain identity/correlation. +export function getPtyChunksAfterOffset( + key: string, + snapshotOffset: number, + snapshotGeneration?: number +): string[] { + const buffer = buffers.get(key) ?? [] + const bufferOffsets = offsets.get(key) ?? [] + const bufferGenerations = generations.get(key) ?? [] + return buffer.filter((_, index) => { + const generation = bufferGenerations[index] + // Offset counters reset with a daemon/PTY generation. Without both sides + // of the generation comparison, the numeric offset is not globally safe; + // deliver the uncertain chunk instead of risking dropped terminal bytes. + if (snapshotGeneration === undefined || generation === undefined) return true + if (generation < snapshotGeneration) return false + if (generation > snapshotGeneration) return true + const offset = bufferOffsets[index] + return offset === undefined || offset > snapshotOffset + }) +} + // Synchronously drain any chunks staged for the next rAF into the buffer. // Used by the terminal runtime right before reading the buffer length as a // snapshot baseline — otherwise pending chunks would be replayed on top of @@ -152,16 +211,40 @@ function __previewChunk(chunk: string): string { return chunk.slice(0, 80).replace(/\n/g, '\\n').replace(/\r/g, '\\r').replace(/\x1b/g, '\\e') } -export function appendPtyChunk(key: string, chunk: string): void { +export function appendPtyChunk( + key: string, + chunk: string, + offset?: number, + generation?: number +): void { if (diagPtyEnabled()) { __appendSeq += 1 console.log(`[diag:pty-append] #${__appendSeq} key=${key} bytes=${chunk.length} preview="${__previewChunk(chunk)}"`) } + const normalizedOffset = typeof offset === 'number' && Number.isFinite(offset) && offset >= 0 + ? offset + : undefined + const normalizedGeneration = + typeof generation === 'number' && Number.isSafeInteger(generation) && generation >= 0 + ? generation + : undefined + if ( + normalizedOffset !== undefined && + normalizedGeneration !== undefined && + isDuplicatePtyChunk(key, chunk, normalizedOffset, normalizedGeneration) + ) { + return + } + const entry: BufferedPtyChunk = { + chunk, + ...(normalizedOffset !== undefined ? { offset: normalizedOffset } : {}), + ...(normalizedGeneration !== undefined ? { generation: normalizedGeneration } : {}) + } const queue = pending.get(key) if (queue) { - queue.push(chunk) + queue.push(entry) } else { - pending.set(key, [chunk]) + pending.set(key, [entry]) } const subscriberCount = listeners.get(key)?.size ?? 0 if (subscriberCount === 0) { @@ -189,6 +272,9 @@ export function appendPtyChunk(key: string, chunk: string): void { export function clearPtyBuffer(key: string): void { cancelPendingFlush(key) buffers.delete(key) + offsets.delete(key) + generations.delete(key) + replayDedupe.delete(key) const keyListeners = listeners.get(key) if (keyListeners) { for (const listener of [...keyListeners]) { @@ -201,6 +287,59 @@ export function clearPtyBuffer(key: string): void { } } +function isDuplicatePtyChunk( + key: string, + chunk: string, + offset: number, + generation: number +): boolean { + let state = replayDedupe.get(key) + if (!state) { + state = { recentByIdentity: new Map(), suppressedReplays: 0 } + replayDedupe.set(key, state) + } + + const identity = `${generation}:${offset}` + const hash = ptyChunkHash(chunk) + if (state.recentByIdentity.get(identity) === hash) { + state.suppressedReplays += 1 + // Exponential debug telemetry keeps replay storms visible without logging + // every suppressed terminal chunk. + if ( + diagPtyEnabled() && + (state.suppressedReplays & (state.suppressedReplays - 1)) === 0 + ) { + console.info('[pty-buffer-store] suppressed replayed PTY chunk', { + key, + generation, + offset, + suppressedTotal: state.suppressedReplays + }) + } + return true + } + + // Same identity with different bytes is not a duplicate. Remember the most + // recently delivered content, but deliver both: identity alone is never a + // safe reason to drop PTY traffic. + state.recentByIdentity.delete(identity) + state.recentByIdentity.set(identity, hash) + if (state.recentByIdentity.size > PTY_REPLAY_WINDOW) { + const oldest = state.recentByIdentity.keys().next().value + if (oldest !== undefined) state.recentByIdentity.delete(oldest) + } + return false +} + +function ptyChunkHash(chunk: string): string { + let hash = 0x811c9dc5 + for (let index = 0; index < chunk.length; index += 1) { + hash ^= chunk.charCodeAt(index) + hash = Math.imul(hash, 0x01000193) + } + return `${chunk.length}:${hash >>> 0}` +} + export function subscribePtyBuffer(key: string, listener: Listener): () => void { let set = listeners.get(key) if (!set) { diff --git a/src/shared/types/ipc.ts b/src/shared/types/ipc.ts index bf6e06dd..1cb43ef2 100644 --- a/src/shared/types/ipc.ts +++ b/src/shared/types/ipc.ts @@ -385,6 +385,10 @@ export interface BrokerAttachTerminalResult { cols: number cursor: [number, number] screen: string + /** Cumulative PTY byte offset represented by this snapshot (Relay v10+). */ + offset?: number + /** Pear event-listener generation that captured this snapshot. */ + generation?: number } } @@ -403,6 +407,8 @@ export interface BrokerTerminalSnapshot { cursor: [number, number] /** Row text for `plain`; decoded ANSI reproduction byte stream for `ansi`. */ screen: string + /** Cumulative PTY byte offset represented by this snapshot (Relay v10+). */ + offset?: number } export interface BrokerSendMessageInput { @@ -976,7 +982,15 @@ export interface PearAPI { shutdown: () => Promise onEvent: (callback: (event: unknown) => void) => () => void onEventStreamDiagnostic: (callback: (event: BrokerEventStreamDiagnostic) => void) => () => void - onPtyChunk: (callback: (projectId: string, name: string, chunk: string) => void) => () => void + onPtyChunk: ( + callback: ( + projectId: string, + name: string, + chunk: string, + offset?: number, + generation?: number + ) => void + ) => () => void onStatus: (callback: (status: BrokerStatusEvent) => void) => () => void } factory: {