webhooks.ts 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147
  1. import { Router } from 'express';
  2. import { chatwootWebhookSchema } from '../../types/chatwoot.js';
  3. import { normalizeChatwootWebhook } from '../../domain/messageNormalizer.js';
  4. import { isTicketMode } from '../../domain/ticketService.js';
  5. import { claimMessage, setMessageStatus } from '../../store/idempotencyStore.js';
  6. import { enqueue } from '../../queue/jobQueue.js';
  7. import { audit } from '../../store/auditLog.js';
  8. import { logger } from '../../logger.js';
  9. import { config } from '../../config.js';
  10. import { drainOnce } from '../../queue/worker.js';
  11. import { evaluateRuntimeSkip } from '../../store/runtimeSettings.js';
  12. export const webhookRouter = Router();
  13. /**
  14. * Chatwoot `message_created` entry point.
  15. *
  16. * Contract: acknowledge fast (202) and never let a slow Flowise call cause a
  17. * Chatwoot-side timeout and retry. Anything that is not an actionable incoming
  18. * customer message is answered 200 with an explicit skip reason, so Chatwoot
  19. * does not keep retrying it.
  20. */
  21. webhookRouter.post('/webhooks/chatwoot', async (req, res, next) => {
  22. try {
  23. const parsed = chatwootWebhookSchema.safeParse(req.body);
  24. if (!parsed.success) {
  25. logger.warn('Malformed Chatwoot webhook payload', {
  26. issues: parsed.error.issues.map((i) => i.path.join('.')),
  27. });
  28. res.status(200).json({ ok: true, skipped: true, reason: 'invalid_payload' });
  29. return;
  30. }
  31. const normalized = normalizeChatwootWebhook(parsed.data);
  32. if (!normalized.ok) {
  33. logger.info('Webhook ignored', { reason: normalized.reason });
  34. res.status(200).json({ ok: true, skipped: true, reason: normalized.reason });
  35. return;
  36. }
  37. const event = normalized.event;
  38. const runtime = await evaluateRuntimeSkip(event.messageCreatedAt ? new Date(event.messageCreatedAt) : null);
  39. if (runtime.decision.skip) {
  40. await setMessageStatus(event.source, event.messageId, 'skipped', runtime.decision.reason).catch(async () => {
  41. await claimMessage(event.source, event.messageId, event.conversationId);
  42. await setMessageStatus(event.source, event.messageId, 'skipped', runtime.decision.skip ? runtime.decision.reason : 'runtime_skip');
  43. });
  44. await audit({
  45. conversationId: event.conversationId,
  46. messageId: event.messageId,
  47. eventType: `skipped_${runtime.decision.reason}`,
  48. summary: `Message not queued: ${runtime.decision.reason}`,
  49. meta: {
  50. stage: 'webhook',
  51. messageCreatedAt: event.messageCreatedAt,
  52. maxMessageAgeMinutes: runtime.settings.maxMessageAgeMinutes,
  53. ignoreMessagesBefore: runtime.settings.ignoreMessagesBefore?.toISOString() ?? null,
  54. },
  55. });
  56. res.status(200).json({ ok: true, skipped: true, reason: runtime.decision.reason, conversationId: event.conversationId });
  57. return;
  58. }
  59. // Idempotency: the unique (source, messageId) insert decides the winner of
  60. // a retry race before any job is created.
  61. const claim = await claimMessage(event.source, event.messageId, event.conversationId);
  62. if (!claim.claimed) {
  63. logger.info('Duplicate message ignored', {
  64. conversationId: event.conversationId,
  65. messageId: event.messageId,
  66. previousStatus: claim.status,
  67. });
  68. res.status(200).json({
  69. ok: true,
  70. duplicate: true,
  71. status: claim.status,
  72. conversationId: event.conversationId,
  73. });
  74. return;
  75. }
  76. // Fast path: when the payload already carries labels and custom attributes,
  77. // a handed-off conversation can be settled here — same data the worker
  78. // would have used, minus a queue round-trip. Conversations that need a
  79. // Chatwoot fetch are still decided in the worker.
  80. if (!event.needsConversationFetch) {
  81. const ticketed = isTicketMode({
  82. labels: event.labels,
  83. custom_attributes: event.customAttributes,
  84. });
  85. if (ticketed) {
  86. await setMessageStatus(event.source, event.messageId, 'skipped', 'ticket_mode');
  87. await audit({
  88. conversationId: event.conversationId,
  89. messageId: event.messageId,
  90. eventType: 'skipped_ticket_mode',
  91. summary: 'Conversation is in manual/ticket mode — Flowise not called (settled at webhook)',
  92. meta: { stage: 'webhook', channel: event.channel, inboxId: event.inboxId },
  93. });
  94. logger.info('Webhook settled without queueing: ticket mode', {
  95. conversationId: event.conversationId,
  96. messageId: event.messageId,
  97. });
  98. res.status(200).json({
  99. ok: true,
  100. skipped: true,
  101. reason: 'ticket_mode',
  102. conversationId: event.conversationId,
  103. });
  104. return;
  105. }
  106. }
  107. const jobId = await enqueue({ type: 'chatwoot_message', payload: { ...event } });
  108. await audit({
  109. conversationId: event.conversationId,
  110. messageId: event.messageId,
  111. jobId,
  112. eventType: 'webhook_accepted',
  113. summary: `Queued job ${jobId} for conversation ${event.conversationId}`,
  114. meta: { channel: event.channel, inboxId: event.inboxId, attachments: event.attachmentCount },
  115. });
  116. res.status(202).json({
  117. ok: true,
  118. accepted: true,
  119. jobId,
  120. conversationId: event.conversationId,
  121. messageId: event.messageId,
  122. });
  123. // With the background worker disabled (one-shot runs) still make progress
  124. // once the response has been flushed. Tests drive the queue explicitly.
  125. const cfg = config();
  126. if (!cfg.WORKER_ENABLED && cfg.NODE_ENV !== 'test') {
  127. void drainOnce().catch((err: unknown) => {
  128. logger.error('Inline drain failed', {
  129. error: err instanceof Error ? err.message : String(err),
  130. });
  131. });
  132. }
  133. } catch (err) {
  134. next(err);
  135. }
  136. });