Download src/services/connectionRegistry.ts from Collabos/chatapi: direct link, hf CLI and curl.
- Browser
- Download file 2.79 kB
-
https://huggingface.co/spaces/Collabos/chatapi/resolve/main/src/services/connectionRegistry.ts
- Command line
-
hf download hf://spaces/Collabos/chatapi/src/services/connectionRegistry.ts
-
curl -L -o connectionRegistry.ts https://huggingface.co/spaces/Collabos/chatapi/resolve/main/src/services/connectionRegistry.ts
2.79 kB
| 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; | |
| }, | |
| }; | |