import { Request, Response } from 'express'; import { sendToN8n } from '../services/n8n.service'; import { sendTinyOrderToGraphs } from '../services/graphs.service'; import axios from 'axios'; import fs from 'fs'; import path from 'path'; type CachedOrderDetails = { expiresAt: number; fullOrderDetails: any; statusProcessamento: string; whatsappVendedor: string; }; type OrderDetailsResult = Omit; const orderDetailsCache = new Map(); const inFlightOrderDetails = new Map>(); const configuredCacheTtlMs = Number(process.env.TINY_ORDER_DETAILS_CACHE_TTL_MS || 120000); const ORDER_DETAILS_CACHE_TTL_MS = Number.isFinite(configuredCacheTtlMs) ? configuredCacheTtlMs : 120000; const numberEnv = (name: string, fallback: number) => { const raw = process.env[name]; if (raw === undefined || raw.trim() === '') return fallback; const value = Number(raw); return Number.isFinite(value) && value >= 0 ? value : fallback; }; const LIVE_TINY_REQUEST_DELAY_MS = numberEnv('TINY_LIVE_REQUEST_DELAY_MS', 6000); const LIVE_TINY_MAX_QUEUE_WAIT_MS = numberEnv('TINY_LIVE_MAX_QUEUE_WAIT_MS', 30000); let liveTinyQueueTail: Promise = Promise.resolve(); let liveTinyNextRequestAt = 0; let liveTinyPendingRequests = 0; class LiveTinyQueueBacklogError extends Error { constructor(orderId: string, estimatedWaitMs: number) { super(`Tiny live queue is too busy for Order ID ${orderId}. Estimated wait ${estimatedWaitMs}ms.`); } } const sleep = (ms: number) => new Promise(resolve => setTimeout(resolve, ms)); function pruneExpiredOrderDetailsCache() { const now = Date.now(); for (const [orderId, cached] of orderDetailsCache.entries()) { if (now > cached.expiresAt) { orderDetailsCache.delete(orderId); } } } function getCachedOrderDetails(orderId: string) { const cached = orderDetailsCache.get(orderId); if (!cached) return null; if (Date.now() > cached.expiresAt) { orderDetailsCache.delete(orderId); return null; } return cached; } function setCachedOrderDetails(orderId: string, cached: Omit) { if (ORDER_DETAILS_CACHE_TTL_MS <= 0) return; pruneExpiredOrderDetailsCache(); orderDetailsCache.set(orderId, { ...cached, expiresAt: Date.now() + ORDER_DETAILS_CACHE_TTL_MS }); } async function runWithLiveTinySlot(orderId: string, servicePhp: string, operation: () => Promise): Promise { if (LIVE_TINY_REQUEST_DELAY_MS <= 0) { return operation(); } const estimatedWaitMs = Math.max(0, liveTinyNextRequestAt - Date.now()) + (liveTinyPendingRequests * LIVE_TINY_REQUEST_DELAY_MS); if (estimatedWaitMs > LIVE_TINY_MAX_QUEUE_WAIT_MS) { throw new LiveTinyQueueBacklogError(orderId, estimatedWaitMs); } liveTinyPendingRequests += 1; const previousTail = liveTinyQueueTail.catch(() => undefined); const runPromise = previousTail.then(async () => { const waitMs = Math.max(0, liveTinyNextRequestAt - Date.now()); if (waitMs > 0) { console.log(`[Tiny Live Queue] Waiting ${waitMs}ms before ${servicePhp} for Order ID: ${orderId}.`); await sleep(waitMs); } liveTinyNextRequestAt = Date.now() + LIVE_TINY_REQUEST_DELAY_MS; console.log(`[Tiny Live Queue] Fetching ${servicePhp} for Order ID: ${orderId}. Pending: ${liveTinyPendingRequests}.`); return operation(); }).finally(() => { liveTinyPendingRequests = Math.max(0, liveTinyPendingRequests - 1); }); liveTinyQueueTail = runPromise.then(() => undefined, () => undefined); return runPromise; } async function tinyLivePost(orderId: string, servicePhp: string, params: URLSearchParams) { return runWithLiveTinySlot(orderId, servicePhp, () => axios.post(`https://api.tiny.com.br/api2/${servicePhp}`, params, { headers: { 'Content-Type': 'application/x-www-form-urlencoded' } })); } async function fetchOrderDetailsFromTiny(orderId: string, tinyApiToken: string): Promise { console.log(`Fetching full details for Order ID: ${orderId} from Tiny API...`); const params = new URLSearchParams(); params.append('token', tinyApiToken); params.append('id', orderId); params.append('formato', 'JSON'); const apiResponse = await tinyLivePost(orderId, 'pedido.obter.php', params); if (apiResponse.data?.retorno?.status !== 'OK') { console.error('Tiny API returned an error:', apiResponse.data?.retorno?.erros || 'Unknown error'); return null; } const fullOrderDetails = apiResponse.data.retorno.pedido; const statusProcessamento = apiResponse.data.retorno.status_processamento || ""; let whatsappVendedor = ""; console.log(`Successfully fetched order details! Found phone: ${fullOrderDetails.cliente?.celular || fullOrderDetails.cliente?.fone || 'None'}`); // Seller details are off by default because this is an extra Tiny API call per order. const idVendedor = fullOrderDetails.id_vendedor; if (process.env.TINY_FETCH_SELLER_DETAILS === 'true' && idVendedor && idVendedor !== "0") { console.log(`Fetching seller details for Seller ID: ${idVendedor}...`); const vendorParams = new URLSearchParams(); vendorParams.append('token', tinyApiToken); vendorParams.append('id', idVendedor); vendorParams.append('formato', 'JSON'); const vendorResponse = await tinyLivePost(orderId, 'contato.obter.php', vendorParams); if (vendorResponse.data?.retorno?.status === 'OK') { const contato = vendorResponse.data.retorno.contato; whatsappVendedor = contato?.celular || contato?.fone || contato?.telefone || ""; console.log(`Successfully fetched seller WhatsApp: ${whatsappVendedor || 'None'}`); } else { console.error('Failed to fetch seller details:', JSON.stringify(vendorResponse.data?.retorno?.erros || 'Unknown API error')); } } return { fullOrderDetails, statusProcessamento, whatsappVendedor }; } async function getOrFetchOrderDetails(orderId: string, tinyApiToken: string): Promise { const cachedDetails = getCachedOrderDetails(orderId); if (cachedDetails) { console.log(`Using cached order details for Order ID: ${orderId}.`); return { fullOrderDetails: cachedDetails.fullOrderDetails, statusProcessamento: cachedDetails.statusProcessamento, whatsappVendedor: cachedDetails.whatsappVendedor }; } const existingFetch = inFlightOrderDetails.get(orderId); if (existingFetch) { console.log(`Waiting for in-flight Tiny fetch for Order ID: ${orderId}.`); return existingFetch; } const fetchPromise = fetchOrderDetailsFromTiny(orderId, tinyApiToken) .then(result => { if (result) { setCachedOrderDetails(orderId, result); } return result; }) .finally(() => { inFlightOrderDetails.delete(orderId); }); inFlightOrderDetails.set(orderId, fetchPromise); return fetchPromise; } export const handleTinyOrderUpdate = async (req: Request, res: Response): Promise => { try { // 1. Security Check: Verify token from Tiny const expectedToken = process.env.TINY_WEBHOOK_SECRET; const providedToken = req.query.token; if (expectedToken && providedToken !== expectedToken) { console.warn('Unauthorized webhook attempt. Invalid or missing token.'); res.status(401).json({ error: 'Unauthorized' }); return; } let payload = req.body || {}; if (payload && typeof payload.dados === 'string') { try { payload.dados = JSON.parse(payload.dados); } catch (e) {} } console.log('Received webhook from Tiny. Acknowledging immediately...'); // Acknowledge Tiny immediately to prevent timeouts res.status(200).send('OK'); const orderId = payload?.dados?.id; if (!orderId) { console.warn('No order ID found in webhook payload. Cannot fetch details.'); return; } const tinyApiToken = process.env.TINY_API_TOKEN; let fullOrderDetails: any = null; let statusProcessamento = ""; let whatsappVendedor = ""; if (tinyApiToken) { try { const details = await getOrFetchOrderDetails(String(orderId), tinyApiToken); if (details) { fullOrderDetails = details.fullOrderDetails; statusProcessamento = details.statusProcessamento; whatsappVendedor = details.whatsappVendedor; } } catch (apiError: any) { if (apiError instanceof LiveTinyQueueBacklogError) { console.warn(`${apiError.message} Forwarding without Tiny enrichment.`); } else { console.error('Failed to fetch from Tiny API:', apiError.message); } } } else { console.warn('TINY_API_TOKEN is not set in environment variables. Skipping full details fetch.'); } // Build the exact flat JSON payload requested by the user const finalPayload = { id: fullOrderDetails?.id || payload.dados?.id || "", numero: fullOrderDetails?.numero || payload.dados?.numero || "", numero_ecommerce: fullOrderDetails?.numero_ecommerce || fullOrderDetails?.ecommerce?.numeroPedidoEcommerce || payload.dados?.idPedidoEcommerce || "", data_pedido: fullOrderDetails?.data_pedido || payload.dados?.data || "", data_prevista: fullOrderDetails?.data_prevista || "", nome: fullOrderDetails?.cliente?.nome || payload.dados?.cliente?.nome || "", valor: parseFloat(fullOrderDetails?.total_pedido || fullOrderDetails?.valor_total || "0"), id_vendedor: fullOrderDetails?.id_vendedor || "", nome_vendedor: fullOrderDetails?.nome_vendedor || "", whatsapp_vendedor: whatsappVendedor, situacao: fullOrderDetails?.situacao || payload.dados?.descricaoSituacao || "", fone: fullOrderDetails?.cliente?.celular || fullOrderDetails?.cliente?.telefone || fullOrderDetails?.cliente?.fone || "", email: fullOrderDetails?.cliente?.email || "", status_processamento: statusProcessamento, forma_envio: fullOrderDetails?.forma_envio || "", codigo_rastreamento: fullOrderDetails?.codigo_rastreamento || "", url_rastreamento: fullOrderDetails?.url_rastreamento || "", itens: fullOrderDetails?.itens || [] // Added so n8n can filter specific products! }; console.log('Forwarding formatted payload to n8n...'); const statusUrl = process.env.N8N_WEBHOOK_STATUS || process.env.N8N_WEBHOOK_URL; const graphsUrl = process.env.N8N_WEBHOOK_GRAPHS; console.log(`[Config Check] Status URL: ${statusUrl ? 'CONFIGURED' : 'MISSING'}`); console.log(`[Config Check] Graphs URL: ${graphsUrl ? 'CONFIGURED' : 'MISSING'} (${graphsUrl})`); const dispatchPromises: Promise[] = []; if (statusUrl) { console.log('-> Preparing dispatch to Status Workflow...'); dispatchPromises.push(sendToN8n(finalPayload, statusUrl)); } if (graphsUrl) { console.log('-> Preparing dispatch to Graphs Workflow...'); dispatchPromises.push(sendToN8n(finalPayload, graphsUrl)); } if (fullOrderDetails && (process.env.GRAPHS_API_URL || process.env.NEXSTAR_GRAPHS_API_URL)) { console.log('-> Preparing direct dispatch to Graphs API...'); dispatchPromises.push(sendTinyOrderToGraphs(fullOrderDetails)); } // Fire them all in parallel so one does not block the other await Promise.allSettled(dispatchPromises); console.log('All n8n dispatch attempts completed.'); } catch (error) { console.error('Error handling Tiny webhook:', error); } }; export const handleTinyStockUpdate = async (req: Request, res: Response): Promise => { try { const expectedToken = process.env.TINY_WEBHOOK_SECRET; const providedToken = req.query.token; if (expectedToken && providedToken !== expectedToken) { console.warn('Unauthorized webhook attempt on stock. Invalid or missing token.'); res.status(401).json({ error: 'Unauthorized' }); return; } res.status(200).send('OK'); const targetUrl = process.env.N8N_WEBHOOK_STOCK; if (!targetUrl) { console.warn('N8N_WEBHOOK_STOCK is not defined in environment variables. Skipping forward.'); return; } console.log('Received stock webhook from Tiny. Calculating delta and stats...'); let payload = req.body || {}; if (payload && typeof payload.dados === 'string') { try { payload.dados = JSON.parse(payload.dados); } catch (e) {} } // --- SMART MEMORY & STATS LOGIC --- const memoryFile = path.join(process.cwd(), 'stock_memory.json'); let memory: Record = {}; // Load memory if it exists if (fs.existsSync(memoryFile)) { try { memory = JSON.parse(fs.readFileSync(memoryFile, 'utf-8')); } catch (e) { console.error("Failed to parse stock_memory.json, starting fresh."); } } const dados = payload.dados || {}; const idProduto = String(dados.idProduto || dados.id || ""); const nomeProduto = String(dados.nome || dados.descricao || "Unknown"); const novoSaldo = Number(dados.saldo || 0); let deltaEstoque = 0; if (idProduto) { // Backwards compatibility for old memory format (if it was just a number) if (typeof memory[idProduto] === 'number') { memory[idProduto] = { saldo: memory[idProduto], addCount: 0, addTotal: 0, addMax: 0, addMin: null }; } const prodMem = memory[idProduto] || { saldo: undefined, addCount: 0, addTotal: 0, addMax: 0, addMin: null }; const estoqueAntigo = prodMem.saldo; if (estoqueAntigo !== undefined) { // Calculate delta deltaEstoque = novoSaldo - estoqueAntigo; // Track statistics if stock was ADDED if (deltaEstoque > 0) { prodMem.addCount += 1; prodMem.addTotal += deltaEstoque; if (deltaEstoque > prodMem.addMax) { prodMem.addMax = deltaEstoque; } if (prodMem.addMin === null || deltaEstoque < prodMem.addMin) { prodMem.addMin = deltaEstoque; } } } else { deltaEstoque = 0; } // Update current balance in memory prodMem.saldo = novoSaldo; memory[idProduto] = prodMem; fs.writeFileSync(memoryFile, JSON.stringify(memory, null, 2)); console.log(`[Stock Tracker] Product ${idProduto}: Old=${estoqueAntigo ?? 'None'} -> New=${novoSaldo}. Delta=${deltaEstoque > 0 ? '+' : ''}${deltaEstoque}`); // --- CSV LOGGING FOR ANALYSIS --- if (deltaEstoque !== 0 || estoqueAntigo === undefined) { const csvFile = path.join(process.cwd(), 'stock_log.csv'); const timestamp = new Date().toISOString().replace(/T/, ' ').replace(/\..+/, ''); const oldStockStr = estoqueAntigo !== undefined ? estoqueAntigo : 'NEW'; const adicionado = deltaEstoque > 0 ? deltaEstoque : 0; const vendido = deltaEstoque < 0 ? Math.abs(deltaEstoque) : 0; // Calculate current average const mediaAdicionado = prodMem.addCount > 0 ? (prodMem.addTotal / prodMem.addCount).toFixed(2) : 0; const minAdicionado = prodMem.addMin !== null ? prodMem.addMin : 0; const maxAdicionado = prodMem.addMax; const csvLine = `"${timestamp}","${idProduto}","${nomeProduto}","${oldStockStr}","${novoSaldo}","${adicionado}","${vendido}","${mediaAdicionado}","${maxAdicionado}","${minAdicionado}"\n`; // Create file with headers if it doesn't exist if (!fs.existsSync(csvFile)) { fs.writeFileSync(csvFile, '"Data","ID_Produto","Nome_Produto","Estoque_Antigo","Estoque_Novo","Adicionado","Vendido","Media_Adicionado","Max_Adicionado","Min_Adicionado"\n'); } fs.appendFileSync(csvFile, csvLine); } } // Inject the delta and stats into the payload so n8n can easily read it if (!payload.dados) payload.dados = {}; payload.dados.delta_estoque = deltaEstoque; if (idProduto && memory[idProduto]) { const pm = memory[idProduto]; payload.dados.media_adicionado = pm.addCount > 0 ? (pm.addTotal / pm.addCount).toFixed(2) : 0; payload.dados.max_adicionado = pm.addMax; payload.dados.min_adicionado = pm.addMin !== null ? pm.addMin : 0; } await sendToN8n(payload, targetUrl); } catch (error) { console.error('Error handling Tiny stock webhook:', error); } }; export const handleTinyGraphsUpdate = async (req: Request, res: Response): Promise => { try { const expectedToken = process.env.TINY_WEBHOOK_SECRET; const providedToken = req.query.token; if (expectedToken && providedToken !== expectedToken) { console.warn('Unauthorized webhook attempt on graphs. Invalid or missing token.'); res.status(401).json({ error: 'Unauthorized' }); return; } res.status(200).send('OK'); const targetUrl = process.env.N8N_WEBHOOK_GRAPHS; if (!targetUrl) { console.warn('N8N_WEBHOOK_GRAPHS is not defined in environment variables. Skipping forward.'); return; } console.log('Received graphs webhook from Tiny. Forwarding to n8n...'); let payload = req.body || {}; if (payload && typeof payload.dados === 'string') { try { payload.dados = JSON.parse(payload.dados); } catch (e) {} } await sendToN8n(payload, targetUrl); } catch (error) { console.error('Error handling Tiny graphs webhook:', error); } };