All checks were successful
Build and Deploy / build-and-deploy (push) Successful in 2m17s
789 lines
32 KiB
JavaScript
789 lines
32 KiB
JavaScript
const crypto = require('crypto');
|
|
const { pool } = require('../db');
|
|
const {
|
|
OLIST_CLIENT_ID,
|
|
OLIST_CLIENT_SECRET,
|
|
OLIST_FRONTEND_URL,
|
|
OLIST_REDIRECT_URI,
|
|
OLIST_SYNC_HOUR,
|
|
OLIST_SYNC_REQUEST_DELAY_MS,
|
|
OLIST_TOKEN_ENCRYPTION_KEY
|
|
} = require('../config');
|
|
const { upsertTinyProductComposition } = require('./productionOrderService');
|
|
|
|
const OLIST_API_URL = 'https://api.tiny.com.br/public-api/v3';
|
|
const OLIST_TOKEN_URL = 'https://accounts.tiny.com.br/realms/tiny/protocol/openid-connect/token';
|
|
const OLIST_AUTHORIZE_URL = 'https://accounts.tiny.com.br/realms/tiny/protocol/openid-connect/auth';
|
|
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 nextScheduledSyncAt = null;
|
|
let nextOlistRequestAt = 0;
|
|
|
|
const normalizeText = (value) => String(value || '').trim();
|
|
const stateHash = (state) => crypto.createHash('sha256').update(state).digest('hex');
|
|
|
|
const getEncryptionKey = () => {
|
|
const raw = normalizeText(OLIST_TOKEN_ENCRYPTION_KEY);
|
|
if (/^[a-f0-9]{64}$/i.test(raw)) return Buffer.from(raw, 'hex');
|
|
try {
|
|
const decoded = Buffer.from(raw, 'base64');
|
|
return decoded.length === 32 ? decoded : null;
|
|
} catch {
|
|
return null;
|
|
}
|
|
};
|
|
|
|
const isOlistConfigured = () => Boolean(
|
|
normalizeText(OLIST_CLIENT_ID)
|
|
&& normalizeText(OLIST_CLIENT_SECRET)
|
|
&& normalizeText(OLIST_REDIRECT_URI)
|
|
&& normalizeText(OLIST_FRONTEND_URL)
|
|
&& getEncryptionKey()
|
|
);
|
|
|
|
const configuredError = () => {
|
|
const error = new Error('Configure as credenciais Olist V3 e a chave de criptografia no ambiente antes de conectar.');
|
|
error.statusCode = 409;
|
|
return error;
|
|
};
|
|
|
|
const assertConfigured = () => {
|
|
if (!isOlistConfigured()) throw configuredError();
|
|
};
|
|
|
|
const encrypt = (value) => {
|
|
const key = getEncryptionKey();
|
|
if (!key) throw configuredError();
|
|
const iv = crypto.randomBytes(12);
|
|
const cipher = crypto.createCipheriv('aes-256-gcm', key, iv);
|
|
const ciphertext = Buffer.concat([cipher.update(value, 'utf8'), cipher.final()]);
|
|
return {
|
|
ciphertext: Buffer.concat([ciphertext, cipher.getAuthTag()]).toString('base64'),
|
|
iv: iv.toString('base64')
|
|
};
|
|
};
|
|
|
|
const decrypt = ({ ciphertext, iv }) => {
|
|
const key = getEncryptionKey();
|
|
if (!key) throw configuredError();
|
|
const payload = Buffer.from(ciphertext, 'base64');
|
|
const authTag = payload.subarray(-16);
|
|
const encrypted = payload.subarray(0, -16);
|
|
const decipher = crypto.createDecipheriv('aes-256-gcm', key, Buffer.from(iv, 'base64'));
|
|
decipher.setAuthTag(authTag);
|
|
return Buffer.concat([decipher.update(encrypted), decipher.final()]).toString('utf8');
|
|
};
|
|
|
|
const mapRun = (row) => ({
|
|
id: row.id,
|
|
trigger: row.trigger,
|
|
mode: row.mode,
|
|
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),
|
|
componentCount: Number(row.component_count || 0),
|
|
referenceCount: Number(row.reference_count || 0),
|
|
skippedReferenceCount: Number(row.skipped_reference_count || 0),
|
|
failed: Number(row.failed || 0),
|
|
errorMessage: row.error_message || '',
|
|
details: row.details || []
|
|
});
|
|
|
|
const getConnection = async () => {
|
|
const result = await pool.query('SELECT * FROM olist_connections WHERE id = 1;');
|
|
return result.rows[0] || null;
|
|
};
|
|
|
|
const listSyncRuns = async ({ page = 1, pageSize = 5 } = {}) => {
|
|
const normalizedPageSize = Math.max(1, Math.min(Number(pageSize) || 5, 20));
|
|
const normalizedPage = Math.max(1, Number(page) || 1);
|
|
const offset = (normalizedPage - 1) * normalizedPageSize;
|
|
const [runsResult, countResult] = await Promise.all([
|
|
pool.query(`
|
|
SELECT * FROM olist_sync_runs
|
|
ORDER BY started_at DESC
|
|
LIMIT $1 OFFSET $2;
|
|
`, [normalizedPageSize, offset]),
|
|
pool.query('SELECT COUNT(*)::int AS count FROM olist_sync_runs;')
|
|
]);
|
|
return {
|
|
items: runsResult.rows.map(mapRun),
|
|
total: Number(countResult.rows[0]?.count || 0),
|
|
page: normalizedPage,
|
|
pageSize: normalizedPageSize
|
|
};
|
|
};
|
|
|
|
const hasScheduledSyncToday = async () => {
|
|
const result = await pool.query(`
|
|
SELECT EXISTS (
|
|
SELECT 1
|
|
FROM olist_sync_runs
|
|
WHERE trigger = 'scheduled'
|
|
AND started_at >= (date_trunc('day', CURRENT_TIMESTAMP AT TIME ZONE 'America/Sao_Paulo') AT TIME ZONE 'America/Sao_Paulo')
|
|
) AS exists;
|
|
`);
|
|
return Boolean(result.rows[0]?.exists);
|
|
};
|
|
|
|
const getTodayOlistRunAt = () => {
|
|
const parts = new Intl.DateTimeFormat('en-CA', {
|
|
timeZone: 'America/Sao_Paulo',
|
|
year: 'numeric',
|
|
month: '2-digit',
|
|
day: '2-digit'
|
|
}).formatToParts(new Date());
|
|
const value = (type) => parts.find(part => part.type === type)?.value;
|
|
return new Date(`${value('year')}-${value('month')}-${value('day')}T${String(OLIST_SYNC_HOUR).padStart(2, '0')}:00:00-03:00`);
|
|
};
|
|
|
|
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;
|
|
};
|
|
|
|
const getProductCacheProgress = async () => {
|
|
const result = await pool.query(`
|
|
SELECT
|
|
COUNT(*)::int AS total,
|
|
COUNT(*) FILTER (WHERE structure_checked_at IS NOT NULL)::int AS checked,
|
|
COUNT(*) FILTER (WHERE structure_checked_at IS NULL)::int AS pending,
|
|
COUNT(*) FILTER (WHERE is_manufactured IS TRUE)::int AS manufactured
|
|
FROM olist_product_cache;
|
|
`);
|
|
const row = result.rows[0];
|
|
return {
|
|
total: Number(row.total || 0),
|
|
checked: Number(row.checked || 0),
|
|
pending: Number(row.pending || 0),
|
|
manufactured: Number(row.manufactured || 0)
|
|
};
|
|
};
|
|
|
|
const markInterruptedOlistRuns = async () => {
|
|
const result = await pool.query(`
|
|
UPDATE olist_sync_runs
|
|
SET status = 'failed',
|
|
completed_at = CURRENT_TIMESTAMP,
|
|
error_message = COALESCE(error_message, 'Sincronização interrompida pelo reinício do servidor.')
|
|
WHERE status = 'running';
|
|
`);
|
|
return result.rowCount || 0;
|
|
};
|
|
|
|
const getOlistStatus = async ({ runsPage = 1, runsPageSize = 5 } = {}) => {
|
|
const [connection, syncRuns, productCache, latestRunResult, completedRunResult] = await Promise.all([
|
|
getConnection(),
|
|
listSyncRuns({ page: runsPage, pageSize: runsPageSize }),
|
|
getProductCacheProgress(),
|
|
pool.query('SELECT * FROM olist_sync_runs ORDER BY started_at DESC LIMIT 1;'),
|
|
pool.query("SELECT EXISTS (SELECT 1 FROM olist_sync_runs WHERE status = 'completed') AS exists;")
|
|
]);
|
|
return {
|
|
configured: isOlistConfigured(),
|
|
connected: Boolean(connection),
|
|
tokenExpiresAt: connection?.token_expires_at || null,
|
|
syncInProgress: Boolean(activeSync),
|
|
nextScheduledSyncAt,
|
|
productCache,
|
|
latestRun: latestRunResult.rows[0] ? mapRun(latestRunResult.rows[0]) : null,
|
|
hasCompletedRun: Boolean(completedRunResult.rows[0]?.exists),
|
|
runs: syncRuns.items,
|
|
runsTotal: syncRuns.total,
|
|
runsPage: syncRuns.page,
|
|
runsPageSize: syncRuns.pageSize
|
|
};
|
|
};
|
|
|
|
const createAuthorizationUrl = async () => {
|
|
assertConfigured();
|
|
const state = crypto.randomBytes(32).toString('base64url');
|
|
await pool.query('DELETE FROM olist_oauth_states WHERE expires_at <= CURRENT_TIMESTAMP;');
|
|
await pool.query(`
|
|
INSERT INTO olist_oauth_states (state_hash, expires_at)
|
|
VALUES ($1, $2);
|
|
`, [stateHash(state), new Date(Date.now() + OAUTH_STATE_TTL_MS)]);
|
|
|
|
const params = new URLSearchParams({
|
|
client_id: OLIST_CLIENT_ID,
|
|
redirect_uri: OLIST_REDIRECT_URI,
|
|
scope: 'openid offline_access',
|
|
response_type: 'code',
|
|
state
|
|
});
|
|
return `${OLIST_AUTHORIZE_URL}?${params.toString()}`;
|
|
};
|
|
|
|
const exchangeToken = async (params) => {
|
|
const response = await fetch(OLIST_TOKEN_URL, {
|
|
method: 'POST',
|
|
headers: { 'Content-Type': 'application/x-www-form-urlencoded' },
|
|
body: new URLSearchParams({
|
|
client_id: OLIST_CLIENT_ID,
|
|
client_secret: OLIST_CLIENT_SECRET,
|
|
...params
|
|
})
|
|
});
|
|
const data = await response.json().catch(() => null);
|
|
if (!response.ok || !data?.access_token) {
|
|
const error = new Error(data?.error_description || data?.error || 'Não foi possível autorizar a conexão com a Olist.');
|
|
error.statusCode = 502;
|
|
throw error;
|
|
}
|
|
return data;
|
|
};
|
|
|
|
const persistTokens = async ({ accessToken, refreshToken, expiresIn, existingConnection }) => {
|
|
const access = encrypt(accessToken);
|
|
const refresh = refreshToken ? encrypt(refreshToken) : {
|
|
ciphertext: existingConnection.refresh_token_ciphertext,
|
|
iv: existingConnection.refresh_token_iv
|
|
};
|
|
const expiresAt = new Date(Date.now() + Math.max(Number(expiresIn) || 14_400, 60) * 1000);
|
|
await pool.query(`
|
|
INSERT INTO olist_connections (
|
|
id, access_token_ciphertext, access_token_iv, refresh_token_ciphertext,
|
|
refresh_token_iv, token_expires_at, connected_at, updated_at
|
|
)
|
|
VALUES (1, $1, $2, $3, $4, $5, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
|
ON CONFLICT (id) DO UPDATE
|
|
SET access_token_ciphertext = EXCLUDED.access_token_ciphertext,
|
|
access_token_iv = EXCLUDED.access_token_iv,
|
|
refresh_token_ciphertext = EXCLUDED.refresh_token_ciphertext,
|
|
refresh_token_iv = EXCLUDED.refresh_token_iv,
|
|
token_expires_at = EXCLUDED.token_expires_at,
|
|
updated_at = CURRENT_TIMESTAMP;
|
|
`, [access.ciphertext, access.iv, refresh.ciphertext, refresh.iv, expiresAt]);
|
|
};
|
|
|
|
const completeAuthorization = async ({ code, state }) => {
|
|
assertConfigured();
|
|
const result = await pool.query(`
|
|
DELETE FROM olist_oauth_states
|
|
WHERE state_hash = $1 AND expires_at > CURRENT_TIMESTAMP
|
|
RETURNING state_hash;
|
|
`, [stateHash(state)]);
|
|
if (!result.rows.length) {
|
|
const error = new Error('A autorização da Olist expirou ou é inválida. Tente conectar novamente.');
|
|
error.statusCode = 400;
|
|
throw error;
|
|
}
|
|
|
|
const token = await exchangeToken({
|
|
grant_type: 'authorization_code',
|
|
redirect_uri: OLIST_REDIRECT_URI,
|
|
code
|
|
});
|
|
if (!token.refresh_token) {
|
|
const error = new Error('A Olist não retornou um refresh token. Verifique as permissões do aplicativo.');
|
|
error.statusCode = 502;
|
|
throw error;
|
|
}
|
|
await persistTokens({ accessToken: token.access_token, refreshToken: token.refresh_token, expiresIn: token.expires_in, existingConnection: {} });
|
|
};
|
|
|
|
const getAccessToken = async (forceRefresh = false) => {
|
|
assertConfigured();
|
|
const connection = await getConnection();
|
|
if (!connection) {
|
|
const error = new Error('Conecte a conta Olist antes de sincronizar.');
|
|
error.statusCode = 409;
|
|
throw error;
|
|
}
|
|
|
|
if (!forceRefresh && new Date(connection.token_expires_at).getTime() > Date.now() + 60_000) {
|
|
return decrypt({ ciphertext: connection.access_token_ciphertext, iv: connection.access_token_iv });
|
|
}
|
|
|
|
const token = await exchangeToken({
|
|
grant_type: 'refresh_token',
|
|
refresh_token: decrypt({ ciphertext: connection.refresh_token_ciphertext, iv: connection.refresh_token_iv })
|
|
});
|
|
await persistTokens({ accessToken: token.access_token, refreshToken: token.refresh_token, expiresIn: token.expires_in, existingConnection: connection });
|
|
return token.access_token;
|
|
};
|
|
|
|
const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms));
|
|
|
|
const waitForOlistRequestSlot = async () => {
|
|
const now = Date.now();
|
|
const waitMs = Math.max(0, nextOlistRequestAt - now);
|
|
if (waitMs) await sleep(waitMs);
|
|
nextOlistRequestAt = Date.now() + OLIST_SYNC_REQUEST_DELAY_MS;
|
|
};
|
|
|
|
const olistRequest = async (path, retry = true) => {
|
|
await waitForOlistRequestSlot();
|
|
const response = await fetch(`${OLIST_API_URL}${path}`, {
|
|
headers: { Authorization: `Bearer ${await getAccessToken(!retry)}` }
|
|
});
|
|
if (response.status === 401 && retry) return olistRequest(path, false);
|
|
const data = await response.json().catch(() => null);
|
|
if (!response.ok) {
|
|
const error = new Error(data?.mensagem || `A Olist retornou HTTP ${response.status}.`);
|
|
error.statusCode = response.status;
|
|
throw error;
|
|
}
|
|
return data;
|
|
};
|
|
|
|
const formatOlistDate = (value) => {
|
|
const date = new Date(value);
|
|
if (Number.isNaN(date.getTime())) return '';
|
|
return date.toISOString().slice(0, 19).replace('T', ' ');
|
|
};
|
|
|
|
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}`);
|
|
products.push(...(Array.isArray(page?.itens) ? page.itens : []));
|
|
offset += page?.paginacao?.limit || 100;
|
|
if (offset >= Number(page?.paginacao?.total || 0)) break;
|
|
} while (true);
|
|
return { products, cancelled: false };
|
|
};
|
|
|
|
const cacheOlistProducts = async (products) => {
|
|
for (const product of products) {
|
|
if (!product?.id || !normalizeText(product.descricao)) continue;
|
|
await pool.query(`
|
|
INSERT INTO olist_product_cache (
|
|
olist_product_id, sku, description, unit, last_changed_at, updated_at
|
|
)
|
|
VALUES ($1, $2, $3, $4, $5, CURRENT_TIMESTAMP)
|
|
ON CONFLICT (olist_product_id) DO UPDATE
|
|
SET sku = EXCLUDED.sku,
|
|
description = EXCLUDED.description,
|
|
unit = EXCLUDED.unit,
|
|
structure_checked_at = CASE
|
|
WHEN olist_product_cache.last_changed_at IS DISTINCT FROM EXCLUDED.last_changed_at
|
|
THEN NULL
|
|
ELSE olist_product_cache.structure_checked_at
|
|
END,
|
|
last_changed_at = EXCLUDED.last_changed_at,
|
|
updated_at = CURRENT_TIMESTAMP;
|
|
`, [
|
|
String(product.id),
|
|
normalizeText(product.sku) || null,
|
|
normalizeText(product.descricao),
|
|
normalizeText(product.unidade) || 'UN',
|
|
product.dataAlteracao || product.dataCriacao || null
|
|
]);
|
|
}
|
|
};
|
|
|
|
const updateCatalogCheckpoint = async () => {
|
|
await pool.query(`
|
|
INSERT INTO olist_sync_state (id, last_catalog_sync_at, updated_at)
|
|
VALUES (1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
|
ON CONFLICT (id) DO UPDATE
|
|
SET last_catalog_sync_at = EXCLUDED.last_catalog_sync_at,
|
|
updated_at = CURRENT_TIMESTAMP;
|
|
`);
|
|
};
|
|
|
|
const listStructureCandidates = async () => {
|
|
const result = await pool.query(`
|
|
SELECT olist_product_id, sku, description, unit
|
|
FROM olist_product_cache
|
|
WHERE structure_checked_at IS NULL
|
|
ORDER BY olist_product_id;
|
|
`);
|
|
return result.rows.map(row => ({
|
|
id: String(row.olist_product_id),
|
|
sku: row.sku || '',
|
|
descricao: row.description,
|
|
unidade: row.unit || 'UN'
|
|
}));
|
|
};
|
|
|
|
const getCachedUnitById = async () => {
|
|
const result = await pool.query('SELECT olist_product_id, unit FROM olist_product_cache;');
|
|
return new Map(result.rows.map(row => [String(row.olist_product_id), row.unit || 'UN']));
|
|
};
|
|
|
|
const markStructureChecked = async (productId, isManufactured) => {
|
|
await pool.query(`
|
|
UPDATE olist_product_cache
|
|
SET structure_checked_at = CURRENT_TIMESTAMP,
|
|
is_manufactured = $2,
|
|
updated_at = CURRENT_TIMESTAMP
|
|
WHERE olist_product_id = $1;
|
|
`, [String(productId), isManufactured]);
|
|
};
|
|
|
|
const createRun = async ({ trigger, mode }) => {
|
|
const result = await pool.query(`
|
|
INSERT INTO olist_sync_runs (trigger, mode, status)
|
|
VALUES ($1, $2, 'running')
|
|
RETURNING *;
|
|
`, [trigger, mode]);
|
|
return mapRun(result.rows[0]);
|
|
};
|
|
|
|
const updateRunProgress = async (id, summary) => {
|
|
const result = await pool.query(`
|
|
UPDATE olist_sync_runs
|
|
SET products_scanned = $2,
|
|
manufactured_products = $3,
|
|
imported = $4,
|
|
component_count = $5,
|
|
reference_count = $6,
|
|
skipped_reference_count = $7,
|
|
failed = $8,
|
|
details = $9::jsonb
|
|
WHERE id = $1
|
|
RETURNING *;
|
|
`, [id, summary.productsScanned, summary.manufacturedProducts, summary.imported, summary.componentCount, summary.referenceCount, summary.skippedReferenceCount, summary.failed, JSON.stringify(summary.details.slice(0, 100))]);
|
|
return mapRun(result.rows[0]);
|
|
};
|
|
|
|
const finishRun = async (id, summary, { status = 'completed', error = null } = {}) => {
|
|
const result = await pool.query(`
|
|
UPDATE olist_sync_runs
|
|
SET status = $2,
|
|
completed_at = CURRENT_TIMESTAMP,
|
|
products_scanned = $3,
|
|
manufactured_products = $4,
|
|
imported = $5,
|
|
component_count = $6,
|
|
reference_count = $7,
|
|
skipped_reference_count = $8,
|
|
failed = $9,
|
|
error_message = $10,
|
|
details = $11::jsonb
|
|
WHERE id = $1
|
|
RETURNING *;
|
|
`, [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]);
|
|
};
|
|
|
|
const runOlistSync = async ({ trigger = 'manual', fullSync = false } = {}) => {
|
|
assertConfigured();
|
|
if (activeSync) return activeSync;
|
|
|
|
activeSync = (async () => {
|
|
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) {
|
|
const error = new Error('Outra sincronização Olist já está em andamento.');
|
|
error.statusCode = 409;
|
|
throw error;
|
|
}
|
|
|
|
const resumePendingFullScan = mode === 'full' && cacheProgress.total > 0 && cacheProgress.pending > 0;
|
|
if (!resumePendingFullScan) {
|
|
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);
|
|
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
|
|
}
|
|
});
|
|
}
|
|
} 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 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) {
|
|
await lockClient.query('SELECT pg_advisory_unlock($1);', [RUN_LOCK_KEY]).catch(() => {});
|
|
lockClient.release();
|
|
}
|
|
}
|
|
})().finally(() => {
|
|
activeSync = null;
|
|
activeRunId = null;
|
|
stopRequested = false;
|
|
});
|
|
|
|
return activeSync;
|
|
};
|
|
|
|
const startOlistSync = async (options) => {
|
|
assertConfigured();
|
|
if (!await getConnection()) {
|
|
const error = new Error('Conecte a conta Olist antes de sincronizar.');
|
|
error.statusCode = 409;
|
|
throw error;
|
|
}
|
|
if (activeSync) return { accepted: false, reason: 'running' };
|
|
runOlistSync(options).catch(error => console.error('Olist sync failed:', error.message));
|
|
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 = async () => {
|
|
if (scheduledSync) clearTimeout(scheduledSync);
|
|
|
|
let scheduledRunAt = getTodayOlistRunAt();
|
|
let delay = Math.max(0, scheduledRunAt.getTime() - Date.now());
|
|
try {
|
|
if (delay === 0 && await hasScheduledSyncToday()) {
|
|
scheduledRunAt.setUTCDate(scheduledRunAt.getUTCDate() + 1);
|
|
delay = scheduledRunAt.getTime() - Date.now();
|
|
}
|
|
} catch (error) {
|
|
console.error('Could not calculate the next daily Olist sync:', error.message);
|
|
}
|
|
|
|
nextScheduledSyncAt = new Date(Date.now() + delay).toISOString();
|
|
scheduledSync = setTimeout(async () => {
|
|
scheduledSync = null;
|
|
nextScheduledSyncAt = null;
|
|
try {
|
|
assertConfigured();
|
|
if (!await getConnection()) {
|
|
console.error('Scheduled Olist sync skipped: no connected Olist account.');
|
|
return;
|
|
}
|
|
await runOlistSync({ trigger: 'scheduled' });
|
|
} catch (error) {
|
|
console.error('Scheduled Olist sync failed:', error.message);
|
|
} finally {
|
|
await scheduleOlistSync();
|
|
}
|
|
}, delay);
|
|
return nextScheduledSyncAt;
|
|
};
|
|
|
|
const getOlistFrontendRedirect = (outcome) => `${OLIST_FRONTEND_URL.replace(/\/$/, '')}/#/admin/olist?olist=${outcome}`;
|
|
|
|
module.exports = {
|
|
completeAuthorization,
|
|
createAuthorizationUrl,
|
|
getOlistFrontendRedirect,
|
|
getOlistRunDetails,
|
|
getOlistStatus,
|
|
isOlistConfigured,
|
|
markInterruptedOlistRuns,
|
|
scheduleOlistSync,
|
|
requestOlistSyncStop,
|
|
startOlistSync
|
|
};
|