fix(input): recover when the seq counter falls behind the server watermark

Browser input is delivered exactly once by (clientId, seq). The server records a
watermark per clientId and discards anything not above it as a duplicate — but
acknowledged it with an ACK indistinguishable from "applied". The client then
dropped the record from its queue, the UI looked perfectly normal, and the
terminal received nothing at all.

The counter is persisted to localStorage through a debounced write. Kill the page
between "sent" and "persisted" and the restored counter is below the server's
watermark, after which every keystroke lands under it, is discarded, and is
ACKed. Reloading does not help: the clientId is restored from localStorage
alongside that stale counter. Measured on a real session — typing into the same
session from a fresh browser (new clientId, no watermark on the server) worked
perfectly, which is what localised the fault to client state.

Three changes:
- on rejection the server replies {"t":"ia",seq,"dup":true,"last":<watermark>}.
  It still ACKs, so the client can drop the record from its queue, but it now
  says the input was not applied and supplies the number needed to climb out.
- on `dup` the client lifts its counter above the watermark and re-queues.
  ⚠️ Only records whose FIRST delivery is being retried are re-sent: a retry
  judged duplicate means the mechanism is working (the original did arrive), and
  re-sending would type the same text twice.
- the counter is now persisted synchronously. The queue payload can stay
  debounced, but the counter is the thing that has to survive a crash, and
  leaving it on the lossiest path cancels the only guarantee there is.

⚠️ Reading the watermark is defensive: the session arrives through a structured
port, and a port missing that method must not take the whole input path down —
a throw inside the handler means the ACK is never sent and the record is stuck in
the client queue forever, which is worse than the ambiguity being fixed. A mock
port's test timeout is what exposed this.
This commit is contained in:
d fei
2026-09-03 02:01:44 -07:00
parent e6dac66a20
commit 05bb7081cc
6 changed files with 198 additions and 8 deletions
+15
View File
@@ -3300,6 +3300,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;
+40 -5
View File
@@ -2794,7 +2794,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
@@ -2972,7 +2972,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);
}
@@ -3078,9 +3082,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. */
+25 -2
View File
@@ -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' &&
+36
View File
@@ -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);
});
});
+78
View File
@@ -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\(\);/);
});
});
+4 -1
View File
@@ -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();