Files
pulse-signage/src/web/bootstrap.js
T
2026-08-03 12:45:27 +01:00

277 lines
9.3 KiB
JavaScript

// Web-to-player websocket bridge and dashboard state synchronizer.
const { WebSocketServer, WebSocket } = require('ws');
const { createDashboardStateService } = require('./lib/dashboard-state');
const { createUploadSyncService } = require('./lib/media');
const { createRequestAuthHeaders } = require('#src/request-auth');
function createWebBootstrap(options) {
const pool = options && options.pool;
const common = options && options.common;
const configuredPlayerInternalBaseUrl = String(options && options.playerInternalBaseUrl || '').replace(/\/$/, '');
const playerPublicBaseUrl = String(options && options.playerPublicBaseUrl || '').replace(/\/$/, '');
const uploadDir = String(options && options.uploadDir || '').trim();
const dashboardRefreshIntervalMs = 5000;
const formatDashboardDate = options && options.formatDashboardDate;
const notifyPlayerScreens = options && options.notifyPlayerScreens;
const backgroundTaskQueue = options && options.backgroundTaskQueue;
if (!pool || !common || !uploadDir || typeof formatDashboardDate !== 'function' || typeof notifyPlayerScreens !== 'function') {
throw new Error('createWebBootstrap requires the web bootstrap dependencies.');
}
const dashboardWs = new WebSocketServer({ noServer: true });
const dashboardClients = new Set();
const playerSnapshotCache = new Map();
const playerSnapshotSockets = new Map();
let dashboardRefreshInFlight = null;
let broadcastDashboardState = null;
let playerInternalBaseUrl = null;
let playerInternalBaseUrlPromise = null;
async function getPlayerInternalBaseUrl() {
if (playerInternalBaseUrl) {
return playerInternalBaseUrl;
}
if (playerInternalBaseUrlPromise) {
return playerInternalBaseUrlPromise;
}
playerInternalBaseUrlPromise = (async function () {
if (!pool) {
return configuredPlayerInternalBaseUrl || null;
}
try {
const [rows] = await pool.query(
`SELECT internal_base_url
FROM d_players
WHERE device_id = '1'
LIMIT 1`
);
const resolvedBaseUrl = String(rows && rows[0] && rows[0].internal_base_url || '').trim().replace(/\/$/, '');
return resolvedBaseUrl || configuredPlayerInternalBaseUrl || null;
} catch (_error) {
return configuredPlayerInternalBaseUrl || null;
}
})().then(function (baseUrl) {
playerInternalBaseUrl = baseUrl || null;
playerInternalBaseUrlPromise = null;
return playerInternalBaseUrl;
}, function () {
playerInternalBaseUrlPromise = null;
return configuredPlayerInternalBaseUrl || null;
});
return playerInternalBaseUrlPromise;
}
async function getPlayerSnapshotSocketUrl(slug) {
const resolvedPlayerInternalBaseUrl = await getPlayerInternalBaseUrl();
if (!resolvedPlayerInternalBaseUrl) {
throw new Error('Unable to resolve the player internal base URL.');
}
const url = new URL(resolvedPlayerInternalBaseUrl.replace(/^http/, 'ws'));
url.pathname = `/ws/screens/${encodeURIComponent(slug)}/events`;
url.search = '';
return url.toString();
}
function storePlayerSnapshot(slug, connections) {
const normalizedSlug = String(slug || '').trim();
const normalizedConnections = Array.isArray(connections) ? connections : [];
playerSnapshotCache.set(normalizedSlug, {
slug: normalizedSlug,
count: normalizedConnections.length,
connections: normalizedConnections
});
}
function clearPlayerSnapshotSocket(slug) {
const key = String(slug || '').trim();
playerSnapshotSockets.delete(key);
}
function ensurePlayerSnapshotSubscription(slug) {
const key = String(slug || '').trim();
if (!key || playerSnapshotSockets.has(key)) {
return;
}
playerSnapshotSockets.set(key, null);
const socketUrlPromise = getPlayerSnapshotSocketUrl(key);
const authHeaders = createRequestAuthHeaders({
method: 'GET',
pathname: `/ws/screens/${encodeURIComponent(key)}/events`
});
socketUrlPromise.then(function (socketUrl) {
const socket = new WebSocket(socketUrl, {
headers: authHeaders
});
playerSnapshotSockets.set(key, socket);
socket.onmessage = function (event) {
try {
const payload = JSON.parse(String(event.data || '{}'));
if (!payload || payload.type !== 'snapshot' || payload.slug !== key) {
return;
}
storePlayerSnapshot(key, payload.connections || []);
if (broadcastDashboardState) {
broadcastDashboardState().catch(function (error) {
console.error(error);
});
}
} catch (_error) {
// Ignore malformed player snapshot payloads.
}
};
socket.onclose = function () {
clearPlayerSnapshotSocket(key);
setTimeout(function () {
ensurePlayerSnapshotSubscription(key);
}, 2000);
};
socket.onerror = function () {
try {
socket.close();
} catch (_error) {
// ignore close errors
}
};
}).catch(function (error) {
clearPlayerSnapshotSocket(key);
console.error(error);
});
}
const dashboardStateService = createDashboardStateService({
pool: pool,
common: common,
playerSnapshotCache: playerSnapshotCache,
playerSnapshotSockets: playerSnapshotSockets,
ensurePlayerSnapshotSubscription: ensurePlayerSnapshotSubscription,
playerPublicBaseUrl: playerPublicBaseUrl,
formatDashboardDate: formatDashboardDate
});
const buildDashboardState = dashboardStateService.buildDashboardState;
const uploadSyncService = createUploadSyncService({
pool: pool,
common: common,
playerInternalBaseUrl: configuredPlayerInternalBaseUrl,
playerSnapshotCache: playerSnapshotCache,
notifyPlayerScreens: notifyPlayerScreens,
backgroundTaskQueue: backgroundTaskQueue
});
const upload = uploadSyncService.createUploadMiddleware(uploadDir);
const collectUploadReferencesFromSlide = uploadSyncService.collectUploadReferencesFromSlide;
const collectUploadReferencesFromTemplate = uploadSyncService.collectUploadReferencesFromTemplate;
const collectUploadReferencesFromPayload = uploadSyncService.collectUploadReferencesFromPayload;
const removeUnusedUploadFiles = uploadSyncService.removeUnusedUploadFiles;
const collectUploadPathsFromDirectory = uploadSyncService.collectUploadPathsFromDirectory;
const syncPlaylistUploadsOnChange = uploadSyncService.syncPlaylistUploadsOnChange;
const runMediaSyncTask = uploadSyncService.runMediaSyncTask;
async function sendDashboardStateToSocket(socket) {
if (!socket || socket.readyState !== WebSocket.OPEN) {
return;
}
const state = await buildDashboardState();
socket.send(JSON.stringify({ type: 'dashboard-state', state: state }));
}
broadcastDashboardState = async function () {
if (dashboardRefreshInFlight) {
return dashboardRefreshInFlight;
}
dashboardRefreshInFlight = (async function () {
const state = await buildDashboardState();
const payload = JSON.stringify({ type: 'dashboard-state', state: state });
for (const socket of dashboardClients) {
if (socket && socket.readyState === WebSocket.OPEN) {
socket.send(payload);
}
}
return state;
})().finally(function () {
dashboardRefreshInFlight = null;
});
return dashboardRefreshInFlight;
};
function installDashboardWebsocket(server, loadCurrentUser) {
server.on('upgrade', async function (request, socket, head) {
let pathname = '';
try {
pathname = new URL(request.url, 'http://localhost').pathname;
} catch (_error) {
socket.destroy();
return;
}
if (pathname !== '/ws/dashboard') {
socket.destroy();
return;
}
try {
const currentUser = await loadCurrentUser(pool, request);
if (!currentUser) {
socket.destroy();
return;
}
} catch (_error) {
socket.destroy();
return;
}
dashboardWs.handleUpgrade(request, socket, head, function (ws) {
dashboardWs.emit('connection', ws, request);
});
});
dashboardWs.on('connection', function (socket) {
dashboardClients.add(socket);
sendDashboardStateToSocket(socket);
socket.on('close', function () {
dashboardClients.delete(socket);
});
socket.on('error', function () {
dashboardClients.delete(socket);
});
});
setInterval(function () {
broadcastDashboardState().catch(function (error) {
console.error(error);
});
}, dashboardRefreshIntervalMs);
}
return {
upload: upload,
uploadSyncService: uploadSyncService,
buildDashboardState: buildDashboardState,
collectUploadReferencesFromSlide: collectUploadReferencesFromSlide,
collectUploadReferencesFromTemplate: collectUploadReferencesFromTemplate,
collectUploadReferencesFromPayload: collectUploadReferencesFromPayload,
removeUnusedUploadFiles: removeUnusedUploadFiles,
collectUploadPathsFromDirectory: collectUploadPathsFromDirectory,
syncPlaylistUploadsOnChange: syncPlaylistUploadsOnChange,
runMediaSyncTask: runMediaSyncTask,
broadcastDashboardState: broadcastDashboardState,
installDashboardWebsocket: installDashboardWebsocket
};
}
module.exports = { createWebBootstrap };