semillero-INVIMA/backend/invima/sync-job-manager.js

193 lines
5.5 KiB
JavaScript

const {
getSyncDatasetKeys,
listInvimaDatasets,
syncInvimaDatasets
} = require('./service');
let activeJob = null;
let lastJob = null;
let sequence = 0;
const cloneJob = (job) => (job ? JSON.parse(JSON.stringify(job)) : null);
const nowIso = () => new Date().toISOString();
const datasetMetadata = () => {
const map = new Map();
for (const item of listInvimaDatasets()) {
map.set(item.key, item);
}
return map;
};
const recalculateTotals = (job) => {
job.completed = job.datasets.filter((item) => item.status === 'completed').length;
job.failed = job.datasets.filter((item) => item.status === 'failed').length;
job.fetched = job.datasets.reduce((total, item) => total + Number(item.fetched || 0), 0);
job.persisted = job.datasets.reduce((total, item) => total + Number(item.persisted || 0), 0);
job.pages = job.datasets.reduce((total, item) => total + Number(item.pages || 0), 0);
job.updatedAt = nowIso();
};
const findDatasetProgress = (job, datasetKey) =>
job.datasets.find((item) => item.key === datasetKey || item.datasetKey === datasetKey);
const applyProgress = (job, event) => {
if (!job || !event) return;
if (event.type === 'sync_started') {
job.requested = Array.isArray(event.keys) ? event.keys.length : job.requested;
job.concurrency = event.concurrency || job.concurrency;
recalculateTotals(job);
return;
}
const dataset = findDatasetProgress(job, event.datasetKey);
if (!dataset) {
recalculateTotals(job);
return;
}
if (event.type === 'dataset_started') {
dataset.status = 'running';
dataset.startedAt = dataset.startedAt || nowIso();
}
if (event.type === 'page_persisted') {
dataset.status = 'running';
dataset.fetched = Number(event.fetched || dataset.fetched || 0);
dataset.persisted = Number(event.persisted || dataset.persisted || 0);
dataset.pages = Number(event.pages || dataset.pages || 0);
dataset.pageSizeUsed = event.pageSizeUsed || dataset.pageSizeUsed || null;
dataset.warnings = Array.isArray(event.warnings) ? event.warnings : dataset.warnings;
}
if (event.type === 'dataset_completed') {
dataset.status = 'completed';
dataset.datasetId = event.datasetId || dataset.datasetId;
dataset.fetched = Number(event.fetched || dataset.fetched || 0);
dataset.persisted = Number(event.persisted || dataset.persisted || 0);
dataset.pages = Number(event.pages || dataset.pages || 0);
dataset.pageSizeUsed = event.pageSizeUsed || dataset.pageSizeUsed || null;
dataset.orderUsed = event.orderUsed || dataset.orderUsed || null;
dataset.warnings = Array.isArray(event.warnings) ? event.warnings : dataset.warnings;
dataset.finishedAt = nowIso();
}
if (event.type === 'dataset_failed') {
dataset.status = 'failed';
dataset.error = event.error || 'sync_failed';
dataset.finishedAt = nowIso();
}
dataset.updatedAt = nowIso();
recalculateTotals(job);
};
const finishJob = (job, status, resultOrError) => {
if (!job) return;
if (resultOrError?.datasets) {
for (const result of resultOrError.datasets) {
applyProgress(job, { type: 'dataset_completed', ...result });
}
}
if (resultOrError?.failed) {
for (const failure of resultOrError.failed) {
applyProgress(job, { type: 'dataset_failed', ...failure });
}
}
job.status = status;
job.summary = resultOrError?.summary || job.summary || null;
job.error = resultOrError instanceof Error ? String(resultOrError.message || resultOrError) : null;
job.finishedAt = nowIso();
recalculateTotals(job);
};
const createJob = (keys, options) => {
const meta = datasetMetadata();
const startedAt = nowIso();
return {
id: `${Date.now()}-${++sequence}`,
status: 'running',
requested: keys.length,
completed: 0,
failed: 0,
fetched: 0,
persisted: 0,
pages: 0,
concurrency: options.concurrency,
pageSize: options.pageSize,
datasets: keys.map((key) => {
const dataset = meta.get(key) || {};
return {
key,
datasetKey: key,
datasetId: dataset.id || null,
name: dataset.name || key,
status: 'pending',
fetched: 0,
persisted: 0,
pages: 0,
warnings: [],
error: null,
startedAt: null,
updatedAt: startedAt,
finishedAt: null
};
}),
summary: null,
error: null,
startedAt,
updatedAt: startedAt,
finishedAt: null
};
};
const startInvimaSyncJob = (pool, options = {}) => {
if (activeJob?.status === 'running') {
return { alreadyRunning: true, job: cloneJob(activeJob) };
}
const keys = getSyncDatasetKeys(options.datasets);
if (!keys.length) {
const error = new Error('No hay datasets validos para sincronizar.');
error.status = 400;
throw error;
}
const job = createJob(keys, options);
activeJob = job;
lastJob = job;
syncInvimaDatasets(pool, {
...options,
datasets: keys,
onProgress: (event) => applyProgress(job, event)
})
.then((result) => {
const hasSuccessfulDatasets = Number(result?.summary?.datasets || 0) > 0;
const status = hasSuccessfulDatasets || !result?.failed?.length ? 'completed' : 'failed';
finishJob(job, status, result);
})
.catch((error) => {
finishJob(job, 'failed', error);
})
.finally(() => {
if (activeJob?.id === job.id) {
activeJob = null;
}
});
return { alreadyRunning: false, job: cloneJob(job) };
};
const getInvimaSyncJob = () => cloneJob(activeJob || lastJob);
module.exports = {
getInvimaSyncJob,
startInvimaSyncJob
};