Add Olist sync logs and stop control
All checks were successful
Build and Deploy / build-and-deploy (push) Successful in 47s
All checks were successful
Build and Deploy / build-and-deploy (push) Successful in 47s
This commit is contained in:
@@ -17,6 +17,8 @@ const OLIST_AUTHORIZE_URL = 'https://accounts.tiny.com.br/realms/tiny/protocol/o
|
||||
const OAUTH_STATE_TTL_MS = 10 * 60 * 1000;
|
||||
const RUN_LOCK_KEY = 8_371_204;
|
||||
let activeSync = null;
|
||||
let activeRunId = null;
|
||||
let stopRequested = false;
|
||||
let scheduledSync = null;
|
||||
let nextOlistRequestAt = 0;
|
||||
|
||||
@@ -82,6 +84,7 @@ const mapRun = (row) => ({
|
||||
status: row.status,
|
||||
startedAt: row.started_at,
|
||||
completedAt: row.completed_at,
|
||||
cancelRequestedAt: row.cancel_requested_at || null,
|
||||
productsScanned: Number(row.products_scanned || 0),
|
||||
manufacturedProducts: Number(row.manufactured_products || 0),
|
||||
imported: Number(row.imported || 0),
|
||||
@@ -107,6 +110,90 @@ const listSyncRuns = async (limit = 5) => {
|
||||
return result.rows.map(mapRun);
|
||||
};
|
||||
|
||||
const logRunEvent = async (runId, {
|
||||
level = 'info',
|
||||
eventType,
|
||||
message,
|
||||
productId = null,
|
||||
sku = '',
|
||||
productDescription = '',
|
||||
metadata = {}
|
||||
}) => {
|
||||
const result = await pool.query(`
|
||||
INSERT INTO olist_sync_events (
|
||||
run_id, level, event_type, message, olist_product_id, sku, product_description, metadata
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8::jsonb)
|
||||
RETURNING *;
|
||||
`, [runId, level, eventType, message, productId ? String(productId) : null, sku || null, productDescription || null, JSON.stringify(metadata)]);
|
||||
return result.rows[0];
|
||||
};
|
||||
|
||||
const isCancellationRequested = async (runId) => {
|
||||
const result = await pool.query('SELECT cancel_requested_at FROM olist_sync_runs WHERE id = $1;', [runId]);
|
||||
return Boolean(result.rows[0]?.cancel_requested_at);
|
||||
};
|
||||
|
||||
const mapEvent = (row) => ({
|
||||
id: Number(row.id),
|
||||
level: row.level,
|
||||
eventType: row.event_type,
|
||||
message: row.message,
|
||||
productId: row.olist_product_id ? String(row.olist_product_id) : null,
|
||||
sku: row.sku || '',
|
||||
productDescription: row.product_description || '',
|
||||
metadata: row.metadata || {},
|
||||
createdAt: row.created_at
|
||||
});
|
||||
|
||||
const getOlistRunDetails = async (runId, { eventsPage = 1, productsPage = 1, pageSize = 20 } = {}) => {
|
||||
const normalizedRunId = Number(runId);
|
||||
const normalizedPageSize = Math.max(5, Math.min(Number(pageSize) || 20, 50));
|
||||
const normalizePage = (page) => Math.max(1, Number(page) || 1);
|
||||
const normalizedEventsPage = normalizePage(eventsPage);
|
||||
const normalizedProductsPage = normalizePage(productsPage);
|
||||
const runResult = await pool.query('SELECT * FROM olist_sync_runs WHERE id = $1;', [normalizedRunId]);
|
||||
if (!runResult.rows[0]) {
|
||||
const error = new Error('Execução Olist não encontrada.');
|
||||
error.statusCode = 404;
|
||||
throw error;
|
||||
}
|
||||
const [eventsResult, eventsCountResult, productsResult, productsCountResult] = await Promise.all([
|
||||
pool.query(`
|
||||
SELECT * FROM olist_sync_events
|
||||
WHERE run_id = $1
|
||||
ORDER BY created_at DESC, id DESC
|
||||
LIMIT $2 OFFSET $3;
|
||||
`, [normalizedRunId, normalizedPageSize, (normalizedEventsPage - 1) * normalizedPageSize]),
|
||||
pool.query('SELECT COUNT(*)::int AS count FROM olist_sync_events WHERE run_id = $1;', [normalizedRunId]),
|
||||
pool.query(`
|
||||
SELECT * FROM olist_sync_events
|
||||
WHERE run_id = $1 AND event_type = 'product_imported'
|
||||
ORDER BY created_at DESC, id DESC
|
||||
LIMIT $2 OFFSET $3;
|
||||
`, [normalizedRunId, normalizedPageSize, (normalizedProductsPage - 1) * normalizedPageSize]),
|
||||
pool.query(`
|
||||
SELECT COUNT(*)::int AS count FROM olist_sync_events
|
||||
WHERE run_id = $1 AND event_type = 'product_imported';
|
||||
`, [normalizedRunId])
|
||||
]);
|
||||
return {
|
||||
run: mapRun(runResult.rows[0]),
|
||||
events: {
|
||||
items: eventsResult.rows.map(mapEvent),
|
||||
total: Number(eventsCountResult.rows[0].count || 0),
|
||||
page: normalizedEventsPage,
|
||||
pageSize: normalizedPageSize
|
||||
},
|
||||
affectedProducts: {
|
||||
items: productsResult.rows.map(mapEvent),
|
||||
total: Number(productsCountResult.rows[0].count || 0),
|
||||
page: normalizedProductsPage,
|
||||
pageSize: normalizedPageSize
|
||||
}
|
||||
};
|
||||
};
|
||||
|
||||
const getCatalogSyncState = async () => {
|
||||
const result = await pool.query('SELECT * FROM olist_sync_state WHERE id = 1;');
|
||||
return result.rows[0] || null;
|
||||
@@ -291,10 +378,11 @@ const formatOlistDate = (value) => {
|
||||
return date.toISOString().slice(0, 19).replace('T', ' ');
|
||||
};
|
||||
|
||||
const listOlistProducts = async (changedAfter = null) => {
|
||||
const listOlistProducts = async (changedAfter = null, shouldStop = async () => false) => {
|
||||
const products = [];
|
||||
let offset = 0;
|
||||
do {
|
||||
if (await shouldStop()) return { products, cancelled: true };
|
||||
const params = new URLSearchParams({ limit: '100', offset: String(offset) });
|
||||
if (changedAfter) params.set('dataAlteracao', formatOlistDate(changedAfter));
|
||||
const page = await olistRequest(`/produtos?${params}`);
|
||||
@@ -302,7 +390,7 @@ const listOlistProducts = async (changedAfter = null) => {
|
||||
offset += page?.paginacao?.limit || 100;
|
||||
if (offset >= Number(page?.paginacao?.total || 0)) break;
|
||||
} while (true);
|
||||
return products;
|
||||
return { products, cancelled: false };
|
||||
};
|
||||
|
||||
const cacheOlistProducts = async (products) => {
|
||||
@@ -400,7 +488,7 @@ const updateRunProgress = async (id, summary) => {
|
||||
return mapRun(result.rows[0]);
|
||||
};
|
||||
|
||||
const finishRun = async (id, summary, error = null) => {
|
||||
const finishRun = async (id, summary, { status = 'completed', error = null } = {}) => {
|
||||
const result = await pool.query(`
|
||||
UPDATE olist_sync_runs
|
||||
SET status = $2,
|
||||
@@ -416,7 +504,7 @@ const finishRun = async (id, summary, error = null) => {
|
||||
details = $11::jsonb
|
||||
WHERE id = $1
|
||||
RETURNING *;
|
||||
`, [id, error ? 'failed' : 'completed', summary.productsScanned, summary.manufacturedProducts, summary.imported, summary.componentCount, summary.referenceCount, summary.skippedReferenceCount, summary.failed, error?.message || null, JSON.stringify(summary.details.slice(0, 100))]);
|
||||
`, [id, status, summary.productsScanned, summary.manufacturedProducts, summary.imported, summary.componentCount, summary.referenceCount, summary.skippedReferenceCount, summary.failed, error?.message || null, JSON.stringify(summary.details.slice(0, 100))]);
|
||||
return mapRun(result.rows[0]);
|
||||
};
|
||||
|
||||
@@ -428,9 +516,11 @@ const runOlistSync = async ({ trigger = 'manual', fullSync = false } = {}) => {
|
||||
const [checkpoint, cacheProgress] = await Promise.all([getCatalogSyncState(), getProductCacheProgress()]);
|
||||
const mode = fullSync || !checkpoint || cacheProgress.pending > 0 ? 'full' : 'incremental';
|
||||
const run = await createRun({ trigger, mode });
|
||||
activeRunId = run.id;
|
||||
const summary = { productsScanned: 0, manufacturedProducts: 0, imported: 0, componentCount: 0, referenceCount: 0, skippedReferenceCount: 0, failed: 0, details: [] };
|
||||
let lockClient;
|
||||
try {
|
||||
await logRunEvent(run.id, { eventType: 'started', message: 'Sincronização iniciada.' });
|
||||
lockClient = await pool.connect();
|
||||
const lock = await lockClient.query('SELECT pg_try_advisory_lock($1) AS locked;', [RUN_LOCK_KEY]);
|
||||
if (!lock.rows[0]?.locked) {
|
||||
@@ -441,52 +531,118 @@ const runOlistSync = async ({ trigger = 'manual', fullSync = false } = {}) => {
|
||||
|
||||
const resumePendingFullScan = mode === 'full' && cacheProgress.total > 0 && cacheProgress.pending > 0;
|
||||
if (!resumePendingFullScan) {
|
||||
const changedProducts = await listOlistProducts(mode === 'incremental' ? checkpoint.last_catalog_sync_at : null);
|
||||
const catalogResult = await listOlistProducts(
|
||||
mode === 'incremental' ? checkpoint.last_catalog_sync_at : null,
|
||||
async () => stopRequested || isCancellationRequested(run.id)
|
||||
);
|
||||
if (catalogResult.cancelled) {
|
||||
await logRunEvent(run.id, { eventType: 'cancelled', message: 'Sincronização interrompida a pedido do usuário antes de concluir a leitura do catálogo.' });
|
||||
return await finishRun(run.id, summary, { status: 'cancelled' });
|
||||
}
|
||||
const changedProducts = catalogResult.products;
|
||||
await cacheOlistProducts(changedProducts);
|
||||
await updateCatalogCheckpoint();
|
||||
await logRunEvent(run.id, {
|
||||
eventType: 'catalog_loaded',
|
||||
message: `${changedProducts.length.toLocaleString('pt-BR')} produtos recebidos da Olist.`,
|
||||
metadata: { productCount: changedProducts.length, mode }
|
||||
});
|
||||
} else {
|
||||
await logRunEvent(run.id, {
|
||||
eventType: 'resumed',
|
||||
message: 'Retomando estruturas pendentes da sincronização anterior.',
|
||||
metadata: { pendingStructures: cacheProgress.pending }
|
||||
});
|
||||
}
|
||||
const [candidates, unitById] = await Promise.all([listStructureCandidates(), getCachedUnitById()]);
|
||||
for (const product of candidates) {
|
||||
if (stopRequested || await isCancellationRequested(run.id)) {
|
||||
await logRunEvent(run.id, { eventType: 'cancelled', message: 'Sincronização interrompida a pedido do usuário.' });
|
||||
return await finishRun(run.id, summary, { status: 'cancelled' });
|
||||
}
|
||||
summary.productsScanned += 1;
|
||||
try {
|
||||
const structure = await olistRequest(`/produtos/${product.id}/fabricado`);
|
||||
const components = Array.isArray(structure?.produtos) ? structure.produtos : [];
|
||||
if (!components.length) {
|
||||
await markStructureChecked(product.id, false);
|
||||
continue;
|
||||
await logRunEvent(run.id, {
|
||||
eventType: 'product_skipped',
|
||||
message: `Produto sem composição fabricada: ${product.descricao}.`,
|
||||
productId: product.id,
|
||||
sku: product.sku,
|
||||
productDescription: product.descricao
|
||||
});
|
||||
} else {
|
||||
const result = await upsertTinyProductComposition({
|
||||
tinyId: String(product.id),
|
||||
productSku: product.sku,
|
||||
productDescription: product.descricao,
|
||||
unit: product.unidade || 'UN',
|
||||
rawPayload: { product, structure }
|
||||
}, components.map(component => ({
|
||||
componentTinyId: String(component?.produto?.id || ''),
|
||||
componentSku: component?.produto?.sku || '',
|
||||
componentName: component?.produto?.descricao || '',
|
||||
quantityPerUnit: component?.quantidade,
|
||||
unit: unitById.get(String(component?.produto?.id)) || 'UN'
|
||||
})));
|
||||
summary.manufacturedProducts += 1;
|
||||
summary.imported += 1;
|
||||
summary.componentCount += result.componentCount || 0;
|
||||
summary.referenceCount += result.referenceCount || 0;
|
||||
summary.skippedReferenceCount += result.skippedReferenceCount || 0;
|
||||
await markStructureChecked(product.id, true);
|
||||
await logRunEvent(run.id, {
|
||||
eventType: 'product_imported',
|
||||
message: `Composição importada: ${product.descricao}.`,
|
||||
productId: product.id,
|
||||
sku: product.sku,
|
||||
productDescription: product.descricao,
|
||||
metadata: {
|
||||
componentCount: result.componentCount || 0,
|
||||
referenceCount: result.referenceCount || 0
|
||||
}
|
||||
});
|
||||
}
|
||||
const result = await upsertTinyProductComposition({
|
||||
tinyId: String(product.id),
|
||||
productSku: product.sku,
|
||||
productDescription: product.descricao,
|
||||
unit: product.unidade || 'UN',
|
||||
rawPayload: { product, structure }
|
||||
}, components.map(component => ({
|
||||
componentTinyId: String(component?.produto?.id || ''),
|
||||
componentSku: component?.produto?.sku || '',
|
||||
componentName: component?.produto?.descricao || '',
|
||||
quantityPerUnit: component?.quantidade,
|
||||
unit: unitById.get(String(component?.produto?.id)) || 'UN'
|
||||
})));
|
||||
summary.manufacturedProducts += 1;
|
||||
summary.imported += 1;
|
||||
summary.componentCount += result.componentCount || 0;
|
||||
summary.referenceCount += result.referenceCount || 0;
|
||||
summary.skippedReferenceCount += result.skippedReferenceCount || 0;
|
||||
await markStructureChecked(product.id, true);
|
||||
} catch (error) {
|
||||
if (error.statusCode === 400 || error.statusCode === 404) {
|
||||
await markStructureChecked(product.id, false);
|
||||
await logRunEvent(run.id, {
|
||||
eventType: 'product_skipped',
|
||||
message: `Produto sem composição fabricada: ${product.descricao}.`,
|
||||
productId: product.id,
|
||||
sku: product.sku,
|
||||
productDescription: product.descricao
|
||||
});
|
||||
} else {
|
||||
summary.failed += 1;
|
||||
summary.details.push({ productId: product.id, sku: product.sku || '', error: error.message });
|
||||
await logRunEvent(run.id, {
|
||||
level: 'error',
|
||||
eventType: 'product_failed',
|
||||
message: `Falha ao consultar a composição: ${error.message}`,
|
||||
productId: product.id,
|
||||
sku: product.sku,
|
||||
productDescription: product.descricao,
|
||||
metadata: { error: error.message }
|
||||
});
|
||||
}
|
||||
}
|
||||
await updateRunProgress(run.id, summary);
|
||||
if (summary.productsScanned % 25 === 0) {
|
||||
await logRunEvent(run.id, {
|
||||
eventType: 'progress',
|
||||
message: `${summary.productsScanned.toLocaleString('pt-BR')} estruturas verificadas; ${summary.imported.toLocaleString('pt-BR')} composições importadas.`,
|
||||
metadata: { productsScanned: summary.productsScanned, imported: summary.imported, failed: summary.failed }
|
||||
});
|
||||
}
|
||||
}
|
||||
await logRunEvent(run.id, { eventType: 'completed', message: 'Sincronização concluída.', metadata: summary });
|
||||
return await finishRun(run.id, summary);
|
||||
} catch (error) {
|
||||
await finishRun(run.id, summary, error);
|
||||
await logRunEvent(run.id, { level: 'error', eventType: 'failed', message: error.message, metadata: { error: error.message } }).catch(() => {});
|
||||
await finishRun(run.id, summary, { status: 'failed', error });
|
||||
throw error;
|
||||
} finally {
|
||||
if (lockClient) {
|
||||
@@ -494,7 +650,11 @@ const runOlistSync = async ({ trigger = 'manual', fullSync = false } = {}) => {
|
||||
lockClient.release();
|
||||
}
|
||||
}
|
||||
})().finally(() => { activeSync = null; });
|
||||
})().finally(() => {
|
||||
activeSync = null;
|
||||
activeRunId = null;
|
||||
stopRequested = false;
|
||||
});
|
||||
|
||||
return activeSync;
|
||||
};
|
||||
@@ -511,6 +671,26 @@ const startOlistSync = async (options) => {
|
||||
return { accepted: true };
|
||||
};
|
||||
|
||||
const requestOlistSyncStop = async () => {
|
||||
if (!activeSync) {
|
||||
const error = new Error('Não há uma sincronização Olist em andamento.');
|
||||
error.statusCode = 409;
|
||||
throw error;
|
||||
}
|
||||
stopRequested = true;
|
||||
if (!activeRunId) return { accepted: true, initializing: true };
|
||||
const result = await pool.query(`
|
||||
UPDATE olist_sync_runs
|
||||
SET cancel_requested_at = COALESCE(cancel_requested_at, CURRENT_TIMESTAMP)
|
||||
WHERE id = $1
|
||||
RETURNING *;
|
||||
`, [activeRunId]);
|
||||
if (result.rows[0]) {
|
||||
await logRunEvent(activeRunId, { eventType: 'cancel_requested', message: 'Parada solicitada pelo usuário.' });
|
||||
}
|
||||
return { accepted: true, initializing: false };
|
||||
};
|
||||
|
||||
const scheduleOlistSync = () => {
|
||||
if (scheduledSync) clearInterval(scheduledSync);
|
||||
scheduledSync = setInterval(() => {
|
||||
@@ -525,9 +705,11 @@ module.exports = {
|
||||
completeAuthorization,
|
||||
createAuthorizationUrl,
|
||||
getOlistFrontendRedirect,
|
||||
getOlistRunDetails,
|
||||
getOlistStatus,
|
||||
isOlistConfigured,
|
||||
markInterruptedOlistRuns,
|
||||
scheduleOlistSync,
|
||||
requestOlistSyncStop,
|
||||
startOlistSync
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user