From 1255e28f6fb698a149ca17a05ec36562b28594d4 Mon Sep 17 00:00:00 2001 From: Codeman maintainer Date: Fri, 19 Jun 2026 16:44:18 +0200 Subject: [PATCH] fix(input): durable exactly-once input delivery so a dropped link can't lose a prompt MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A "sent" prompt could vanish with no trace on a flaky connection (e.g. a train): with local echo on, Enter cleared the overlay then sent over the WebSocket fire-and-forget. On a half-open socket (readyState===OPEN, dead TCP) ws.send() doesn't throw, so the frame was silently discarded, nothing was enqueued, and navigator.onLine stayed true — the prompt was lost and never resent. Replace the best-effort offline queue with a durable, acknowledged delivery layer: - Client (app.js): every input frame is recorded with a stable clientId + monotonic per-session seq and persisted to localStorage BEFORE delivery, and only dropped on a server ACK. Delivered over WS (acked via {t:'ia',seq}) or, when the socket is down, POST in seq order (HTTP 2xx = ACK). A 2s sweep force-reconnects a WS whose oldest frame is unacked past 4s (half-open sockets never recover on their own); on reconnect/reload all pending frames re-deliver. Survives reconnects AND page reloads. Connection indicator shows pending count. - Server: Session.shouldApplyInput(clientId, seq) applies each frame exactly once (bounded MRU map); ws-routes + POST /input dedup a redelivered seq but still ACK it (200 / {t:'ia'}), so an at-least-once resend can never type the prompt twice. Untagged input (curl/legacy) applies unconditionally — no behavior change. - terminal-ui.js sendInput() (voice / keyboard-accessory / paste) now routes through the same durable layer. Tests: test/reliable-input-dedup.test.ts (exactly-once semantics on the real Session) + POST /input dedup route tests. Design: docs/reliable-input-delivery.md. Co-Authored-By: Claude Opus 4.8 (1M context) --- docs/reliable-input-delivery.md | 72 ++++++ src/session.ts | 36 +++ src/web/public/app.js | 373 +++++++++++++++++++++++------ src/web/public/terminal-ui.js | 11 +- src/web/routes/session-routes.ts | 9 +- src/web/routes/ws-routes.ts | 25 +- src/web/schemas.ts | 9 + test/mocks/mock-session.ts | 10 + test/reliable-input-dedup.test.ts | 73 ++++++ test/routes/session-routes.test.ts | 37 +++ 10 files changed, 567 insertions(+), 88 deletions(-) create mode 100644 docs/reliable-input-delivery.md create mode 100644 test/reliable-input-dedup.test.ts diff --git a/docs/reliable-input-delivery.md b/docs/reliable-input-delivery.md new file mode 100644 index 00000000..accc8032 --- /dev/null +++ b/docs/reliable-input-delivery.md @@ -0,0 +1,72 @@ +# Reliable input delivery (exactly-once, durable) + +## The bug this fixes + +With local echo on, pressing Enter cleared the overlay and then sent the prompt +over the WebSocket **fire-and-forget** (`ws.send({t:'i',d})`). On a flaky link +(e.g. a moving train) the socket is frequently *half-open*: `readyState === OPEN` +so `ws.send()` does **not** throw, but the underlying TCP is dead, so the frame is +silently discarded. Nothing was enqueued (the send "succeeded"), the on-screen +prompt was already wiped, and `navigator.onLine` stays `true` — so a long typed +prompt vanished with no trace and no resend. + +## The guarantee + +Every byte of user input is **recorded durably before delivery** and **only +dropped once the server ACKs it** — so a half-open socket, a reconnect, or a page +reload can never lose input. Redelivery is **exactly-once**: the server applies +each `(clientId, seq)` at most once, so a resend can't type the prompt twice. + +## How it works + +### Client (`app.js`) + +- A stable **`clientId`** (`localStorage['codeman:clientId']`) identifies this + browser to the server's dedup across reconnects and reloads. +- Each input frame gets a **monotonic per-session `seq`**. Frame records + (`{seq,data,useMux,ts,tries,sentAt}`) live in `_pendingDeliveries` + (`Map`), persisted (debounced, + flushed on `pagehide`/ + `visibilitychange`) to `localStorage['codeman:pendingInput']`. The seq counters + persist too, so seqs stay monotonic across reloads (never reset — a reset would + let the server treat fresh input as an already-applied duplicate). +- **Delivery** (`_drainSession`): + - **WS path** — when the socket is `OPEN` for the session, send each not-yet-sent + record (`sentAt === 0`) in seq order over the single ordered stream. Records + stay pending until the server's `{t:'ia',seq}` ACK removes them. + - **POST path** — when no WS, POST records in order, awaiting each (the HTTP 2xx + *is* the ACK). A 404/410 (session gone) drops the record rather than retry + forever. +- **Half-open recovery** (`_redeliverSweep`, every 2s): if the active WS session's + oldest record is unacked past `_reliableAckTimeoutMs` (4s), the socket is assumed + dead — `ws.close()` forces a fast reconnect; `onopen` (`_onWsReady`) resets + `sentAt = 0` and re-sends everything pending. Also re-drains background sessions + over POST, and fires on SSE-reconnect / `online`. +- The connection indicator shows pending count/bytes (`_pendingBytes`). + +### Server + +- **`Session.shouldApplyInput(clientId, seq)`** — returns `true` exactly once per + `(clientId, seq)`: the first time a seq strictly greater than that client's + last-applied is seen. A replayed/lower seq returns `false`. Bounded MRU map + (`MAX_INPUT_DEDUP_CLIENTS = 256`). +- **WS route** (`ws-routes.ts`) — parses optional `cid`/`seq` on `{t:'i'}`; applies + via `shouldApplyInput` (skips a duplicate, still ACKs with `{t:'ia',seq}` so the + client drops it). Untagged frames apply unconditionally (no behavior change). +- **POST route** (`/api/sessions/:id/input`) — optional `seq`/`clientId` in + `SessionInputWithLimitSchema`; a deduped duplicate returns 200 without writing + (the 200 is the client's ACK). `curl`/legacy callers omit the fields and always + apply. + +## Known limitation + +Dedup state is in-memory on the server. A **server restart** between a write and +the client's redelivery of that same seq could re-apply it (a rare duplicate). +This is a deliberate trade-off: favor *never losing input* over a rare duplicate +across the narrow restart window. + +## Tests + +- `test/reliable-input-dedup.test.ts` — `Session.shouldApplyInput` exactly-once + semantics (monotonic, per-client, gap-tolerant, eviction-safe). +- `test/routes/session-routes.test.ts` — POST `/input` applies a tagged + `(clientId, seq)` once on redelivery; untagged input always applies. diff --git a/src/session.ts b/src/session.ts index 6b115bc3..a6281069 100644 --- a/src/session.ts +++ b/src/session.ts @@ -2213,6 +2213,42 @@ export class Session extends EventEmitter { } } + /** + * Per-client highest-applied input sequence, for exactly-once input delivery. + * Keyed by the web client's stable `clientId`. Bounded so many devices over a + * long-lived session can't grow it without limit (insertion order = MRU, so + * eviction drops the least-recently-active client). + */ + private _appliedInputSeq = new Map(); + private static readonly MAX_INPUT_DEDUP_CLIENTS = 256; + + /** + * Decide whether an input frame should be applied to the PTY or skipped as a + * duplicate redelivery. Returns true exactly once per (clientId, seq): the + * first time a seq strictly greater than the client's last-applied is seen. + * A redelivery of an already-applied seq (the client never got our ACK and + * resent) returns false. Callers should ACK regardless — a duplicate is, from + * the client's view, "delivered" — and only `write()` the PTY when this is + * true. Relies on the client delivering one client's frames in seq order over + * a single ordered stream, so `seq <= last` ⇒ already applied. + * + * Without this, the client's at-least-once redelivery (needed because a + * half-open socket silently drops frames with no error) would type a prompt + * twice whenever an ACK is lost after the write landed. + */ + shouldApplyInput(clientId: string, seq: number): boolean { + const last = this._appliedInputSeq.get(clientId); + if (last !== undefined && seq <= last) return false; + // Re-insert to move this client to the MRU end for fair eviction. + if (last !== undefined) this._appliedInputSeq.delete(clientId); + this._appliedInputSeq.set(clientId, seq); + if (this._appliedInputSeq.size > Session.MAX_INPUT_DEDUP_CLIENTS) { + const oldest = this._appliedInputSeq.keys().next().value; + if (oldest !== undefined) this._appliedInputSeq.delete(oldest); + } + return true; + } + /** * Sends input via the terminal multiplexer's direct input mechanism. * diff --git a/src/web/public/app.js b/src/web/public/app.js index e5153a25..59a79a92 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -468,13 +468,28 @@ class CodemanApp { this.maxReconnectAttempts = 10; this.isOnline = navigator.onLine; - // Offline input queue - this._inputQueue = new Map(); // Map - this._inputQueueMaxBytes = 64 * 1024; // 64KB cap per session + // Reliable, durable input delivery (replaces the old best-effort queue). + // Every input byte is recorded with a stable clientId + a monotonic + // per-session seq, persisted to localStorage, and only dropped once the + // server ACKs that exact seq — so a half-open socket silently dropping a + // frame, a reconnect, or a page reload can never lose a typed prompt. + // Exactly-once: the server applies each (clientId, seq) at most once. this._connectionStatus = 'connected'; - - // Sequential input send chain — ensures keystroke ordering across async fetches - this._inputSendChain = Promise.resolve(); + this._clientId = ''; + this._seqCounters = new Map(); // sessionId -> last issued seq + this._pendingDeliveries = new Map(); // sessionId -> [{seq,data,useMux,ts,tries,sentAt}] + this._postDraining = new Set(); // sessionIds with an in-flight POST drainer + this._persistReliableTimer = null; + this._reliableAckTimeoutMs = 4000; // unacked WS frame older than this ⇒ socket likely dead + this._reliableMaxBytes = 256 * 1024; // cap on the persisted backlog + this._loadReliableState(); + this._reliableSweepTimer = setInterval(() => this._redeliverSweep(), 2000); + // Flush the durable queue synchronously when the page is hidden/closed — + // debounced persistence may have a pending write we mustn't lose on reload. + window.addEventListener('pagehide', () => this._persistReliableNow()); + document.addEventListener('visibilitychange', () => { + if (document.visibilityState === 'hidden') this._persistReliableNow(); + }); // Local echo overlay — DOM overlay positioned at the visible ❯ prompt // (not at buffer.cursorY, which reflects Ink's internal cursor position) @@ -1931,8 +1946,10 @@ class CodemanApp { setConnectionStatus(status) { this._connectionStatus = status; this._updateConnectionIndicator(); - if (status === 'connected' && this._inputQueue.size > 0) { - this._drainInputQueues(); + if (status === 'connected') { + // Reconnected (SSE) — push any durably-queued input out immediately + // instead of waiting for the next 2s sweep. + this._redeliverSweep(); } } @@ -1965,6 +1982,9 @@ class CodemanApp { // went over HTTP, which never claims (see ws-routes sizingToken). this.sendResize(sessionId)?.catch?.(() => {}); this._startMobileResizeRetry(sessionId); + // Flush any durably-queued input over the fresh socket (covers frames a + // prior half-open socket silently dropped, and input typed while offline). + this._onWsReady(sessionId); } }; @@ -1979,6 +1999,10 @@ class CodemanApp { this._onSessionClearTerminal({ id: sessionId }); } else if (msg.t === 'r') { this._onSessionNeedsRefresh({ id: sessionId }); + } 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); } } catch { // Ignore malformed messages @@ -2062,79 +2086,270 @@ class CodemanApp { } /** - * Send input to server without blocking the keystroke flush cycle. - * Uses a sequential promise chain to preserve character ordering - * across concurrent async fetches. + * Public input entry point — name/signature kept for all call sites. + * Records the input durably, then delivers it reliably (exactly-once). Never + * blocks the keystroke flush; never silently drops on a half-open socket. + * @param {string} sessionId + * @param {string} input + * @param {{useMux?: boolean}} [opts] - useMux only affects the POST fallback. */ - _sendInputAsync(sessionId, input) { - // Queue immediately if offline - if (!this.isOnline || this._connectionStatus === 'disconnected') { - this._enqueueInput(sessionId, input); + _sendInputAsync(sessionId, input, opts) { + if (!sessionId || !input) return; + this._reliableSend(sessionId, input, opts?.useMux === true); + } + + /** Record one input frame and kick delivery. The record lives until ACKed. */ + _reliableSend(sessionId, data, useMux) { + const seq = this._nextSeq(sessionId); + const rec = { seq, data, useMux: !!useMux, ts: Date.now(), tries: 0, sentAt: 0 }; + let list = this._pendingDeliveries.get(sessionId); + if (!list) { + list = []; + this._pendingDeliveries.set(sessionId, list); + } + list.push(rec); + this._persistReliableState(); + this._updateConnectionIndicator(); + this._drainSession(sessionId); + } + + _nextSeq(sessionId) { + const next = (this._seqCounters.get(sessionId) || 0) + 1; + this._seqCounters.set(sessionId, next); + return next; + } + + /** Deliver all unacked records for a session, in seq order. */ + _drainSession(sessionId) { + const list = this._pendingDeliveries.get(sessionId); + if (!list || list.length === 0) return; + + // Fast path: WebSocket open for this session — fire each not-yet-sent record + // over the single ordered stream. They stay pending until the server ACKs + // them ({t:'ia'}); a frame swallowed by a half-open socket is re-sent after + // the sweep force-reconnects (which resets sentAt=0 in _onWsReady). + if (this._ws && this._ws.readyState === WebSocket.OPEN && this._wsSessionId === sessionId) { + for (const rec of list) { + if (rec.sentAt !== 0) continue; + try { + this._ws.send(JSON.stringify({ t: 'i', d: rec.data, seq: rec.seq, cid: this._clientId })); + rec.sentAt = Date.now(); + rec.tries++; + } catch { + break; // socket died mid-send — reconnect/POST drainer retries + } + } return; } - // Fast path: WebSocket — fire-and-forget, inherently ordered (single TCP stream). - if (this._wsReady && this._wsSessionId === sessionId) { + // Slow path: no WS — POST records in order, awaiting each (the HTTP 2xx is + // the ACK). Serialized per session so seq order survives async fetches. + if (this._postDraining.has(sessionId)) return; + this._postDraining.add(sessionId); + (async () => { try { - this._ws.send(JSON.stringify({ t: 'i', d: input })); - 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(() => { - const fetchPromise = fetch(`/api/sessions/${sessionId}/input`, { - method: 'POST', - headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify({ input }), - keepalive: input.length < 65536, - }); - - // Handle response asynchronously — don't block next keystroke on response - fetchPromise.then(resp => { - if (!resp.ok) { - this._enqueueInput(sessionId, input); - } else { - this.clearPendingHooks(sessionId); + for (;;) { + const cur = this._pendingDeliveries.get(sessionId); + if (!cur || cur.length === 0) break; + // If the WebSocket came back mid-drain, yield to it (the acked stream) + // so we don't redundantly re-POST what onopen is already re-sending. + if (this._ws && this._ws.readyState === WebSocket.OPEN && this._wsSessionId === sessionId) { + break; + } + const rec = cur[0]; + rec.tries++; + rec.sentAt = Date.now(); + let resp = null; + try { + const body = { input: rec.data, seq: rec.seq, clientId: this._clientId }; + if (rec.useMux) body.useMux = true; + resp = await fetch(`/api/sessions/${sessionId}/input`, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify(body), + keepalive: rec.data.length < 65536, + }); + } catch { + resp = null; + } + if (resp && resp.ok) { + this._ackDelivery(sessionId, rec.seq); + } else if (resp && (resp.status === 404 || resp.status === 410)) { + // Session no longer exists — the input can never land. Drop it + // rather than retry forever (not a "lost" prompt: the target is gone). + this._ackDelivery(sessionId, rec.seq); + } else { + break; // offline / 5xx — leave queued; sweep + reconnect retry later + } } - }).catch(() => { - this._enqueueInput(sessionId, input); - }); - - // Return immediately after fetch is dispatched (don't await response) - }); + } finally { + this._postDraining.delete(sessionId); + } + })(); } - - _enqueueInput(sessionId, input) { - const existing = this._inputQueue.get(sessionId) || ''; - let combined = existing + input; - // Enforce 64KB cap — keep most recent keystrokes - if (combined.length > this._inputQueueMaxBytes) { - combined = combined.slice(combined.length - this._inputQueueMaxBytes); - } - this._inputQueue.set(sessionId, combined); - this._updateConnectionIndicator(); - } - - async _drainInputQueues() { - if (this._inputQueue.size === 0) return; - // Snapshot and clear - const queued = new Map(this._inputQueue); - this._inputQueue.clear(); - this._updateConnectionIndicator(); - - for (const [sessionId, input] of queued) { - const resp = await this._apiPost(`/api/sessions/${sessionId}/input`, { input }); - if (!resp?.ok) { - this._enqueueInput(sessionId, input); + /** Drop an ACKed record (by exact seq) and persist. */ + _ackDelivery(sessionId, seq) { + const list = this._pendingDeliveries.get(sessionId); + if (list) { + const idx = list.findIndex((r) => r.seq === seq); + if (idx !== -1) { + list.splice(idx, 1); + if (list.length === 0) this._pendingDeliveries.delete(sessionId); + // When nothing is left pending anywhere, flush durable state immediately + // (not debounced) so a reload in the next 250ms can't redeliver an + // already-delivered frame — otherwise localStorage briefly still shows it. + if (this._pendingDeliveries.size === 0) this._persistReliableNow(); + else this._persistReliableState(); + this._updateConnectionIndicator(); } } - this._updateConnectionIndicator(); + this.clearPendingHooks?.(sessionId); + } + + /** Server input-ACK frame ({t:'ia',seq}) over the WebSocket. */ + _onWsInputAck(seq) { + if (this._wsSessionId && Number.isInteger(seq)) this._ackDelivery(this._wsSessionId, seq); + } + + /** Called from ws.onopen — flush everything pending over the fresh socket. */ + _onWsReady(sessionId) { + const list = this._pendingDeliveries.get(sessionId); + if (list) for (const r of list) r.sentAt = 0; // fresh socket ⇒ re-send all + this._drainSession(sessionId); + } + + /** + * Periodic retry. For the active WS session, an oldest frame unacked past the + * timeout means the socket is (half-)dead — close it to force a fast reconnect + * (onclose → reconnect → onopen → _onWsReady re-sends). Other sessions just + * (re)drain over POST. + */ + _redeliverSweep() { + if (this._pendingDeliveries.size === 0) return; + for (const sessionId of [...this._pendingDeliveries.keys()]) { + const list = this._pendingDeliveries.get(sessionId); + if (!list || list.length === 0) continue; + const isActiveWs = + this._ws && this._ws.readyState === WebSocket.OPEN && this._wsSessionId === sessionId; + if (isActiveWs) { + const oldest = list[0]; + if (oldest && oldest.sentAt && Date.now() - oldest.sentAt > this._reliableAckTimeoutMs) { + try { + this._ws.close(); // half-open: never recovers on its own — force reconnect + } catch { + /* ignore */ + } + continue; + } + } + this._drainSession(sessionId); + } + } + + /** Total bytes/count still awaiting ACK across all sessions (for the indicator). */ + _pendingBytes() { + let bytes = 0; + let count = 0; + for (const list of this._pendingDeliveries.values()) { + for (const r of list) { + bytes += r.data.length; + count++; + } + } + return { bytes, count }; + } + + // ---- durable persistence (localStorage; quota- and disabled-storage-safe) -- + + _loadReliableState() { + // Stable client identity for server-side dedup across reconnects/reloads. + try { + this._clientId = localStorage.getItem('codeman:clientId') || ''; + } catch { + this._clientId = ''; + } + if (!this._clientId) { + this._clientId = 'c-' + Math.random().toString(36).slice(2) + '-' + Date.now().toString(36); + try { + localStorage.setItem('codeman:clientId', this._clientId); + } catch { + /* storage disabled — dedup degrades to per-load, still no loss */ + } + } + try { + const raw = localStorage.getItem('codeman:pendingInput'); + if (!raw) return; + const saved = JSON.parse(raw); + if (saved && saved.seqs) { + for (const [s, n] of Object.entries(saved.seqs)) { + if (Number.isFinite(n)) this._seqCounters.set(s, n); + } + } + if (saved && saved.pending) { + for (const [s, recs] of Object.entries(saved.pending)) { + if (Array.isArray(recs) && recs.length) { + // Reset sentAt so they re-deliver promptly on this fresh load. + this._pendingDeliveries.set( + s, + recs + .filter((r) => r && typeof r.data === 'string' && Number.isInteger(r.seq)) + .map((r) => ({ + seq: r.seq, + data: r.data, + useMux: !!r.useMux, + ts: r.ts || Date.now(), + tries: 0, + sentAt: 0, + })) + ); + } + } + } + } catch { + /* corrupt/parse error — start clean rather than throw */ + } + } + + _persistReliableState() { + // Debounced — typing without local echo calls this per keystroke. + if (this._persistReliableTimer) return; + this._persistReliableTimer = setTimeout(() => { + this._persistReliableTimer = null; + this._persistReliableNow(); + }, 250); + } + + _persistReliableNow() { + if (this._persistReliableTimer) { + clearTimeout(this._persistReliableTimer); + this._persistReliableTimer = null; + } + try { + const seqs = {}; + for (const [s, n] of this._seqCounters) seqs[s] = n; + const pending = {}; + let bytes = 0; + for (const [s, list] of this._pendingDeliveries) { + if (!list.length) continue; + pending[s] = list.map((r) => ({ + seq: r.seq, + data: r.data, + useMux: r.useMux, + ts: r.ts, + tries: r.tries, + })); + for (const r of list) bytes += r.data.length; + } + // Bound the persisted backlog. On extreme overflow keep the seq counters + // (so future input stays monotonic and dedup-safe) but skip the payloads — + // the in-memory queue still delivers; only cross-reload durability is lost. + const payload = + bytes > this._reliableMaxBytes ? { seqs } : { seqs, pending }; + localStorage.setItem('codeman:pendingInput', JSON.stringify(payload)); + } catch { + /* QuotaExceeded or disabled storage — in-memory delivery is unaffected */ + } } _updateConnectionIndicator() { @@ -2143,11 +2358,9 @@ class CodemanApp { const text = this.$('connectionText'); if (!indicator || !dot || !text) return; - let totalBytes = 0; - for (const v of this._inputQueue.values()) totalBytes += v.length; - + const { bytes: totalBytes, count } = this._pendingBytes(); const status = this._connectionStatus; - const hasQueue = totalBytes > 0; + const hasQueue = count > 0; // Connected with empty queue — hide if ((status === 'connected' || status === 'connecting') && !hasQueue) { @@ -2179,6 +2392,8 @@ class CodemanApp { this.isOnline = true; this.reconnectAttempts = 0; this.connectSSE(); + // Network came back — drain durably-queued input right away. + this._redeliverSweep(); }); window.addEventListener('offline', () => { this.isOnline = false; @@ -3686,7 +3901,13 @@ class CodemanApp { this._flushedOffsets?.delete(sessionId); this._flushedTexts?.delete(sessionId); - this._inputQueue.delete(sessionId); + // Drop any durably-queued input for a session that's actually gone (deleted/ + // exited). Not a lost prompt — the target no longer exists. Only reached on + // real session removal, never on a tab switch. + this._pendingDeliveries?.delete(sessionId); + this._seqCounters?.delete(sessionId); + this._postDraining?.delete(sessionId); + this._persistReliableState(); this.ralphStates.delete(sessionId); this.ralphClosedSessions.delete(sessionId); this.projectInsights.delete(sessionId); diff --git a/src/web/public/terminal-ui.js b/src/web/public/terminal-ui.js index 65e767ee..c131c4c4 100644 --- a/src/web/public/terminal-ui.js +++ b/src/web/public/terminal-ui.js @@ -2195,12 +2195,11 @@ Object.assign(CodemanApp.prototype, { * @returns {Promise} */ async sendInput(input) { - if (!this.activeSessionId) return; - await fetch(`/api/sessions/${this.activeSessionId}/input`, { - method: 'POST', - headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify({ input, useMux: true }), - }); + if (!this.activeSessionId || !input) return; + // Route through the durable, exactly-once delivery layer (useMux for the + // POST fallback) so voice / keyboard-accessory / paste input also survives a + // dropped link instead of being lost in a single best-effort fetch. + this._sendInputAsync(this.activeSessionId, input, { useMux: true }); }, // ═══════════════════════════════════════════════════════════════ diff --git a/src/web/routes/session-routes.ts b/src/web/routes/session-routes.ts index 247e1ebd..c3ec4e33 100644 --- a/src/web/routes/session-routes.ts +++ b/src/web/routes/session-routes.ts @@ -662,7 +662,7 @@ export function registerSessionRoutes( app.post('/api/sessions/:id/input', async (req) => { const { id } = req.params as { id: string }; - const { input, useMux } = parseBody(SessionInputWithLimitSchema, req.body); + const { input, useMux, seq, clientId } = parseBody(SessionInputWithLimitSchema, req.body); const session = findSessionOrFail(ctx, id); const inputStr = String(input); @@ -673,6 +673,13 @@ export function registerSessionRoutes( ); } + // Reliable delivery (POST fallback when the WebSocket is down): a 2xx IS the + // client's ACK, so a tagged duplicate redelivery must still return 200 but + // skip the write. Untagged requests (curl/legacy) always apply. + if (typeof clientId === 'string' && typeof seq === 'number' && !session.shouldApplyInput(clientId, seq)) { + return {}; + } + // Write input to PTY. Direct write is synchronous; writeViaMux // (tmux send-keys) is fire-and-forget to avoid blocking the HTTP response. if (useMux) { diff --git a/src/web/routes/ws-routes.ts b/src/web/routes/ws-routes.ts index 8bcca517..8aaa1a21 100644 --- a/src/web/routes/ws-routes.ts +++ b/src/web/routes/ws-routes.ts @@ -21,8 +21,11 @@ * {"t":"o","d":"..."} — terminal output * {"t":"c"} — clear terminal * {"t":"r"} — needs refresh (reload buffer) + * {"t":"ia","seq":N} — input ACK (echoes the seq of an applied/deduped input frame) * Client -> Server: - * {"t":"i","d":"..."} — input (keystroke or paste) + * {"t":"i","d":"...","seq":N,"cid":"..."} — input (keystroke or paste). seq+cid are + * optional reliable-delivery tags: the server applies each + * (cid,seq) at-most-once and ACKs with {"t":"ia","seq":N}. * {"t":"z","c":N,"r":N,"f":bool} — resize terminal (f=true forces SIGWINCH even if dims unchanged) */ @@ -123,10 +126,22 @@ export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort, getHost const msg = JSON.parse(String(raw)); if (msg.t === 'i' && typeof msg.d === 'string') { if (msg.d.length > MAX_INPUT_LENGTH) return; - // Typed input from a claim-holding desktop keeps the claim "hot" - // and re-asserts the desktop layout after a mobile override. - if (holdsDesktopClaim) session.noteDesktopActivity(); - session.write(msg.d); + // Reliable delivery: when the frame carries a clientId + seq, apply it + // exactly once (skip a duplicate redelivery) but ACK it regardless so + // the client can drop it from its durable queue. Frames without seq + // (legacy/other tools) are applied as-is — no behavior change. + const cid = typeof msg.cid === 'string' ? msg.cid : null; + const seq = Number.isInteger(msg.seq) ? (msg.seq as number) : null; + const apply = cid && seq !== null ? session.shouldApplyInput(cid, seq) : true; + if (apply) { + // Typed input from a claim-holding desktop keeps the claim "hot" + // and re-asserts the desktop layout after a mobile override. + if (holdsDesktopClaim) session.noteDesktopActivity(); + session.write(msg.d); + } + if (seq !== null && socket.readyState === 1) { + socket.send(`{"t":"ia","seq":${seq}}`); + } } else if ( msg.t === 'z' && Number.isInteger(msg.c) && diff --git a/src/web/schemas.ts b/src/web/schemas.ts index 29fb0e6b..06457799 100644 --- a/src/web/schemas.ts +++ b/src/web/schemas.ts @@ -469,6 +469,15 @@ export const SettingsUpdateSchema = z export const SessionInputWithLimitSchema = z.object({ input: z.string().max(100000), // 100KB max input useMux: z.boolean().optional(), + // Reliable-delivery dedup (optional; absent for curl/legacy clients). The web + // client tags each input with a stable clientId + a monotonic per-session seq + // and redelivers anything it hasn't seen ACKed (e.g. a frame silently dropped + // by a half-open WebSocket on a flaky link). The server applies each (clientId, + // seq) at-most-once via Session.shouldApplyInput so a redelivery can't type the + // prompt twice. `.optional()` (not `.nullish()`) — the client omits them when + // unset rather than sending null. See docs/reliable-input-delivery.md. + seq: z.number().int().nonnegative().optional(), + clientId: z.string().max(128).optional(), }); // ========== Session Mutation Routes ========== diff --git a/test/mocks/mock-session.ts b/test/mocks/mock-session.ts index c1923dec..a16154c6 100644 --- a/test/mocks/mock-session.ts +++ b/test/mocks/mock-session.ts @@ -39,6 +39,16 @@ export class MockSession extends EventEmitter { return true; } + /** Exactly-once input dedup — mirrors Session.shouldApplyInput so route tests + * exercising the reliable-delivery path behave like production. */ + private _appliedInputSeq = new Map(); + shouldApplyInput(clientId: string, seq: number): boolean { + const last = this._appliedInputSeq.get(clientId); + if (last !== undefined && seq <= last) return false; + this._appliedInputSeq.set(clientId, seq); + return true; + } + /** Get the last written data */ get lastWrite(): string | undefined { return this.writeBuffer[this.writeBuffer.length - 1]; diff --git a/test/reliable-input-dedup.test.ts b/test/reliable-input-dedup.test.ts new file mode 100644 index 00000000..d77bef2d --- /dev/null +++ b/test/reliable-input-dedup.test.ts @@ -0,0 +1,73 @@ +/** + * @fileoverview Exactly-once input delivery — Session.shouldApplyInput dedup. + * + * Guards the server half of the reliable-input-delivery feature: the web client + * tags each input frame with a stable clientId + a monotonic per-session seq and + * redelivers anything it hasn't seen ACKed (a half-open socket silently drops + * frames on a flaky link). shouldApplyInput must apply each (clientId, seq) + * exactly once so a redelivery can never type the prompt twice — while still + * applying untagged input (curl/legacy) unconditionally at the call sites. + * + * See docs/reliable-input-delivery.md. + */ +import { describe, it, expect } from 'vitest'; +import { Session } from '../src/session.js'; + +function makeSession(): Session { + // workingDir is the only required field; no PTY is spawned until start(), + // and TmuxManager no-ops under VITEST — so this is a cheap, side-effect-free + // instance for exercising the pure dedup bookkeeping. + return new Session({ workingDir: '/tmp' }); +} + +describe('Session.shouldApplyInput (exactly-once input dedup)', () => { + it('applies a fresh (clientId, seq) exactly once', () => { + const s = makeSession(); + expect(s.shouldApplyInput('clientA', 1)).toBe(true); + // Same seq redelivered (lost ACK) — must NOT apply again. + expect(s.shouldApplyInput('clientA', 1)).toBe(false); + }); + + it('applies strictly increasing seqs and rejects stale ones', () => { + const s = makeSession(); + expect(s.shouldApplyInput('c', 1)).toBe(true); + expect(s.shouldApplyInput('c', 2)).toBe(true); + expect(s.shouldApplyInput('c', 3)).toBe(true); + // Out-of-order / replayed lower seqs are duplicates. + expect(s.shouldApplyInput('c', 2)).toBe(false); + expect(s.shouldApplyInput('c', 1)).toBe(false); + // The next genuinely-new seq still applies. + expect(s.shouldApplyInput('c', 4)).toBe(true); + }); + + it('tracks each client independently', () => { + const s = makeSession(); + expect(s.shouldApplyInput('a', 5)).toBe(true); + // A different client at seq 1 is not shadowed by client a's higher seq. + expect(s.shouldApplyInput('b', 1)).toBe(true); + expect(s.shouldApplyInput('b', 1)).toBe(false); + expect(s.shouldApplyInput('a', 6)).toBe(true); + }); + + it('tolerates a seq gap (skips never collapse a new seq to a duplicate)', () => { + const s = makeSession(); + expect(s.shouldApplyInput('c', 1)).toBe(true); + // Client jumped seq (e.g. resumed after a reload that kept the counter). + expect(s.shouldApplyInput('c', 100)).toBe(true); + expect(s.shouldApplyInput('c', 100)).toBe(false); + expect(s.shouldApplyInput('c', 50)).toBe(false); + expect(s.shouldApplyInput('c', 101)).toBe(true); + }); + + it('keeps recent clients dedup-correct past the eviction bound', () => { + const s = makeSession(); + // Far exceed MAX_INPUT_DEDUP_CLIENTS (256) with one-shot clients, then prove + // a freshly-active client is still deduped correctly (MRU eviction). + for (let i = 0; i < 400; i++) { + expect(s.shouldApplyInput(`oneshot-${i}`, 1)).toBe(true); + } + expect(s.shouldApplyInput('recent', 1)).toBe(true); + expect(s.shouldApplyInput('recent', 1)).toBe(false); + expect(s.shouldApplyInput('recent', 2)).toBe(true); + }); +}); diff --git a/test/routes/session-routes.test.ts b/test/routes/session-routes.test.ts index 39cac61b..7f0d8a39 100644 --- a/test/routes/session-routes.test.ts +++ b/test/routes/session-routes.test.ts @@ -355,6 +355,43 @@ describe('session-routes', () => { const body = JSON.parse(res.body); expect(body.success).toBe(false); }); + + it('applies a tagged (clientId, seq) input exactly once on redelivery', async () => { + const url = `/api/sessions/${harness.ctx._sessionId}/input`; + // eslint-disable-next-line @typescript-eslint/no-explicit-any + const session = harness.ctx.sessions.get(harness.ctx._sessionId) as any; + session.writeBuffer.length = 0; + + const post = (payload: unknown) => harness.app.inject({ method: 'POST', url, payload }); + + // First delivery of seq 1 — applied (200, written once). + const first = await post({ input: 'prompt', seq: 1, clientId: 'cid-1' }); + expect(first.statusCode).toBe(200); + + // Redelivery of the SAME seq (client never saw the ACK) — still 200, but + // must NOT write again. + const dup = await post({ input: 'prompt', seq: 1, clientId: 'cid-1' }); + expect(dup.statusCode).toBe(200); + + // A genuinely new seq — applied. + const next = await post({ input: '\r', seq: 2, clientId: 'cid-1' }); + expect(next.statusCode).toBe(200); + + expect(session.writeBuffer).toEqual(['prompt', '\r']); + }); + + it('always applies untagged input (curl/legacy, no dedup)', async () => { + const url = `/api/sessions/${harness.ctx._sessionId}/input`; + // eslint-disable-next-line @typescript-eslint/no-explicit-any + const session = harness.ctx.sessions.get(harness.ctx._sessionId) as any; + session.writeBuffer.length = 0; + + const post = () => harness.app.inject({ method: 'POST', url, payload: { input: 'x' } }); + await post(); + await post(); + // No seq/clientId ⇒ no dedup ⇒ both writes land. + expect(session.writeBuffer).toEqual(['x', 'x']); + }); }); // ========== POST /api/sessions/:id/resize ==========