File size: 2,786 Bytes
87cb242
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
import WebSocket from 'ws';

/**
 * In-process map of uid → WebSocket.
 * This is NOT distributed state (that lives in Redis Bucket 1).
 * This map is only valid for THIS process instance.
 *
 * On HF Spaces (single container), this is sufficient.
 * For multi-instance horizontal scale: replace with Redis pub/sub routing.
 */

interface SocketMeta {
  ws: WebSocket;
  uid: string;
  socketId: string;
  connectedAt: number;
  lastHeartbeat: number;
}

// Two-way map for O(1) lookups in both directions
const byUid = new Map<string, SocketMeta>();
const bySocketId = new Map<string, SocketMeta>();

export const registry = {
  /** Register a new connection */
  add(socketId: string, uid: string, ws: WebSocket): void {
    const meta: SocketMeta = {
      ws,
      uid,
      socketId,
      connectedAt: Date.now(),
      lastHeartbeat: Date.now(),
    };
    // Remove any stale entry for this uid (reconnect scenario)
    const existing = byUid.get(uid);
    if (existing) {
      bySocketId.delete(existing.socketId);
      // Close the old zombie socket if it is still "open"
      if (existing.ws.readyState === WebSocket.OPEN) {
        existing.ws.terminate();
      }
    }
    byUid.set(uid, meta);
    bySocketId.set(socketId, meta);
  },

  /** Remove a connection by socketId */
  remove(socketId: string): SocketMeta | undefined {
    const meta = bySocketId.get(socketId);
    if (!meta) return undefined;
    bySocketId.delete(socketId);
    // Only remove from byUid if this socket is still the current one
    if (byUid.get(meta.uid)?.socketId === socketId) {
      byUid.delete(meta.uid);
    }
    return meta;
  },

  /** Get WebSocket for a uid */
  getByUid(uid: string): SocketMeta | undefined {
    return byUid.get(uid);
  },

  /** Get WebSocket by its socketId */
  getBySocketId(socketId: string): SocketMeta | undefined {
    return bySocketId.get(socketId);
  },

  /** Update heartbeat timestamp */
  heartbeat(socketId: string): void {
    const meta = bySocketId.get(socketId);
    if (meta) meta.lastHeartbeat = Date.now();
  },

  /** How many sockets are active right now */
  size(): number {
    return bySocketId.size;
  },

  /**
   * Zombie pruner — iterates all connections and terminates any that
   * haven't sent a heartbeat within 3× the expected interval.
   * Called by a setInterval in index.ts.
   */
  pruneZombies(maxSilenceMs: number): number {
    const now = Date.now();
    let pruned = 0;
    for (const [socketId, meta] of bySocketId) {
      if (now - meta.lastHeartbeat > maxSilenceMs) {
        meta.ws.terminate();
        bySocketId.delete(socketId);
        if (byUid.get(meta.uid)?.socketId === socketId) {
          byUid.delete(meta.uid);
        }
        pruned++;
      }
    }
    return pruned;
  },
};