feat(split): Pane B input goes through the exactly-once queue

TerminalTile used to send every xterm onData chunk as a bare {t:'i'}
frame: no seq, no ACK, silently dropped while its socket was down, and
typing never acknowledged the session's idle alert. Keystrokes and
pastes now go through app._sendInputAsync over the tile's own socket,
registered in the input-socket map while open (HTTP fallback while
not), so they are ACKed, persisted until delivered, redelivered after a
drop, and the ACK clears the idle alert.

What xterm generates on its own stays out of that persisted queue: a
query reply (DA/CPR/OSC) is dropped, as the primary pane drops it, and a
focus or mouse report goes out once via _sendInputEphemeral. The tile's
socket carries the tab identity with a :tile suffix so it can never
evict the primary pane's socket.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Codeman maintainer
2026-10-06 09:41:10 +02:00
parent d1bbb4cc26
commit 5e3dbf2057
2 changed files with 385 additions and 6 deletions
+64 -6
View File
@@ -91,6 +91,10 @@
this._liveQueue = null; this._liveQueue = null;
this._markerOwed = false; this._markerOwed = false;
this._onWheel = null; this._onWheel = null;
// `{ ws, lastRecvAt }`, registered with the app's input-socket map while
// this pane's socket is open, so the exactly-once input queue delivers this
// session's keystrokes over it (app.js _inputSocketFor). Null otherwise.
this._inputHandle = null;
} }
async connect() { async connect() {
@@ -116,11 +120,7 @@
this._installWheelListener(); this._installWheelListener();
this.terminal.onData((data) => { this.terminal.onData((data) => this._onTerminalData(data));
if (this.ws && this.ws.readyState === WebSocket.OPEN) {
this.ws.send(JSON.stringify({ t: 'i', d: data }));
}
});
// Pane B has no gates of its own by default, so every app-level chord // Pane B has no gates of its own by default, so every app-level chord
// that the document capture-phase handler (app.js) only preventDefault()s // that the document capture-phase handler (app.js) only preventDefault()s
@@ -266,15 +266,25 @@
if (this._destroyed) return; if (this._destroyed) return;
const proto = location.protocol === 'https:' ? 'wss:' : 'ws:'; const proto = location.protocol === 'https:' ? 'wss:' : 'ws:';
const url = `${proto}//${location.host}${window.CodemanBase.base}/ws/sessions/${this.sessionId}/terminal`; // The tab's own connection identity plus a `:tile` suffix. The server
// supersedes a socket that reuses a cid on the same session (4010), so a
// pane must never share the primary pane's exact cid: were both ever on
// one session they would evict each other in a loop. Input frames still
// carry the BARE clientId, which is what the server dedups on.
const app = global.app;
const cid = app?._clientId ? `${app._clientId}:${app._wsTabNonce}:tile` : '';
const cidQuery = cid ? `?cid=${encodeURIComponent(cid)}` : '';
const url = `${proto}//${location.host}${window.CodemanBase.base}/ws/sessions/${this.sessionId}/terminal${cidQuery}`;
this.ws = new WebSocket(url); this.ws = new WebSocket(url);
this.ws.onopen = () => { this.ws.onopen = () => {
this._wsReady = true; this._wsReady = true;
this._registerInputSocket();
this._sendResize(); this._sendResize();
}; };
this.ws.onmessage = (event) => { this.ws.onmessage = (event) => {
if (this._inputHandle) this._inputHandle.lastRecvAt = Date.now();
try { try {
const msg = JSON.parse(event.data); const msg = JSON.parse(event.data);
if (msg.t === 'o') { if (msg.t === 'o') {
@@ -287,6 +297,9 @@
// _onSessionNeedsRefresh (app.js:2990) — Pane B has its own // _onSessionNeedsRefresh (app.js:2990) — Pane B has its own
// buffer loader for the same reason connect() does. // buffer loader for the same reason connect() does.
this._refreshBuffer(); this._refreshBuffer();
} else if (msg.t === 'ia') {
// Input ACK. The frame names no session, so it is this pane's.
global.app?._onWsInputAck?.(msg.seq, msg, this.sessionId);
} }
} catch { } catch {
/* Malformed frame — ignore, matches primary pane's tolerance. */ /* Malformed frame — ignore, matches primary pane's tolerance. */
@@ -322,10 +335,51 @@
_onSocketClosed() { _onSocketClosed() {
this._wsReady = false; this._wsReady = false;
this._wsClosed = true; this._wsClosed = true;
this._unregisterInputSocket();
if (this._bufferLoading) this._markerOwed = true; if (this._bufferLoading) this._markerOwed = true;
else this._writeDisconnectedMarker(); else this._writeDisconnectedMarker();
} }
// Keystrokes and pastes go through the app's exactly-once input queue (seq,
// ACK, persisted until delivered, redelivered after a drop), over this
// pane's own socket while it is open and the HTTP fallback while it is not.
// What xterm GENERATES must not be queued: a query reply (DA/CPR/OSC) is
// dropped, as the primary pane drops it, because forwarding it types
// "0;276;0c" into the CLI, and replaying one after a reload would do so
// again into a later screen. A focus or mouse report is real input the
// program asked for, but nobody typed it: it goes out once, never
// persisted. Same predicates as the primary pane (terminal-ui.js onData).
_onTerminalData(data) {
const input = global.CodemanTerminalInput;
if (input?.shouldSuppressTerminalQueryResponse?.(data)) return;
const app = global.app;
if (input?.isTerminalFocusOrMouseReport?.(data)) {
app?._sendInputEphemeral?.(this.sessionId, data);
return;
}
app?._sendInputAsync?.(this.sessionId, data);
}
// Joins the app's input-socket map for this session and flushes anything
// already queued for it (typed while the socket was down, or left over from
// a reload) over the fresh socket. Called from onopen.
_registerInputSocket() {
const app = global.app;
if (!this.ws || !app?._registerInputSocket) return;
this._unregisterInputSocket();
this._inputHandle = { ws: this.ws, lastRecvAt: 0 };
app._registerInputSocket(this.sessionId, this._inputHandle);
app._onWsReady?.(this.sessionId);
}
// Leaves the map; only this pane's own handle is removed (a replacement
// socket's registration survives a late close of the old one).
_unregisterInputSocket() {
if (!this._inputHandle) return;
global.app?._unregisterInputSocket?.(this.sessionId, this._inputHandle);
this._inputHandle = null;
}
// Settles a marker the pane owes: set when a close lands during a load (the // Settles a marker the pane owes: set when a close lands during a load (the
// replay would otherwise sit below it) or when a load wipes the terminal on // replay would otherwise sit below it) or when a load wipes the terminal on
// a closed socket. Called from each load's own finally, just before // a closed socket. Called from each load's own finally, just before
@@ -615,6 +669,10 @@
destroy() { destroy() {
this._destroyed = true; this._destroyed = true;
// Anything still queued for this session stays in the app's queue and is
// delivered over HTTP by the redelivery sweep, so closing the pane mid-
// keystroke loses nothing.
this._unregisterInputSocket();
if (this._onWheel) { if (this._onWheel) {
this.mountEl?.removeEventListener('wheel', this._onWheel, { capture: true }); this.mountEl?.removeEventListener('wheel', this._onWheel, { capture: true });
this._onWheel = null; this._onWheel = null;
+321
View File
@@ -0,0 +1,321 @@
/**
* @fileoverview TerminalTile input rides the app's exactly-once queue, and only
* what a human typed is queued.
*
* Before, the split pane's second terminal sent every xterm `onData` chunk as a
* bare `{t:'i', d}` frame: no seq, no ACK, dropped while the socket was down,
* never acknowledged an idle alert. It now goes through `app._sendInputAsync`,
* over the pane's own socket registered in the app's input-socket map. That
* queue PERSISTS and REDELIVERS, so it must never hold what xterm generates on
* its own: a query reply (DA/CPR/OSC) is dropped, exactly as the primary pane
* drops it, and a focus or mouse report goes out once, ephemeral.
*
* Real code under test: constants.js + app.js (the queue) + terminal-ui.js (the
* shared input predicates) + terminal-tile.js, in one `vm` context. xterm, the
* fit addon and WebSocket are fakes; `connect()` runs for real.
*/
import { readFileSync } from 'node:fs';
import { performance } from 'node:perf_hooks';
import { resolve } from 'node:path';
import vm from 'node:vm';
import { beforeEach, describe, expect, it, vi } from 'vitest';
type Frame = { t: string; d?: string; seq?: number; cid?: string; c?: number; r?: number };
class FakeSocket {
static OPEN = 1;
static instances: FakeSocket[] = [];
readyState = 0;
sent: Frame[] = [];
onopen: (() => void) | null = null;
onmessage: ((ev: { data: string }) => void) | null = null;
onclose: ((ev?: { code: number }) => void) | null = null;
onerror: (() => void) | null = null;
constructor(public url: string) {
FakeSocket.instances.push(this);
}
send(data: string) {
this.sent.push(JSON.parse(data) as Frame);
}
close = vi.fn(() => {
this.readyState = 3;
});
open() {
this.readyState = 1;
this.onopen?.();
}
receive(msg: object) {
this.onmessage?.({ data: JSON.stringify(msg) });
}
inputFrames() {
return this.sent.filter((f) => f.t === 'i');
}
}
class FakeTerminal {
static last: FakeTerminal | null = null;
options: Record<string, unknown>;
cols = 80;
rows = 24;
dataCb: ((data: string) => void) | null = null;
buffer = { active: { type: 'normal', viewportY: 0, length: 24 } };
constructor(options: Record<string, unknown>) {
this.options = { ...options };
FakeTerminal.last = this;
}
loadAddon() {}
open() {}
onData(cb: (data: string) => void) {
this.dataCb = cb;
}
attachCustomKeyEventHandler() {}
registerLinkProvider() {}
write(_data: string, cb?: () => void) {
cb?.();
}
clear() {}
resize(cols: number, rows: number) {
this.cols = cols;
this.rows = rows;
}
dispose() {}
type(data: string) {
this.dataCb?.(data);
}
}
const fetchMock = vi.fn();
function loadContext() {
const read = (f: string) => readFileSync(resolve(import.meta.dirname, `../src/web/public/${f}`), 'utf8');
const windowStub: Record<string, unknown> = {
addEventListener: vi.fn(),
removeEventListener: vi.fn(),
CodemanBase: { base: '' },
};
const context = vm.createContext({
console: { ...console, log: vi.fn(), debug: vi.fn() },
performance,
setInterval: vi.fn(),
clearInterval: vi.fn(),
setTimeout,
clearTimeout,
requestAnimationFrame: vi.fn(),
HTMLCanvasElement: class HTMLCanvasElement {},
WebSocket: FakeSocket,
Terminal: FakeTerminal,
FitAddon: {
FitAddon: class {
fit() {}
proposeDimensions() {
return { cols: 80, rows: 24 };
}
},
},
fetch: (...args: unknown[]) => fetchMock(...args),
location: { protocol: 'http:', host: 'codeman.test' },
document: { addEventListener: vi.fn(), documentElement: { dataset: {} } },
localStorage: { length: 0, key: vi.fn(), getItem: vi.fn(), setItem: vi.fn(), removeItem: vi.fn() },
window: windowStub,
MobileDetection: {
isTouchDevice: () => false,
isHandheldDevice: () => false,
getDeviceType: () => 'desktop',
},
});
vm.runInContext(
`${read('constants.js')}\n${read('app.js')}\n${read('terminal-ui.js')}\n${read('terminal-tile.js')}\n` +
'globalThis.__CodemanApp = CodemanApp;',
context
);
return {
CodemanApp: (context as unknown as { __CodemanApp: { prototype: object } }).__CodemanApp,
windowStub,
};
}
const { CodemanApp, windowStub } = loadContext();
type App = Record<string, unknown> & {
_pendingDeliveries: Map<string, Array<{ seq: number; data: string }>>;
markIdleAlertSeen: ReturnType<typeof vi.fn>;
};
function makeApp(): App {
const app = Object.create(CodemanApp.prototype) as App;
app._clientId = 'c-test';
app._wsTabNonce = 'nonce-1';
app._seqCounters = new Map();
app._pendingDeliveries = new Map();
app._postDraining = new Set();
app._extraInputSockets = new Map();
app._persistReliableState = vi.fn();
app._persistReliableNow = vi.fn();
app._updateConnectionIndicator = vi.fn();
app.markIdleAlertSeen = vi.fn();
app.showToast = vi.fn();
app._ws = null;
app._wsSessionId = null;
app._estimateReplayRows = (text: string) => text.split('\n').length;
return app;
}
type Tile = {
connect(): Promise<void>;
destroy(): void;
ws: FakeSocket | null;
};
const TerminalTile = windowStub.TerminalTile as new (id: string, mount: unknown, opts?: object) => Tile;
async function connectTile(app: App) {
windowStub.app = app;
const tile = new TerminalTile('s-tile', { addEventListener: vi.fn(), removeEventListener: vi.fn() }, { mode: 'claude' });
await tile.connect();
const ws = FakeSocket.instances.at(-1)!;
return { tile, ws, term: FakeTerminal.last! };
}
beforeEach(() => {
FakeSocket.instances = [];
fetchMock.mockReset();
fetchMock.mockImplementation(async () => ({
ok: true,
status: 200,
json: async () => ({ data: { terminalBuffer: '' } }),
}));
});
describe('TerminalTile socket identity', () => {
it('connects with the tab identity plus a :tile suffix, never the primary pane cid', async () => {
const { ws } = await connectTile(makeApp());
const cid = new URL(ws.url).searchParams.get('cid');
expect(cid).toBe('c-test:nonce-1:tile');
});
});
describe('TerminalTile input through the exactly-once queue', () => {
it('sends typed input as seq-tagged frames with the bare clientId once the socket opens', async () => {
const app = makeApp();
const { ws, term } = await connectTile(app);
ws.open();
term.type('h');
term.type('i');
expect(ws.inputFrames().map((f) => [f.d, f.seq, f.cid])).toEqual([
['h', 1, 'c-test'],
['i', 2, 'c-test'],
]);
});
it('drops a record on its ACK and acknowledges the session idle alert', async () => {
const app = makeApp();
const { ws, term } = await connectTile(app);
ws.open();
term.type('x');
expect(app._pendingDeliveries.get('s-tile')).toHaveLength(1);
ws.receive({ t: 'ia', seq: 1 });
expect(app._pendingDeliveries.get('s-tile')).toBeUndefined();
expect(app.markIdleAlertSeen).toHaveBeenCalledWith('s-tile');
});
it('flushes input typed before the socket opened, in order, once it does', async () => {
const app = makeApp();
// POSTs fail, as they would while the server restarts, so the input waits.
const { ws, term } = await connectTile(app);
fetchMock.mockImplementation(async () => ({ ok: false, status: 503, json: async () => ({}) }));
term.type('a');
term.type('b');
await new Promise((r) => setTimeout(r, 0));
expect(ws.inputFrames()).toEqual([]);
ws.open();
expect(ws.inputFrames().map((f) => [f.d, f.seq])).toEqual([
['a', 1],
['b', 2],
]);
});
it('keeps queuing after the socket closes (HTTP fallback), instead of dropping keystrokes', async () => {
const app = makeApp();
const { ws, term } = await connectTile(app);
ws.open();
ws.readyState = 3;
ws.onclose?.({ code: 1006 });
const posts: Array<{ input: string; seq: number }> = [];
fetchMock.mockImplementation(async (_url: string, init?: { body?: string }) => {
if (init?.body) posts.push(JSON.parse(init.body));
return { ok: true, status: 200, json: async () => ({}) };
});
term.type('z');
await new Promise((r) => setTimeout(r, 0));
expect(posts.map((p) => [p.input, p.seq])).toEqual([['z', 1]]);
expect(ws.inputFrames()).toEqual([]);
});
});
describe('what xterm generates never enters the durable queue', () => {
it('drops a DA query reply entirely', async () => {
const app = makeApp();
const { ws, term } = await connectTile(app);
ws.open();
term.type('\x1b[?1;2c');
expect(ws.inputFrames()).toEqual([]);
expect(app._pendingDeliveries.get('s-tile')).toBeUndefined();
});
it('sends a mouse report once, without a seq, and never persists it', async () => {
const app = makeApp();
const { ws, term } = await connectTile(app);
ws.open();
term.type('\x1b[<0;10;5M');
expect(ws.inputFrames()).toEqual([{ t: 'i', d: '\x1b[<0;10;5M' }]);
expect(app._pendingDeliveries.get('s-tile')).toBeUndefined();
});
it('sends a focus report once, without a seq', async () => {
const app = makeApp();
const { ws, term } = await connectTile(app);
ws.open();
term.type('\x1b[I');
expect(ws.inputFrames()).toEqual([{ t: 'i', d: '\x1b[I' }]);
expect(app._pendingDeliveries.get('s-tile')).toBeUndefined();
});
});
describe('TerminalTile leaves the input-socket map', () => {
it('on close, so a later keystroke cannot be written into a dead socket', async () => {
const app = makeApp();
const { ws } = await connectTile(app);
ws.open();
expect((app._extraInputSockets as Map<string, unknown>).has('s-tile')).toBe(true);
ws.onclose?.({ code: 1006 });
expect((app._extraInputSockets as Map<string, unknown>).has('s-tile')).toBe(false);
});
it('on destroy, while pending input stays queued for the HTTP sweep', async () => {
const app = makeApp();
const { tile, ws, term } = await connectTile(app);
ws.open();
term.type('q');
tile.destroy();
expect((app._extraInputSockets as Map<string, unknown>).has('s-tile')).toBe(false);
expect(app._pendingDeliveries.get('s-tile')?.map((r) => r.data)).toEqual(['q']);
});
});