module.exports = function registerAdminDataSourceRoutes(app, deps) { const pool = deps.pool; const common = deps.common; const pages = deps.pages; const backgroundTaskQueue = deps.backgroundTaskQueue; const { hasAnyPermission } = require('../../../rbac'); const formatDashboardDate = deps.formatDashboardDate || function (value) { return value ? String(value) : ''; }; const getAuditUserId = deps.getAuditUserId; const redirectAfterSave = deps.redirectAfterSave; const requirePermission = deps.requirePermission; const fetchRssFeedItems = deps.fetchRssFeedItems || (common && common.fetchRssFeedItems) || null; const replaceRssFeedItems = deps.replaceRssFeedItems || (common && common.replaceRssFeedItems) || null; const fetchApiSourceResponse = common && common.fetchApiSourceResponse ? common.fetchApiSourceResponse : null; const { buildPagination } = require('../../lib/pagination'); const LIST_PAGE_SIZE = 10; function toIsoTimestamp(value) { if (!value) { return ''; } const date = new Date(value); return Number.isNaN(date.getTime()) ? '' : date.toISOString(); } async function refreshApiSourceInBackground(apiSourceId, apiUrl, actorId) { const connection = await pool.getConnection(); try { let responseDetails = null; let pullError = ''; try { responseDetails = await loadApiSourceResponse(apiUrl); } catch (error) { pullError = String(error && error.message ? error.message : 'Unable to load API response.'); } await connection.beginTransaction(); await connection.query( 'UPDATE api_sources SET last_pulled_at = ?, last_pull_error = ?, last_response_status = ?, last_response_content_type = ?, last_response_json = ?, modified_by = ? WHERE id = ?', [new Date(), pullError || null, responseDetails ? responseDetails.responseStatus : null, responseDetails ? responseDetails.responseContentType : null, responseDetails ? responseDetails.responseJson : null, actorId, apiSourceId] ); await connection.commit(); } catch (error) { try { await connection.rollback(); } catch (_rollbackError) { // Ignore rollback failures and surface the original error. } throw error; } finally { connection.release(); } } async function refreshRssFeedInBackground(rssFeedId, feedUrl, itemLimit, actorId) { const connection = await pool.getConnection(); try { let updatedItems = []; let pullError = ''; try { updatedItems = await loadRssFeedItems(feedUrl, itemLimit); } catch (error) { pullError = String(error && error.message ? error.message : 'Unable to load feed items.'); } await connection.beginTransaction(); await connection.query( 'UPDATE rss_feeds SET modified_by = ? WHERE id = ?', [actorId, rssFeedId] ); if (replaceRssFeedItems) { await replaceRssFeedItems(connection, rssFeedId, updatedItems); } await connection.commit(); if (pullError) { console.error('[admin-data-sources] RSS feed refresh completed with an error for feed ' + rssFeedId + ': ' + pullError); } } catch (error) { try { await connection.rollback(); } catch (_rollbackError) { // Ignore rollback failures and surface the original error. } throw error; } finally { connection.release(); } } if (!pool || !common || !pages || typeof getAuditUserId !== 'function' || typeof redirectAfterSave !== 'function' || typeof requirePermission !== 'function' || !backgroundTaskQueue) { throw new Error('registerAdminDataSourceRoutes requires the data source route dependencies.'); } function formatRecurringKey(sourceType, id) { return sourceType + '-refresh:' + Number(id); } function buildRecurringTitle(sourceType) { return sourceType === 'rss-feed' ? 'RSS feed refresh' : 'API source refresh'; } function registerRecurringRefresh(sourceType, id, name, intervalValue, intervalUnit, run) { backgroundTaskQueue.registerRecurringTask({ key: formatRecurringKey(sourceType, id), title: buildRecurringTitle(sourceType), category: 'data-source', intervalMs: require('../../lib/background-task-queue').normalizeIntervalMs(intervalValue, intervalUnit), metadata: { sourceType: sourceType, sourceId: Number(id), sourceName: name }, run: run }); } function removeRecurringRefresh(sourceType, id) { backgroundTaskQueue.removeRecurringTask(formatRecurringKey(sourceType, id)); } function getTaskStatusById(taskId) { if (!backgroundTaskQueue || typeof backgroundTaskQueue.getTaskById !== 'function') { return null; } const task = backgroundTaskQueue.getTaskById(taskId); if (!task) { return null; } return { id: task.id, key: task.key, status: task.status, finishedAt: task.finishedAt || '', errorMessage: task.errorMessage || '' }; } async function loadRssFeedItems(feedUrl, itemLimit) { if (typeof fetchRssFeedItems === 'function') { return fetchRssFeedItems(feedUrl, itemLimit); } return []; } async function loadApiSourceResponse(apiUrl) { if (typeof fetchApiSourceResponse === 'function') { return fetchApiSourceResponse(apiUrl); } return { responseJson: null, responseStatus: null, responseContentType: null }; } function sendRefreshTaskState(req, res, sourceType, sourceId) { const taskId = Number(req.query.refresh_task_id); if (!Number.isFinite(taskId) || taskId <= 0) { return res.status(400).json({ error: 'Missing refresh task id.' }); } const task = getTaskStatusById(taskId); const expectedKey = sourceType + '-refresh:' + Number(sourceId); if (!task || task.key !== expectedKey) { return res.status(404).json({ error: 'Refresh task not found.' }); } res.json(task); } app.get('/data-sources', function (req, res, next) { if (!req.currentUser) { return res.redirect('/login?message=' + encodeURIComponent('Please sign in to continue.')); } if (hasAnyPermission(req.currentUser, ['rss-feeds.read', 'api-sources.read'])) { if (hasAnyPermission(req.currentUser, ['rss-feeds.read'])) { return res.redirect('/data-sources/rss-feeds'); } return res.redirect('/data-sources/api-sources'); } const error = new Error('You do not have permission to access this area.'); error.statusCode = 403; error.expose = true; next(error); }); app.get('/data-sources/api-sources', requirePermission('api-sources.read'), async function (req, res, next) { try { const page = Math.max(1, Math.floor(Number(req.query.page) || 1)); const data = await common.fetchApiSourcesPage(pool, page, LIST_PAGE_SIZE); const apiSources = (data.apiSources || []).map(function (apiSource) { return Object.assign({}, apiSource, { intervalLabel: apiSource.update_interval_unit === 'seconds' ? (Math.max(1, Number(apiSource.update_interval_value) || 0) === 1 ? 'Every second' : `Every ${Math.max(1, Number(apiSource.update_interval_value) || 0)} seconds`) : (Math.max(1, Number(apiSource.update_interval_value) || 0) === 1 ? 'Every minute' : `Every ${Math.max(1, Number(apiSource.update_interval_value) || 0)} minutes`), lastPullLabel: apiSource.last_pulled_at ? formatDashboardDate(apiSource.last_pulled_at) : 'Never', lastPulledAtValue: apiSource.last_pulled_at ? new Date(apiSource.last_pulled_at).toISOString() : '' }); }); res.send(pages.renderApiSourcesPage({ apiSources: apiSources, pagination: buildPagination(data.totalItems, data.currentPage, 'page', {}, LIST_PAGE_SIZE, 'API sources', 'API source pages') }, req.query.message ? String(req.query.message) : '', req.currentUser, formatDashboardDate)); } catch (error) { next(error); } }); app.get('/data-sources/api-sources/new', requirePermission('api-sources.create'), function (req, res) { res.send(pages.renderApiSourceFormPage(null, 'create', req.query.message ? String(req.query.message) : '', req.currentUser)); }); app.post('/data-sources/api-sources', requirePermission('api-sources.create'), async function (req, res, next) { const connection = await pool.getConnection(); try { const payload = common.buildApiSourcePayload(req, null); if (await common.fetchDuplicateName(pool, 'api_sources', payload.name)) { return res.redirect('/data-sources/api-sources/new?message=' + encodeURIComponent('An API source with that name already exists.')); } const actorId = getAuditUserId(req); await connection.beginTransaction(); const [result] = await connection.query( 'INSERT INTO api_sources (name, api_url, update_interval_value, update_interval_unit, last_pulled_at, last_pull_error, last_response_status, last_response_content_type, last_response_json, created_by, modified_by) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)', [payload.name, payload.apiUrl, payload.updateIntervalValue, payload.updateIntervalUnit, null, null, null, null, null, actorId, actorId] ); await connection.commit(); registerRecurringRefresh('api-source', result.insertId, payload.name, payload.updateIntervalValue, payload.updateIntervalUnit, function () { return refreshApiSourceInBackground(result.insertId, payload.apiUrl, actorId); }); const refreshTask = await backgroundTaskQueue.enqueueTask({ key: 'api-source-refresh:' + result.insertId, title: 'API source refresh', category: 'data-source', taskType: 'data-source-refresh', payload: { sourceType: 'api-source', sourceId: result.insertId, sourceName: payload.name, actorId: actorId }, metadata: { sourceType: 'api-source', sourceId: result.insertId, sourceName: payload.name } }); const message = 'API source created. Refresh is running in the background.'; res.redirect('/data-sources/api-sources/' + result.insertId + '/edit?message=' + encodeURIComponent(message) + '&refresh_task_id=' + encodeURIComponent(refreshTask && refreshTask.id ? refreshTask.id : '')); } catch (error) { try { await connection.rollback(); } catch (_rollbackError) { // Ignore rollback failures and surface the original error. } next(error); } finally { connection.release(); } }); app.get('/data-sources/api-sources/:id/edit', requirePermission('api-sources.update'), async function (req, res, next) { try { const apiSource = await common.fetchApiSourceById(pool, Number(req.params.id)); if (!apiSource) { return res.status(404).send('API source not found'); } res.send(pages.renderApiSourceEditPage(Object.assign({}, apiSource, { apiUrl: apiSource.api_url, updateIntervalValue: apiSource.update_interval_value, updateIntervalUnit: apiSource.update_interval_unit || 'minutes', lastPulledAtValue: toIsoTimestamp(apiSource.last_pulled_at), lastPulledAtLabel: apiSource.last_pulled_at ? String(apiSource.last_pulled_at) : '', lastPullError: apiSource.last_pull_error || '', lastResponseStatus: apiSource.last_response_status, lastResponseContentType: apiSource.last_response_content_type, lastResponseJson: apiSource.last_response_json || '' }), { lastResponseJson: apiSource.last_response_json || '', lastPullError: apiSource.last_pull_error || '' }, req.query.message ? String(req.query.message) : '', req.currentUser)); } catch (error) { next(error); } }); app.get('/data-sources/api-sources/:id/state', requirePermission('api-sources.update'), function (req, res) { sendRefreshTaskState(req, res, 'api-source', Number(req.params.id)); }); app.post('/data-sources/api-sources/:id', requirePermission('api-sources.update'), async function (req, res, next) { const connection = await pool.getConnection(); try { const apiSource = await common.fetchApiSourceById(pool, Number(req.params.id)); if (!apiSource) { return res.status(404).send('API source not found'); } const payload = common.buildApiSourcePayload(req, apiSource); if (await common.fetchDuplicateName(pool, 'api_sources', payload.name, apiSource.id)) { return res.redirect('/data-sources/api-sources/' + apiSource.id + '/edit?message=' + encodeURIComponent('An API source with that name already exists.')); } const actorId = getAuditUserId(req); await connection.beginTransaction(); await connection.query( 'UPDATE api_sources SET name = ?, api_url = ?, update_interval_value = ?, update_interval_unit = ?, last_pulled_at = ?, last_pull_error = ?, last_response_status = ?, last_response_content_type = ?, last_response_json = ?, modified_by = ? WHERE id = ?', [payload.name, payload.apiUrl, payload.updateIntervalValue, payload.updateIntervalUnit, null, null, null, null, null, actorId, apiSource.id] ); await connection.commit(); registerRecurringRefresh('api-source', apiSource.id, payload.name, payload.updateIntervalValue, payload.updateIntervalUnit, function () { return refreshApiSourceInBackground(apiSource.id, payload.apiUrl, actorId); }); const message = 'API source updated. Refresh is running in the background.'; const refreshTask = await backgroundTaskQueue.enqueueTask({ key: 'api-source-refresh:' + apiSource.id, title: 'API source refresh', category: 'data-source', taskType: 'data-source-refresh', payload: { sourceType: 'api-source', sourceId: apiSource.id, sourceName: payload.name, actorId: actorId }, metadata: { sourceType: 'api-source', sourceId: apiSource.id, sourceName: payload.name } }); res.redirect('/data-sources/api-sources/' + apiSource.id + '/edit?message=' + encodeURIComponent(message) + '&refresh_task_id=' + encodeURIComponent(refreshTask && refreshTask.id ? refreshTask.id : '')); } catch (error) { try { await connection.rollback(); } catch (_rollbackError) { // Ignore rollback failures and surface the original error. } next(error); } finally { connection.release(); } }); app.post('/data-sources/api-sources/:id/delete', requirePermission('api-sources.delete'), async function (req, res, next) { try { const apiSource = await common.fetchApiSourceById(pool, Number(req.params.id)); if (!apiSource) { return res.status(404).send('API source not found'); } const connection = await pool.getConnection(); try { await connection.beginTransaction(); await connection.query('DELETE FROM api_sources WHERE id = ?', [apiSource.id]); await connection.commit(); removeRecurringRefresh('api-source', apiSource.id); } catch (error) { await connection.rollback(); throw error; } finally { connection.release(); } res.redirect('/data-sources/api-sources?message=' + encodeURIComponent('API source deleted.')); } catch (error) { next(error); } }); app.get('/data-sources/rss-feeds', requirePermission('rss-feeds.read'), async function (req, res, next) { try { const page = Math.max(1, Math.floor(Number(req.query.page) || 1)); const data = await common.fetchRssFeedsPage(pool, page, LIST_PAGE_SIZE); res.send(pages.renderRssFeedsPage({ rssFeeds: data.rssFeeds || [], pagination: buildPagination(data.totalItems, data.currentPage, 'page', {}, LIST_PAGE_SIZE, 'RSS feeds', 'RSS feed pages') }, req.query.message ? String(req.query.message) : '', req.currentUser)); } catch (error) { next(error); } }); app.get('/data-sources/rss-feeds/new', requirePermission('rss-feeds.create'), function (req, res) { res.send(pages.renderRssFeedFormPage(null, 'create', req.query.message ? String(req.query.message) : '', req.currentUser)); }); app.post('/data-sources/rss-feeds', requirePermission('rss-feeds.create'), async function (req, res, next) { const connection = await pool.getConnection(); try { const payload = common.buildRssFeedPayload(req, null); if (await common.fetchDuplicateName(pool, 'rss_feeds', payload.name)) { return res.redirect('/data-sources/rss-feeds/new?message=' + encodeURIComponent('An RSS feed with that name already exists.')); } const actorId = getAuditUserId(req); await connection.beginTransaction(); const [result] = await connection.query( 'INSERT INTO rss_feeds (name, feed_url, update_interval_value, update_interval_unit, item_limit, created_by, modified_by) VALUES (?, ?, ?, ?, ?, ?, ?)', [payload.name, payload.feedUrl, payload.updateIntervalValue, payload.updateIntervalUnit, payload.itemLimit, actorId, actorId] ); await connection.commit(); registerRecurringRefresh('rss-feed', result.insertId, payload.name, payload.updateIntervalValue, payload.updateIntervalUnit, function () { return refreshRssFeedInBackground(result.insertId, payload.feedUrl, payload.itemLimit, actorId); }); const refreshTask = await backgroundTaskQueue.enqueueTask({ key: 'rss-feed-refresh:' + result.insertId, title: 'RSS feed refresh', category: 'data-source', taskType: 'data-source-refresh', payload: { sourceType: 'rss-feed', sourceId: result.insertId, sourceName: payload.name, actorId: actorId }, metadata: { sourceType: 'rss-feed', sourceId: result.insertId, sourceName: payload.name } }); const message = 'RSS feed created. Refresh is running in the background.'; res.redirect('/data-sources/rss-feeds/' + result.insertId + '/edit?message=' + encodeURIComponent(message) + '&refresh_task_id=' + encodeURIComponent(refreshTask && refreshTask.id ? refreshTask.id : '')); } catch (error) { try { await connection.rollback(); } catch (_rollbackError) { // Ignore rollback failures and surface the original error. } next(error); } finally { connection.release(); } }); app.get('/data-sources/rss-feeds/:id/edit', requirePermission('rss-feeds.update'), async function (req, res, next) { try { const rssFeed = await common.fetchRssFeedById(pool, Number(req.params.id)); if (!rssFeed) { return res.status(404).send('RSS feed not found'); } const pulledItems = typeof common.fetchRssFeedItemsByFeedId === 'function' ? await common.fetchRssFeedItemsByFeedId(pool, rssFeed.id) : []; res.send(pages.renderRssFeedEditPage(Object.assign({}, rssFeed, { feedUrl: rssFeed.feed_url, updateIntervalValue: rssFeed.update_interval_value, updateIntervalUnit: rssFeed.update_interval_unit || 'minutes', itemLimit: rssFeed.item_limit }), { pulledItems: pulledItems, pullError: '' }, req.query.message ? String(req.query.message) : '', req.currentUser)); } catch (error) { next(error); } }); app.get('/data-sources/rss-feeds/:id/state', requirePermission('rss-feeds.update'), function (req, res) { sendRefreshTaskState(req, res, 'rss-feed', Number(req.params.id)); }); app.post('/data-sources/rss-feeds/:id', requirePermission('rss-feeds.update'), async function (req, res, next) { const connection = await pool.getConnection(); try { const rssFeed = await common.fetchRssFeedById(pool, Number(req.params.id)); if (!rssFeed) { return res.status(404).send('RSS feed not found'); } const payload = common.buildRssFeedPayload(req, rssFeed); if (await common.fetchDuplicateName(pool, 'rss_feeds', payload.name, rssFeed.id)) { return res.redirect('/data-sources/rss-feeds/' + rssFeed.id + '/edit?message=' + encodeURIComponent('An RSS feed with that name already exists.')); } const actorId = getAuditUserId(req); await connection.beginTransaction(); await connection.query( 'UPDATE rss_feeds SET name = ?, feed_url = ?, update_interval_value = ?, update_interval_unit = ?, item_limit = ?, modified_by = ? WHERE id = ?', [payload.name, payload.feedUrl, payload.updateIntervalValue, payload.updateIntervalUnit, payload.itemLimit, actorId, rssFeed.id] ); await connection.commit(); registerRecurringRefresh('rss-feed', rssFeed.id, payload.name, payload.updateIntervalValue, payload.updateIntervalUnit, function () { return refreshRssFeedInBackground(rssFeed.id, payload.feedUrl, payload.itemLimit, actorId); }); const message = 'RSS feed updated. Refresh is running in the background.'; const refreshTask = await backgroundTaskQueue.enqueueTask({ key: 'rss-feed-refresh:' + rssFeed.id, title: 'RSS feed refresh', category: 'data-source', taskType: 'data-source-refresh', payload: { sourceType: 'rss-feed', sourceId: rssFeed.id, sourceName: payload.name, actorId: actorId }, metadata: { sourceType: 'rss-feed', sourceId: rssFeed.id, sourceName: payload.name } }); res.redirect('/data-sources/rss-feeds/' + rssFeed.id + '/edit?message=' + encodeURIComponent(message) + '&refresh_task_id=' + encodeURIComponent(refreshTask && refreshTask.id ? refreshTask.id : '')); } catch (error) { try { await connection.rollback(); } catch (_rollbackError) { // Ignore rollback failures and surface the original error. } next(error); } finally { connection.release(); } }); app.post('/data-sources/rss-feeds/:id/delete', requirePermission('rss-feeds.delete'), async function (req, res, next) { try { const rssFeed = await common.fetchRssFeedById(pool, Number(req.params.id)); if (!rssFeed) { return res.status(404).send('RSS feed not found'); } const connection = await pool.getConnection(); try { await connection.beginTransaction(); await connection.query('DELETE FROM rss_feeds WHERE id = ?', [rssFeed.id]); await connection.commit(); removeRecurringRefresh('rss-feed', rssFeed.id); } catch (error) { await connection.rollback(); throw error; } finally { connection.release(); } res.redirect('/data-sources/rss-feeds?message=' + encodeURIComponent('RSS feed deleted.')); } catch (error) { next(error); } }); };