import { storage } from '../storage'; export interface SendWebhookOptions { retries?: number; taskId?: number; conversationId?: number; organizationId?: number; botId?: number; eventType?: string; } /** * Simple counting semaphore to cap concurrent outbound webhook requests. */ class Semaphore { private slots: number; private queue: Array<() => void> = []; constructor(slots: number) { this.slots = slots; } acquire(): Promise { if (this.slots > 0) { this.slots--; return Promise.resolve(); } return new Promise(resolve => this.queue.push(resolve)); } release(): void { if (this.queue.length > 0) { const next = this.queue.shift()!; next(); } else { this.slots++; } } } // Max 5 concurrent outbound webhook requests across the whole process const webhookSemaphore = new Semaphore(5); /** * Send a webhook with automatic retry (exponential backoff) and concurrency limiting. * * @param url - Target webhook URL * @param payload - JSON-serialisable request body * @param headers - Additional HTTP headers (Content-Type is always set) * @param options - Retry count and optional audit-log context */ export async function sendWebhook( url: string, payload: unknown, headers: Record = {}, options: SendWebhookOptions = {} ): Promise { const { retries = 3, taskId, organizationId, botId, eventType } = options; const body = JSON.stringify(payload); await webhookSemaphore.acquire(); try { let lastError: unknown; for (let attempt = 0; attempt < retries; attempt++) { if (attempt > 0) { // Exponential backoff: 2 s, 4 s, 8 s, … const delayMs = Math.pow(2, attempt) * 1000; await new Promise(resolve => setTimeout(resolve, delayMs)); } try { const response = await fetch(url, { method: 'POST', headers: { 'Content-Type': 'application/json', ...headers }, body, }); if (response.ok) { return; // success } const text = await response.text().catch(() => ''); lastError = new Error(`HTTP ${response.status}: ${text}`); console.error( `[webhook] Attempt ${attempt + 1}/${retries} failed (HTTP ${response.status}) — ${url}` ); } catch (err) { lastError = err; console.error( `[webhook] Attempt ${attempt + 1}/${retries} failed (network) — ${url}:`, err ); } } // All retries exhausted const errorMessage = lastError instanceof Error ? lastError.message : String(lastError); console.error( `[webhook] Delivery permanently failed after ${retries} attempts — ${url}: ${errorMessage}` ); // Record failure in task audit log when task context is available if (taskId != null && organizationId != null) { storage .addTaskAuditLog({ taskId, organizationId, action: 'webhook.failed', fieldName: 'webhook', newValue: url, changedBy: null, changedByName: 'system', metadata: { url, botId: botId ?? null, eventType: eventType ?? null, error: errorMessage, attempts: retries, }, }) .catch(err => console.error('[webhook] Failed to write audit log entry:', err) ); } } finally { webhookSemaphore.release(); } }