fix(BUG-8): ws adapter auto-reconnect after drop

WS adapter had no reconnect. WS dies (idle/error/close) → wsReady=null,
subscribers dead forever, display frozen until full reload.

Changes (src/storage/ws.js):
- onClose: schedule reconnect via setTimeout(500ms), ensureWs re-arms.
  Guard: disposed flag stops reconnect after dispose.
- onOpen: resubscribe all existing doc/coll subscribers (backend state
  may have changed). Re-fetch current values on RECONNECT only (skip
  first connect — initial REST fetch in subscribe* already did). Added
  everConnected flag to distinguish first vs reconnect.
- reconnectTimer unref'd (Node) to avoid hanging event loop.
- dispose(cb): set disposed, clear timer, close ws, then cb.

Also fixed test teardown leaks:
- server/index.js close(): terminate all wss.clients before wss.close().
  Reconnect test spawned new ws to server; old close hung on live conn.
- both ws test factories: port 0 (OS picks free) instead of module-local
  nextPort counter. Parallel jest workers collided on EADDRINUSE.

Tests: ws-reconnect GREEN (1.7s), ws-contract 23 GREEN. No regression.
server suite 24/24. shared 90/90.
This commit is contained in:
david raistrick
2026-07-01 18:26:42 -04:00
parent afdd72e829
commit c1d982b4a4
5 changed files with 61 additions and 17 deletions
+47 -3
View File
@@ -39,13 +39,51 @@ function createWsStorage({ baseUrl, wsUrl } = {}) {
let ws = null;
let wsReady = null;
let disposed = false;
let reconnectTimer = null;
let everConnected = false;
const RECONNECT_DELAY = 500;
function ensureWs() {
if (wsReady) return wsReady;
wsReady = new Promise((resolve, reject) => {
ws = new WebSocketImpl(WS);
const onOpen = () => resolve(ws);
const onOpen = () => {
const isReconnect = everConnected;
everConnected = true;
// resubscribe all existing subscribers after (re)connect
for (const p of docSubs.keys()) {
ws.send(JSON.stringify({ type: 'subscribe', kind: 'doc', path: p }));
}
for (const p of collSubs.keys()) {
ws.send(JSON.stringify({ type: 'subscribe', kind: 'collection', path: p }));
}
// On RECONNECT only: re-fetch current values — catches writes that
// happened while disconnected (broadcast missed). Skip on first connect
// (initial REST fetch in subscribeDoc/subscribeCollection already did).
if (isReconnect) {
for (const [p, cbs] of docSubs) {
storage.getDoc(p).then(doc => { cbs.forEach(cb => cb(doc)); }).catch(() => {});
}
for (const [p, cbs] of collSubs) {
storage.getCollection(p).then(docs => { cbs.forEach(cb => cb(docs)); }).catch(() => {});
}
}
resolve(ws);
};
const onError = (err) => { wsReady = null; reject(err instanceof Event ? new Error('ws error') : err); };
const onClose = () => { wsReady = null; };
const onClose = () => {
wsReady = null;
ws = null;
if (disposed) return;
// auto-reconnect (BUG-8): try again after delay. ensureWs() re-arms.
if (reconnectTimer) clearTimeout(reconnectTimer);
reconnectTimer = setTimeout(() => {
reconnectTimer = null;
if (!disposed) ensureWs().catch(() => {});
}, RECONNECT_DELAY);
if (reconnectTimer && typeof reconnectTimer.unref === 'function') reconnectTimer.unref();
};
const onMessage = (ev) => {
const raw = typeof ev === 'string' ? ev : (ev.data !== undefined ? ev.data : ev);
let msg; try { msg = JSON.parse(typeof raw === 'string' ? raw : raw.toString()); } catch { return; }
@@ -168,7 +206,13 @@ function createWsStorage({ baseUrl, wsUrl } = {}) {
return () => { collSubs.get(p)?.delete(cb); };
},
dispose() { if (ws) ws.close(); docSubs.clear(); collSubs.clear(); },
dispose(cb) {
disposed = true;
if (reconnectTimer) { clearTimeout(reconnectTimer); reconnectTimer = null; }
if (ws) ws.close();
docSubs.clear(); collSubs.clear();
if (typeof cb === 'function') cb();
},
_api: api,
_test: {