mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 12:39:42 +02:00
fix(input): durable exactly-once input delivery so a dropped link can't lose a prompt
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) <noreply@anthropic.com>
This commit is contained in:
@@ -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<sessionId, record[]>`), 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.
|
||||
@@ -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<string, number>();
|
||||
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.
|
||||
*
|
||||
|
||||
+297
-76
@@ -468,13 +468,28 @@ class CodemanApp {
|
||||
this.maxReconnectAttempts = 10;
|
||||
this.isOnline = navigator.onLine;
|
||||
|
||||
// Offline input queue
|
||||
this._inputQueue = new Map(); // Map<sessionId, string>
|
||||
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);
|
||||
|
||||
@@ -2195,12 +2195,11 @@ Object.assign(CodemanApp.prototype, {
|
||||
* @returns {Promise<void>}
|
||||
*/
|
||||
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 });
|
||||
},
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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) &&
|
||||
|
||||
@@ -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 ==========
|
||||
|
||||
@@ -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<string, number>();
|
||||
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];
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -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 ==========
|
||||
|
||||
Reference in New Issue
Block a user