diff --git a/.github/workflows/deploy.yml b/.github/workflows/deploy.yml index 0549286d6..1519a8b19 100644 --- a/.github/workflows/deploy.yml +++ b/.github/workflows/deploy.yml @@ -162,6 +162,7 @@ jobs: node test-channel-issue-1101.js node test-observer-iata-1188.js node test-pull-to-reconnect-1091.js + node test-issue-117-ws-watchdog.js node test-issue-111-drawer-version.js node test-channel-fluid-layout.js node test-issue-1279-p2-code-filter.js diff --git a/cmd/server/websocket.go b/cmd/server/websocket.go index d470b850a..7b492e622 100644 --- a/cmd/server/websocket.go +++ b/cmd/server/websocket.go @@ -19,6 +19,7 @@ type Hub struct { upgrader websocket.Upgrader allowedOrigins []string // exact-match allowlist for /ws CheckOrigin (see SetAllowedOrigins) limits *wsLimiter // #1794: per-IP caps and deny list; nil allows everything + pingInterval time.Duration // writePump tick: protocol ping + app heartbeat (#117) } // SetAllowedOrigins configures the exact-match origin allowlist consulted by @@ -74,6 +75,15 @@ func (h *Hub) checkOrigin(r *http.Request) bool { return false } +// wsHeartbeat is written to every client on each ping tick (#117). Browser +// JS never sees protocol ping frames, so without a frame the page receives a +// quiet mesh and a silently dead (half-open) socket look the same to it. +// public/app.js matches these exact bytes (WS_HEARTBEAT) and drops a socket +// that has been silent for WS_STALE_MS; ws_heartbeat_117_test.go keeps the +// two in step. Tabs still running an app.js from before this change treat +// it as an ordinary message of an unknown type (see the #117 PR). +var wsHeartbeat = []byte(`{"type":"heartbeat"}`) + // Client is a single WebSocket connection. type Client struct { conn *websocket.Conn @@ -109,7 +119,8 @@ func (h *Hub) ConfigureLimits(maxConnsPerIP, upgradesPerMin int, trustedProxies, func NewHub() *Hub { h := &Hub{ - clients: make(map[*Client]bool), + clients: make(map[*Client]bool), + pingInterval: 30 * time.Second, } h.upgrader = websocket.Upgrader{ ReadBufferSize: 1024, @@ -214,7 +225,7 @@ func (h *Hub) ServeWS(w http.ResponseWriter, r *http.Request) { } h.Register(client) - go client.writePump() + go client.writePump(h.pingInterval) go client.readPump(h) } @@ -248,8 +259,8 @@ func (c *Client) readPump(hub *Hub) { } } -func (c *Client) writePump() { - ticker := time.NewTicker(30 * time.Second) +func (c *Client) writePump(pingInterval time.Duration) { + ticker := time.NewTicker(pingInterval) defer func() { ticker.Stop() c.conn.Close() @@ -270,6 +281,10 @@ func (c *Client) writePump() { if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil { return } + // #117: same tick, same (only) writer goroutine; see wsHeartbeat. + if err := c.conn.WriteMessage(websocket.TextMessage, wsHeartbeat); err != nil { + return + } } } } diff --git a/cmd/server/ws_heartbeat_117_test.go b/cmd/server/ws_heartbeat_117_test.go new file mode 100644 index 000000000..0ade44745 --- /dev/null +++ b/cmd/server/ws_heartbeat_117_test.go @@ -0,0 +1,154 @@ +package main + +import ( + "bytes" + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "regexp" + "strconv" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/gorilla/websocket" +) + +// Issue #117: browsers cannot see WebSocket ping frames, so the server also +// writes a small application-level heartbeat on the existing ping tick, from +// the single writer goroutine. public/app.js replaces a socket that has been +// silent for longer than its threshold. + +func dialHeartbeatHub(t *testing.T, interval time.Duration) (*Hub, *websocket.Conn, *atomic.Int32) { + t.Helper() + hub := NewHub() + hub.pingInterval = interval + srv := httptest.NewServer(http.HandlerFunc(hub.ServeWS)) + t.Cleanup(srv.Close) + conn, _, err := websocket.DefaultDialer.Dial("ws"+srv.URL[4:], nil) + if err != nil { + t.Fatalf("dial: %v", err) + } + t.Cleanup(func() { conn.Close() }) + var pings atomic.Int32 + conn.SetPingHandler(func(data string) error { + pings.Add(1) + return conn.WriteControl(websocket.PongMessage, []byte(data), time.Now().Add(time.Second)) + }) + return hub, conn, &pings +} + +func TestHubDefaultPingInterval_117(t *testing.T) { + // public/app.js WS_STALE_MS is derived from this interval (two missed + // heartbeats plus slack); TestClientStaleThresholdCoversTwoHeartbeats_117 + // checks the pair. + if got := NewHub().pingInterval; got != 30*time.Second { + t.Fatalf("default ping interval = %v, want 30s", got) + } +} + +func TestWritePumpSendsHeartbeatOnEveryPingTick_117(t *testing.T) { + _, conn, pings := dialHeartbeatHub(t, 25*time.Millisecond) + const want = 4 + got := 0 + deadline := time.Now().Add(5 * time.Second) + for got < want { + conn.SetReadDeadline(deadline) + typ, msg, err := conn.ReadMessage() + if err != nil { + t.Fatalf("after %d heartbeats: read: %v (no application heartbeat reaches the page)", got, err) + } + if typ != websocket.TextMessage || !bytes.Equal(msg, wsHeartbeat) { + t.Fatalf("unexpected frame %d %q", typ, msg) + } + got++ + } + // The protocol-level ping is still sent on the same tick; it is written + // before the heartbeat, so by now at least as many pings have arrived. + if p := int(pings.Load()); p < want { + t.Fatalf("%d pings for %d heartbeats: the protocol ping must stay", p, want) + } +} + +// The heartbeat bytes are what app.js matches exactly, and a plain JSON +// object with a type other than "packet" for anything that parses it. +func TestHeartbeatFrameShape_117(t *testing.T) { + var v struct { + Type string `json:"type"` + } + if err := json.Unmarshal(wsHeartbeat, &v); err != nil || v.Type != "heartbeat" { + t.Fatalf("heartbeat %q: type=%q err=%v", wsHeartbeat, v.Type, err) + } + if len(wsHeartbeat) > 32 { + t.Fatalf("heartbeat is %d bytes; keep it small", len(wsHeartbeat)) + } +} + +// Broadcasts still arrive byte-for-byte, interleaved with heartbeats. +func TestBroadcastsUnchangedAlongsideHeartbeats_117(t *testing.T) { + hub, conn, _ := dialHeartbeatHub(t, 20*time.Millisecond) + for hub.ClientCount() == 0 { + time.Sleep(time.Millisecond) + } + type packetMsg struct { + Type string `json:"type"` + Data map[string]string `json:"data"` + } + want, _ := json.Marshal(packetMsg{Type: "packet", Data: map[string]string{"hash": "abc"}}) + hub.Broadcast(packetMsg{Type: "packet", Data: map[string]string{"hash": "abc"}}) + deadline := time.Now().Add(5 * time.Second) + for { + conn.SetReadDeadline(deadline) + _, msg, err := conn.ReadMessage() + if err != nil { + t.Fatalf("read: %v", err) + } + if bytes.Equal(msg, wsHeartbeat) { + continue + } + if !bytes.Equal(msg, want) { + t.Fatalf("broadcast changed: got %q, want %q", msg, want) + } + return + } +} + +// Dead-client detection is unchanged: the read deadline is refreshed only by +// pongs (a client never sends heartbeats), and all writes, the heartbeat +// included, stay on the one writer goroutine per client. +func TestHeartbeatDoesNotChangeDeadClientHandling_117(t *testing.T) { + src, err := os.ReadFile("websocket.go") + if err != nil { + t.Fatal(err) + } + s := string(src) + if strings.Count(s, "SetReadDeadline(time.Now().Add(60 * time.Second))") != 2 { + t.Fatal("readPump's 60s read deadline / pong refresh changed") + } + if strings.Count(s, "go client.writePump(") != 1 { + t.Fatal("there must be exactly one writer goroutine per client") + } +} + +// app.js must treat two missed heartbeats (plus slack) as stale, never less, +// and match the server's heartbeat bytes exactly. +func TestClientStaleThresholdCoversTwoHeartbeats_117(t *testing.T) { + app, err := os.ReadFile("../../public/app.js") + if err != nil { + t.Fatal(err) + } + m := regexp.MustCompile(`const WS_STALE_MS = (\d+);`).FindSubmatch(app) + if m == nil { + t.Fatal("WS_STALE_MS not found in public/app.js") + } + staleMs, _ := strconv.Atoi(string(m[1])) + if stale := time.Duration(staleMs) * time.Millisecond; stale <= 2*NewHub().pingInterval { + t.Fatalf("WS_STALE_MS %v does not tolerate one lost heartbeat at %v", stale, NewHub().pingInterval) + } + h := regexp.MustCompile("const WS_HEARTBEAT = '([^']*)';").FindSubmatch(app) + if h == nil || !bytes.Equal(h[1], wsHeartbeat) { + t.Fatalf("app.js WS_HEARTBEAT %q != server %q", h, wsHeartbeat) + } +} diff --git a/public/app.js b/public/app.js index 149798478..7b9fef991 100644 --- a/public/app.js +++ b/public/app.js @@ -678,6 +678,20 @@ function buildHexLegend(ranges) { let ws = null; let wsListeners = []; +// #117: a half-open connection (a proxy or NAT dropping state, a laptop that +// slept, a handshake that never completes) can leave the shared socket +// looking OPEN but silent, with no onclose, so every view stops updating. +// The server writes WS_HEARTBEAT on each 30 s ping tick +// (cmd/server/websocket.go wsHeartbeat), and any received frame counts as +// life. A socket silent for WS_STALE_MS, measured from its creation, is +// replaced once. 75 s = two heartbeat intervals plus slack, so one late or +// lost heartbeat is tolerated. +const WS_STALE_MS = 75000; +const WS_HEARTBEAT = '{"type":"heartbeat"}'; +let wsLastFrameAt = 0; +let wsWatchdogTimer = null; +let wsReconnectTimer = null; // the one pending reconnect, if any + // --- Brand-logo packet-driven pulse (#1173) --- // Replaces the legacy live-dot indicator. Class-toggle only (CSS animations); colors come from // --logo-accent / --logo-accent-hi tokens. Test seam at window.__corescopeLogo. @@ -821,19 +835,69 @@ const Logo = (function () { return api; })(); +// Detach the current socket's handlers, then close it. A replaced socket's +// close event can arrive much later (a half-open connection waits out the +// closing handshake) and must not schedule another connection. +function dropWS() { + clearTimeout(wsWatchdogTimer); + wsWatchdogTimer = null; + if (!ws) return; + const old = ws; + ws = null; + old.onopen = old.onclose = old.onerror = old.onmessage = null; + try { old.close(); } catch (_) {} +} + +// Watchdog and resume check. Inert while a reconnect is already scheduled +// (ordinary close keeps its configured delay) or with no socket. +function checkWSLiveness() { + clearTimeout(wsWatchdogTimer); + wsWatchdogTimer = null; + if (!ws || wsReconnectTimer) return; + const silentMs = Date.now() - wsLastFrameAt; + // Date.now(), not performance.now(): a tab resumed from sleep must be + // measured against real elapsed time. A negative reading means the wall + // clock stepped back, so the silence cannot be measured; treat it as stale + // instead of re-arming for the size of the step. Either step direction + // costs at most one extra reconnect. + if (silentMs >= 0 && silentMs < WS_STALE_MS) { + wsWatchdogTimer = setTimeout(checkWSLiveness, WS_STALE_MS - silentMs); + return; + } + Logo.setConnected(false); + connectWS(); +} + +// Opens the shared socket, replacing (detaching + closing) any previous one +// and cancelling a pending reconnect, so onclose, the watchdog, resume +// checks and pull-to-reconnect can never leave two live sockets. function connectWS() { + clearTimeout(wsReconnectTimer); + wsReconnectTimer = null; + dropWS(); const proto = location.protocol === 'https:' ? 'wss:' : 'ws:'; - ws = new WebSocket(`${proto}//${location.host}`); - ws.onopen = () => Logo.setConnected(true); - ws.onclose = () => { + const sock = new WebSocket(`${proto}//${location.host}`); + ws = sock; + wsLastFrameAt = Date.now(); // from creation: a stuck handshake is caught too + wsWatchdogTimer = setTimeout(checkWSLiveness, WS_STALE_MS); + sock.onopen = () => Logo.setConnected(true); + sock.onclose = () => { + if (ws !== sock) return; + clearTimeout(wsWatchdogTimer); + wsWatchdogTimer = null; Logo.setConnected(false); // WS_RECONNECT_MS comes from roles.js (operator setting `wsReconnectMs`). // It used to be honoured only by the live map's private socket; now that // every view shares this one, the setting applies here or nowhere. - setTimeout(connectWS, window.WS_RECONNECT_MS || 3000); + clearTimeout(wsReconnectTimer); + wsReconnectTimer = setTimeout(connectWS, window.WS_RECONNECT_MS || 3000); }; - ws.onerror = () => ws.close(); - ws.onmessage = (e) => { + sock.onerror = () => sock.close(); + sock.onmessage = (e) => { + wsLastFrameAt = Date.now(); + // Consumed here, before the logo pulse, cache invalidation and every + // onWS listener (and so every pause buffer). + if (e.data === WS_HEARTBEAT) return; Logo.pulse(e); try { const msg = JSON.parse(e.data); @@ -850,6 +914,16 @@ function connectWS() { }; } +// Timers in a hidden or sleeping tab can run late (or not at all), so check +// as soon as the page is visible or back online instead of waiting out a +// watchdog that may be minutes behind. +function setupWSResumeCheck() { + document.addEventListener('visibilitychange', () => { + if (!document.hidden) checkWSLiveness(); + }); + window.addEventListener('online', checkWSLiveness); +} + function onWS(fn) { wsListeners.push(fn); } function offWS(fn) { wsListeners = wsListeners.filter(f => f !== fn); } @@ -908,16 +982,12 @@ function pullReconnect() { // If WS is connected (readyState OPEN), give a brief "Connected" // confirmation but still cycle so the user sees fresh data. const wasOpen = ws && ws.readyState === 1; - if (wasOpen) { - _showPullToast('Connected', true); - // Fast cycle: close and let onclose reconnect immediately - try { ws.close(); } catch (e) {} - } else { - _showPullToast('Reconnecting…', true); - try { if (ws) ws.close(); } catch (e) {} - // onclose handler schedules reconnect; force one now in case ws was null - try { connectWS(); } catch (e) {} - } + _showPullToast(wasOpen ? 'Connected' : 'Reconnecting…', true); + // #117: replace the socket now in both cases. An OPEN socket may be + // half-open, and its close event can take about a minute after close(); + // connectWS() detaches and closes the old socket and cancels a pending + // reconnect itself, so this cannot stack sockets. + try { connectWS(); } catch (e) {} } function _isTouchDevice() { @@ -1349,6 +1419,7 @@ window.addEventListener('timestamp-mode-changed', () => { }); window.addEventListener('DOMContentLoaded', () => { connectWS(); + setupWSResumeCheck(); setupPullToReconnect(); // --- Dark Mode --- diff --git a/test-all.sh b/test-all.sh index 773184657..e54213f91 100755 --- a/test-all.sh +++ b/test-all.sh @@ -114,6 +114,7 @@ node test-network-digest-tool.js node test-position-gaps-tool.js node test-gps-sanity-tool.js node test-map-scope-filter.js +node test-issue-117-ws-watchdog.js node test-issue-111-drawer-version.js echo "" diff --git a/test-issue-117-ws-watchdog.js b/test-issue-117-ws-watchdog.js new file mode 100644 index 000000000..51b9ca018 --- /dev/null +++ b/test-issue-117-ws-watchdog.js @@ -0,0 +1,352 @@ +/* test-issue-117-ws-watchdog.js — shared WebSocket liveness (#117). + * + * Loads the real public/app.js in a vm with a fake clock, fake timers and a + * fake WebSocket, boots it through its DOMContentLoaded listeners, and checks: + * - a silent OPEN socket and a handshake that never opens are replaced + * once after WS_STALE_MS, measured from socket creation; + * - any frame (heartbeat or traffic) keeps the socket; + * - heartbeats are consumed before the logo pulse, cache invalidation and + * onWS listeners (so before every pause buffer); packets are unchanged; + * - wall-clock steps, online and visible-tab resume; + * - ordinary close keeps the configured reconnect delay; + * - onclose, watchdog, resume and pull-to-reconnect never leave more than + * one live socket or more than one pending reconnect, and a replaced + * socket's late events are detached. + */ +'use strict'; +const vm = require('vm'); +const fs = require('fs'); +const assert = require('assert'); + +console.log('--- test-issue-117-ws-watchdog.js ---'); +let passed = 0, failed = 0; +function test(name, fn) { + try { fn(); passed++; console.log(' ✅ ' + name); } + catch (e) { failed++; console.log(' ❌ ' + name + ': ' + e.message); } +} + +const APP = fs.readFileSync(__dirname + '/public/app.js', 'utf8'); +const STALE_MATCH = APP.match(/const WS_STALE_MS = (\d+);/); +const STALE = STALE_MATCH ? Number(STALE_MATCH[1]) : 75000; // master has none: use the documented value +const HEARTBEAT = '{"type":"heartbeat"}'; +const PACKET = JSON.stringify({ type: 'packet', data: { packet: { hash: 'abc' } } }); + +function makeBox(opts) { + opts = opts || {}; + // Fake clock: `mono` drives timers, `wall` is what Date.now() returns. + // They move together unless a test steps the wall clock. + const clock = { mono: 0, wall: 1790000000000 }; + let seq = 0; + const timers = new Map(); + function setTimeout_(fn, ms) { + const id = ++seq; + timers.set(id, { fn, at: clock.mono + Math.max(0, Number(ms) || 0), id }); + return id; + } + function clearTimeout_(id) { timers.delete(id); } + function advance(ms) { + const end = clock.mono + ms; + for (;;) { + let next = null; + for (const t of timers.values()) if (t.at <= end && (!next || t.at < next.at || (t.at === next.at && t.id < next.id))) next = t; + if (!next) break; + clock.wall += next.at - clock.mono; + clock.mono = next.at; + timers.delete(next.id); + next.fn(); + } + clock.wall += end - clock.mono; + clock.mono = end; + } + class FakeDate extends Date { + constructor(...a) { if (a.length) super(...a); else super(clock.wall); } + static now() { return clock.wall; } + } + + const sockets = []; + function FakeWS(url) { + this.url = url; + this.readyState = 0; + this.closeCalled = false; + this.onopen = this.onclose = this.onerror = this.onmessage = null; + sockets.push(this); + } + FakeWS.prototype.close = function () { + if (this.closeCalled) return; + this.closeCalled = true; + this.readyState = 2; + // The close event arrives later (after the closing handshake; much + // later on a half-open connection). + const self = this; + setTimeout_(function () { self.readyState = 3; if (self.onclose) self.onclose({}); }, opts.closeEventMs == null ? 50 : opts.closeEventMs); + }; + FakeWS.prototype.send = function () {}; + FakeWS.prototype.open = function () { this.readyState = 1; if (this.onopen) this.onopen({}); }; + FakeWS.prototype.recv = function (data) { if (this.onmessage) this.onmessage({ data }); }; + FakeWS.prototype.serverClose = function () { this.readyState = 3; if (this.onclose) this.onclose({}); }; + + const docListeners = {}, winListeners = {}; + function el(id) { + return { + id, style: {}, dataset: {}, textContent: '', innerHTML: '', value: '', + classList: { add() {}, remove() {}, toggle() {}, contains() { return false; } }, + addEventListener() {}, removeEventListener() {}, setAttribute() {}, getAttribute() { return null; }, + appendChild(c) { return c; }, remove() {}, querySelector() { return null; }, querySelectorAll() { return []; }, + contains() { return false; }, + }; + } + const doc = { + hidden: false, readyState: 'loading', + documentElement: Object.assign(el('html'), { scrollTop: 0, style: { setProperty() {} } }), + body: el('body'), head: el('head'), + createElement: (t) => el(t), getElementById: () => null, + querySelector: () => null, querySelectorAll: () => [], + addEventListener(ev, fn) { (docListeners[ev] = docListeners[ev] || []).push(fn); }, + removeEventListener() {}, + }; + const win = { + addEventListener(ev, fn) { (winListeners[ev] = winListeners[ev] || []).push(fn); }, + removeEventListener() {}, dispatchEvent() { return true; }, + matchMedia: () => ({ matches: false, addEventListener() {}, addListener() {} }), + }; + if (opts.reconnectMs) win.WS_RECONNECT_MS = opts.reconnectMs; + const ctx = { + console: { log() {}, warn() {}, error() {}, info() {}, debug() {} }, + setTimeout: setTimeout_, clearTimeout: clearTimeout_, + setInterval: () => 0, clearInterval() {}, + Date: FakeDate, Math, JSON, Object, Array, String, Number, Boolean, Error, RegExp, Map, Set, WeakMap, Symbol, Promise, + requestAnimationFrame: () => 0, cancelAnimationFrame() {}, + performance: { now: () => clock.mono }, + location: { protocol: 'http:', host: 'localhost', hash: '', pathname: '/', search: '' }, + navigator: { userAgent: 'test', maxTouchPoints: 0 }, + WebSocket: FakeWS, + fetch: () => new Promise(() => {}), + localStorage: { getItem: () => null, setItem() {}, removeItem() {} }, + sessionStorage: { getItem: () => null, setItem() {}, removeItem() {} }, + document: doc, window: win, + CustomEvent: function (type, init) { this.type = type; this.detail = (init || {}).detail; }, + Event: function (type) { this.type = type; }, + MutationObserver: function () { this.observe = function () {}; this.disconnect = function () {}; }, + ResizeObserver: function () { this.observe = function () {}; this.disconnect = function () {}; }, + }; + Object.assign(win, { location: ctx.location, localStorage: ctx.localStorage, document: doc, navigator: ctx.navigator }); + ctx.self = win; ctx.globalThis = ctx; + vm.createContext(ctx); + vm.runInContext(APP, ctx); + vm.runInContext('window.connectWS = connectWS; window.onWS = onWS; window.offWS = offWS; window.pullReconnect = pullReconnect; window.__api = api;', ctx); + // Boot through the page's own startup listeners; a stubbed-DOM failure + // after the socket wiring is irrelevant here. + for (const fn of docListeners.DOMContentLoaded || []) { try { fn({}); } catch (_) {} } + for (const fn of winListeners.DOMContentLoaded || []) { try { fn({}); } catch (_) {} } + + const box = { + ctx, clock, timers, sockets, advance, + live() { return sockets.filter((s) => !s.closeCalled && s.readyState !== 3); }, + current() { return sockets[sockets.length - 1]; }, + stepWall(ms) { clock.wall += ms; }, + setHidden(h) { doc.hidden = h; for (const fn of docListeners.visibilitychange || []) fn({}); }, + online() { for (const fn of winListeners.online || []) fn({}); }, + // pending timers whose callback is connectWS itself (the reconnect) + pendingReconnects() { + let n = 0; + for (const t of timers.values()) if (t.fn === ctx.connectWS) n++; + return n; + }, + }; + return box; +} + +function booted(opts) { + const b = makeBox(opts); + assert.strictEqual(b.sockets.length, 1, 'startup should open exactly one socket'); + return b; +} + +console.log('\n=== silent sockets are replaced once ==='); + +test('a silent OPEN socket is replaced exactly once after WS_STALE_MS', () => { + const b = booted(); + const s = b.current(); s.open(); + b.advance(10000); s.recv(PACKET); + b.advance(STALE - 1); + assert.strictEqual(b.sockets.length, 1, 'replaced before the threshold'); + b.advance(2); + assert.strictEqual(b.sockets.length, 2, 'a socket silent for WS_STALE_MS was not replaced'); + assert(s.closeCalled, 'the stale socket was not closed'); + assert.strictEqual(s.onmessage, null, 'the stale socket still has its handlers'); + assert.strictEqual(b.live().length, 1); + b.advance(1000); // the old socket's close event arrives: nothing more + assert.strictEqual(b.sockets.length, 2, 'the replaced socket\'s close event scheduled another connection'); + assert.strictEqual(b.pendingReconnects(), 0); +}); + +test('a handshake that never opens is replaced after WS_STALE_MS from creation', () => { + const b = booted(); + b.advance(STALE - 1); + assert.strictEqual(b.sockets.length, 1); + b.advance(2); + assert.strictEqual(b.sockets.length, 2, 'a stuck handshake was never replaced'); + assert.strictEqual(b.live().length, 1); +}); + +test('heartbeats keep a quiet socket for 10 minutes', () => { + const b = booted(); + const s = b.current(); s.open(); + for (let t = 0; t < 600000; t += 30000) { b.advance(30000); s.recv(HEARTBEAT); } + assert.strictEqual(b.sockets.length, 1, 'a socket with regular heartbeats was replaced'); +}); + +test('packet traffic alone keeps the socket', () => { + const b = booted(); + const s = b.current(); s.open(); + for (let t = 0; t < 600000; t += STALE / 2) { b.advance(STALE / 2); s.recv(PACKET); } + assert.strictEqual(b.sockets.length, 1); +}); + +console.log('\n=== heartbeats are consumed first ==='); + +test('a heartbeat reaches no onWS listener, no logo pulse and no cache invalidation', () => { + const b = booted(); + const s = b.current(); s.open(); + const got = []; + b.ctx.window.onWS((m) => got.push(m)); + const logo = b.ctx.window.__corescopeLogo; + const before = logo.stats.triggered + logo.stats.dropped; + s.recv(HEARTBEAT); + assert.strictEqual(got.length, 0, 'the heartbeat was dispatched to a listener (and so to pause buffers)'); + assert.strictEqual(logo.stats.triggered + logo.stats.dropped, before, 'the heartbeat pulsed the logo'); + assert(!b.ctx.window.__api._invalidateTimer, 'the heartbeat scheduled cache invalidation'); +}); + +test('packet messages are dispatched unchanged', () => { + const b = booted(); + const s = b.current(); s.open(); + const got = []; + b.ctx.window.onWS((m) => got.push(m)); + s.recv(PACKET); + assert.strictEqual(got.length, 1); + assert.deepStrictEqual(JSON.parse(JSON.stringify(got[0])), JSON.parse(PACKET)); +}); + +console.log('\n=== clock steps and resume ==='); + +test('the wall clock stepping back does not postpone detection by the size of the step', () => { + const b = booted(); + const s = b.current(); s.open(); + b.advance(1000); s.recv(PACKET); + b.stepWall(-3600000); // one hour back + b.advance(STALE + 1); + assert.strictEqual(b.sockets.length, 2, 'after a backward clock step the silent socket was not replaced within WS_STALE_MS'); +}); + +test('a forward step (sleep) is caught on visible-tab resume at once', () => { + const b = booted(); + const s = b.current(); s.open(); s.recv(PACKET); + b.setHidden(true); + b.stepWall(10 * 60000); // asleep: timers did not run + assert.strictEqual(b.sockets.length, 1, 'hiding the tab must not trigger a check'); + b.setHidden(false); + assert.strictEqual(b.sockets.length, 2, 'resume after a long silence did not replace the socket'); + assert.strictEqual(b.live().length, 1); +}); + +test('online after a long silence reconnects at once; with recent traffic it does not', () => { + const b = booted(); + const s = b.current(); s.open(); s.recv(PACKET); + b.advance(5000); + b.online(); + assert.strictEqual(b.sockets.length, 1, 'online with recent traffic replaced a healthy socket'); + b.stepWall(STALE); + b.online(); + assert.strictEqual(b.sockets.length, 2, 'online after silence did not reconnect'); +}); + +test('repeated resume events open one socket', () => { + const b = booted(); + const s = b.current(); s.open(); + b.stepWall(STALE * 3); + b.setHidden(false); b.online(); b.setHidden(false); b.online(); + assert.strictEqual(b.sockets.length, 2, b.sockets.length + ' sockets after four resume events'); + b.advance(STALE - 1); + assert.strictEqual(b.sockets.length, 2); + assert.strictEqual(b.live().length, 1); +}); + +console.log('\n=== close, pull and races ==='); + +test('an ordinary close keeps the configured reconnect delay, and schedules one reconnect', () => { + const b = booted({ reconnectMs: 5000 }); + const s = b.current(); s.open(); s.recv(PACKET); + s.serverClose(); + assert.strictEqual(b.pendingReconnects(), 1, 'onclose must leave exactly one pending reconnect'); + b.advance(4999); + assert.strictEqual(b.sockets.length, 1, 'reconnected before the configured delay'); + b.advance(1); + assert.strictEqual(b.sockets.length, 2); + b.advance(STALE * 2); // no watchdog on the closed socket + assert.strictEqual(b.live().length, 1); +}); + +test('resume and the watchdog during a pending reconnect do not add a socket', () => { + const b = booted({ reconnectMs: 5000 }); + const s = b.current(); s.open(); + b.advance(STALE - 1000); + s.serverClose(); + b.stepWall(STALE); + b.setHidden(false); b.online(); + b.advance(4999); + assert.strictEqual(b.sockets.length, 1, 'a resume check bypassed the pending reconnect'); + b.advance(1); + assert.strictEqual(b.sockets.length, 2); + assert.strictEqual(b.live().length, 1); + assert.strictEqual(b.pendingReconnects(), 0); +}); + +test('a pull during the reconnect delay cancels the pending reconnect', () => { + const b = booted({ reconnectMs: 5000 }); + const s = b.current(); s.open(); + s.serverClose(); + b.advance(1000); + b.ctx.window.pullReconnect(); + assert.strictEqual(b.pendingReconnects(), 0, 'the reconnect scheduled by onclose is still pending'); + b.current().open(); + b.advance(10000); + assert.strictEqual(b.sockets.length, 2, b.sockets.length + ' sockets: the old reconnect replaced the pulled socket'); + assert.strictEqual(b.live().length, 1); +}); + +test('pull-to-reconnect on a socket that is not open leaves one socket', () => { + const b = booted(); + const s = b.current(); // still CONNECTING + b.ctx.window.pullReconnect(); + b.advance(10000); + assert.strictEqual(b.live().length, 1, b.live().length + ' live sockets after pull-to-reconnect'); + assert(s.closeCalled); + assert.strictEqual(b.sockets.length, 2, b.sockets.length + ' sockets constructed'); +}); + +test('pull-to-reconnect on an OPEN (possibly half-open) socket replaces it at once', () => { + const b = booted({ closeEventMs: 60000 }); + const s = b.current(); s.open(); + b.ctx.window.pullReconnect(); + assert.strictEqual(b.sockets.length, 2, 'pull waited for the old socket\'s close event'); + b.advance(120000); + assert.strictEqual(b.live().length, 1); +}); + +test('pull, watchdog and close racing each other end with one socket and no pending reconnect', () => { + const b = booted({ reconnectMs: 3000 }); + let s = b.current(); s.open(); + b.ctx.window.pullReconnect(); + b.ctx.window.pullReconnect(); + s = b.current(); s.serverClose(); + b.ctx.window.pullReconnect(); + b.advance(STALE + 10); // watchdog of a never-opened socket + b.setHidden(false); b.online(); + b.advance(20000); + assert.strictEqual(b.live().length, 1, b.live().length + ' live sockets'); + assert(b.pendingReconnects() <= 1); +}); + +console.log(`\n${passed} passed, ${failed} failed`); +if (failed > 0) process.exit(1);