Refine player control-plane flow

This commit is contained in:
2026-08-07 13:10:16 +01:00
parent 74318eb34e
commit 433be06bc7
35 changed files with 1604 additions and 233 deletions
@@ -18,50 +18,57 @@ function registerFontSweepTask(options) {
const uploadSyncService = options && options.uploadSyncService;
const pushUploadFileToPlayer = uploadSyncService && uploadSyncService.pushUploadFileToPlayer;
const removeUploadFileFromPlayer = uploadSyncService && uploadSyncService.removeUploadFileFromPlayer;
const getPlayerTaskMetadata = uploadSyncService && uploadSyncService.getPlayerTaskMetadata;
const mediaDir = String(options && options.mediaDir || '').trim();
if (!backgroundTaskQueue || typeof pushUploadFileToPlayer !== 'function' || typeof removeUploadFileFromPlayer !== 'function' || !mediaDir) {
throw new Error('registerFontSweepTask requires the font sweep dependencies.');
}
backgroundTaskQueue.registerRecurringTask({
key: TASK.key,
title: TASK.title,
category: TASK.category,
intervalMs: TASK.intervalMs,
metadata: {
mediaDir: mediaDir
},
run: async function () {
const desiredOperations = collectFontLibrarySyncOperations(mediaDir);
const desiredUploadPaths = new Set(desiredOperations.map(function (operation) {
return operation && operation.uploadPath ? operation.uploadPath : '';
}).filter(Boolean));
const currentUploadPaths = await collectFontLibraryDirectoryUploadPaths(mediaDir);
const metadataPromise = typeof getPlayerTaskMetadata === 'function'
? Promise.resolve(getPlayerTaskMetadata())
: Promise.resolve({});
for (let i = 0; i < desiredOperations.length; i += 1) {
const operation = desiredOperations[i] || {};
const uploadPath = String(operation.uploadPath || '').trim();
if (!uploadPath) {
continue;
return metadataPromise.then(function (metadata) {
backgroundTaskQueue.registerRecurringTask({
key: TASK.key,
title: TASK.title,
category: TASK.category,
intervalMs: TASK.intervalMs,
metadata: Object.assign({
mediaDir: mediaDir
}, metadata || {}),
run: async function () {
const desiredOperations = collectFontLibrarySyncOperations(mediaDir);
const desiredUploadPaths = new Set(desiredOperations.map(function (operation) {
return operation && operation.uploadPath ? operation.uploadPath : '';
}).filter(Boolean));
const currentUploadPaths = await collectFontLibraryDirectoryUploadPaths(mediaDir);
for (let i = 0; i < desiredOperations.length; i += 1) {
const operation = desiredOperations[i] || {};
const uploadPath = String(operation.uploadPath || '').trim();
if (!uploadPath) {
continue;
}
if (String(operation.type || '').trim().toLowerCase() === 'delete') {
await removeUploadFileFromPlayer(uploadPath, mediaDir);
} else {
await pushUploadFileToPlayer(uploadPath, mediaDir);
}
}
if (String(operation.type || '').trim().toLowerCase() === 'delete') {
for (let i = 0; i < currentUploadPaths.length; i += 1) {
const uploadPath = String(currentUploadPaths[i] || '').trim();
if (!uploadPath || desiredUploadPaths.has(uploadPath)) {
continue;
}
await removeUploadFileFromPlayer(uploadPath, mediaDir);
} else {
await pushUploadFileToPlayer(uploadPath, mediaDir);
}
}
for (let i = 0; i < currentUploadPaths.length; i += 1) {
const uploadPath = String(currentUploadPaths[i] || '').trim();
if (!uploadPath || desiredUploadPaths.has(uploadPath)) {
continue;
}
await removeUploadFileFromPlayer(uploadPath, mediaDir);
}
}
});
});
}
@@ -8,22 +8,30 @@ const TASK = {
function registerInitialFontSyncTask(options) {
const backgroundTaskQueue = options && options.backgroundTaskQueue;
const mediaDir = String(options && options.mediaDir || '').trim();
const uploadSyncService = options && options.uploadSyncService;
if (!backgroundTaskQueue || !mediaDir) {
throw new Error('registerInitialFontSyncTask requires the initial font sync dependencies.');
}
return backgroundTaskQueue.enqueueTask({
key: TASK.key,
title: 'Initial font sync',
category: TASK.category,
taskType: 'font-sync',
payload: {
mode: 'initial',
uploadDir: mediaDir,
operations: collectFontLibrarySyncOperations(mediaDir)
},
persist: true
const metadataPromise = uploadSyncService && typeof uploadSyncService.getPlayerTaskMetadata === 'function'
? uploadSyncService.getPlayerTaskMetadata()
: Promise.resolve({});
return Promise.resolve(metadataPromise).then(function (metadata) {
return backgroundTaskQueue.enqueueTask({
key: TASK.key,
title: 'Initial font sync',
category: TASK.category,
taskType: 'font-sync',
metadata: Object.assign({}, metadata || {}),
payload: {
mode: 'initial',
uploadDir: mediaDir,
operations: collectFontLibrarySyncOperations(mediaDir)
},
persist: true
});
}).catch(function (error) {
console.warn('Unable to queue initial font sync:', error);
});
@@ -6,21 +6,29 @@ const TASK = {
function registerInitialMediaSyncTask(options) {
const backgroundTaskQueue = options && options.backgroundTaskQueue;
const mediaDir = String(options && options.mediaDir || '').trim();
const uploadSyncService = options && options.uploadSyncService;
if (!backgroundTaskQueue || !mediaDir) {
throw new Error('registerInitialMediaSyncTask requires the initial media sync dependencies.');
}
return backgroundTaskQueue.enqueueTask({
key: TASK.key,
title: 'Initial media sync',
category: TASK.category,
taskType: 'media-sync',
payload: {
mode: 'initial',
uploadDir: mediaDir
},
persist: true
const metadataPromise = uploadSyncService && typeof uploadSyncService.getPlayerTaskMetadata === 'function'
? uploadSyncService.getPlayerTaskMetadata()
: Promise.resolve({});
return Promise.resolve(metadataPromise).then(function (metadata) {
return backgroundTaskQueue.enqueueTask({
key: TASK.key,
title: 'Initial media sync',
category: TASK.category,
taskType: 'media-sync',
metadata: Object.assign({}, metadata || {}),
payload: {
mode: 'initial',
uploadDir: mediaDir
},
persist: true
});
}).catch(function (error) {
console.warn('Unable to queue initial media sync:', error);
});
+146 -24
View File
@@ -6,12 +6,40 @@ const crypto = require('crypto');
const multer = require('multer');
const { createRequestAuthHeaders } = require('#src/request-auth');
const { collectFontLibrarySyncOperations } = require('./font-library');
const { getConfiguredPlayerIdentifier, resolvePlayerRegistration } = require('#src/data/player-registry');
const { fetchPlayerRegistrations, getConfiguredPlayerIdentifier } = require('#src/data/player-registry');
function normalizeUploadRoot(uploadDir) {
return path.resolve(String(uploadDir || '').trim());
}
function isLocalLikeBaseUrl(value) {
let host = '';
try {
host = new URL(String(value || '').trim().replace(/\/$/, '')).hostname.toLowerCase();
} catch (_error) {
return false;
}
return host === 'localhost'
|| host === '127.0.0.1'
|| host === '::1'
|| host === 'host.docker.internal'
|| host === 'player'
|| host === 'web'
|| host === 'player-bridge'
|| host.endsWith('.local')
|| host.endsWith('.internal')
|| host.endsWith('.docker.internal');
}
function normalizeBaseUrl(value) {
return String(value || '').trim().replace(/\/$/, '');
}
function normalizePlayerRowBaseUrl(player) {
return normalizeBaseUrl(player && player.internal_base_url);
}
function createUploadSyncService(options) {
const pool = options && options.pool;
const common = options && options.common;
@@ -21,54 +49,120 @@ function createUploadSyncService(options) {
const backgroundTaskQueue = options && options.backgroundTaskQueue;
const MAX_UPLOAD_BYTES = 1024 * 1024 * 1024;
const MAX_FIELD_BYTES = 10 * 1024 * 1024;
const PLAYER_UPLOAD_SYNC_RETRY_LOG_INTERVAL_MS = 60000;
const pendingPlayerUploadSyncs = new Map();
let pendingPlayerUploadSyncFlushTimer = null;
let pendingPlayerUploadSyncFlushInFlight = null;
let pendingPlayerUploadSyncRetryLogAt = 0;
const pendingPlaylistUploadSyncs = new Map();
let pendingPlaylistUploadSyncFlushTimer = null;
let pendingPlaylistUploadSyncFlushInFlight = null;
let playerInternalBaseUrl = null;
let playerInternalBaseUrlPromise = null;
let playerTaskMetadata = null;
let playerTaskMetadataPromise = null;
if (!common || !playerSnapshotCache || typeof notifyPlayerScreens !== 'function') {
throw new Error('createUploadSyncService requires the upload dependencies.');
}
async function getPlayerInternalBaseUrl() {
if (playerInternalBaseUrl) {
return playerInternalBaseUrl;
const metadata = await getPlayerTaskMetadata();
return metadata && metadata.playerInternalBaseUrl ? metadata.playerInternalBaseUrl : null;
}
async function getPlayerTaskMetadata() {
if (playerTaskMetadata) {
return playerTaskMetadata;
}
if (playerInternalBaseUrlPromise) {
return playerInternalBaseUrlPromise;
if (playerTaskMetadataPromise) {
return playerTaskMetadataPromise;
}
playerInternalBaseUrlPromise = (async function () {
if (!pool) {
return configuredPlayerInternalBaseUrl || null;
}
playerTaskMetadataPromise = (async function () {
try {
const configuredPlayerIdentifier = getConfiguredPlayerIdentifier();
const player = await resolvePlayerRegistration(pool, configuredPlayerIdentifier);
const resolvedBaseUrl = String(player && player.internal_base_url || '').trim().replace(/\/$/, '');
if (resolvedBaseUrl) {
playerInternalBaseUrl = resolvedBaseUrl;
return resolvedBaseUrl;
if (pool && typeof fetchPlayerRegistrations === 'function') {
const configuredPlayerIdentifier = getConfiguredPlayerIdentifier();
const players = await fetchPlayerRegistrations(pool);
const exactPlayer = Array.isArray(players)
? players.find(function (player) {
return String(player && player.identifier || '').trim() === configuredPlayerIdentifier;
})
: null;
const registeredPlayers = Array.isArray(players) ? players : [];
const preferredPlayer = registeredPlayers.find(function (player) {
const internalBaseUrl = normalizePlayerRowBaseUrl(player);
return internalBaseUrl && !isLocalLikeBaseUrl(internalBaseUrl);
}) || exactPlayer || registeredPlayers[0] || null;
const resolvedInternalBaseUrl = normalizePlayerRowBaseUrl(preferredPlayer);
const resolvedPublicBaseUrl = normalizeBaseUrl(preferredPlayer && preferredPlayer.public_base_url);
const resolvedIdentifier = String(preferredPlayer && preferredPlayer.identifier || '').trim();
if (resolvedInternalBaseUrl || resolvedPublicBaseUrl || resolvedIdentifier) {
playerInternalBaseUrl = resolvedInternalBaseUrl || null;
playerTaskMetadata = {
playerIdentifier: resolvedIdentifier || null,
playerPublicBaseUrl: resolvedPublicBaseUrl || null,
playerInternalBaseUrl: resolvedInternalBaseUrl || null,
playerLabel: resolvedIdentifier || resolvedPublicBaseUrl || resolvedInternalBaseUrl || null
};
return playerTaskMetadata;
}
}
} catch (_error) {
}
return configuredPlayerInternalBaseUrl || null;
})().then(function (baseUrl) {
playerInternalBaseUrlPromise = null;
return baseUrl || null;
playerInternalBaseUrl = configuredPlayerInternalBaseUrl || null;
playerTaskMetadata = {
playerIdentifier: getConfiguredPlayerIdentifier() || null,
playerPublicBaseUrl: null,
playerInternalBaseUrl: playerInternalBaseUrl,
playerLabel: getConfiguredPlayerIdentifier() || playerInternalBaseUrl || null
};
return playerTaskMetadata;
})().then(function (metadata) {
playerTaskMetadataPromise = null;
return metadata || null;
}, function () {
playerInternalBaseUrlPromise = null;
return configuredPlayerInternalBaseUrl || null;
playerTaskMetadataPromise = null;
return {
playerIdentifier: getConfiguredPlayerIdentifier() || null,
playerPublicBaseUrl: null,
playerInternalBaseUrl: configuredPlayerInternalBaseUrl || null,
playerLabel: getConfiguredPlayerIdentifier() || configuredPlayerInternalBaseUrl || null
};
});
return playerInternalBaseUrlPromise;
return playerTaskMetadataPromise;
}
function formatPlayerTaskLabel(metadata) {
const playerLabel = String(metadata && metadata.playerLabel || '').trim();
if (playerLabel) {
return playerLabel;
}
const playerIdentifier = String(metadata && metadata.playerIdentifier || '').trim();
if (playerIdentifier) {
return playerIdentifier;
}
const playerPublicBaseUrl = String(metadata && metadata.playerPublicBaseUrl || '').trim();
if (playerPublicBaseUrl) {
return playerPublicBaseUrl;
}
return '';
}
function logMediaSyncSummary(level, message, metadata, details) {
const suffix = formatPlayerTaskLabel(metadata);
const logger = level === 'warn' ? console.warn : console.info;
if (details !== undefined) {
logger(`[media-sync] ${message}${suffix ? ` for ${suffix}` : ''}`, details);
return;
}
logger(`[media-sync] ${message}${suffix ? ` for ${suffix}` : ''}`);
}
function createUploadMiddleware(uploadDir) {
@@ -561,6 +655,10 @@ function createUploadSyncService(options) {
pendingPlaylistUploadSyncs.delete(operation.key);
}
if (pendingEntries.length) {
logMediaSyncSummary('info', `Playlist sync flushed ${pendingEntries.length} task${pendingEntries.length === 1 ? '' : 's'}`, pendingEntries[0] && pendingEntries[0].metadata);
}
})().finally(function () {
pendingPlaylistUploadSyncFlushInFlight = null;
if (pendingPlaylistUploadSyncs.size) {
@@ -582,6 +680,11 @@ function createUploadSyncService(options) {
pendingPlayerUploadSyncFlushInFlight = (async function () {
const pendingEntries = Array.from(pendingPlayerUploadSyncs.values());
const playerMetadata = pendingEntries.length && pendingEntries[0] && pendingEntries[0].metadata
? pendingEntries[0].metadata
: await getPlayerTaskMetadata();
let successCount = 0;
let failureCount = 0;
for (let i = 0; i < pendingEntries.length; i += 1) {
const operation = pendingEntries[i];
let success = false;
@@ -591,9 +694,25 @@ function createUploadSyncService(options) {
success = await pushUploadFileToPlayer(operation.uploadPath, operation.uploadDir);
}
if (success) {
successCount += 1;
pendingPlayerUploadSyncs.delete(operation.uploadPath);
} else {
failureCount += 1;
}
}
if (successCount) {
logMediaSyncSummary('info', `Media sync completed ${successCount} upload${successCount === 1 ? '' : 's'}`, playerMetadata);
}
if (failureCount) {
const now = Date.now();
if (!pendingPlayerUploadSyncRetryLogAt || now - pendingPlayerUploadSyncRetryLogAt >= PLAYER_UPLOAD_SYNC_RETRY_LOG_INTERVAL_MS) {
pendingPlayerUploadSyncRetryLogAt = now;
logMediaSyncSummary('warn', `Player unavailable, retry queued for ${failureCount} upload${failureCount === 1 ? '' : 's'}`, playerMetadata);
}
} else if (!pendingPlayerUploadSyncs.size) {
pendingPlayerUploadSyncRetryLogAt = 0;
}
})().finally(function () {
pendingPlayerUploadSyncFlushInFlight = null;
if (pendingPlayerUploadSyncs.size) {
@@ -689,12 +808,14 @@ function createUploadSyncService(options) {
safePayload.operation = Object.assign({}, safePayload.operation);
delete safePayload.operation.pool;
}
const playerMetadata = await getPlayerTaskMetadata();
const definition = {
key: taskKey,
title: title,
category: 'media-sync',
taskType: 'media-sync',
metadata: Object.assign({}, playerMetadata || {}),
payload: safePayload,
persist: true
};
@@ -732,7 +853,8 @@ function createUploadSyncService(options) {
flushPendingPlaylistUploadSyncs: flushPendingPlaylistUploadSyncs,
flushPendingPlayerUploadSyncs: flushPendingPlayerUploadSyncs,
runMediaSyncTask: runMediaSyncTask,
queueMediaSyncTask: queueMediaSyncTask
queueMediaSyncTask: queueMediaSyncTask,
getPlayerTaskMetadata: getPlayerTaskMetadata
};
}
+49 -16
View File
@@ -1,5 +1,29 @@
const { createRequestAuthHeaders } = require('#src/request-auth');
const { getConfiguredPlayerIdentifier, resolvePlayerRegistration } = require('#src/data/player-registry');
const { fetchPlayerRegistrations, getConfiguredPlayerIdentifier } = require('#src/data/player-registry');
function isLocalLikeBaseUrl(value) {
let host = '';
try {
host = new URL(String(value || '').trim().replace(/\/$/, '')).hostname.toLowerCase();
} catch (_error) {
return false;
}
return host === 'localhost'
|| host === '127.0.0.1'
|| host === '::1'
|| host === 'host.docker.internal'
|| host === 'player'
|| host === 'web'
|| host === 'player-bridge'
|| host.endsWith('.local')
|| host.endsWith('.internal')
|| host.endsWith('.docker.internal');
}
function normalizeBaseUrl(value) {
return String(value || '').trim().replace(/\/$/, '');
}
function createPlayerActionService(options) {
const pool = options && options.pool;
@@ -14,10 +38,6 @@ function createPlayerActionService(options) {
let playerInternalBaseUrlPromise = null;
async function getPlayerInternalBaseUrl() {
if (configuredPlayerInternalBaseUrl) {
return configuredPlayerInternalBaseUrl;
}
if (playerInternalBaseUrl) {
return playerInternalBaseUrl;
}
@@ -27,22 +47,35 @@ function createPlayerActionService(options) {
}
playerInternalBaseUrlPromise = (async function () {
if (!pool) {
return configuredPlayerInternalBaseUrl || null;
}
try {
const configuredPlayerIdentifier = getConfiguredPlayerIdentifier();
const player = await resolvePlayerRegistration(pool, configuredPlayerIdentifier);
const resolvedBaseUrl = String(player && player.internal_base_url || '').trim().replace(/\/$/, '');
if (resolvedBaseUrl) {
playerInternalBaseUrl = resolvedBaseUrl;
return resolvedBaseUrl;
if (pool && typeof fetchPlayerRegistrations === 'function') {
const configuredPlayerIdentifier = getConfiguredPlayerIdentifier();
const players = await fetchPlayerRegistrations(pool);
const exactPlayer = Array.isArray(players)
? players.find(function (player) {
return String(player && player.identifier || '').trim() === configuredPlayerIdentifier;
})
: null;
const registeredPlayers = Array.isArray(players) ? players : [];
const preferredPlayer = registeredPlayers.find(function (player) {
const internalBaseUrl = normalizeBaseUrl(player && player.internal_base_url);
return internalBaseUrl && !isLocalLikeBaseUrl(internalBaseUrl);
}) || exactPlayer || registeredPlayers[0] || null;
const resolvedBaseUrl = normalizeBaseUrl(preferredPlayer && preferredPlayer.internal_base_url);
if (resolvedBaseUrl) {
playerInternalBaseUrl = resolvedBaseUrl;
return resolvedBaseUrl;
}
}
} catch (_error) {
}
return configuredPlayerInternalBaseUrl || null;
if (configuredPlayerInternalBaseUrl) {
playerInternalBaseUrl = configuredPlayerInternalBaseUrl;
return configuredPlayerInternalBaseUrl;
}
return null;
})().then(function (baseUrl) {
playerInternalBaseUrlPromise = null;
return baseUrl || null;