When a WhatsApp contact starts a new conversation by sending multiple images at once (an album), each image arrives as a separate webhook. Because no conversation exists yet, the concurrent workers each pass the "does a conversation exist?" check and each create their own conversation — producing one conversation per image instead of one grouped conversation. This fix serializes webhook processing per `(inbox, contact)` using a Redis lock at the job level, so only one webhook at a time can create the initial conversation for a given contact. Concurrent workers retry with backoff and append to the same conversation once the lock is released. ## Closes - Closes #13261 ## How to test 1. On a WhatsApp inbox, ensure there is no active (open) conversation with a specific test contact — resolve or delete any existing one. 2. From a phone, select 6+ images in the WhatsApp gallery and send them as a single album to the Chatwoot-connected number. 3. Open the Chatwoot dashboard and confirm exactly **one** new conversation is created, with all images grouped under it. 4. Repeat the test with a mix of attachment types (XMLs, PDFs, images) sent in rapid succession — still one conversation. ## What changed - New Redis key `WHATSAPP_MESSAGE_CREATE_LOCK::<inbox_id>::<sender_id>` in `lib/redis/redis_keys.rb`. - `Webhooks::WhatsappEventsJob` now inherits from `MutexApplicationJob` and wraps event processing in `with_lock(key)`, matching the pattern already used by `FacebookEventsJob`, `InstagramEventsJob`, and `TiktokEventsJob`. - Uses `retry_on LockAcquisitionError, wait: 1.second, attempts: 8` so concurrent webhooks retry until the lock is free instead of poll-waiting inside the service. - Sender ID is derived from the webhook payload (contact's `from`, or `to` for SMB echo events); status-only webhooks bypass the lock. - Issue 1 from the report (same `source_id` redelivery) was already handled previously by `Whatsapp::MessageDedupLock` (atomic `SET NX EX`); no changes needed there. --------- Co-authored-by: Muhsin <12408980+muhsin-k@users.noreply.github.com> Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
61 lines
3.2 KiB
Ruby
61 lines
3.2 KiB
Ruby
module Redis::RedisKeys
|
|
## Inbox Keys
|
|
# Array storing the ordered ids for agent round robin assignment
|
|
ROUND_ROBIN_AGENTS = 'ROUND_ROBIN_AGENTS:%<inbox_id>d'.freeze
|
|
|
|
## Conversation keys
|
|
# Detect whether to send an email reply to the conversation
|
|
CONVERSATION_MAILER_KEY = 'CONVERSATION::%<conversation_id>d'.freeze
|
|
# Whether a conversation is muted ?
|
|
CONVERSATION_MUTE_KEY = 'CONVERSATION::%<id>d::MUTED'.freeze
|
|
CONVERSATION_DRAFT_MESSAGE = 'CONVERSATION::%<id>d::DRAFT_MESSAGE'.freeze
|
|
|
|
## User Keys
|
|
# SSO Auth Tokens
|
|
USER_SSO_AUTH_TOKEN = 'USER_SSO_AUTH_TOKEN::%<user_id>d::%<token>s'.freeze
|
|
|
|
## Online Status Keys
|
|
# hash containing user_id key and status as value
|
|
ONLINE_STATUS = 'ONLINE_STATUS::%<account_id>d'.freeze
|
|
# sorted set storing online presense of account contacts
|
|
ONLINE_PRESENCE_CONTACTS = 'ONLINE_PRESENCE::%<account_id>d::CONTACTS'.freeze
|
|
# sorted set storing online presense of account users
|
|
ONLINE_PRESENCE_USERS = 'ONLINE_PRESENCE::%<account_id>d::USERS'.freeze
|
|
|
|
## Authorization Status Keys
|
|
# Used to track token expiry and such issues for facebook slack integrations etc
|
|
AUTHORIZATION_ERROR_COUNT = 'AUTHORIZATION_ERROR_COUNT:%<obj_type>s:%<obj_id>d'.freeze
|
|
REAUTHORIZATION_REQUIRED = 'REAUTHORIZATION_REQUIRED:%<obj_type>s:%<obj_id>d'.freeze
|
|
|
|
## Internal Installation related keys
|
|
CHATWOOT_INSTALLATION_ONBOARDING = 'CHATWOOT_INSTALLATION_ONBOARDING'.freeze
|
|
CHATWOOT_INSTALLATION_CONFIG_RESET_WARNING = 'CHATWOOT_CONFIG_RESET_WARNING'.freeze
|
|
LATEST_CHATWOOT_VERSION = 'LATEST_CHATWOOT_VERSION'.freeze
|
|
# Check if a message create with same source-id is in progress?
|
|
MESSAGE_SOURCE_KEY = 'MESSAGE_SOURCE_KEY::%<id>s'.freeze
|
|
OPENAI_CONVERSATION_KEY = 'OPEN_AI_CONVERSATION_KEY::V1::%<event_name>s::%<conversation_id>d::%<updated_at>d'.freeze
|
|
|
|
## Sempahores / Locks
|
|
# We don't want to process messages from the same sender concurrently to prevent creating double conversations
|
|
FACEBOOK_MESSAGE_MUTEX = 'FB_MESSAGE_CREATE_LOCK::%<sender_id>s::%<recipient_id>s'.freeze
|
|
IG_MESSAGE_MUTEX = 'IG_MESSAGE_CREATE_LOCK::%<sender_id>s::%<ig_account_id>s'.freeze
|
|
TIKTOK_MESSAGE_MUTEX = 'TIKTOK_MESSAGE_CREATE_LOCK::%<business_id>s::%<conversation_id>s'.freeze
|
|
TIKTOK_REFRESH_TOKEN_MUTEX = 'TIKTOK_REFRESH_TOKEN_LOCK::%<channel_id>s'.freeze
|
|
SLACK_MESSAGE_MUTEX = 'SLACK_MESSAGE_LOCK::%<conversation_id>s::%<reference_id>s'.freeze
|
|
EMAIL_MESSAGE_MUTEX = 'EMAIL_CHANNEL_LOCK::%<inbox_id>s'.freeze
|
|
WHATSAPP_MESSAGE_MUTEX = 'WHATSAPP_MESSAGE_CREATE_LOCK::%<inbox_id>s::%<sender_id>s'.freeze
|
|
CRM_PROCESS_MUTEX = 'CRM_PROCESS_MUTEX::%<hook_id>s'.freeze
|
|
CAPTAIN_DOCUMENT_SYNC_MUTEX = 'CAPTAIN_DOCUMENT_SYNC_LOCK::%<document_id>s'.freeze
|
|
|
|
## Auto Assignment Keys
|
|
# Track conversation assignments to agents for rate limiting
|
|
ASSIGNMENT_KEY = 'ASSIGNMENT::%<inbox_id>d::AGENT::%<agent_id>d::CONVERSATION::%<conversation_id>d'.freeze
|
|
ASSIGNMENT_KEY_PATTERN = 'ASSIGNMENT::%<inbox_id>d::AGENT::%<agent_id>d::*'.freeze
|
|
|
|
## Account Onboarding
|
|
ACCOUNT_ONBOARDING_ENRICHMENT = 'ONBOARDING_ENRICHMENT::%<account_id>d'.freeze
|
|
|
|
## Account Email Rate Limiting
|
|
ACCOUNT_OUTBOUND_EMAIL_COUNT_KEY = 'OUTBOUND_EMAIL_COUNT::%<account_id>d::%<date>s'.freeze
|
|
end
|