diff --git a/README.md b/README.md index 6d07f46..e173c1b 100644 --- a/README.md +++ b/README.md @@ -67,7 +67,7 @@ OLIST_SYNC_ENABLED=true `OLIST_TOKEN_ENCRYPTION_KEY` must be stable between deploys because Graphs uses it to encrypt OAuth tokens stored in PostgreSQL. The initial sync is full; later scheduled runs are incremental and use the product update date. Manual JSON import remains available as a fallback. -Graphs waits at least 2.5 seconds between every Olist API request (24 requests/minute). This keeps capacity available for other Tiny/Olist integrations. `OLIST_SYNC_REQUEST_DELAY_MS` can increase that delay but cannot lower it below 2500. +Graphs waits at least 2.5 seconds between every Olist API request (24 requests/minute). This keeps capacity available for other Tiny/Olist integrations. The first catalog scan caches product IDs, SKUs, descriptions, and units; each checked BOM is persisted immediately, so a restarted full sync resumes from the remaining products. Later runs use Olist's `dataAlteracao` filter instead of listing the full catalog. `OLIST_SYNC_REQUEST_DELAY_MS` can increase that delay but cannot lower it below 2500. ### Cutting Workbook Import diff --git a/backend/db.js b/backend/db.js index 63212e9..450a53c 100644 --- a/backend/db.js +++ b/backend/db.js @@ -318,6 +318,27 @@ const initDB = async () => { ); `); + await pool.query(` + CREATE TABLE IF NOT EXISTS olist_product_cache ( + olist_product_id BIGINT PRIMARY KEY, + sku VARCHAR(100), + description TEXT NOT NULL, + unit VARCHAR(30) NOT NULL DEFAULT 'UN', + last_changed_at TIMESTAMPTZ, + structure_checked_at TIMESTAMPTZ, + is_manufactured BOOLEAN, + updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP + ); + `); + + await pool.query(` + CREATE TABLE IF NOT EXISTS olist_sync_state ( + id SMALLINT PRIMARY KEY CHECK (id = 1), + last_catalog_sync_at TIMESTAMPTZ, + updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP + ); + `); + await pool.query(` CREATE TABLE IF NOT EXISTS consumption_references ( id SERIAL PRIMARY KEY, @@ -596,6 +617,7 @@ const initDB = async () => { await pool.query(`CREATE INDEX IF NOT EXISTS idx_catalog_products_type ON catalog_products (type);`); await pool.query(`CREATE INDEX IF NOT EXISTS idx_catalog_products_category_id ON catalog_products (category_id);`); await pool.query(`CREATE INDEX IF NOT EXISTS idx_olist_sync_runs_started_at ON olist_sync_runs (started_at DESC);`); + await pool.query(`CREATE INDEX IF NOT EXISTS idx_olist_product_cache_pending_structure ON olist_product_cache (structure_checked_at) WHERE structure_checked_at IS NULL;`); await pool.query(`CREATE INDEX IF NOT EXISTS idx_consumption_references_product_id ON consumption_references (product_id);`); await pool.query(`CREATE INDEX IF NOT EXISTS idx_consumption_references_material_product_id ON consumption_references (material_product_id);`); await pool.query(` diff --git a/backend/services/olistService.js b/backend/services/olistService.js index dd426db..4477cb7 100644 --- a/backend/services/olistService.js +++ b/backend/services/olistService.js @@ -107,13 +107,37 @@ const listSyncRuns = async (limit = 5) => { 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] = await Promise.all([getConnection(), listSyncRuns()]); + 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 }; }; @@ -250,11 +274,19 @@ const olistRequest = async (path, retry = true) => { return data; }; -const listOlistProducts = async () => { +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 page = await olistRequest(`/produtos?${new URLSearchParams({ limit: '100', offset: String(offset) })}`); + 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; @@ -262,6 +294,75 @@ const listOlistProducts = async () => { 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) @@ -296,8 +397,8 @@ const runOlistSync = async ({ trigger = 'manual', fullSync = false } = {}) => { if (activeSync) return activeSync; activeSync = (async () => { - const lastSuccessful = (await listSyncRuns(20)).find(run => run.status === 'completed'); - const mode = fullSync || !lastSuccessful ? 'full' : 'incremental'; + 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; @@ -310,18 +411,22 @@ const runOlistSync = async ({ trigger = 'manual', fullSync = false } = {}) => { throw error; } - const allProducts = await listOlistProducts(); - const unitById = new Map(allProducts.map(product => [String(product.id), product.unidade || 'UN'])); - const changedAfter = mode === 'incremental' ? new Date(lastSuccessful.completedAt || lastSuccessful.startedAt) : null; - const candidates = changedAfter - ? allProducts.filter(product => !product.dataAlteracao || new Date(product.dataAlteracao).getTime() >= changedAfter.getTime()) - : allProducts; + 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) continue; + if (!components.length) { + await markStructureChecked(product.id, false); + continue; + } const result = await upsertTinyProductComposition({ tinyId: String(product.id), productSku: product.sku, @@ -340,8 +445,11 @@ const runOlistSync = async ({ trigger = 'manual', fullSync = false } = {}) => { 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) { + 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 }); } diff --git a/src/pages/Registrations.tsx b/src/pages/Registrations.tsx index 06d2f07..093da7c 100644 --- a/src/pages/Registrations.tsx +++ b/src/pages/Registrations.tsx @@ -593,6 +593,9 @@ const Registrations = () => { {olistStatus.runs[0].mode === 'full' ? 'Carga completa' : 'Atualização incremental'} {olistStatus.runs[0].imported} composições {olistStatus.runs[0].referenceCount} referências + {olistStatus.productCache.total > 0 && ( + {olistStatus.productCache.checked.toLocaleString('pt-BR')} / {olistStatus.productCache.total.toLocaleString('pt-BR')} estruturas verificadas + )} {olistStatus.runs[0].failed > 0 && {olistStatus.runs[0].failed} falhas} )} diff --git a/src/types.ts b/src/types.ts index 958b90c..7ed1547 100644 --- a/src/types.ts +++ b/src/types.ts @@ -227,6 +227,12 @@ export interface OlistStatus { connected: boolean; tokenExpiresAt: string | null; syncInProgress: boolean; + productCache: { + total: number; + checked: number; + pending: number; + manufactured: number; + }; runs: OlistSyncRun[]; }