Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/workflows/deploy.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
23 changes: 19 additions & 4 deletions cmd/server/websocket.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)
}

Expand Down Expand Up @@ -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()
Expand All @@ -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
}
}
}
}
Expand Down
154 changes: 154 additions & 0 deletions cmd/server/ws_heartbeat_117_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
103 changes: 87 additions & 16 deletions public/app.js
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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);
Expand All @@ -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); }

Expand Down Expand Up @@ -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() {
Expand Down Expand Up @@ -1349,6 +1419,7 @@ window.addEventListener('timestamp-mode-changed', () => {
});
window.addEventListener('DOMContentLoaded', () => {
connectWS();
setupWSResumeCheck();
setupPullToReconnect();

// --- Dark Mode ---
Expand Down
1 change: 1 addition & 0 deletions test-all.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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 ""
Expand Down
Loading
Loading