mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-10-02 13:39:41 +02:00
Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
abbbf9e90a | ||
|
|
3a41de7b57 | ||
|
|
78c568e5f7 | ||
|
|
cc624d2575 | ||
|
|
5844720525 | ||
|
|
ceaf4624a1 | ||
|
|
a6597e4a9a |
Generated
+103
-25
@@ -1,12 +1,12 @@
|
|||||||
{
|
{
|
||||||
"name": "aicodeman",
|
"name": "aicodeman",
|
||||||
"version": "0.3.8",
|
"version": "0.3.11",
|
||||||
"lockfileVersion": 3,
|
"lockfileVersion": 3,
|
||||||
"requires": true,
|
"requires": true,
|
||||||
"packages": {
|
"packages": {
|
||||||
"": {
|
"": {
|
||||||
"name": "aicodeman",
|
"name": "aicodeman",
|
||||||
"version": "0.3.8",
|
"version": "0.3.11",
|
||||||
"hasInstallScript": true,
|
"hasInstallScript": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"workspaces": [
|
"workspaces": [
|
||||||
@@ -17,6 +17,7 @@
|
|||||||
"@fastify/compress": "^8.3.1",
|
"@fastify/compress": "^8.3.1",
|
||||||
"@fastify/cookie": "^11.0.2",
|
"@fastify/cookie": "^11.0.2",
|
||||||
"@fastify/static": "^8.0.0",
|
"@fastify/static": "^8.0.0",
|
||||||
|
"@fastify/websocket": "^11.2.0",
|
||||||
"@xterm/addon-fit": "^0.11.0",
|
"@xterm/addon-fit": "^0.11.0",
|
||||||
"@xterm/addon-unicode11": "^0.9.0",
|
"@xterm/addon-unicode11": "^0.9.0",
|
||||||
"@xterm/addon-webgl": "^0.19.0",
|
"@xterm/addon-webgl": "^0.19.0",
|
||||||
@@ -45,6 +46,7 @@
|
|||||||
"@types/react": "^19.2.14",
|
"@types/react": "^19.2.14",
|
||||||
"@types/uuid": "^10.0.0",
|
"@types/uuid": "^10.0.0",
|
||||||
"@types/web-push": "^3.6.4",
|
"@types/web-push": "^3.6.4",
|
||||||
|
"@types/ws": "^8.18.1",
|
||||||
"@vitest/coverage-v8": "^4.0.18",
|
"@vitest/coverage-v8": "^4.0.18",
|
||||||
"agent-browser": "^0.6.0",
|
"agent-browser": "^0.6.0",
|
||||||
"esbuild": "^0.27.3",
|
"esbuild": "^0.27.3",
|
||||||
@@ -945,6 +947,53 @@
|
|||||||
"glob": "^11.0.0"
|
"glob": "^11.0.0"
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"node_modules/@fastify/websocket": {
|
||||||
|
"version": "11.2.0",
|
||||||
|
"resolved": "https://registry.npmjs.org/@fastify/websocket/-/websocket-11.2.0.tgz",
|
||||||
|
"integrity": "sha512-3HrDPbAG1CzUCqnslgJxppvzaAZffieOVbLp1DAy1huCSynUWPifSvfdEDUR8HlJLp3sp1A36uOM2tJogADS8w==",
|
||||||
|
"funding": [
|
||||||
|
{
|
||||||
|
"type": "github",
|
||||||
|
"url": "https://github.com/sponsors/fastify"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"type": "opencollective",
|
||||||
|
"url": "https://opencollective.com/fastify"
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"license": "MIT",
|
||||||
|
"dependencies": {
|
||||||
|
"duplexify": "^4.1.3",
|
||||||
|
"fastify-plugin": "^5.0.0",
|
||||||
|
"ws": "^8.16.0"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"node_modules/@fastify/websocket/node_modules/duplexify": {
|
||||||
|
"version": "4.1.3",
|
||||||
|
"resolved": "https://registry.npmjs.org/duplexify/-/duplexify-4.1.3.tgz",
|
||||||
|
"integrity": "sha512-M3BmBhwJRZsSx38lZyhE53Csddgzl5R7xGJNk7CVddZD6CcmwMCH8J+7AprIrQKH7TonKxaCjcv27Qmf+sQ+oA==",
|
||||||
|
"license": "MIT",
|
||||||
|
"dependencies": {
|
||||||
|
"end-of-stream": "^1.4.1",
|
||||||
|
"inherits": "^2.0.3",
|
||||||
|
"readable-stream": "^3.1.1",
|
||||||
|
"stream-shift": "^1.0.2"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"node_modules/@fastify/websocket/node_modules/readable-stream": {
|
||||||
|
"version": "3.6.2",
|
||||||
|
"resolved": "https://registry.npmjs.org/readable-stream/-/readable-stream-3.6.2.tgz",
|
||||||
|
"integrity": "sha512-9u/sniCrY3D5WdsERHzHE4G2YCXqoG5FTHUiCC4SIbr6XcLZBY05ya9EKjYek9O5xOAwjGq+1JdGBAS7Q9ScoA==",
|
||||||
|
"license": "MIT",
|
||||||
|
"dependencies": {
|
||||||
|
"inherits": "^2.0.3",
|
||||||
|
"string_decoder": "^1.1.1",
|
||||||
|
"util-deprecate": "^1.0.1"
|
||||||
|
},
|
||||||
|
"engines": {
|
||||||
|
"node": ">= 6"
|
||||||
|
}
|
||||||
|
},
|
||||||
"node_modules/@humanfs/core": {
|
"node_modules/@humanfs/core": {
|
||||||
"version": "0.19.1",
|
"version": "0.19.1",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
@@ -2118,6 +2167,16 @@
|
|||||||
"@types/node": "*"
|
"@types/node": "*"
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"node_modules/@types/ws": {
|
||||||
|
"version": "8.18.1",
|
||||||
|
"resolved": "https://registry.npmjs.org/@types/ws/-/ws-8.18.1.tgz",
|
||||||
|
"integrity": "sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==",
|
||||||
|
"dev": true,
|
||||||
|
"license": "MIT",
|
||||||
|
"dependencies": {
|
||||||
|
"@types/node": "*"
|
||||||
|
}
|
||||||
|
},
|
||||||
"node_modules/@types/yauzl": {
|
"node_modules/@types/yauzl": {
|
||||||
"version": "2.10.3",
|
"version": "2.10.3",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
@@ -3019,7 +3078,9 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/basic-ftp": {
|
"node_modules/basic-ftp": {
|
||||||
"version": "5.1.0",
|
"version": "5.2.0",
|
||||||
|
"resolved": "https://registry.npmjs.org/basic-ftp/-/basic-ftp-5.2.0.tgz",
|
||||||
|
"integrity": "sha512-VoMINM2rqJwJgfdHq6RiUudKt2BV+FY5ZFezP/ypmwayk68+NzzAQy4XXLlqsGD4MCzq3DrmNFD/uUmBJuGoXw==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"engines": {
|
"engines": {
|
||||||
@@ -4416,7 +4477,9 @@
|
|||||||
"license": "BSD-3-Clause"
|
"license": "BSD-3-Clause"
|
||||||
},
|
},
|
||||||
"node_modules/fastify": {
|
"node_modules/fastify": {
|
||||||
"version": "5.7.4",
|
"version": "5.8.2",
|
||||||
|
"resolved": "https://registry.npmjs.org/fastify/-/fastify-5.8.2.tgz",
|
||||||
|
"integrity": "sha512-lZmt3navvZG915IE+f7/TIVamxIwmBd+OMB+O9WBzcpIwOo6F0LTh0sluoMFk5VkrKTvvrwIaoJPkir4Z+jtAg==",
|
||||||
"funding": [
|
"funding": [
|
||||||
{
|
{
|
||||||
"type": "github",
|
"type": "github",
|
||||||
@@ -4438,7 +4501,7 @@
|
|||||||
"fast-json-stringify": "^6.0.0",
|
"fast-json-stringify": "^6.0.0",
|
||||||
"find-my-way": "^9.0.0",
|
"find-my-way": "^9.0.0",
|
||||||
"light-my-request": "^6.0.0",
|
"light-my-request": "^6.0.0",
|
||||||
"pino": "^10.1.0",
|
"pino": "^9.14.0 || ^10.1.0",
|
||||||
"process-warning": "^5.0.0",
|
"process-warning": "^5.0.0",
|
||||||
"rfdc": "^1.3.1",
|
"rfdc": "^1.3.1",
|
||||||
"secure-json-parse": "^4.0.0",
|
"secure-json-parse": "^4.0.0",
|
||||||
@@ -4601,6 +4664,20 @@
|
|||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "Unlicense"
|
"license": "Unlicense"
|
||||||
},
|
},
|
||||||
|
"node_modules/fsevents": {
|
||||||
|
"version": "2.3.3",
|
||||||
|
"resolved": "https://registry.npmjs.org/fsevents/-/fsevents-2.3.3.tgz",
|
||||||
|
"integrity": "sha512-5xoDfX+fL7faATnagmWPpbFtwh/R77WmMMqqHGS65C3vvB0YHrgF+B1YmZ3441tMj5n63k0212XNoJwzlhffQw==",
|
||||||
|
"hasInstallScript": true,
|
||||||
|
"license": "MIT",
|
||||||
|
"optional": true,
|
||||||
|
"os": [
|
||||||
|
"darwin"
|
||||||
|
],
|
||||||
|
"engines": {
|
||||||
|
"node": "^8.16.0 || ^10.6.0 || >=11.0.0"
|
||||||
|
}
|
||||||
|
},
|
||||||
"node_modules/function-bind": {
|
"node_modules/function-bind": {
|
||||||
"version": "1.1.2",
|
"version": "1.1.2",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
@@ -5641,7 +5718,9 @@
|
|||||||
"license": "ISC"
|
"license": "ISC"
|
||||||
},
|
},
|
||||||
"node_modules/minimatch": {
|
"node_modules/minimatch": {
|
||||||
"version": "10.2.2",
|
"version": "10.2.4",
|
||||||
|
"resolved": "https://registry.npmjs.org/minimatch/-/minimatch-10.2.4.tgz",
|
||||||
|
"integrity": "sha512-oRjTw/97aTBN0RHbYCdtF1MQfvusSIBQM0IZEgzl6426+8jSC0nF1a/GmnVLpfB9yyr6g6FTqWqiZVbxrtaCIg==",
|
||||||
"license": "BlueOak-1.0.0",
|
"license": "BlueOak-1.0.0",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"brace-expansion": "^5.0.2"
|
"brace-expansion": "^5.0.2"
|
||||||
@@ -6163,6 +6242,21 @@
|
|||||||
"node": ">=18"
|
"node": ">=18"
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"node_modules/playwright/node_modules/fsevents": {
|
||||||
|
"version": "2.3.2",
|
||||||
|
"resolved": "https://registry.npmjs.org/fsevents/-/fsevents-2.3.2.tgz",
|
||||||
|
"integrity": "sha512-xiqMQR4xAeHTuB9uWm+fFRcIOgKBMiOBP+eXiyT7jsgVCq1bkVygt00oASowB7EdtpOHaaPgKt812P9ab+DDKA==",
|
||||||
|
"dev": true,
|
||||||
|
"hasInstallScript": true,
|
||||||
|
"license": "MIT",
|
||||||
|
"optional": true,
|
||||||
|
"os": [
|
||||||
|
"darwin"
|
||||||
|
],
|
||||||
|
"engines": {
|
||||||
|
"node": "^8.16.0 || ^10.6.0 || >=11.0.0"
|
||||||
|
}
|
||||||
|
},
|
||||||
"node_modules/pngjs": {
|
"node_modules/pngjs": {
|
||||||
"version": "7.0.0",
|
"version": "7.0.0",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
@@ -6625,14 +6719,6 @@
|
|||||||
"version": "4.0.4",
|
"version": "4.0.4",
|
||||||
"license": "MIT"
|
"license": "MIT"
|
||||||
},
|
},
|
||||||
"node_modules/randombytes": {
|
|
||||||
"version": "2.1.0",
|
|
||||||
"dev": true,
|
|
||||||
"license": "MIT",
|
|
||||||
"dependencies": {
|
|
||||||
"safe-buffer": "^5.1.0"
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"node_modules/react": {
|
"node_modules/react": {
|
||||||
"version": "19.2.4",
|
"version": "19.2.4",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
@@ -7023,14 +7109,6 @@
|
|||||||
"node": ">=10"
|
"node": ">=10"
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/serialize-javascript": {
|
|
||||||
"version": "6.0.2",
|
|
||||||
"dev": true,
|
|
||||||
"license": "BSD-3-Clause",
|
|
||||||
"dependencies": {
|
|
||||||
"randombytes": "^2.1.0"
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"node_modules/set-blocking": {
|
"node_modules/set-blocking": {
|
||||||
"version": "2.0.0",
|
"version": "2.0.0",
|
||||||
"license": "ISC"
|
"license": "ISC"
|
||||||
@@ -7401,14 +7479,15 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/terser-webpack-plugin": {
|
"node_modules/terser-webpack-plugin": {
|
||||||
"version": "5.3.16",
|
"version": "5.4.0",
|
||||||
|
"resolved": "https://registry.npmjs.org/terser-webpack-plugin/-/terser-webpack-plugin-5.4.0.tgz",
|
||||||
|
"integrity": "sha512-Bn5vxm48flOIfkdl5CaD2+1CiUVbonWQ3KQPyP7/EuIl9Gbzq/gQFOzaMFUEgVjB1396tcK0SG8XcNJ/2kDH8g==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@jridgewell/trace-mapping": "^0.3.25",
|
"@jridgewell/trace-mapping": "^0.3.25",
|
||||||
"jest-worker": "^27.4.5",
|
"jest-worker": "^27.4.5",
|
||||||
"schema-utils": "^4.3.0",
|
"schema-utils": "^4.3.0",
|
||||||
"serialize-javascript": "^6.0.2",
|
|
||||||
"terser": "^5.31.1"
|
"terser": "^5.31.1"
|
||||||
},
|
},
|
||||||
"engines": {
|
"engines": {
|
||||||
@@ -8488,7 +8567,6 @@
|
|||||||
},
|
},
|
||||||
"node_modules/ws": {
|
"node_modules/ws": {
|
||||||
"version": "8.19.0",
|
"version": "8.19.0",
|
||||||
"dev": true,
|
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"engines": {
|
"engines": {
|
||||||
"node": ">=10.0.0"
|
"node": ">=10.0.0"
|
||||||
|
|||||||
@@ -51,6 +51,7 @@
|
|||||||
"@fastify/compress": "^8.3.1",
|
"@fastify/compress": "^8.3.1",
|
||||||
"@fastify/cookie": "^11.0.2",
|
"@fastify/cookie": "^11.0.2",
|
||||||
"@fastify/static": "^8.0.0",
|
"@fastify/static": "^8.0.0",
|
||||||
|
"@fastify/websocket": "^11.2.0",
|
||||||
"@xterm/addon-fit": "^0.11.0",
|
"@xterm/addon-fit": "^0.11.0",
|
||||||
"@xterm/addon-unicode11": "^0.9.0",
|
"@xterm/addon-unicode11": "^0.9.0",
|
||||||
"@xterm/addon-webgl": "^0.19.0",
|
"@xterm/addon-webgl": "^0.19.0",
|
||||||
@@ -76,6 +77,7 @@
|
|||||||
"@types/react": "^19.2.14",
|
"@types/react": "^19.2.14",
|
||||||
"@types/uuid": "^10.0.0",
|
"@types/uuid": "^10.0.0",
|
||||||
"@types/web-push": "^3.6.4",
|
"@types/web-push": "^3.6.4",
|
||||||
|
"@types/ws": "^8.18.1",
|
||||||
"@vitest/coverage-v8": "^4.0.18",
|
"@vitest/coverage-v8": "^4.0.18",
|
||||||
"agent-browser": "^0.6.0",
|
"agent-browser": "^0.6.0",
|
||||||
"esbuild": "^0.27.3",
|
"esbuild": "^0.27.3",
|
||||||
|
|||||||
+119
-8
@@ -172,9 +172,9 @@ const _SSE_HANDLER_MAP = [
|
|||||||
[SSE_EVENTS.SESSION_CREATED, '_onSessionCreated'],
|
[SSE_EVENTS.SESSION_CREATED, '_onSessionCreated'],
|
||||||
[SSE_EVENTS.SESSION_UPDATED, '_onSessionUpdated'],
|
[SSE_EVENTS.SESSION_UPDATED, '_onSessionUpdated'],
|
||||||
[SSE_EVENTS.SESSION_DELETED, '_onSessionDeleted'],
|
[SSE_EVENTS.SESSION_DELETED, '_onSessionDeleted'],
|
||||||
[SSE_EVENTS.SESSION_TERMINAL, '_onSessionTerminal'],
|
[SSE_EVENTS.SESSION_TERMINAL, '_onSSETerminal'],
|
||||||
[SSE_EVENTS.SESSION_NEEDS_REFRESH, '_onSessionNeedsRefresh'],
|
[SSE_EVENTS.SESSION_NEEDS_REFRESH, '_onSSENeedsRefresh'],
|
||||||
[SSE_EVENTS.SESSION_CLEAR_TERMINAL, '_onSessionClearTerminal'],
|
[SSE_EVENTS.SESSION_CLEAR_TERMINAL, '_onSSEClearTerminal'],
|
||||||
[SSE_EVENTS.SESSION_COMPLETION, '_onSessionCompletion'],
|
[SSE_EVENTS.SESSION_COMPLETION, '_onSessionCompletion'],
|
||||||
[SSE_EVENTS.SESSION_ERROR, '_onSessionError'],
|
[SSE_EVENTS.SESSION_ERROR, '_onSessionError'],
|
||||||
[SSE_EVENTS.SESSION_EXIT, '_onSessionExit'],
|
[SSE_EVENTS.SESSION_EXIT, '_onSessionExit'],
|
||||||
@@ -361,6 +361,11 @@ class CodemanApp {
|
|||||||
// Tracks pending hook events that need resolution (permission_prompt, elicitation_dialog, idle_prompt)
|
// Tracks pending hook events that need resolution (permission_prompt, elicitation_dialog, idle_prompt)
|
||||||
this.pendingHooks = new Map();
|
this.pendingHooks = new Map();
|
||||||
|
|
||||||
|
// WebSocket terminal I/O (low-latency bypass of HTTP POST + SSE)
|
||||||
|
this._ws = null; // WebSocket instance for active session
|
||||||
|
this._wsSessionId = null; // Session ID the WS is connected to
|
||||||
|
this._wsReady = false; // True when WS is open and ready for I/O
|
||||||
|
|
||||||
// Terminal write batching with DEC 2026 sync support
|
// Terminal write batching with DEC 2026 sync support
|
||||||
this.pendingWrites = [];
|
this.pendingWrites = [];
|
||||||
this.writeFrameScheduled = false;
|
this.writeFrameScheduled = false;
|
||||||
@@ -1864,6 +1869,7 @@ class CodemanApp {
|
|||||||
}
|
}
|
||||||
|
|
||||||
_onSessionDeleted(data) {
|
_onSessionDeleted(data) {
|
||||||
|
if (this._wsSessionId === data.id) this._disconnectWs();
|
||||||
this._cleanupSessionData(data.id);
|
this._cleanupSessionData(data.id);
|
||||||
if (this.activeSessionId === data.id) {
|
if (this.activeSessionId === data.id) {
|
||||||
this.activeSessionId = null;
|
this.activeSessionId = null;
|
||||||
@@ -1878,6 +1884,21 @@ class CodemanApp {
|
|||||||
if (this.sessions.size === 0) this.stopSystemStatsPolling();
|
if (this.sessions.size === 0) this.stopSystemStatsPolling();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SSE wrappers — skip terminal events when WebSocket is delivering for this session.
|
||||||
|
// WS handler calls the underlying _onSession* methods directly.
|
||||||
|
_onSSETerminal(data) {
|
||||||
|
if (this._wsReady && this._wsSessionId === data.id) return;
|
||||||
|
this._onSessionTerminal(data);
|
||||||
|
}
|
||||||
|
_onSSENeedsRefresh(data) {
|
||||||
|
if (this._wsReady && this._wsSessionId === data?.id) return;
|
||||||
|
this._onSessionNeedsRefresh(data);
|
||||||
|
}
|
||||||
|
_onSSEClearTerminal(data) {
|
||||||
|
if (this._wsReady && this._wsSessionId === data?.id) return;
|
||||||
|
this._onSessionClearTerminal(data);
|
||||||
|
}
|
||||||
|
|
||||||
_onSessionTerminal(data) {
|
_onSessionTerminal(data) {
|
||||||
if (data.id === this.activeSessionId) {
|
if (data.id === this.activeSessionId) {
|
||||||
if (data.data.length > 32768) _crashDiag.log(`TERMINAL: ${(data.data.length/1024).toFixed(0)}KB`);
|
if (data.data.length > 32768) _crashDiag.log(`TERMINAL: ${(data.data.length/1024).toFixed(0)}KB`);
|
||||||
@@ -1982,6 +2003,7 @@ class CodemanApp {
|
|||||||
}
|
}
|
||||||
|
|
||||||
_onSessionExit(data) {
|
_onSessionExit(data) {
|
||||||
|
if (this._wsSessionId === data.id) this._disconnectWs();
|
||||||
const session = this.sessions.get(data.id);
|
const session = this.sessions.get(data.id);
|
||||||
if (session) {
|
if (session) {
|
||||||
session.status = 'stopped';
|
session.status = 'stopped';
|
||||||
@@ -2835,6 +2857,72 @@ class CodemanApp {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ═══════════════════════════════════════════════════════════════
|
||||||
|
// WebSocket Terminal I/O
|
||||||
|
// ═══════════════════════════════════════════════════════════════
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Open a WebSocket for terminal I/O on the given session.
|
||||||
|
* Replaces HTTP POST input and SSE terminal output with a single
|
||||||
|
* bidirectional connection. Falls back to SSE+POST if WS fails.
|
||||||
|
*/
|
||||||
|
_connectWs(sessionId) {
|
||||||
|
this._disconnectWs();
|
||||||
|
|
||||||
|
const proto = location.protocol === 'https:' ? 'wss:' : 'ws:';
|
||||||
|
const url = `${proto}//${location.host}/ws/sessions/${sessionId}/terminal`;
|
||||||
|
const ws = new WebSocket(url);
|
||||||
|
this._ws = ws;
|
||||||
|
this._wsSessionId = sessionId;
|
||||||
|
|
||||||
|
ws.onopen = () => {
|
||||||
|
// Only mark ready if this is still the intended session
|
||||||
|
if (this._ws === ws) {
|
||||||
|
this._wsReady = true;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
ws.onmessage = (event) => {
|
||||||
|
if (this._ws !== ws) return;
|
||||||
|
try {
|
||||||
|
const msg = JSON.parse(event.data);
|
||||||
|
if (msg.t === 'o') {
|
||||||
|
// Terminal output — route through the same batching pipeline as SSE
|
||||||
|
this._onSessionTerminal({ id: sessionId, data: msg.d });
|
||||||
|
} else if (msg.t === 'c') {
|
||||||
|
this._onSessionClearTerminal({ id: sessionId });
|
||||||
|
} else if (msg.t === 'r') {
|
||||||
|
this._onSessionNeedsRefresh({ id: sessionId });
|
||||||
|
}
|
||||||
|
} catch {
|
||||||
|
// Ignore malformed messages
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
ws.onclose = () => {
|
||||||
|
if (this._ws === ws) {
|
||||||
|
this._ws = null;
|
||||||
|
this._wsSessionId = null;
|
||||||
|
this._wsReady = false;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
ws.onerror = () => {
|
||||||
|
// onclose will fire after onerror — cleanup happens there
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Close the active WebSocket connection (if any). */
|
||||||
|
_disconnectWs() {
|
||||||
|
if (this._ws) {
|
||||||
|
this._ws.onclose = null; // Prevent re-entrant cleanup
|
||||||
|
this._ws.close();
|
||||||
|
this._ws = null;
|
||||||
|
this._wsSessionId = null;
|
||||||
|
this._wsReady = false;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Send input to server without blocking the keystroke flush cycle.
|
* Send input to server without blocking the keystroke flush cycle.
|
||||||
* Uses a sequential promise chain to preserve character ordering
|
* Uses a sequential promise chain to preserve character ordering
|
||||||
@@ -2847,11 +2935,19 @@ class CodemanApp {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Chain on dispatch only — wait for the previous request to be sent before
|
// Fast path: WebSocket — fire-and-forget, inherently ordered (single TCP stream).
|
||||||
// dispatching the next one (preserves keystroke ordering), but don't wait
|
if (this._wsReady && this._wsSessionId === sessionId) {
|
||||||
// for the server's response. The server handles writeViaMux as
|
try {
|
||||||
// fire-and-forget anyway, so the HTTP response carries no useful data
|
this._ws.send(JSON.stringify({ t: 'i', d: input }));
|
||||||
// beyond success/failure for retry purposes.
|
this.clearPendingHooks(sessionId);
|
||||||
|
return;
|
||||||
|
} catch {
|
||||||
|
// WS send failed — fall through to HTTP POST
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Slow path: HTTP POST — chain on dispatch only, don't wait for response.
|
||||||
|
// The server handles writeViaMux as fire-and-forget anyway.
|
||||||
this._inputSendChain = this._inputSendChain.then(() => {
|
this._inputSendChain = this._inputSendChain.then(() => {
|
||||||
const fetchPromise = fetch(`/api/sessions/${sessionId}/input`, {
|
const fetchPromise = fetch(`/api/sessions/${sessionId}/input`, {
|
||||||
method: 'POST',
|
method: 'POST',
|
||||||
@@ -3632,6 +3728,9 @@ class CodemanApp {
|
|||||||
|
|
||||||
if (selectGen !== this._selectGeneration) return; // newer tab switch won
|
if (selectGen !== this._selectGeneration) return; // newer tab switch won
|
||||||
|
|
||||||
|
// Close WebSocket for previous session (new one opens after buffer load)
|
||||||
|
this._disconnectWs();
|
||||||
|
|
||||||
// Clean up flicker filter state when switching sessions
|
// Clean up flicker filter state when switching sessions
|
||||||
if (this.flickerFilterTimeout) {
|
if (this.flickerFilterTimeout) {
|
||||||
clearTimeout(this.flickerFilterTimeout);
|
clearTimeout(this.flickerFilterTimeout);
|
||||||
@@ -3932,6 +4031,9 @@ class CodemanApp {
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// Open WebSocket for low-latency terminal I/O (after buffer load completes)
|
||||||
|
this._connectWs(sessionId);
|
||||||
|
|
||||||
_crashDiag.log('FOCUS');
|
_crashDiag.log('FOCUS');
|
||||||
this.terminal.focus();
|
this.terminal.focus();
|
||||||
this.terminal.scrollToBottom();
|
this.terminal.scrollToBottom();
|
||||||
@@ -5254,6 +5356,15 @@ class CodemanApp {
|
|||||||
async sendResize(sessionId) {
|
async sendResize(sessionId) {
|
||||||
const dims = this.getTerminalDimensions();
|
const dims = this.getTerminalDimensions();
|
||||||
if (!dims) return;
|
if (!dims) return;
|
||||||
|
// Fast path: WebSocket resize
|
||||||
|
if (this._wsReady && this._wsSessionId === sessionId) {
|
||||||
|
try {
|
||||||
|
this._ws.send(JSON.stringify({ t: 'z', c: dims.cols, r: dims.rows }));
|
||||||
|
return;
|
||||||
|
} catch {
|
||||||
|
// Fall through to HTTP POST
|
||||||
|
}
|
||||||
|
}
|
||||||
await fetch(`/api/sessions/${sessionId}/resize`, {
|
await fetch(`/api/sessions/${sessionId}/resize`, {
|
||||||
method: 'POST',
|
method: 'POST',
|
||||||
headers: { 'Content-Type': 'application/json' },
|
headers: { 'Content-Type': 'application/json' },
|
||||||
|
|||||||
@@ -14,3 +14,4 @@ export { registerSessionRoutes } from './session-routes.js';
|
|||||||
export { registerRespawnRoutes } from './respawn-routes.js';
|
export { registerRespawnRoutes } from './respawn-routes.js';
|
||||||
export { registerRalphRoutes } from './ralph-routes.js';
|
export { registerRalphRoutes } from './ralph-routes.js';
|
||||||
export { registerPlanRoutes } from './plan-routes.js';
|
export { registerPlanRoutes } from './plan-routes.js';
|
||||||
|
export { registerWsRoutes } from './ws-routes.js';
|
||||||
|
|||||||
@@ -0,0 +1,179 @@
|
|||||||
|
/**
|
||||||
|
* @fileoverview WebSocket terminal I/O route.
|
||||||
|
*
|
||||||
|
* Provides a low-latency bidirectional channel for terminal input/output,
|
||||||
|
* bypassing the HTTP POST + SSE path that adds per-request middleware overhead.
|
||||||
|
* Auth is checked once on the WebSocket upgrade handshake (cookies are included
|
||||||
|
* automatically by the browser). After upgrade, the connection is raw — no
|
||||||
|
* per-message middleware processing.
|
||||||
|
*
|
||||||
|
* Additive: the existing HTTP POST /api/sessions/:id/input and SSE session:terminal
|
||||||
|
* paths remain fully functional. The frontend opts into WS when available and
|
||||||
|
* falls back transparently.
|
||||||
|
*
|
||||||
|
* Terminal output is micro-batched at 8ms to group Ink's rapid cursor-up redraws
|
||||||
|
* into single frames, preventing flicker from split ANSI sequences. This matches
|
||||||
|
* the SSE path's server-side batching (16-50ms) but at a shorter interval since
|
||||||
|
* WS has no Traefik buffering overhead.
|
||||||
|
*
|
||||||
|
* Protocol (all JSON text frames):
|
||||||
|
* Server -> Client:
|
||||||
|
* {"t":"o","d":"..."} — terminal output
|
||||||
|
* {"t":"c"} — clear terminal
|
||||||
|
* {"t":"r"} — needs refresh (reload buffer)
|
||||||
|
* Client -> Server:
|
||||||
|
* {"t":"i","d":"..."} — input (keystroke or paste)
|
||||||
|
* {"t":"z","c":N,"r":N} — resize terminal
|
||||||
|
*/
|
||||||
|
|
||||||
|
import { FastifyInstance } from 'fastify';
|
||||||
|
import type { WebSocket } from 'ws';
|
||||||
|
import type { SessionPort } from '../ports/session-port.js';
|
||||||
|
import { MAX_INPUT_LENGTH } from '../../config/terminal-limits.js';
|
||||||
|
|
||||||
|
/** Micro-batch interval for terminal output (ms). Short enough for low latency,
|
||||||
|
* long enough to group Ink's rapid cursor-up redraw sequences into single frames. */
|
||||||
|
const WS_BATCH_INTERVAL_MS = 8;
|
||||||
|
|
||||||
|
/** Flush immediately when batch exceeds this size (bytes) for responsiveness. */
|
||||||
|
const WS_BATCH_FLUSH_THRESHOLD = 16384;
|
||||||
|
|
||||||
|
/** How often to ping each WebSocket client (ms). Detects stale connections that
|
||||||
|
* TCP keepalive won't catch for minutes, especially through tunnels/proxies. */
|
||||||
|
const WS_PING_INTERVAL_MS = 30_000;
|
||||||
|
|
||||||
|
/** If pong isn't received within this window after a ping, terminate the socket. */
|
||||||
|
const WS_PONG_TIMEOUT_MS = 10_000;
|
||||||
|
|
||||||
|
/** DEC 2026 synchronized update markers. Wrapping output in these tells xterm.js
|
||||||
|
* to buffer all content and render atomically in a single frame — eliminates
|
||||||
|
* flicker from cursor-up redraws that Ink sends without its own sync markers
|
||||||
|
* (DA capability negotiation fails through the PTY→server→WS proxy chain). */
|
||||||
|
const DEC_2026_START = '\x1b[?2026h';
|
||||||
|
const DEC_2026_END = '\x1b[?2026l';
|
||||||
|
|
||||||
|
export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort): void {
|
||||||
|
app.get<{ Params: { id: string } }>('/ws/sessions/:id/terminal', { websocket: true }, (socket: WebSocket, req) => {
|
||||||
|
const { id } = req.params;
|
||||||
|
const session = ctx.sessions.get(id);
|
||||||
|
|
||||||
|
if (!session) {
|
||||||
|
socket.close(4004, 'Session not found');
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Per-connection micro-batch state
|
||||||
|
let batchChunks: string[] = [];
|
||||||
|
let batchSize = 0;
|
||||||
|
let batchTimer: ReturnType<typeof setTimeout> | null = null;
|
||||||
|
|
||||||
|
const flushBatch = () => {
|
||||||
|
batchTimer = null;
|
||||||
|
if (batchChunks.length === 0 || socket.readyState !== 1) {
|
||||||
|
batchChunks = [];
|
||||||
|
batchSize = 0;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
const data = batchChunks.join('');
|
||||||
|
batchChunks = [];
|
||||||
|
batchSize = 0;
|
||||||
|
socket.send(`{"t":"o","d":${JSON.stringify(DEC_2026_START + data + DEC_2026_END)}}`);
|
||||||
|
};
|
||||||
|
|
||||||
|
// Attach message handler synchronously BEFORE any async work
|
||||||
|
// (@fastify/websocket requirement to avoid dropped messages).
|
||||||
|
socket.on('message', (raw) => {
|
||||||
|
try {
|
||||||
|
const msg = JSON.parse(String(raw));
|
||||||
|
if (msg.t === 'i' && typeof msg.d === 'string') {
|
||||||
|
if (msg.d.length > MAX_INPUT_LENGTH) return;
|
||||||
|
session.write(msg.d);
|
||||||
|
} else if (
|
||||||
|
msg.t === 'z' &&
|
||||||
|
Number.isInteger(msg.c) &&
|
||||||
|
Number.isInteger(msg.r) &&
|
||||||
|
msg.c >= 1 &&
|
||||||
|
msg.c <= 500 &&
|
||||||
|
msg.r >= 1 &&
|
||||||
|
msg.r <= 200
|
||||||
|
) {
|
||||||
|
session.resize(msg.c, msg.r);
|
||||||
|
}
|
||||||
|
} catch {
|
||||||
|
// Ignore malformed messages
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// Terminal output -> micro-batched WS send
|
||||||
|
const onTerminal = (data: string) => {
|
||||||
|
batchChunks.push(data);
|
||||||
|
batchSize += data.length;
|
||||||
|
|
||||||
|
// Flush immediately for large batches (responsiveness during bulk output)
|
||||||
|
if (batchSize > WS_BATCH_FLUSH_THRESHOLD) {
|
||||||
|
if (batchTimer) {
|
||||||
|
clearTimeout(batchTimer);
|
||||||
|
}
|
||||||
|
flushBatch();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Start timer if not already running
|
||||||
|
if (!batchTimer) {
|
||||||
|
batchTimer = setTimeout(flushBatch, WS_BATCH_INTERVAL_MS);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
const onClearTerminal = () => {
|
||||||
|
if (socket.readyState === 1) {
|
||||||
|
socket.send('{"t":"c"}');
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
const onNeedsRefresh = () => {
|
||||||
|
if (socket.readyState === 1) {
|
||||||
|
socket.send('{"t":"r"}');
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
session.on('terminal', onTerminal);
|
||||||
|
session.on('clearTerminal', onClearTerminal);
|
||||||
|
session.on('needsRefresh', onNeedsRefresh);
|
||||||
|
|
||||||
|
// Heartbeat: detect stale connections (especially through tunnels where
|
||||||
|
// TCP RST can take minutes to propagate).
|
||||||
|
let pongTimeout: ReturnType<typeof setTimeout> | null = null;
|
||||||
|
let alive = true;
|
||||||
|
|
||||||
|
socket.on('pong', () => {
|
||||||
|
alive = true;
|
||||||
|
if (pongTimeout) {
|
||||||
|
clearTimeout(pongTimeout);
|
||||||
|
pongTimeout = null;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
const pingInterval = setInterval(() => {
|
||||||
|
if (!alive) {
|
||||||
|
// Previous ping never got a pong — connection is dead
|
||||||
|
socket.terminate();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
alive = false;
|
||||||
|
socket.ping();
|
||||||
|
pongTimeout = setTimeout(() => {
|
||||||
|
socket.terminate();
|
||||||
|
}, WS_PONG_TIMEOUT_MS);
|
||||||
|
}, WS_PING_INTERVAL_MS);
|
||||||
|
|
||||||
|
socket.on('close', () => {
|
||||||
|
clearInterval(pingInterval);
|
||||||
|
if (pongTimeout) clearTimeout(pongTimeout);
|
||||||
|
if (batchTimer) clearTimeout(batchTimer);
|
||||||
|
batchChunks = [];
|
||||||
|
session.off('terminal', onTerminal);
|
||||||
|
session.off('clearTerminal', onClearTerminal);
|
||||||
|
session.off('needsRefresh', onNeedsRefresh);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
}
|
||||||
@@ -31,6 +31,7 @@ import Fastify, { FastifyInstance, FastifyReply } from 'fastify';
|
|||||||
import fastifyCompress from '@fastify/compress';
|
import fastifyCompress from '@fastify/compress';
|
||||||
import fastifyCookie from '@fastify/cookie';
|
import fastifyCookie from '@fastify/cookie';
|
||||||
import fastifyStatic from '@fastify/static';
|
import fastifyStatic from '@fastify/static';
|
||||||
|
import fastifyWebsocket from '@fastify/websocket';
|
||||||
import { join, dirname } from 'node:path';
|
import { join, dirname } from 'node:path';
|
||||||
import { fileURLToPath } from 'node:url';
|
import { fileURLToPath } from 'node:url';
|
||||||
import { existsSync, mkdirSync, readFileSync, chmodSync } from 'node:fs';
|
import { existsSync, mkdirSync, readFileSync, chmodSync } from 'node:fs';
|
||||||
@@ -103,6 +104,7 @@ import {
|
|||||||
registerRespawnRoutes,
|
registerRespawnRoutes,
|
||||||
registerRalphRoutes,
|
registerRalphRoutes,
|
||||||
registerPlanRoutes,
|
registerPlanRoutes,
|
||||||
|
registerWsRoutes,
|
||||||
} from './routes/index.js';
|
} from './routes/index.js';
|
||||||
|
|
||||||
const __dirname = dirname(fileURLToPath(import.meta.url));
|
const __dirname = dirname(fileURLToPath(import.meta.url));
|
||||||
@@ -563,6 +565,9 @@ export class WebServer extends EventEmitter {
|
|||||||
this.qrAuthFailures = authState.qrAuthFailures;
|
this.qrAuthFailures = authState.qrAuthFailures;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// WebSocket support (terminal I/O — low-latency bidirectional channel)
|
||||||
|
await this.app.register(fastifyWebsocket);
|
||||||
|
|
||||||
// Security headers + CORS
|
// Security headers + CORS
|
||||||
registerSecurityHeaders(this.app, this.https);
|
registerSecurityHeaders(this.app, this.https);
|
||||||
// Service worker must never be cached — browsers check for SW updates on navigation
|
// Service worker must never be cached — browsers check for SW updates on navigation
|
||||||
@@ -700,6 +705,7 @@ export class WebServer extends EventEmitter {
|
|||||||
registerRespawnRoutes(this.app, ctx);
|
registerRespawnRoutes(this.app, ctx);
|
||||||
registerRalphRoutes(this.app, ctx);
|
registerRalphRoutes(this.app, ctx);
|
||||||
registerPlanRoutes(this.app, ctx);
|
registerPlanRoutes(this.app, ctx);
|
||||||
|
registerWsRoutes(this.app, ctx);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -0,0 +1,362 @@
|
|||||||
|
/**
|
||||||
|
* @fileoverview Tests for WebSocket terminal I/O route.
|
||||||
|
*
|
||||||
|
* Unlike other route tests that use app.inject(), WebSocket testing requires
|
||||||
|
* a real listening server since inject() doesn't support upgrade requests.
|
||||||
|
* Uses the `ws` package (transitive dep of @fastify/websocket) as the client.
|
||||||
|
*
|
||||||
|
* Port: 3170 (ws-routes tests)
|
||||||
|
*/
|
||||||
|
|
||||||
|
import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest';
|
||||||
|
import Fastify, { type FastifyInstance } from 'fastify';
|
||||||
|
import fastifyWebsocket from '@fastify/websocket';
|
||||||
|
import WebSocket from 'ws';
|
||||||
|
import { createMockRouteContext, type MockRouteContext } from '../mocks/index.js';
|
||||||
|
import { registerWsRoutes } from '../../src/web/routes/ws-routes.js';
|
||||||
|
|
||||||
|
const PORT = 3170;
|
||||||
|
|
||||||
|
/** Helper: open a WebSocket connection and wait for it to reach OPEN state. */
|
||||||
|
function connectWs(path: string): Promise<WebSocket> {
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
const ws = new WebSocket(`ws://127.0.0.1:${PORT}${path}`);
|
||||||
|
ws.on('open', () => resolve(ws));
|
||||||
|
ws.on('error', reject);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Helper: wait for the next WS message, parsed as JSON. */
|
||||||
|
function nextMessage(ws: WebSocket, timeoutMs = 2000): Promise<unknown> {
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
const timer = setTimeout(() => reject(new Error('WS message timeout')), timeoutMs);
|
||||||
|
ws.once('message', (raw) => {
|
||||||
|
clearTimeout(timer);
|
||||||
|
resolve(JSON.parse(String(raw)));
|
||||||
|
});
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Helper: wait for WS close event and return { code, reason }. */
|
||||||
|
function waitForClose(ws: WebSocket, timeoutMs = 2000): Promise<{ code: number; reason: string }> {
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
const timer = setTimeout(() => reject(new Error('WS close timeout')), timeoutMs);
|
||||||
|
ws.on('close', (code, reason) => {
|
||||||
|
clearTimeout(timer);
|
||||||
|
resolve({ code, reason: reason.toString() });
|
||||||
|
});
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('ws-routes', () => {
|
||||||
|
let app: FastifyInstance;
|
||||||
|
let ctx: MockRouteContext;
|
||||||
|
|
||||||
|
beforeEach(async () => {
|
||||||
|
app = Fastify({ logger: false });
|
||||||
|
await app.register(fastifyWebsocket);
|
||||||
|
|
||||||
|
ctx = createMockRouteContext({ sessionId: 'ws-test-session' });
|
||||||
|
registerWsRoutes(app, ctx as never);
|
||||||
|
|
||||||
|
await app.listen({ port: PORT, host: '127.0.0.1' });
|
||||||
|
});
|
||||||
|
|
||||||
|
afterEach(async () => {
|
||||||
|
await app.close();
|
||||||
|
});
|
||||||
|
|
||||||
|
// ========== Session not found ==========
|
||||||
|
|
||||||
|
describe('session not found', () => {
|
||||||
|
it('closes with 4004 when session does not exist', async () => {
|
||||||
|
const ws = new WebSocket(`ws://127.0.0.1:${PORT}/ws/sessions/nonexistent/terminal`);
|
||||||
|
const { code, reason } = await waitForClose(ws);
|
||||||
|
expect(code).toBe(4004);
|
||||||
|
expect(reason).toBe('Session not found');
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
// ========== Terminal output ==========
|
||||||
|
|
||||||
|
describe('terminal output', () => {
|
||||||
|
it('receives terminal output via WS with DEC 2026 sync markers', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
|
||||||
|
// Emit terminal data from mock session
|
||||||
|
session.emit('terminal', 'hello world');
|
||||||
|
|
||||||
|
// Wait for the micro-batched message (8ms batch interval + margin)
|
||||||
|
const msg = (await nextMessage(ws)) as { t: string; d: string };
|
||||||
|
expect(msg.t).toBe('o');
|
||||||
|
// Should contain DEC 2026 sync markers wrapping the data
|
||||||
|
expect(msg.d).toContain('hello world');
|
||||||
|
expect(msg.d).toMatch(/^\x1b\[\?2026h/); // starts with DEC 2026 start
|
||||||
|
expect(msg.d).toMatch(/\x1b\[\?2026l$/); // ends with DEC 2026 end
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('sends clearTerminal event as {"t":"c"}', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
ctx._session.emit('clearTerminal');
|
||||||
|
|
||||||
|
const msg = (await nextMessage(ws)) as { t: string };
|
||||||
|
expect(msg.t).toBe('c');
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('sends needsRefresh event as {"t":"r"}', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
ctx._session.emit('needsRefresh');
|
||||||
|
|
||||||
|
const msg = (await nextMessage(ws)) as { t: string };
|
||||||
|
expect(msg.t).toBe('r');
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
// ========== Client input ==========
|
||||||
|
|
||||||
|
describe('client input', () => {
|
||||||
|
it('forwards input messages to session.write()', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
session.writeBuffer = [];
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'i', d: 'ls -la\r' }));
|
||||||
|
|
||||||
|
// Give the message handler time to process
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.writeBuffer).toContain('ls -la\r');
|
||||||
|
});
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('ignores input exceeding MAX_INPUT_LENGTH', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
session.writeBuffer = [];
|
||||||
|
|
||||||
|
// MAX_INPUT_LENGTH is 64KB
|
||||||
|
const hugeInput = 'x'.repeat(65 * 1024);
|
||||||
|
ws.send(JSON.stringify({ t: 'i', d: hugeInput }));
|
||||||
|
|
||||||
|
// Send a valid message after to confirm the connection still works
|
||||||
|
ws.send(JSON.stringify({ t: 'i', d: 'ok' }));
|
||||||
|
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.writeBuffer).toContain('ok');
|
||||||
|
});
|
||||||
|
|
||||||
|
// The oversized input should not have been written
|
||||||
|
expect(session.writeBuffer).not.toContain(hugeInput);
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('ignores malformed JSON messages', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
session.writeBuffer = [];
|
||||||
|
|
||||||
|
ws.send('not-json{{{');
|
||||||
|
|
||||||
|
// Send valid input to verify connection still alive
|
||||||
|
ws.send(JSON.stringify({ t: 'i', d: 'after-bad' }));
|
||||||
|
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.writeBuffer).toContain('after-bad');
|
||||||
|
});
|
||||||
|
|
||||||
|
// Only 'after-bad' should be in the buffer
|
||||||
|
expect(session.writeBuffer).toHaveLength(1);
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
// ========== Resize validation ==========
|
||||||
|
|
||||||
|
describe('resize validation', () => {
|
||||||
|
it('accepts valid resize within bounds', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'z', c: 120, r: 40 }));
|
||||||
|
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.resize).toHaveBeenCalledWith(120, 40);
|
||||||
|
});
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('accepts resize at minimum bounds (1x1)', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'z', c: 1, r: 1 }));
|
||||||
|
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.resize).toHaveBeenCalledWith(1, 1);
|
||||||
|
});
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('accepts resize at maximum bounds (500x200)', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'z', c: 500, r: 200 }));
|
||||||
|
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.resize).toHaveBeenCalledWith(500, 200);
|
||||||
|
});
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects resize with cols out of bounds (0 cols)', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'z', c: 0, r: 40 }));
|
||||||
|
|
||||||
|
// Send a valid message to confirm processing continues
|
||||||
|
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.writeBuffer).toContain('sentinel');
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(session.resize).not.toHaveBeenCalled();
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects resize with cols exceeding 500', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'z', c: 501, r: 40 }));
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.writeBuffer).toContain('sentinel');
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(session.resize).not.toHaveBeenCalled();
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects resize with rows exceeding 200', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'z', c: 80, r: 201 }));
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.writeBuffer).toContain('sentinel');
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(session.resize).not.toHaveBeenCalled();
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects resize with non-integer values', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'z', c: 80.5, r: 24 }));
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.writeBuffer).toContain('sentinel');
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(session.resize).not.toHaveBeenCalled();
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects resize with negative values', async () => {
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
try {
|
||||||
|
const session = ctx._session;
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'z', c: -1, r: 24 }));
|
||||||
|
|
||||||
|
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
expect(session.writeBuffer).toContain('sentinel');
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(session.resize).not.toHaveBeenCalled();
|
||||||
|
} finally {
|
||||||
|
ws.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
// ========== Connection cleanup ==========
|
||||||
|
|
||||||
|
describe('connection cleanup', () => {
|
||||||
|
it('removes session event listeners on close', async () => {
|
||||||
|
const session = ctx._session;
|
||||||
|
const listenersBefore = session.listenerCount('terminal');
|
||||||
|
|
||||||
|
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
|
||||||
|
|
||||||
|
// A listener was added for 'terminal'
|
||||||
|
expect(session.listenerCount('terminal')).toBe(listenersBefore + 1);
|
||||||
|
expect(session.listenerCount('clearTerminal')).toBeGreaterThanOrEqual(1);
|
||||||
|
expect(session.listenerCount('needsRefresh')).toBeGreaterThanOrEqual(1);
|
||||||
|
|
||||||
|
// Close the WS connection
|
||||||
|
ws.close();
|
||||||
|
await waitForClose(ws);
|
||||||
|
|
||||||
|
// Give the server-side close handler time to run
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 50));
|
||||||
|
|
||||||
|
// Listeners should be cleaned up
|
||||||
|
expect(session.listenerCount('terminal')).toBe(listenersBefore);
|
||||||
|
expect(session.listenerCount('clearTerminal')).toBe(0);
|
||||||
|
expect(session.listenerCount('needsRefresh')).toBe(0);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user