Files
pulse-signage/src/web/lib/upload-sync.js
T

716 lines
24 KiB
JavaScript

const fs = require('fs');
const path = require('path');
const crypto = require('crypto');
const multer = require('multer');
const { createRequestAuthHeaders } = require('../../request-auth');
function normalizeUploadRoot(uploadDir) {
return path.resolve(String(uploadDir || '').trim());
}
function createUploadSyncService(options) {
const pool = options && options.pool;
const common = options && options.common;
const playerInternalBaseUrl = String(options && options.playerInternalBaseUrl || '').replace(/\/$/, '');
const playerSnapshotCache = options && options.playerSnapshotCache;
const notifyPlayerScreens = options && options.notifyPlayerScreens;
const backgroundTaskQueue = options && options.backgroundTaskQueue;
const MAX_UPLOAD_BYTES = 1024 * 1024 * 1024;
let playerUploadSyncMode = null;
let playerUploadSyncModePromise = null;
const pendingPlayerUploadSyncs = new Map();
let pendingPlayerUploadSyncFlushTimer = null;
let pendingPlayerUploadSyncFlushInFlight = null;
const pendingPlaylistUploadSyncs = new Map();
let pendingPlaylistUploadSyncFlushTimer = null;
let pendingPlaylistUploadSyncFlushInFlight = null;
if (!common || !playerSnapshotCache || typeof notifyPlayerScreens !== 'function') {
throw new Error('createUploadSyncService requires the upload dependencies.');
}
function createUploadMiddleware(uploadDir) {
const storage = multer.diskStorage({
destination: function (_req, _file, cb) {
cb(null, uploadDir);
},
filename: function (_req, file, cb) {
const safeExt = path.extname(file.originalname || '').toLowerCase();
const stamp = `${Date.now()}-${crypto.randomUUID()}`;
cb(null, `${stamp}${safeExt}`);
}
});
return multer({
storage: storage,
limits: {
fileSize: MAX_UPLOAD_BYTES
}
});
}
function normalizeUploadReference(uploadPath) {
const value = String(uploadPath || '').trim();
if (!value || !value.startsWith('/media/')) {
return null;
}
return value;
}
function getUploadRelativePath(uploadPath) {
const value = normalizeUploadReference(uploadPath);
if (!value) {
return null;
}
return value.replace(/^\/media\//, '');
}
function resolveUploadFilePath(uploadDir, uploadPath) {
const relativePath = getUploadRelativePath(uploadPath);
if (!relativePath) {
return null;
}
const normalizedUploadDir = normalizeUploadRoot(uploadDir);
if (!normalizedUploadDir) {
return null;
}
if (relativePath.startsWith('uploads/')) {
return path.join(normalizedUploadDir, relativePath.slice('uploads/'.length));
}
return path.join(path.dirname(normalizedUploadDir), relativePath);
}
function collectUploadReferencesFromValue(value, refs) {
if (!value) {
return refs;
}
const stack = [value];
while (stack.length) {
const current = stack.pop();
if (Array.isArray(current)) {
current.forEach(function (item) {
stack.push(item);
});
continue;
}
if (current && typeof current === 'object') {
Object.keys(current).forEach(function (key) {
stack.push(current[key]);
});
continue;
}
if (typeof current === 'string') {
const reference = normalizeUploadReference(current);
if (reference) {
refs.add(reference);
}
}
}
return refs;
}
function collectUploadReferencesFromSlide(slide) {
const refs = new Set();
if (!slide) {
return refs;
}
collectUploadReferencesFromValue(slide.media_path, refs);
collectUploadReferencesFromValue(common.parseJsonSafe(slide.content_json), refs);
return refs;
}
function collectUploadReferencesFromTemplate(template) {
const refs = new Set();
if (!template) {
return refs;
}
collectUploadReferencesFromValue(template.background_image_path, refs);
return refs;
}
function collectUploadReferencesFromPayload(payload) {
const refs = new Set();
if (!payload) {
return refs;
}
collectUploadReferencesFromValue(payload.mediaPath, refs);
collectUploadReferencesFromValue(common.parseJsonSafe(payload.contentJson), refs);
collectUploadReferencesFromValue(payload.backgroundImagePath, refs);
return refs;
}
async function countUploadReferences(pool, uploadPath) {
const [slideRows] = await pool.query(
`SELECT COUNT(*) AS ref_count
FROM slides
WHERE media_path = ?
OR JSON_SEARCH(COALESCE(content_json, JSON_OBJECT()), 'one', ?) IS NOT NULL`,
[uploadPath, uploadPath]
);
const [templateRows] = await pool.query(
'SELECT COUNT(*) AS ref_count FROM slide_templates WHERE background_image_path = ?',
[uploadPath]
);
return Number(slideRows[0].ref_count || 0) + Number(templateRows[0].ref_count || 0);
}
async function removeUnusedUploadFiles(pool, uploadDir, uploadPaths) {
const uniquePaths = Array.from(new Set((uploadPaths || []).map(normalizeUploadReference).filter(Boolean)));
for (let i = 0; i < uniquePaths.length; i += 1) {
const uploadPath = uniquePaths[i];
const referenceCount = await countUploadReferences(pool, uploadPath);
if (referenceCount > 0) {
continue;
}
const filePath = resolveUploadFilePath(uploadDir, uploadPath);
try {
await fs.promises.unlink(filePath);
} catch (error) {
if (error && error.code !== 'ENOENT') {
console.warn('Unable to remove unused upload file:', filePath, error);
}
}
queuePlayerUploadSync({
type: 'delete',
uploadPath: uploadPath,
uploadDir: uploadDir
});
}
}
async function collectUploadPathsFromDirectory(uploadDir) {
const normalizedUploadDir = normalizeUploadRoot(uploadDir);
if (!normalizedUploadDir) {
return [];
}
const uploadPaths = [];
async function walkDirectory(currentDir, relativeDir) {
let entries = [];
try {
entries = await fs.promises.readdir(currentDir, { withFileTypes: true });
} catch (error) {
if (error && error.code !== 'ENOENT') {
console.warn('Unable to read upload directory:', currentDir, error);
}
return;
}
for (const entry of entries) {
const entryName = String(entry && entry.name || '').trim();
if (!entryName || entryName === '.' || entryName === '..') {
continue;
}
const nextRelativePath = relativeDir ? path.posix.join(relativeDir, entryName) : entryName;
const nextAbsolutePath = path.join(currentDir, entryName);
if (entry.isDirectory && entry.isDirectory()) {
await walkDirectory(nextAbsolutePath, nextRelativePath);
continue;
}
if (entry.isFile && !entry.isFile()) {
continue;
}
uploadPaths.push('/media/uploads/' + nextRelativePath.replace(/\\/g, '/'));
}
}
await walkDirectory(normalizedUploadDir, '');
return uploadPaths;
}
async function getPlayerUploadSyncMode(localUploadDir) {
if (playerUploadSyncMode) {
return playerUploadSyncMode;
}
if (playerUploadSyncModePromise) {
return playerUploadSyncModePromise;
}
playerUploadSyncModePromise = (async function () {
try {
const authHeaders = createRequestAuthHeaders({
method: 'GET',
pathname: '/api/media/config'
});
const response = await fetch(`${playerInternalBaseUrl}/api/media/config`, {
headers: {
Accept: 'application/json',
...authHeaders
}
});
if (!response.ok) {
return null;
}
const data = await response.json();
const playerUploadDir = data && (data.uploadDir || data.mediaDir) ? normalizeUploadRoot(data.uploadDir || data.mediaDir) : null;
if (!playerUploadDir) {
return null;
}
return playerUploadDir === normalizeUploadRoot(localUploadDir) ? 'shared' : 'different';
} catch (_error) {
return null;
}
})().then(function (mode) {
if (mode) {
playerUploadSyncMode = mode;
}
playerUploadSyncModePromise = null;
return mode;
}, function () {
playerUploadSyncModePromise = null;
return null;
});
return playerUploadSyncModePromise;
}
async function shouldMirrorUploads(localUploadDir) {
return Boolean(localUploadDir);
}
function queuePlayerUploadSync(operation) {
if (!operation || !operation.uploadPath) {
return;
}
pendingPlayerUploadSyncs.set(normalizeUploadReference(operation.uploadPath), {
type: operation.type === 'delete' ? 'delete' : 'put',
uploadPath: normalizeUploadReference(operation.uploadPath),
uploadDir: operation.uploadDir || null
});
schedulePendingPlayerUploadSyncFlush();
}
function schedulePendingPlayerUploadSyncFlush() {
if (pendingPlayerUploadSyncFlushTimer) {
return;
}
pendingPlayerUploadSyncFlushTimer = setTimeout(function () {
pendingPlayerUploadSyncFlushTimer = null;
flushPendingPlayerUploadSyncs().catch(function (error) {
console.warn('Unable to flush pending upload syncs:', error);
});
}, 5000);
}
async function pushUploadFileToPlayer(uploadPath, localUploadDir) {
if (!uploadPath || !(await shouldMirrorUploads(localUploadDir))) {
return false;
}
const relativePath = getUploadRelativePath(uploadPath);
const sourcePath = resolveUploadFilePath(localUploadDir, uploadPath);
if (!relativePath || !sourcePath) {
return false;
}
let fileBuffer = null;
try {
fileBuffer = await fs.promises.readFile(sourcePath);
} catch (error) {
if (!error || error.code !== 'ENOENT') {
console.warn('Unable to read upload for player sync:', sourcePath, error);
}
return;
}
try {
const authHeaders = createRequestAuthHeaders({
method: 'PUT',
pathname: `/api/media/${encodeURIComponent(relativePath)}`,
body: fileBuffer
});
const response = await fetch(`${playerInternalBaseUrl}/api/media/${encodeURIComponent(relativePath)}`, {
method: 'PUT',
headers: {
'Content-Type': 'application/octet-stream',
...authHeaders
},
body: fileBuffer
});
if (!response.ok) {
console.warn('Unable to sync upload to player:', relativePath, response.status, response.statusText);
return false;
}
return true;
} catch (error) {
console.warn('Unable to sync upload to player:', relativePath, error);
return false;
}
}
async function removeUploadFileFromPlayer(uploadPath, localUploadDir) {
if (!uploadPath || !(await shouldMirrorUploads(localUploadDir))) {
return false;
}
const relativePath = getUploadRelativePath(uploadPath);
if (!relativePath) {
return false;
}
try {
const authHeaders = createRequestAuthHeaders({
method: 'DELETE',
pathname: `/api/media/${encodeURIComponent(relativePath)}`
});
const response = await fetch(`${playerInternalBaseUrl}/api/media/${encodeURIComponent(relativePath)}`, {
method: 'DELETE',
headers: {
Accept: 'application/json',
...authHeaders
}
});
if (!response.ok && response.status !== 404) {
console.warn('Unable to remove upload from player:', relativePath, response.status, response.statusText);
return false;
}
return true;
} catch (error) {
console.warn('Unable to remove upload from player:', relativePath, error);
return false;
}
}
async function syncUploadRefsToPlayer(uploadRefs, localUploadDir) {
if (!(await shouldMirrorUploads(localUploadDir))) {
return;
}
const uniqueRefs = Array.from(new Set((uploadRefs || []).map(normalizeUploadReference).filter(Boolean)));
for (let i = 0; i < uniqueRefs.length; i += 1) {
const success = await pushUploadFileToPlayer(uniqueRefs[i], localUploadDir);
if (!success) {
queuePlayerUploadSync({
type: 'put',
uploadPath: uniqueRefs[i],
uploadDir: localUploadDir
});
}
}
}
async function syncExistingUploadsToPlayer(pool, localUploadDir) {
return queueMediaSyncTask('media-sync:initial', 'Initial media sync', {
mode: 'initial',
uploadDir: localUploadDir
});
}
function getVisibleCurrentSlideIds() {
const visibleSlideIds = new Set();
playerSnapshotCache.forEach(function (snapshot) {
const connections = snapshot && Array.isArray(snapshot.connections) ? snapshot.connections : [];
connections.forEach(function (connection) {
const currentSlide = connection && connection.currentSlide && typeof connection.currentSlide === 'object'
? connection.currentSlide
: null;
const slideId = currentSlide && currentSlide.id !== undefined && currentSlide.id !== null
? String(currentSlide.id).trim()
: '';
if (slideId) {
visibleSlideIds.add(slideId);
}
});
});
return visibleSlideIds;
}
function isScreenRefreshBlocked(screenSlug, blockedSlideIds, screenSlideCounts) {
const slideIds = Array.isArray(blockedSlideIds)
? blockedSlideIds.map(function (value) {
return String(value || '').trim();
}).filter(Boolean)
: [];
if (!slideIds.length) {
return false;
}
const normalizedScreenSlug = String(screenSlug || '').trim();
const slideCount = screenSlideCounts && Object.prototype.hasOwnProperty.call(screenSlideCounts, normalizedScreenSlug)
? Number(screenSlideCounts[normalizedScreenSlug])
: null;
if (Number.isFinite(slideCount) && slideCount <= 1) {
return false;
}
const snapshot = playerSnapshotCache.get(String(screenSlug || '').trim());
const connections = snapshot && Array.isArray(snapshot.connections) ? snapshot.connections : [];
return connections.some(function (connection) {
const currentSlide = connection && connection.currentSlide && typeof connection.currentSlide === 'object'
? connection.currentSlide
: null;
const currentSlideId = currentSlide && currentSlide.id !== undefined && currentSlide.id !== null
? String(currentSlide.id).trim()
: '';
return currentSlideId && slideIds.includes(currentSlideId);
});
}
function splitRefreshScreenSlugsByVisibility(screenSlugs, blockedSlideIds, screenSlideCounts) {
const ready = [];
const blocked = [];
Array.from(new Set(Array.isArray(screenSlugs) ? screenSlugs : [])).forEach(function (screenSlug) {
const normalizedScreenSlug = String(screenSlug || '').trim();
if (!normalizedScreenSlug) {
return;
}
if (isScreenRefreshBlocked(normalizedScreenSlug, blockedSlideIds, screenSlideCounts)) {
blocked.push(normalizedScreenSlug);
} else {
ready.push(normalizedScreenSlug);
}
});
return { ready: ready, blocked: blocked };
}
function normalizePlaylistUploadSyncOperation(options) {
return {
key: String(options && options.key ? options.key : '').trim(),
pool: options && options.pool ? options.pool : null,
localUploadDir: options && options.localUploadDir ? options.localUploadDir : null,
previousUploadRefs: Array.from(new Set(options && options.previousUploadRefs ? options.previousUploadRefs : [])),
nextUploadRefs: Array.from(new Set(options && options.nextUploadRefs ? options.nextUploadRefs : [])),
blockedSlideIds: Array.from(new Set(options && options.blockedSlideIds ? options.blockedSlideIds : [])).map(function (value) {
return String(value || '').trim();
}).filter(Boolean),
refreshScreenSlugs: Array.from(new Set(options && options.refreshScreenSlugs ? options.refreshScreenSlugs : [])).map(function (value) {
return String(value || '').trim();
}).filter(Boolean)
,
screenSlideCounts: options && options.screenSlideCounts && typeof options.screenSlideCounts === 'object'
? options.screenSlideCounts
: {}
};
}
function queuePlaylistUploadSync(operation) {
if (!operation || !operation.key) {
return;
}
pendingPlaylistUploadSyncs.set(operation.key, normalizePlaylistUploadSyncOperation(operation));
schedulePendingPlaylistUploadSyncFlush();
}
function schedulePendingPlaylistUploadSyncFlush() {
if (pendingPlaylistUploadSyncFlushTimer) {
return;
}
pendingPlaylistUploadSyncFlushTimer = setTimeout(function () {
pendingPlaylistUploadSyncFlushTimer = null;
flushPendingPlaylistUploadSyncs().catch(function (error) {
console.warn('Unable to flush pending playlist upload syncs:', error);
});
}, 5000);
}
async function syncPlaylistUploadsOnChange(options) {
const operation = normalizePlaylistUploadSyncOperation(options);
if (!operation.key) {
return;
}
return queueMediaSyncTask('media-sync:' + operation.key, 'Media sync', {
mode: 'playlist',
operation: operation
});
}
async function flushPendingPlaylistUploadSyncs() {
if (pendingPlaylistUploadSyncFlushInFlight) {
return pendingPlaylistUploadSyncFlushInFlight;
}
if (!pendingPlaylistUploadSyncs.size) {
return null;
}
pendingPlaylistUploadSyncFlushInFlight = (async function () {
const pendingEntries = Array.from(pendingPlaylistUploadSyncs.values());
for (let i = 0; i < pendingEntries.length; i += 1) {
const operation = pendingEntries[i];
const refreshTargets = splitRefreshScreenSlugsByVisibility(operation.refreshScreenSlugs, operation.blockedSlideIds, operation.screenSlideCounts);
if (!refreshTargets.ready.length) {
continue;
}
if (refreshTargets.ready.length) {
await notifyPlayerScreens(refreshTargets.ready, 'refresh');
}
pendingPlaylistUploadSyncs.delete(operation.key);
}
})().finally(function () {
pendingPlaylistUploadSyncFlushInFlight = null;
if (pendingPlaylistUploadSyncs.size) {
schedulePendingPlaylistUploadSyncFlush();
}
});
return pendingPlaylistUploadSyncFlushInFlight;
}
async function flushPendingPlayerUploadSyncs() {
if (pendingPlayerUploadSyncFlushInFlight) {
return pendingPlayerUploadSyncFlushInFlight;
}
if (!pendingPlayerUploadSyncs.size) {
return null;
}
pendingPlayerUploadSyncFlushInFlight = (async function () {
const pendingEntries = Array.from(pendingPlayerUploadSyncs.values());
for (let i = 0; i < pendingEntries.length; i += 1) {
const operation = pendingEntries[i];
let success = false;
if (operation.type === 'delete') {
success = await removeUploadFileFromPlayer(operation.uploadPath, operation.uploadDir);
} else {
success = await pushUploadFileToPlayer(operation.uploadPath, operation.uploadDir);
}
if (success) {
pendingPlayerUploadSyncs.delete(operation.uploadPath);
}
}
})().finally(function () {
pendingPlayerUploadSyncFlushInFlight = null;
if (pendingPlayerUploadSyncs.size) {
schedulePendingPlayerUploadSyncFlush();
}
});
return pendingPlayerUploadSyncFlushInFlight;
}
async function runMediaSyncTask(payload) {
const taskPayload = payload || {};
const mode = String(taskPayload.mode || '').trim();
if (mode === 'initial') {
const uploadDir = String(taskPayload.uploadDir || '').trim();
if (!(await shouldMirrorUploads(uploadDir))) {
return;
}
const data = await common.fetchAdminData(pool);
const uploadRefs = new Set();
(data.slides || []).forEach(function (slide) {
collectUploadReferencesFromSlide(slide).forEach(function (reference) {
uploadRefs.add(reference);
});
});
(data.templates || []).forEach(function (template) {
collectUploadReferencesFromTemplate(template).forEach(function (reference) {
uploadRefs.add(reference);
});
});
Array.from(uploadRefs).forEach(function (uploadPath) {
queuePlayerUploadSync({
type: 'put',
uploadPath: uploadPath,
uploadDir: uploadDir
});
});
await flushPendingPlayerUploadSyncs();
return;
}
if (mode === 'playlist') {
const operation = normalizePlaylistUploadSyncOperation(taskPayload.operation || taskPayload);
if (operation.nextUploadRefs.length) {
await syncUploadRefsToPlayer(operation.nextUploadRefs, operation.localUploadDir);
}
if (operation.previousUploadRefs.length) {
const nextUploadRefSet = new Set(operation.nextUploadRefs);
await removeUnusedUploadFiles(pool, operation.localUploadDir, operation.previousUploadRefs.filter(function (reference) {
return !nextUploadRefSet.has(reference);
}));
}
if (operation.refreshScreenSlugs.length) {
const refreshTargets = splitRefreshScreenSlugsByVisibility(operation.refreshScreenSlugs, operation.blockedSlideIds, operation.screenSlideCounts);
if (refreshTargets.ready.length) {
await notifyPlayerScreens(refreshTargets.ready, 'refresh');
}
refreshTargets.blocked.forEach(function (screenSlug) {
queuePlaylistUploadSync({
key: operation.key + ':refresh:' + screenSlug,
blockedSlideIds: operation.blockedSlideIds,
refreshScreenSlugs: [screenSlug]
});
});
}
return;
}
throw new Error('Unknown media sync task mode.');
}
async function queueMediaSyncTask(taskKey, title, payload) {
const safePayload = Object.assign({}, payload || {});
delete safePayload.pool;
if (safePayload.operation && typeof safePayload.operation === 'object') {
safePayload.operation = Object.assign({}, safePayload.operation);
delete safePayload.operation.pool;
}
const definition = {
key: taskKey,
title: title,
category: 'media-sync',
taskType: 'media-sync',
payload: safePayload,
persist: true
};
if (backgroundTaskQueue && typeof backgroundTaskQueue.enqueueTaskAndWait === 'function') {
return backgroundTaskQueue.enqueueTaskAndWait(definition);
}
return runMediaSyncTask(safePayload);
}
return {
createUploadMiddleware: createUploadMiddleware,
normalizeUploadReference: normalizeUploadReference,
collectUploadReferencesFromValue: collectUploadReferencesFromValue,
collectUploadReferencesFromSlide: collectUploadReferencesFromSlide,
collectUploadReferencesFromTemplate: collectUploadReferencesFromTemplate,
collectUploadReferencesFromPayload: collectUploadReferencesFromPayload,
countUploadReferences: countUploadReferences,
removeUnusedUploadFiles: removeUnusedUploadFiles,
collectUploadPathsFromDirectory: collectUploadPathsFromDirectory,
getPlayerUploadSyncMode: getPlayerUploadSyncMode,
shouldMirrorUploads: shouldMirrorUploads,
queuePlayerUploadSync: queuePlayerUploadSync,
schedulePendingPlayerUploadSyncFlush: schedulePendingPlayerUploadSyncFlush,
pushUploadFileToPlayer: pushUploadFileToPlayer,
removeUploadFileFromPlayer: removeUploadFileFromPlayer,
syncUploadRefsToPlayer: syncUploadRefsToPlayer,
syncExistingUploadsToPlayer: syncExistingUploadsToPlayer,
getVisibleCurrentSlideIds: getVisibleCurrentSlideIds,
isScreenRefreshBlocked: isScreenRefreshBlocked,
splitRefreshScreenSlugsByVisibility: splitRefreshScreenSlugsByVisibility,
normalizePlaylistUploadSyncOperation: normalizePlaylistUploadSyncOperation,
queuePlaylistUploadSync: queuePlaylistUploadSync,
schedulePendingPlaylistUploadSyncFlush: schedulePendingPlaylistUploadSyncFlush,
syncPlaylistUploadsOnChange: syncPlaylistUploadsOnChange,
flushPendingPlaylistUploadSyncs: flushPendingPlaylistUploadSyncs,
flushPendingPlayerUploadSyncs: flushPendingPlayerUploadSyncs,
runMediaSyncTask: runMediaSyncTask,
queueMediaSyncTask: queueMediaSyncTask
};
}
module.exports = { createUploadSyncService };