diff --git a/src/session.ts b/src/session.ts index b99d403e..5486b722 100644 --- a/src/session.ts +++ b/src/session.ts @@ -3572,6 +3572,21 @@ export class Session extends EventEmitter { * half-open socket silently drops frames with no error) would type a prompt * twice whenever an ACK is lost after the write landed. */ + /** + * The highest input seq recorded for `clientId`, or 0 when this session has + * never seen it. + * + * Reported back on a REJECTED (duplicate) frame so the client can lift its own + * counter above this watermark. Without that number a client whose persisted + * counter fell behind ours has no way to find its way out: every fresh + * keystroke it sends lands at or below the watermark, is dropped as a + * duplicate, and is ACKed anyway — so the UI looks healthy while nothing is + * delivered, and a reload restores the same stale counter from localStorage. + */ + lastInputSeq(clientId: string): number { + return this._appliedInputSeq.get(clientId) ?? 0; + } + shouldApplyInput(clientId: string, seq: number): boolean { const last = this._appliedInputSeq.get(clientId); if (last !== undefined && seq <= last) return false; diff --git a/src/web/public/app.js b/src/web/public/app.js index f91f748b..c1fdc2e0 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -2895,7 +2895,7 @@ class CodemanApp { } else if (msg.t === 'ia') { // Input ACK — the server applied (or deduped) this seq; drop it from // the durable queue so it can never be re-delivered/lost. - this._onWsInputAck(msg.seq); + this._onWsInputAck(msg.seq, msg); } } catch { // Ignore malformed messages @@ -3073,7 +3073,11 @@ class CodemanApp { this._pendingDeliveries.set(sessionId, list); } list.push(rec); - this._persistReliableState(); + // ⚠️ SYNCHRONOUS, not the debounced writer: the seq counter is precisely the + // thing that must survive a crash, and a debounce puts it on the path most + // likely to be lost. A counter that comes back BELOW the server's watermark + // makes every later keystroke a silently-dropped duplicate (see _onWsInputAck). + this._persistReliableNow(); this._updateConnectionIndicator(); this._drainSession(sessionId); } @@ -3179,9 +3183,40 @@ class CodemanApp { this.markIdleAlertSeen?.(sessionId); } - /** Server input-ACK frame ({t:'ia',seq}) over the WebSocket. */ - _onWsInputAck(seq) { - if (this._wsSessionId && Number.isInteger(seq)) this._ackDelivery(this._wsSessionId, seq); + /** + * Server input-ACK frame ({t:'ia',seq}) over the WebSocket. + * + * `dup:true` means the server REJECTED the frame as already-seen rather than + * applying it, and `last` is its watermark for this clientId. That combination + * is the escape hatch from a rolled-back counter: our seqs persist on a + * debounced write, so a tab killed between a send and that write comes back + * counting from BELOW the server's watermark, and from then on every keystroke + * is dropped-but-ACKed — a silently dead terminal that a reload cannot fix, + * because the stale counter is restored from localStorage too. + * + * ⚠️ Only a FIRST-attempt record is re-queued. A retry (`tries > 1`) being + * called a duplicate is the mechanism working as designed — the original did + * land — and re-sending it would type the same thing twice. + */ + _onWsInputAck(seq, msg) { + const sessionId = this._wsSessionId; + if (!sessionId || !Number.isInteger(seq)) return; + if (msg && msg.dup) { + const list = this._pendingDeliveries.get(sessionId); + const rec = list && list.find((r) => r.seq === seq); + const watermark = Number.isInteger(msg.last) ? msg.last : seq; + // Lift the counter clear of the server's watermark before anything else, so + // the re-queue below (and every later keystroke) gets an acceptable seq. + if ((this._seqCounters.get(sessionId) || 0) <= watermark) { + this._seqCounters.set(sessionId, watermark); + this._persistReliableNow(); + } + const lost = rec && rec.tries <= 1 ? rec.data : null; + this._ackDelivery(sessionId, seq); + if (lost !== null) this._reliableSend(sessionId, lost, rec.useMux); + return; + } + this._ackDelivery(sessionId, seq); } /** Called from ws.onopen — flush everything pending over the fresh socket. */ diff --git a/src/web/routes/ws-routes.ts b/src/web/routes/ws-routes.ts index fc6e4c61..ed60f5d3 100644 --- a/src/web/routes/ws-routes.ts +++ b/src/web/routes/ws-routes.ts @@ -192,8 +192,31 @@ export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort, getHost // a duplicate: the input was lost for good. if (!delivered && cid && seq !== null) session.forgetInputSeq(cid, seq); } - if (delivered && seq !== null && socket.readyState === 1) { - socket.send(`{"t":"ia","seq":${seq}}`); + if (seq !== null && socket.readyState === 1) { + if (apply) { + if (delivered) socket.send(`{"t":"ia","seq":${seq}}`); + } else { + // REJECTED as a duplicate. ACK it — the client must still drop it + // from its durable queue — but say so, and hand back our watermark. + // + // A plain ACK here is indistinguishable from "applied", which is + // what made a client with a rolled-back counter unrecoverable: its + // seqs persist to localStorage on a DEBOUNCED write, so a tab killed + // between a send and that write comes back with a counter BELOW this + // watermark, every later keystroke lands at or under it, and each one + // is dropped-but-ACKed. The UI stays clean, nothing is delivered, and + // a reload restores the same stale counter. `last` is what lets the + // client lift itself out. + // ⚠️ Defensive: the session arrives through a structural port, and an + // implementation without this method must not take the whole input + // path down with it — a throw here aborts the message handler and the + // frame is never ACKed at all, which strands it in the client's queue. + const watermark = + typeof (session as { lastInputSeq?: (c: string) => number }).lastInputSeq === 'function' + ? (session as { lastInputSeq: (c: string) => number }).lastInputSeq(cid as string) + : seq; + socket.send(`{"t":"ia","seq":${seq},"dup":true,"last":${watermark}}`); + } } } else if ( msg.t === 'z' && diff --git a/test/reliable-input-dedup.test.ts b/test/reliable-input-dedup.test.ts index d77bef2d..f65f3c84 100644 --- a/test/reliable-input-dedup.test.ts +++ b/test/reliable-input-dedup.test.ts @@ -71,3 +71,39 @@ describe('Session.shouldApplyInput (exactly-once input dedup)', () => { expect(s.shouldApplyInput('recent', 2)).toBe(true); }); }); + +describe('Session.lastInputSeq (the watermark a stuck client needs)', () => { + it('reports 0 for a client it has never seen', () => { + expect(makeSession().lastInputSeq('c-new')).toBe(0); + }); + + it('reports the highest seq applied for that client', () => { + const s = makeSession(); + s.shouldApplyInput('c-1', 7); + expect(s.lastInputSeq('c-1')).toBe(7); + }); + + it('is what a rolled-back client must clear to be heard again', () => { + // The failure this exists for: the client's seq counter persists on a + // DEBOUNCED write, so a tab killed between a send and that write comes back + // counting from below the watermark. Every later keystroke then lands at or + // under it and is rejected — silently, because a rejected frame is ACKed too. + const s = makeSession(); + for (let i = 1; i <= 40; i++) s.shouldApplyInput('c-1', i); + // Restored counter starts over at 1: dropped, and every subsequent one too. + expect(s.shouldApplyInput('c-1', 1)).toBe(false); + expect(s.shouldApplyInput('c-1', 2)).toBe(false); + // The watermark it is handed back is exactly what makes it recoverable. + const watermark = s.lastInputSeq('c-1'); + expect(watermark).toBe(40); + expect(s.shouldApplyInput('c-1', watermark + 1)).toBe(true); + }); + + it('does not resurrect a seq that forgetInputSeq rolled back', () => { + const s = makeSession(); + s.shouldApplyInput('c-1', 5); + s.forgetInputSeq('c-1', 5); + expect(s.lastInputSeq('c-1')).toBe(4); + expect(s.shouldApplyInput('c-1', 5)).toBe(true); + }); +}); diff --git a/test/reliable-input-recovery.test.ts b/test/reliable-input-recovery.test.ts new file mode 100644 index 00000000..fe56168a --- /dev/null +++ b/test/reliable-input-recovery.test.ts @@ -0,0 +1,78 @@ +/** + * @fileoverview A client whose seq counter rolled back must heal itself. + * + * The failure: the browser tags input with (clientId, seq) and persists the + * counters to localStorage on a DEBOUNCED write. A tab killed between a send and + * that write comes back counting from BELOW the server's watermark, so every + * later keystroke is rejected as a duplicate — and, because a rejected frame was + * ACKed exactly like an applied one, the client dropped it from its queue and the + * UI looked perfectly healthy while the terminal took no input at all. Reloading + * could not help: clientId and the stale counter both come back from localStorage. + * + * Observed live on a server session: a fresh browser (new clientId, no watermark) + * typed into the same session fine, which is what isolated it to client state. + */ +import { readFileSync } from 'node:fs'; +import { resolve } from 'node:path'; +import { describe, it, expect } from 'vitest'; + +const appSource = readFileSync(resolve(import.meta.dirname, '../src/web/public/app.js'), 'utf8'); +const wsSource = readFileSync(resolve(import.meta.dirname, '../src/web/routes/ws-routes.ts'), 'utf8'); + +describe('the duplicate ACK carries what the client needs', () => { + it('marks a rejected frame as dup and reports the watermark', () => { + // A bare ACK is indistinguishable from "applied" — that ambiguity is the bug. + expect(wsSource).toMatch(/"dup":true,"last":\$\{watermark\}/); + expect(wsSource).toContain('lastInputSeq'); + }); + + it('reads the watermark defensively, so a port without it still ACKs', () => { + // The session arrives through a structural port. A throw inside the message + // handler aborts it before the ACK is sent, stranding the frame in the + // client's durable queue — which is worse than the ambiguity being fixed here. + expect(wsSource).toMatch(/typeof \(session as \{ lastInputSeq\?/); + }); + + it('still ACKs a rejected frame, so the client can drop it from its queue', () => { + // Silence would strand the record and the redelivery sweep would spin on it. + const block = wsSource.slice(wsSource.indexOf('if (seq !== null && socket.readyState === 1)')); + expect(block.slice(0, 1200)).toContain('"t":"ia"'); + }); +}); + +describe('the client lifts itself over the watermark', () => { + const handler = appSource.slice( + appSource.indexOf('_onWsInputAck(seq, msg)'), + appSource.indexOf('/** Called from ws.onopen') + ); + + it('raises the counter to the watermark it was handed', () => { + expect(handler).toMatch(/_seqCounters\.set\(sessionId, watermark\)/); + }); + + it('re-queues a FIRST-attempt frame, whose input was genuinely lost', () => { + expect(handler).toMatch(/rec\.tries <= 1/); + expect(handler).toMatch(/this\._reliableSend\(sessionId, lost/); + }); + + it('does NOT re-queue a retry, which the dedup correctly suppressed', () => { + // A retry called a duplicate means the original DID land; re-sending it would + // type the same thing twice — the exact thing exactly-once delivery prevents. + expect(handler).toMatch(/const lost = rec && rec\.tries <= 1 \? rec\.data : null;/); + }); + + it('persists the raised counter immediately, not on the debounce', () => { + expect(handler).toContain('this._persistReliableNow()'); + }); +}); + +describe('the seq counter is persisted synchronously on every send', () => { + it('_reliableSend uses the immediate writer, never the debounced one', () => { + // The counter is precisely what must survive a crash, so it cannot ride the + // path most likely to be lost. (The queue PAYLOAD may still be debounced.) + const send = appSource.slice(appSource.indexOf('_reliableSend(sessionId, data, useMux)')); + const body = send.slice(0, send.indexOf('_nextSeq(sessionId) {')); + expect(body).toContain('this._persistReliableNow();'); + expect(body).not.toMatch(/list\.push\(rec\);\s*\n\s*this\._persistReliableState\(\);/); + }); +}); diff --git a/test/routes/ws-routes.test.ts b/test/routes/ws-routes.test.ts index 76efc37d..0d7c3bd0 100644 --- a/test/routes/ws-routes.test.ts +++ b/test/routes/ws-routes.test.ts @@ -250,7 +250,10 @@ describe('ws-routes', () => { ws.send(JSON.stringify({ t: 'i', d: 'again\r', cid: 'c1', seq: 7 })); - expect(await nextMessage(ws)).toEqual({ t: 'ia', seq: 7 }); + // The ACK now SAYS it was a duplicate and hands back the watermark: a bare + // ACK is indistinguishable from "applied", and that ambiguity left a client + // whose seq counter had rolled back silently unable to type at all. + expect(await nextMessage(ws)).toEqual({ t: 'ia', seq: 7, dup: true, last: 7 }); expect(session.writeBuffer).not.toContain('again\r'); } finally { ws.close();