mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 12:39:42 +02:00
COD-133 fix header WS status indicator + typing lag from v1.1.15 merge
The upstream v1.1.15 merge spliced upstream's transport-object indicator
body onto local's _connectionStatus-based _updateConnectionIndicator()
without defining `transport`, so every transport.* reference threw
ReferenceError on any queued state. That hid the "WS" status and, because
_reliableSend() updates the indicator before _drainSession(), made every
keystroke skip immediate delivery (input flushed only on the 2s sweep =
typing lag).
- Rewrite _updateConnectionIndicator() to show the terminal WebSocket
transport from _wsState (WS / HTTP / WS… / Offline), falling back to the
SSE _connectionStatus only on the idle dashboard.
- Only annotate a backlog (· N queued) above 4 bytes so normal typing no
longer flickers "sending 1B" on each key press.
- test/connection-indicator.test.ts (new): transport labels, the >4B
threshold, an exhaustive never-throws guard for the ReferenceError, and
the _reliableSend -> _drainSession invariant (typing-lag guard).
- test/input-send-order.test.ts: reconcile to local's durable input layer
(the prior coalescing-fallback tests had been failing since 1255e28).
This commit is contained in:
+57
-23
@@ -2358,34 +2358,68 @@ class CodemanApp {
|
|||||||
const text = this.$('connectionText');
|
const text = this.$('connectionText');
|
||||||
if (!indicator || !dot || !text) return;
|
if (!indicator || !dot || !text) return;
|
||||||
|
|
||||||
const status = this._connectionStatus;
|
const { bytes: totalBytes, count } = this._pendingBytes();
|
||||||
|
const hasQueue = count > 0;
|
||||||
|
// Only surface a backlog once it's more than a few bytes. A single keystroke
|
||||||
|
// (1B) ACKs in milliseconds, so without this the label flickered "sending 1B"
|
||||||
|
// on every key press. Above this threshold means input is genuinely backing up.
|
||||||
|
const BACKLOG_HINT_BYTES = 4;
|
||||||
|
const showBacklog = totalBytes > BACKLOG_HINT_BYTES;
|
||||||
|
const formatBytes = (b) => (b < 1024 ? `${b}B` : `${(b / 1024).toFixed(1)}KB`);
|
||||||
|
const queuedSuffix = showBacklog ? ` · ${formatBytes(totalBytes)} queued` : '';
|
||||||
|
|
||||||
// While the connection is healthy, never surface the input queue. With the
|
// Hard offline (browser reports no network) dominates everything.
|
||||||
// reliable-delivery layer every keystroke is briefly "pending" until its ACK
|
if (!this.isOnline || this._connectionStatus === 'offline') {
|
||||||
// lands a few ms later — showing that flashed "Sending 1B…" on every single
|
indicator.style.display = 'flex';
|
||||||
// character. The indicator is only meaningful for an actual connection
|
dot.className = 'connection-dot offline';
|
||||||
// problem (reconnecting / offline), where the queued byte count reassures
|
text.textContent = showBacklog ? `Offline (${formatBytes(totalBytes)} queued)` : 'Offline';
|
||||||
// the user their typing is safely buffered and will be sent.
|
indicator.title = 'No network connection';
|
||||||
if (status === 'connected' || status === 'connecting') {
|
|
||||||
indicator.style.display = 'none';
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
const { bytes: totalBytes, count } = this._pendingBytes();
|
// With an active terminal, show its transport (WebSocket vs HTTP fallback).
|
||||||
const hasQueue = count > 0;
|
if (this.activeSessionId) {
|
||||||
indicator.style.display = 'flex';
|
let cls, label, detail;
|
||||||
dot.className = 'connection-dot';
|
switch (this._wsState) {
|
||||||
|
case 'connected':
|
||||||
const formatBytes = (b) => (b < 1024 ? `${b}B` : `${(b / 1024).toFixed(1)}KB`);
|
cls = 'connected'; label = 'WS'; detail = 'Terminal connected over WebSocket';
|
||||||
|
break;
|
||||||
if (status === 'reconnecting') {
|
case 'fallback':
|
||||||
dot.classList.add('reconnecting');
|
cls = 'fallback'; label = 'HTTP'; detail = 'WebSocket unavailable — input sent over HTTP';
|
||||||
text.textContent = hasQueue ? `Reconnecting (${formatBytes(totalBytes)} queued)` : 'Reconnecting...';
|
break;
|
||||||
} else {
|
case 'reconnecting':
|
||||||
// Offline or disconnected
|
cls = 'reconnecting'; label = 'WS…'; detail = 'Reconnecting WebSocket';
|
||||||
dot.classList.add('offline');
|
break;
|
||||||
text.textContent = hasQueue ? `Offline (${formatBytes(totalBytes)} queued)` : 'Offline';
|
case 'connecting':
|
||||||
|
default:
|
||||||
|
cls = 'reconnecting'; label = 'WS…'; detail = 'Connecting WebSocket';
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
indicator.style.display = 'flex';
|
||||||
|
dot.className = `connection-dot ${cls}`;
|
||||||
|
text.textContent = `${label}${queuedSuffix}`;
|
||||||
|
indicator.title = detail;
|
||||||
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// No active terminal — reflect the SSE event stream only when it needs attention.
|
||||||
|
if (this._connectionStatus === 'reconnecting' || this._connectionStatus === 'disconnected') {
|
||||||
|
indicator.style.display = 'flex';
|
||||||
|
dot.className = 'connection-dot reconnecting';
|
||||||
|
text.textContent = showBacklog ? `Reconnecting (${formatBytes(totalBytes)} queued)` : 'Reconnecting...';
|
||||||
|
indicator.title = 'Reconnecting to server';
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Idle dashboard, healthy stream — hide unless input is genuinely queued.
|
||||||
|
if (!hasQueue) {
|
||||||
|
indicator.style.display = 'none';
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
indicator.style.display = 'flex';
|
||||||
|
dot.className = 'connection-dot draining';
|
||||||
|
text.textContent = showBacklog ? `Sending ${formatBytes(totalBytes)}...` : 'Sending...';
|
||||||
|
indicator.title = 'Delivering queued input';
|
||||||
}
|
}
|
||||||
|
|
||||||
setupOnlineDetection() {
|
setupOnlineDetection() {
|
||||||
|
|||||||
@@ -0,0 +1,204 @@
|
|||||||
|
/**
|
||||||
|
* @fileoverview Regression tests for the header connection indicator
|
||||||
|
* (`CodemanApp._updateConnectionIndicator`) and the invariant that updating it
|
||||||
|
* never aborts durable input delivery.
|
||||||
|
*
|
||||||
|
* Guards two regressions from the upstream v1.1.15 merge (fixed in COD-133):
|
||||||
|
* 1. The indicator body referenced an undefined `transport` object, throwing
|
||||||
|
* `ReferenceError: transport is not defined` on every queued state. Because
|
||||||
|
* `_reliableSend()` calls `_updateConnectionIndicator()` *before*
|
||||||
|
* `_drainSession()`, the throw skipped immediate delivery on every keystroke
|
||||||
|
* → input only flushed on the 2s sweep (large typing lag) and the indicator
|
||||||
|
* never rendered (missing "WS" status).
|
||||||
|
* 2. The restored body only read the SSE `_connectionStatus`, so it never
|
||||||
|
* surfaced the terminal WebSocket transport ("WS" / "HTTP"), and it flashed
|
||||||
|
* "sending 1B" on every single keystroke.
|
||||||
|
*
|
||||||
|
* Loaded via `vm` with a stubbed context (no jsdom — see input-send-order.test.ts).
|
||||||
|
*/
|
||||||
|
import { readFileSync } from 'node:fs';
|
||||||
|
import { performance } from 'node:perf_hooks';
|
||||||
|
import { resolve } from 'node:path';
|
||||||
|
import vm from 'node:vm';
|
||||||
|
import { describe, expect, it, vi } from 'vitest';
|
||||||
|
|
||||||
|
function loadCodemanAppClass() {
|
||||||
|
const constants = readFileSync(resolve(import.meta.dirname, '../src/web/public/constants.js'), 'utf8');
|
||||||
|
const source = readFileSync(resolve(import.meta.dirname, '../src/web/public/app.js'), 'utf8');
|
||||||
|
const context = vm.createContext({
|
||||||
|
console,
|
||||||
|
performance,
|
||||||
|
setInterval: vi.fn(),
|
||||||
|
clearInterval: vi.fn(),
|
||||||
|
setTimeout,
|
||||||
|
clearTimeout,
|
||||||
|
requestAnimationFrame: vi.fn(),
|
||||||
|
HTMLCanvasElement: class HTMLCanvasElement {},
|
||||||
|
fetch: (...args: Parameters<typeof fetch>) => global.fetch(...args),
|
||||||
|
document: { addEventListener: vi.fn() },
|
||||||
|
localStorage: {
|
||||||
|
length: 0,
|
||||||
|
key: vi.fn(),
|
||||||
|
getItem: vi.fn(),
|
||||||
|
setItem: vi.fn(),
|
||||||
|
removeItem: vi.fn(),
|
||||||
|
},
|
||||||
|
window: { addEventListener: vi.fn(), removeEventListener: vi.fn() },
|
||||||
|
MobileDetection: {},
|
||||||
|
});
|
||||||
|
vm.runInContext(`${constants}\n${source}\nglobalThis.__CodemanApp = CodemanApp;`, context);
|
||||||
|
return (context as { __CodemanApp: new () => unknown }).__CodemanApp;
|
||||||
|
}
|
||||||
|
|
||||||
|
const CodemanApp = loadCodemanAppClass();
|
||||||
|
|
||||||
|
function fakeElement() {
|
||||||
|
return { style: { display: '' }, title: '', textContent: '', className: '' };
|
||||||
|
}
|
||||||
|
|
||||||
|
type Indicator = {
|
||||||
|
$: (id: string) => unknown;
|
||||||
|
_pendingDeliveries: Map<string, Array<{ seq: number; data: string }>>;
|
||||||
|
_connectionStatus: string;
|
||||||
|
_wsState: string;
|
||||||
|
activeSessionId: string | null;
|
||||||
|
isOnline: boolean;
|
||||||
|
_updateConnectionIndicator: () => void;
|
||||||
|
};
|
||||||
|
|
||||||
|
function makeApp(overrides: Partial<Indicator> & { queuedBytes?: number } = {}) {
|
||||||
|
const app = Object.create((CodemanApp as { prototype: object }).prototype) as Indicator & {
|
||||||
|
queuedBytes?: number;
|
||||||
|
};
|
||||||
|
const els: Record<string, ReturnType<typeof fakeElement>> = {
|
||||||
|
connectionIndicator: fakeElement(),
|
||||||
|
connectionDot: fakeElement(),
|
||||||
|
connectionText: fakeElement(),
|
||||||
|
};
|
||||||
|
app.$ = (id: string) => els[id];
|
||||||
|
app._pendingDeliveries = new Map();
|
||||||
|
app._connectionStatus = 'connected';
|
||||||
|
app._wsState = 'disconnected';
|
||||||
|
app.activeSessionId = null;
|
||||||
|
app.isOnline = true;
|
||||||
|
Object.assign(app, overrides);
|
||||||
|
const queued = overrides.queuedBytes ?? 0;
|
||||||
|
if (queued > 0) {
|
||||||
|
app._pendingDeliveries.set('s1', [{ seq: 1, data: 'x'.repeat(queued) }]);
|
||||||
|
}
|
||||||
|
return { app, els };
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('connection indicator — transport display', () => {
|
||||||
|
it('shows "WS" with a connected dot when the terminal WebSocket is open', () => {
|
||||||
|
const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected' });
|
||||||
|
app._updateConnectionIndicator();
|
||||||
|
expect(els.connectionIndicator.style.display).toBe('flex');
|
||||||
|
expect(els.connectionText.textContent).toBe('WS');
|
||||||
|
expect(els.connectionDot.className).toContain('connected');
|
||||||
|
});
|
||||||
|
|
||||||
|
it('shows "HTTP" with a fallback dot when the socket dropped to HTTP POST', () => {
|
||||||
|
const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'fallback' });
|
||||||
|
app._updateConnectionIndicator();
|
||||||
|
expect(els.connectionText.textContent).toBe('HTTP');
|
||||||
|
expect(els.connectionDot.className).toContain('fallback');
|
||||||
|
});
|
||||||
|
|
||||||
|
it('shows "WS…" while connecting or reconnecting the socket', () => {
|
||||||
|
for (const state of ['connecting', 'reconnecting']) {
|
||||||
|
const { app, els } = makeApp({ activeSessionId: 's1', _wsState: state });
|
||||||
|
app._updateConnectionIndicator();
|
||||||
|
expect(els.connectionText.textContent).toBe('WS…');
|
||||||
|
expect(els.connectionDot.className).toContain('reconnecting');
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it('shows "Offline" when the browser reports no network', () => {
|
||||||
|
const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected', isOnline: false });
|
||||||
|
app._updateConnectionIndicator();
|
||||||
|
expect(els.connectionText.textContent).toBe('Offline');
|
||||||
|
expect(els.connectionDot.className).toContain('offline');
|
||||||
|
});
|
||||||
|
|
||||||
|
it('hides on an idle dashboard (no active session, healthy stream, no queue)', () => {
|
||||||
|
const { app, els } = makeApp({ activeSessionId: null, _connectionStatus: 'connected' });
|
||||||
|
app._updateConnectionIndicator();
|
||||||
|
expect(els.connectionIndicator.style.display).toBe('none');
|
||||||
|
});
|
||||||
|
|
||||||
|
it('surfaces SSE reconnecting on the dashboard when there is no active terminal', () => {
|
||||||
|
const { app, els } = makeApp({ activeSessionId: null, _connectionStatus: 'reconnecting' });
|
||||||
|
app._updateConnectionIndicator();
|
||||||
|
expect(els.connectionText.textContent).toContain('Reconnecting');
|
||||||
|
expect(els.connectionDot.className).toContain('reconnecting');
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
describe('connection indicator — keystroke backlog threshold', () => {
|
||||||
|
it('does NOT annotate a single-keystroke (1B) queue — no "sending 1B" flicker', () => {
|
||||||
|
const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected', queuedBytes: 1 });
|
||||||
|
app._updateConnectionIndicator();
|
||||||
|
expect(els.connectionText.textContent).toBe('WS');
|
||||||
|
expect(els.connectionText.textContent).not.toMatch(/queued/);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('does NOT annotate at the 4B threshold boundary', () => {
|
||||||
|
const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected', queuedBytes: 4 });
|
||||||
|
app._updateConnectionIndicator();
|
||||||
|
expect(els.connectionText.textContent).toBe('WS');
|
||||||
|
});
|
||||||
|
|
||||||
|
it('annotates a genuine backlog (>4B) with a queued byte count', () => {
|
||||||
|
const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected', queuedBytes: 40 });
|
||||||
|
app._updateConnectionIndicator();
|
||||||
|
expect(els.connectionText.textContent).toMatch(/^WS · 40B queued$/);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
describe('connection indicator — never throws (the ReferenceError regression)', () => {
|
||||||
|
it('renders every transport × stream × queue combination without throwing', () => {
|
||||||
|
const wsStates = ['disconnected', 'connecting', 'connected', 'reconnecting', 'fallback'];
|
||||||
|
const sseStates = ['connected', 'connecting', 'reconnecting', 'disconnected', 'offline'];
|
||||||
|
for (const ws of wsStates) {
|
||||||
|
for (const sse of sseStates) {
|
||||||
|
for (const active of ['s1', null] as const) {
|
||||||
|
for (const queuedBytes of [0, 1, 4, 200]) {
|
||||||
|
for (const isOnline of [true, false]) {
|
||||||
|
const { app } = makeApp({
|
||||||
|
activeSessionId: active,
|
||||||
|
_wsState: ws,
|
||||||
|
_connectionStatus: sse,
|
||||||
|
isOnline,
|
||||||
|
queuedBytes,
|
||||||
|
});
|
||||||
|
expect(() => app._updateConnectionIndicator()).not.toThrow();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
describe('durable input delivery is not aborted by the indicator (typing-lag regression)', () => {
|
||||||
|
it('_reliableSend reaches _drainSession after updating the indicator', () => {
|
||||||
|
const { app } = makeApp({ activeSessionId: 's1', _wsState: 'connected' });
|
||||||
|
const a = app as unknown as {
|
||||||
|
_seqCounters: Map<string, number>;
|
||||||
|
_persistReliableState: () => void;
|
||||||
|
_drainSession: (id: string) => void;
|
||||||
|
_reliableSend: (id: string, data: string, useMux: boolean) => void;
|
||||||
|
};
|
||||||
|
a._seqCounters = new Map();
|
||||||
|
a._persistReliableState = vi.fn();
|
||||||
|
const drain = vi.fn();
|
||||||
|
a._drainSession = drain;
|
||||||
|
|
||||||
|
// A single keystroke. The indicator runs first; if it throws, drain is skipped.
|
||||||
|
a._reliableSend('s1', 'x', false);
|
||||||
|
|
||||||
|
expect(drain).toHaveBeenCalledWith('s1');
|
||||||
|
expect(app._pendingDeliveries.get('s1')).toHaveLength(1);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -0,0 +1,159 @@
|
|||||||
|
/**
|
||||||
|
* @fileoverview Input dispatch ordering for the durable, acknowledged delivery
|
||||||
|
* layer (`CodemanApp._sendInputAsync` → `_reliableSend` → `_drainSession`).
|
||||||
|
*
|
||||||
|
* Local replaced the upstream best-effort "coalescing fallback queue" with the
|
||||||
|
* durable per-(clientId, seq) layer in commit 1255e28 (docs/reliable-input-
|
||||||
|
* delivery.md). This suite verifies the client-side ordering guarantees of that
|
||||||
|
* layer: each input is a distinct seq-tagged frame, delivered in order over the
|
||||||
|
* WebSocket when open, serialized over HTTP POST when not, and only dropped on a
|
||||||
|
* server ACK (HTTP 2xx). Exactly-once application is covered server-side in
|
||||||
|
* test/reliable-input-dedup.test.ts; the header transport indicator ("WS"/"HTTP")
|
||||||
|
* is covered in test/connection-indicator.test.ts.
|
||||||
|
*
|
||||||
|
* Loaded via `vm` with a stubbed context (no jsdom).
|
||||||
|
*/
|
||||||
|
import { readFileSync } from 'node:fs';
|
||||||
|
import { performance } from 'node:perf_hooks';
|
||||||
|
import { resolve } from 'node:path';
|
||||||
|
import vm from 'node:vm';
|
||||||
|
import { describe, expect, it, vi } from 'vitest';
|
||||||
|
|
||||||
|
function loadCodemanAppClass() {
|
||||||
|
const constants = readFileSync(resolve(import.meta.dirname, '../src/web/public/constants.js'), 'utf8');
|
||||||
|
const source = readFileSync(resolve(import.meta.dirname, '../src/web/public/app.js'), 'utf8');
|
||||||
|
const context = vm.createContext({
|
||||||
|
console,
|
||||||
|
performance,
|
||||||
|
setInterval: vi.fn(),
|
||||||
|
clearInterval: vi.fn(),
|
||||||
|
setTimeout,
|
||||||
|
clearTimeout,
|
||||||
|
requestAnimationFrame: vi.fn(),
|
||||||
|
HTMLCanvasElement: class HTMLCanvasElement {},
|
||||||
|
WebSocket: { OPEN: 1 },
|
||||||
|
fetch: (...args: Parameters<typeof fetch>) => global.fetch(...args),
|
||||||
|
document: { addEventListener: vi.fn() },
|
||||||
|
localStorage: {
|
||||||
|
length: 0,
|
||||||
|
key: vi.fn(),
|
||||||
|
getItem: vi.fn(),
|
||||||
|
setItem: vi.fn(),
|
||||||
|
removeItem: vi.fn(),
|
||||||
|
},
|
||||||
|
window: { addEventListener: vi.fn(), removeEventListener: vi.fn() },
|
||||||
|
MobileDetection: {},
|
||||||
|
});
|
||||||
|
vm.runInContext(`${constants}\n${source}\nglobalThis.__CodemanApp = CodemanApp;`, context);
|
||||||
|
return (context as { __CodemanApp: new () => unknown }).__CodemanApp;
|
||||||
|
}
|
||||||
|
|
||||||
|
const CodemanApp = loadCodemanAppClass();
|
||||||
|
|
||||||
|
async function waitForCalls(calls: unknown[], count: number) {
|
||||||
|
for (let i = 0; i < 50; i++) {
|
||||||
|
if (calls.length >= count) return;
|
||||||
|
await new Promise((r) => setTimeout(r, 0));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type Frame = { t: string; d: string; seq: number; cid: string };
|
||||||
|
type PostBody = { input: string; seq: number; clientId: string };
|
||||||
|
|
||||||
|
type App = {
|
||||||
|
_sendInputAsync: (sessionId: string, input: string, opts?: { useMux?: boolean }) => void;
|
||||||
|
_pendingDeliveries: Map<string, Array<{ seq: number; data: string; sentAt: number }>>;
|
||||||
|
_ws: { readyState: number; send: (data: string) => void } | null;
|
||||||
|
_wsSessionId: string | null;
|
||||||
|
activeSessionId: string | null;
|
||||||
|
};
|
||||||
|
|
||||||
|
function makeApp(): App {
|
||||||
|
const app = Object.create((CodemanApp as { prototype: object }).prototype) as App & Record<string, unknown>;
|
||||||
|
app._clientId = 'c-test';
|
||||||
|
app._seqCounters = new Map();
|
||||||
|
app._pendingDeliveries = new Map();
|
||||||
|
app._postDraining = new Set();
|
||||||
|
app._persistReliableState = vi.fn();
|
||||||
|
app._persistReliableNow = vi.fn();
|
||||||
|
app._updateConnectionIndicator = vi.fn();
|
||||||
|
app.clearPendingHooks = vi.fn();
|
||||||
|
app.activeSessionId = 'session-1';
|
||||||
|
app.isOnline = true;
|
||||||
|
app._connectionStatus = 'connected';
|
||||||
|
app._ws = null;
|
||||||
|
app._wsSessionId = null;
|
||||||
|
return app as unknown as App;
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('durable input delivery — send ordering', () => {
|
||||||
|
it('delivers rapid input as distinct ordered seq frames over an open WebSocket (no coalescing)', () => {
|
||||||
|
const app = makeApp();
|
||||||
|
const frames: Frame[] = [];
|
||||||
|
app._ws = { readyState: 1, send: (d: string) => frames.push(JSON.parse(d)) };
|
||||||
|
app._wsSessionId = 'session-1';
|
||||||
|
|
||||||
|
app._sendInputAsync('session-1', 'a');
|
||||||
|
app._sendInputAsync('session-1', 'b');
|
||||||
|
app._sendInputAsync('session-1', 'c');
|
||||||
|
|
||||||
|
// Each keystroke is its own frame, in seq order — never merged into "abc".
|
||||||
|
expect(frames.map((f) => f.d)).toEqual(['a', 'b', 'c']);
|
||||||
|
expect(frames.map((f) => f.seq)).toEqual([1, 2, 3]);
|
||||||
|
expect(frames.every((f) => f.t === 'i' && f.cid === 'c-test')).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('POSTs queued input one frame at a time in seq order when no socket is open', async () => {
|
||||||
|
const app = makeApp();
|
||||||
|
const calls: PostBody[] = [];
|
||||||
|
const completions: Array<() => void> = [];
|
||||||
|
global.fetch = vi.fn(async (_url, init) => {
|
||||||
|
calls.push(JSON.parse(String(init?.body)) as PostBody);
|
||||||
|
await new Promise<void>((r) => completions.push(r));
|
||||||
|
return new Response('{}', { status: 200 });
|
||||||
|
});
|
||||||
|
|
||||||
|
app._sendInputAsync('session-1', 'a');
|
||||||
|
app._sendInputAsync('session-1', 'b');
|
||||||
|
|
||||||
|
// Serialized: only the first frame is in flight until its 2xx ACK lands.
|
||||||
|
await waitForCalls(calls, 1);
|
||||||
|
expect(calls.map((c) => c.input)).toEqual(['a']);
|
||||||
|
|
||||||
|
completions.shift()?.(); // ACK 'a'
|
||||||
|
await waitForCalls(calls, 2);
|
||||||
|
expect(calls.map((c) => c.input)).toEqual(['a', 'b']);
|
||||||
|
expect(calls.map((c) => c.seq)).toEqual([1, 2]);
|
||||||
|
|
||||||
|
completions.shift()?.();
|
||||||
|
await waitForCalls(calls, 2);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('leaves a frame queued (unacked) when HTTP delivery fails', async () => {
|
||||||
|
const app = makeApp();
|
||||||
|
const calls: PostBody[] = [];
|
||||||
|
global.fetch = vi.fn(async (_url, init) => {
|
||||||
|
calls.push(JSON.parse(String(init?.body)) as PostBody);
|
||||||
|
return new Response('busy', { status: 503 });
|
||||||
|
});
|
||||||
|
|
||||||
|
app._sendInputAsync('session-1', 'a');
|
||||||
|
await waitForCalls(calls, 1);
|
||||||
|
await new Promise((r) => setTimeout(r, 0));
|
||||||
|
|
||||||
|
// 5xx is not an ACK — the frame must survive for the sweep/reconnect to retry.
|
||||||
|
expect(app._pendingDeliveries.get('session-1')).toHaveLength(1);
|
||||||
|
expect(app._pendingDeliveries.get('session-1')?.[0].data).toBe('a');
|
||||||
|
});
|
||||||
|
|
||||||
|
it('drops a frame addressed to a vanished session (404) instead of retrying forever', async () => {
|
||||||
|
const app = makeApp();
|
||||||
|
global.fetch = vi.fn(async () => new Response('gone', { status: 404 }));
|
||||||
|
|
||||||
|
app._sendInputAsync('session-1', 'a');
|
||||||
|
await new Promise((r) => setTimeout(r, 0));
|
||||||
|
await new Promise((r) => setTimeout(r, 0));
|
||||||
|
|
||||||
|
expect(app._pendingDeliveries.get('session-1')).toBeUndefined();
|
||||||
|
});
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user