Files
graphs/backend/services/olistService.js
2026-08-10 13:38:03 -03:00

504 lines
20 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_INTERVAL_MS,
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 scheduledSync = 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,
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 (limit = 5) => {
const result = await pool.query(`
SELECT * FROM olist_sync_runs
ORDER BY started_at DESC
LIMIT $1;
`, [Math.max(1, Math.min(Number(limit) || 5, 20))]);
return result.rows.map(mapRun);
};
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 getOlistStatus = async () => {
const [connection, runs, productCache] = await Promise.all([getConnection(), listSyncRuns(), getProductCacheProgress()]);
return {
configured: isOlistConfigured(),
connected: Boolean(connection),
tokenExpiresAt: connection?.token_expires_at || null,
syncInProgress: Boolean(activeSync),
productCache,
runs
};
};
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',
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) => {
const products = [];
let offset = 0;
do {
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;
};
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 finishRun = async (id, summary, 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, 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))]);
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 });
const summary = { productsScanned: 0, manufacturedProducts: 0, imported: 0, componentCount: 0, referenceCount: 0, skippedReferenceCount: 0, failed: 0, details: [] };
let lockClient;
try {
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 changedProducts = await listOlistProducts(mode === 'incremental' ? checkpoint.last_catalog_sync_at : null);
await cacheOlistProducts(changedProducts);
await updateCatalogCheckpoint();
}
const [candidates, unitById] = await Promise.all([listStructureCandidates(), getCachedUnitById()]);
for (const product of candidates) {
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;
}
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);
} else {
summary.failed += 1;
summary.details.push({ productId: product.id, sku: product.sku || '', error: error.message });
}
}
}
return await finishRun(run.id, summary);
} catch (error) {
await finishRun(run.id, summary, error);
throw error;
} finally {
if (lockClient) {
await lockClient.query('SELECT pg_advisory_unlock($1);', [RUN_LOCK_KEY]).catch(() => {});
lockClient.release();
}
}
})().finally(() => { activeSync = null; });
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 scheduleOlistSync = () => {
if (scheduledSync) clearInterval(scheduledSync);
scheduledSync = setInterval(() => {
startOlistSync({ trigger: 'scheduled' }).catch(error => console.error('Scheduled Olist sync skipped:', error.message));
}, OLIST_SYNC_INTERVAL_MS);
return scheduledSync;
};
const getOlistFrontendRedirect = (outcome) => `${OLIST_FRONTEND_URL.replace(/\/$/, '')}/#/registrations?olist=${outcome}`;
module.exports = {
completeAuthorization,
createAuthorizationUrl,
getOlistFrontendRedirect,
getOlistStatus,
isOlistConfigured,
scheduleOlistSync,
startOlistSync
};