Release v2.10.3
This commit is contained in:
+129
-5
@@ -2,7 +2,7 @@
|
||||
|
||||
const crypto = require('crypto');
|
||||
const { WebSocketServer, WebSocket } = require('ws');
|
||||
const { isClientNameAvailable } = require('#src/data/client-name-check');
|
||||
const { findAvailableClientName, isClientNameAvailable } = require('#src/data/client-name-check');
|
||||
const { verifyPageAuthToken, verifyRequestAuth } = require('#src/request-auth');
|
||||
const PAGE_AUTH_COOKIE_NAME = 'pulse_page_auth';
|
||||
|
||||
@@ -22,6 +22,8 @@ function normalizePlayerPublicBaseUrl(pageUrl) {
|
||||
function createPlayerRuntime(options) {
|
||||
const pool = options && options.pool ? options.pool : null;
|
||||
const notifySnapshot = typeof options.notifySnapshot === 'function' ? options.notifySnapshot : null;
|
||||
const persistClientName = typeof options.persistClientName === 'function' ? options.persistClientName : null;
|
||||
const touchClientLastSeen = typeof options.touchClientLastSeen === 'function' ? options.touchClientLastSeen : null;
|
||||
const normalizeDeviceId = typeof options.normalizeDeviceId === 'function'
|
||||
? options.normalizeDeviceId
|
||||
: function (value) {
|
||||
@@ -30,7 +32,11 @@ function createPlayerRuntime(options) {
|
||||
const connectionsBySlug = new Map();
|
||||
const dashboardListenersBySlug = new Map();
|
||||
const announcementListenersBySlug = new Map();
|
||||
const pendingCommandAcks = new Map();
|
||||
const wss = new WebSocketServer({ noServer: true });
|
||||
const staleConnectionMs = Number(options && options.staleConnectionMs) > 0
|
||||
? Number(options.staleConnectionMs)
|
||||
: 3 * 60 * 1000;
|
||||
|
||||
function hasActiveClientId(clientId, currentConnection) {
|
||||
const normalizedClientId = String(clientId || '').trim();
|
||||
@@ -180,6 +186,57 @@ function createPlayerRuntime(options) {
|
||||
}
|
||||
}
|
||||
|
||||
function removeStaleConnections() {
|
||||
const cutoff = Date.now() - staleConnectionMs;
|
||||
for (const [slug, bucket] of connectionsBySlug.entries()) {
|
||||
for (const [connectionId, connection] of bucket.entries()) {
|
||||
if (connection.lastSeenAt && connection.lastSeenAt.getTime() > cutoff) {
|
||||
continue;
|
||||
}
|
||||
removeConnection(slug, connectionId);
|
||||
try {
|
||||
connection.socket.close(1000, 'Connection heartbeat expired.');
|
||||
} catch (_error) {
|
||||
}
|
||||
broadcastConnectionSnapshot(slug);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const staleConnectionSweep = setInterval(removeStaleConnections, Math.min(staleConnectionMs, 60 * 1000));
|
||||
if (typeof staleConnectionSweep.unref === 'function') {
|
||||
staleConnectionSweep.unref();
|
||||
}
|
||||
|
||||
function checkWebsocketHealth() {
|
||||
for (const [slug, bucket] of connectionsBySlug.entries()) {
|
||||
for (const [connectionId, connection] of bucket.entries()) {
|
||||
if (!connection.isAlive) {
|
||||
removeConnection(slug, connectionId);
|
||||
try {
|
||||
connection.socket.terminate();
|
||||
} catch (_error) {
|
||||
}
|
||||
broadcastConnectionSnapshot(slug);
|
||||
continue;
|
||||
}
|
||||
|
||||
connection.isAlive = false;
|
||||
try {
|
||||
connection.socket.ping();
|
||||
} catch (_error) {
|
||||
removeConnection(slug, connectionId);
|
||||
broadcastConnectionSnapshot(slug);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const websocketHealthCheck = setInterval(checkWebsocketHealth, 30 * 1000);
|
||||
if (typeof websocketHealthCheck.unref === 'function') {
|
||||
websocketHealthCheck.unref();
|
||||
}
|
||||
|
||||
function getDashboardListenerBucket(slug) {
|
||||
const key = String(slug || '').trim();
|
||||
if (!key) {
|
||||
@@ -370,7 +427,10 @@ function createPlayerRuntime(options) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
const target = bucket.get(String(connectionId || '').trim());
|
||||
const normalizedConnectionId = String(connectionId || '').trim();
|
||||
const target = bucket.get(normalizedConnectionId) || Array.from(bucket.values()).find(function (connection) {
|
||||
return String(connection && connection.clientId || '').trim() === normalizedConnectionId;
|
||||
});
|
||||
if (!target || target.socket.readyState !== WebSocket.OPEN) {
|
||||
return 0;
|
||||
}
|
||||
@@ -382,8 +442,28 @@ function createPlayerRuntime(options) {
|
||||
payload.targetConnectionId = target.id;
|
||||
payload.sentAt = new Date().toISOString();
|
||||
|
||||
const requestId = String(payload.requestId || '').trim();
|
||||
if (!requestId) {
|
||||
target.socket.send(JSON.stringify(payload));
|
||||
return 1;
|
||||
}
|
||||
|
||||
const acknowledgement = new Promise(function (resolve) {
|
||||
const timeout = setTimeout(function () {
|
||||
pendingCommandAcks.delete(requestId);
|
||||
resolve(0);
|
||||
}, 5000);
|
||||
pendingCommandAcks.set(requestId, {
|
||||
connection: target,
|
||||
resolve: function (acknowledged) {
|
||||
clearTimeout(timeout);
|
||||
pendingCommandAcks.delete(requestId);
|
||||
resolve(acknowledged ? 1 : 0);
|
||||
}
|
||||
});
|
||||
});
|
||||
target.socket.send(JSON.stringify(payload));
|
||||
return 1;
|
||||
return await acknowledgement;
|
||||
}
|
||||
|
||||
async function broadcastCommand(slug, commandOrPayload) {
|
||||
@@ -527,6 +607,7 @@ function createPlayerRuntime(options) {
|
||||
connectedAt: new Date(),
|
||||
lastSeenAt: new Date()
|
||||
};
|
||||
connection.isAlive = true;
|
||||
const bucket = getConnectionBucket(slug);
|
||||
|
||||
if (!bucket) {
|
||||
@@ -536,7 +617,11 @@ function createPlayerRuntime(options) {
|
||||
|
||||
bucket.set(connectionId, connection);
|
||||
|
||||
socket.on('message', function (rawMessage) {
|
||||
socket.on('pong', function () {
|
||||
connection.isAlive = true;
|
||||
});
|
||||
|
||||
socket.on('message', async function (rawMessage) {
|
||||
connection.lastSeenAt = new Date();
|
||||
let payload = null;
|
||||
try {
|
||||
@@ -546,6 +631,13 @@ function createPlayerRuntime(options) {
|
||||
}
|
||||
|
||||
if (!payload || payload.type !== 'state') {
|
||||
if (payload && payload.type === 'command-ack') {
|
||||
const requestId = String(payload.requestId || '').trim();
|
||||
const pendingAck = pendingCommandAcks.get(requestId);
|
||||
if (pendingAck && pendingAck.connection === connection) {
|
||||
pendingAck.resolve(payload.ok !== false);
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -561,6 +653,26 @@ function createPlayerRuntime(options) {
|
||||
if (!connection.clientName && connection.clientId) {
|
||||
connection.clientName = connection.clientId;
|
||||
}
|
||||
if (connection.clientName) {
|
||||
const connectionIdentity = connection.deviceId || connection.clientId;
|
||||
const selectedClientName = await findAvailableClientName(pool, connection.clientName, connectionIdentity, snapshotAllConnections());
|
||||
if (selectedClientName && selectedClientName !== connection.clientName) {
|
||||
connection.clientName = selectedClientName;
|
||||
if (persistClientName) {
|
||||
await persistClientName(connection.deviceId, selectedClientName);
|
||||
}
|
||||
if (socket.readyState === WebSocket.OPEN) {
|
||||
socket.send(JSON.stringify({ type: 'client-name-updated', clientName: selectedClientName }));
|
||||
}
|
||||
}
|
||||
}
|
||||
if (touchClientLastSeen) {
|
||||
try {
|
||||
await touchClientLastSeen(connection.clientId);
|
||||
} catch (_error) {
|
||||
// A heartbeat failure must not interrupt playback or websocket state.
|
||||
}
|
||||
}
|
||||
connection.userAgent = payload.userAgent ? String(payload.userAgent).trim() : connection.userAgent;
|
||||
connection.viewport = payload.viewport && typeof payload.viewport === 'object' ? payload.viewport : connection.viewport;
|
||||
connection.page = payload.page ? String(payload.page).trim() : connection.page;
|
||||
@@ -581,12 +693,24 @@ function createPlayerRuntime(options) {
|
||||
broadcastConnectionSnapshot(slug);
|
||||
});
|
||||
|
||||
socket.on('close', function () {
|
||||
socket.on('close', function (code, reason) {
|
||||
pendingCommandAcks.forEach(function (pendingAck, requestId) {
|
||||
if (pendingAck.connection === connection) {
|
||||
pendingAck.resolve(false);
|
||||
pendingCommandAcks.delete(requestId);
|
||||
}
|
||||
});
|
||||
removeConnection(slug, connectionId);
|
||||
broadcastConnectionSnapshot(slug);
|
||||
});
|
||||
|
||||
socket.on('error', function () {
|
||||
pendingCommandAcks.forEach(function (pendingAck, requestId) {
|
||||
if (pendingAck.connection === connection) {
|
||||
pendingAck.resolve(false);
|
||||
pendingCommandAcks.delete(requestId);
|
||||
}
|
||||
});
|
||||
removeConnection(slug, connectionId);
|
||||
broadcastConnectionSnapshot(slug);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user