iachat/app/services/whatsapp/providers/whatsapp_baileys_service.rb
Gabriel Jablonski 72c9821270
feat(whatsapp): add emoji reactions UI (#276)
* feat(whatsapp): add emoji reactions UI

Adds end-to-end agent UI for emoji reactions on WhatsApp inboxes
(Cloud API, Baileys, Z-API). Reactions arrive as Messages with
is_reaction=true; this PR exposes them in the bubble UI and lets
agents react with toggle/replace/remove semantics.

- Add POST /reactions endpoint with toggle/replace logic that handles
  multi-device echoes from the same connected number
- Add Channel::Whatsapp#supports_reactions? capability
- Add Message.hide_removed_reactions scope and use it in conversation
  card preview / last_non_activity_message
- Enrich last_non_activity_message with in_reply_to_snippet for
  reaction previews in chat list
- Frontend: hover EmojiReactionPicker (8 quick + full picker) with
  alignment-aware positioning, single ReactionDisplay chip aggregating
  emojis with total count, conversation card preview shows "Você
  reagiu" for own/multi-device echoes

* fix: address CodeRabbit review feedback

- MessagePreview: render "Você" for outgoing reaction echoes that have no
  sender (multi-device echoes from the connected number)
- MessagesView#findCurrentUserReaction: prefer active reactions over
  deleted rows so a stale deleted echo cannot hijack the toggle target
- conversationHelper: drop removed reactions up-front so the activity
  fallback never returns null when older non-removed messages exist
- imap_import rake: wrap IMAP work in begin/ensure so the session is
  closed even when uid_search/scan_new_email_uids raises
- ReactionDisplay: include reaction.id in the user row so v-for keys
  stay stable across re-renders

* fix: address CodeRabbit round 2 feedback

- enterprise Message override of mark_pending_conversation_as_open_for_human_response
  now early-returns on reaction? so reactions can no longer auto-open Captain-pending
  conversations (matches the OSS guard)
- whatsapp incoming reaction-removal handlers (Cloud/Baileys/Z-API) look up the
  reaction Message globally by sender instead of through the inbound conversation
  scope, then operate on existing.conversation; otherwise an old/resolved thread
  could be silently no-op'd while the inbound flow created a stray empty thread
- EmojiReactionPicker: localize quick-emoji tooltip labels via i18n keys
- Message.vue: track pendingTimeouts and clear them on unmount so the cooldown
  setTimeout no longer touches state after teardown
- toggleMessageReaction action returns the API promise so callers can reconcile
  if the cable echo is delayed

* fix: address CodeRabbit round 3 feedback

- MessageFinder#page_window: pluck the 20-row window IDs before taking .min
  so the latest page honors PAGE_LIMIT (ActiveRecord's .minimum(:id) ignores
  .limit and aggregates over the whole relation)
- ReactionsController#current_user_reaction: rank active reactions ahead of
  deleted rows (same invariant as the frontend lookup) so a stale deleted
  echo can no longer hijack the toggle target and resurrect itself
- Whatsapp incoming handlers (Cloud, Baileys individual & group, Z-API) now
  branch on reaction_removal? BEFORE set_conversation / find_or_create_group_
  conversation, so a blank reaction-removal webhook can never open or create
  a stray thread just to no-op
- Message#reaction?: strict-boolean cast (via ActiveModel::Type::Boolean)
  so a stored string "false" no longer leaks through .present? as truthy

* fix: address CodeRabbit round 4 feedback

- MessageList: anchor unread divider on the filtered visibleMessages
  (firstUnreadId can land on a reaction that's filtered out, otherwise
  the separator silently disappears)
- ReactionDisplay: render the removable user row as a real <button> when
  it's the current user's reaction so keyboard users can focus/activate it
- MessagesView#findCurrentUserReaction: read sender_type from m.sender_type
  OR m.sender?.type so REST-loaded messages match the same row instead of
  spawning a duplicate optimistic reaction
- Whatsapp incoming reaction-removal lookup (Cloud, Baileys, Z-API): pick
  the newest active row first and only fall back to the newest deleted row
  when no active reaction exists, mirroring the controller invariant
- CardMessagePreview: use MESSAGE_TYPE.OUTGOING constant in place of the
  literal 1 for the multi-device reaction echo check

* fix: address CodeRabbit round 5 feedback

- ReactionsController#ensure_target_is_reactable: reject activity,
  template, failed, is_unsupported and missing-source_id targets so the
  API mirrors the client toolbar gate and refuses reactions that could
  never land on WhatsApp
- MessageList reaction aggregator: treat "agent reacted via Chatwoot"
  and "agent reacted via the connected phone" as the same self bucket
  so the chip no longer double-counts the current user when both shapes
  coexist for one target
- internalChat ReactionDisplay: render the removable user row as a real
  <button> so keyboard users can focus and trigger removal (mirrors the
  fix already applied to components-next/message/ReactionDisplay)
- EventDataPresenter#push_last_non_activity_message: reorder
  created_at: :desc before .first so the cable snapshot publishes the
  latest preview instead of the oldest row
- Z-API mark_existing_reaction_as_removed: drop the blanket
  `return unless incoming_message?` and route the lookup by direction
  (contact sender for incoming removals, senderless outgoing row for
  multi-device removals from the connected phone). Chatwoot-originated
  echoes stay idempotent because the active-first guard finds nothing
  once the controller has flipped deleted=true locally
- spec: assert reaction removal does not change messages.count on the
  in-place Cloud path

* fix: address CodeRabbit round 6 feedback

- ReactionsController: validate the emoji payload is a single grapheme
  cluster containing a Unicode Emoji codepoint (not just <=32 bytes), so
  arbitrary short strings like "ok" or "123" can no longer be persisted
  as a reaction or enqueued as a WhatsApp reaction send
- target_unreactable_error: add the content_attributes['deleted'] guard
  to mirror the frontend picker gate on deleted messages
- IncomingMessageBaseService: move contact_processable? AFTER the
  reaction_removal? early-return so a blocked contact's removal webhook
  can still reconcile an existing reaction row instead of leaving a
  stale chip/preview
- imap_import rake: add safe_close_imap(imap) that falls back to
  disconnect when logout raises Net::IMAP::Error, mirroring
  terminate_imap_connection in BaseFetchEmailService, and replace the
  three ensure-block imap&.logout sites with it

* fix: address CodeRabbit round 7 feedback

- CardMessagePreview: resolve `lastNonActivityMessage` against the live
  `messages` array by id before rendering, so the chat-card preview
  picks up the freshest copy instead of the (possibly stale) snapshot
  that was mutated in place by a reaction toggle / multi-device echo
- Message + ReactionDisplay: thread an `inboxSupportsReactions` →
  `read-only` prop into the chip so non-supported channels (eg.
  360Dialog) render historical reactions without a clickable
  toggle/remove path that would only hit a 422
- conversations/index.js: replace the truthiness `&&` guard around the
  out-of-order MESSAGE_UPDATED check with `Number.isFinite` parsing so
  a malformed/missing `updated_at` is treated as stale instead of
  silently overwriting a fresher local row
- Baileys mark_existing_reaction_as_removed: drop the blanket
  `return unless incoming?` and split the lookup by direction
  (sender for incoming, sender-less outgoing for multi-device removals)
  to mirror the Z-API/Cloud handlers
- Whatsapp reaction-removal lookup (Cloud, Baileys, Z-API): drop the
  fallback to the newest deleted row so a Chatwoot-originated removal
  echo no-ops cleanly instead of bumping `updated_at` and dispatching
  another `conversation.updated`
- conversation jbuilder: explicit `reorder(created_at: :desc)` on
  `last_non_activity_message` so REST and cable both serialize the
  same most-recent preview

* fix: address CodeRabbit round 8 feedback

- ReactionsController#current_user_reaction: also match on
  content_attributes.in_reply_to_external_id = @target_message.source_id
  (via OR with the existing in_reply_to check), so WhatsApp-echoed
  reactions persisted by the incoming handlers — where in_reply_to could
  be blanked if the target wasn't resolvable at save time — are found and
  toggled instead of stacking a duplicate self-reaction
- Mirror the same defensive OR check in the frontend
  MessagesView#findCurrentUserReaction, and thread the target's
  source_id through the toggleReaction event from Message.vue so the
  lookup sees it

* fix: address CodeRabbit round 9 feedback

- emoji_payload_valid?: tighten the final property check from \p{Emoji}
  to \p{Extended_Pictographic} so plain "1", "#", "*" (which Unicode
  tags as Emoji because they're valid keycap bases) are rejected as
  reaction payloads
- EmojiReactionPicker: mirror the translated `title` into `aria-label`
  on the icon-only smile-plus / plus buttons so screen readers announce
  a meaningful action name
- internalChat ReactionDisplay: close the popover when the post-removal
  state would leave ≤1 reactions, so a singleton-user popover never
  lingers after removing one of a pair
- EventDataPresenter + conversation jbuilder: strip HTML before
  truncating `in_reply_to_snippet` so reactions to email/HTML bubbles
  don't surface literal "<p>..." markup in the chat-list preview

* fix: address CodeRabbit round 10 feedback

- MessageList#reactionsByMessageId: break createdAt ties with `<=` so a
  later iteration wins on second-resolution tie; two toggles in the same
  second no longer leave the chip pointing at a stale row
- MessagePreview: require a non-empty `message.attachments` array (via
  `?.length`) before taking the attachment preview branch, so a removed
  reaction with `[]` no longer renders the attachment placeholder
- MessagesView#findCurrentUserReaction: replace the sort-based pick with
  a reduce that deterministically takes the last element on tie, so a
  fast toggle can't hit a stale/deleted row with the same created_at
- Baileys group handler: guard against `@sender_contact.blank?` before
  dispatching mark_existing_reaction_as_removed, otherwise a nil sender
  would fall into the senderless-outgoing branch and match the wrong row
- WhatsApp reaction-removal lookups (Cloud, Baileys, Z-API): scope the
  base query to `inbox_id: inbox.id` so a colliding WhatsApp message id
  across inboxes can never mutate a reaction row from another inbox

* fix: address CodeRabbit round 11 feedback

- ReactionsController#emoji_payload_valid?: broaden the final property
  check to accept flag and keycap emoji. `\p{Extended_Pictographic}` by
  itself is per-codepoint, so 🇧🇷 (two Regional Indicators) and 1️⃣
  (digit + VS16 + U+20E3) failed validation. Allow any grapheme cluster
  that contains at least one pictographic codepoint, a Regional
  Indicator, or the combining keycap, while still rejecting plain
  ASCII like "ok", "1", "#"
- Message.vue#canShowReactionToolbar: hide the picker when the target
  has no provider source_id, mirroring the server guard in
  ReactionsController#target_unreactable_error instead of letting the
  click fall through to a 422
- MessageList#reactionsByMessageId: fall back to a sourceId → id
  lookup when a reaction only carries `inReplyToExternalId` (WhatsApp
  echo / phone-originated), so its chip still renders against the
  target bubble after reload
- getLastMessage: merge the fresher store fields onto the API
  snapshot instead of replacing it, so jbuilder-only enrichments like
  `in_reply_to_snippet` survive the store refresh

* fix(reactions): preserve API fields on card preview and expose a11y state on quick picker

* fix(reactions): consistent originalId resolution, natural PT-BR snippet phrasing, accurate outgoing-echo spec

* fix(reactions): reject requests missing emoji param and align zapi outgoing-echo spec fixture

* fix(reactions): activity preview fallback, camelCase event listener, EN snippet quoting, fromMe group removals, REST chat-only preview

* fix(reactions): reject non-string emoji, scope page reactions to window, exempt reactions from human_response, add cloud multi-device removal

* test(message): isolate hide_removed_reactions deleted-branch from blank-content branch

* fix(reactions): coerce in_reply_to_snippet to plain String

strip_tags returns an ActiveSupport::SafeBuffer; truncate preserves the
class. When this snippet flowed into ActionCableBroadcastJob via the
CONVERSATION_UPDATED dispatch, Sidekiq's strict-args check rejected the
non-native JSON type, raising synchronously through the dispatcher and
turning the reactions controller response into a 500 even though the row
had already persisted. The UI then surfaced the generic 'failed to update
reaction' toast despite the chip rendering correctly.

Wrap with String.new so the broadcast payload contains plain Strings.

* fix(reactions): don't auto-scroll to bottom on reaction add

ADD_MESSAGE emits SCROLL_TO_MESSAGE for every new push, which makes
sense for regular outgoing messages (the user just hit send and wants
to see it). Reactions render as chips on the parent bubble, so the
auto-scroll yanked the agent away from whichever older message they
were reacting to. Skip the emit when the incoming message is flagged
as a reaction.

* fix(reactions): skip scroll on conversation-only updates triggered by reactions

The reactions controller dispatches CONVERSATION_UPDATED so the chat list
preview can refresh in place. UPDATE_CONVERSATION mutation always emitted
SCROLL_TO_MESSAGE for the open conversation, so every toggle yanked the
viewport back to the bottom even after the previous fix in ADD_MESSAGE.
When the refreshed preview row is itself a reaction the update is
preview-only and the scroll is unwanted; for a regular incoming message
the latest non-activity row is the message itself, which still triggers
the scroll as before.

* fix(reactions): anchor compact picker to button side instead of centering

The compact picker was centered on the smile button, so half its width
always extended toward the bubble side and overflowed past the chat edge
on short messages. Anchor it to the button's outer side and nudge 4px
toward the bubble so it lines up with the trigger.

* test(reactions): regression coverage for safebuffer + scroll skip

The previous CodeRabbit rounds shipped three bugs the existing specs
didn't catch: a SafeBuffer return from `strip_tags` that 500'd the
reactions controller via Sidekiq strict-args, and two SCROLL_TO_MESSAGE
emits (one per mutation) that yanked the open conversation to the
bottom on every emoji toggle. Lock all three behaviors.

Also tighten the spec policy in AGENTS.md so new features default to
having specs instead of skipping them.

* test(baileys): align send_message_body helper with id:updated_at format

The reactions PR switched chatwootMessageId to "<id>:<updated_at>" so
toggle/replace cycles get a fresh idempotency key against baileys-api,
but the shared spec helper still merged the bare integer id. 18 baileys
provider specs were silently broken on CI as a result.

* fix(reactions): skip set_contact for unknown reaction-removal webhooks

Reaction-removal cloud webhooks were unconditionally creating a contact
even when the sender was unknown and there was nothing to remove,
because set_contact ran before the reaction_removal? short-circuit
(needed earlier so blocked-contact reconciliation works). Add a
sender-agnostic existence check on the inbox/in_reply_to scope and bail
out before set_contact when no candidate row exists.

Also realign two specs that were not updated when the chatwootMessageId
format gained an `:updated_at` suffix and when zapi reaction-removal
short-circuited instead of creating a Message.

* test(conversation): include last_non_activity_message in push_data fixture

Reactions PR added last_non_activity_message to the push_data payload
but conversation_spec's exact-match expectation wasn't updated, so the
sharded CI shard that landed on this file flipped red.
2026-04-30 21:09:12 -03:00

885 lines
27 KiB
Ruby

class Whatsapp::Providers::WhatsappBaileysService < Whatsapp::Providers::BaseService # rubocop:disable Metrics/ClassLength
include BaileysHelper
class MessageContentTypeNotSupported < StandardError; end
class ProviderUnavailableError < StandardError; end
class GroupParticipantNotAllowedError < StandardError; end
class MessageAlreadyProcessingError < StandardError; end
DEFAULT_CLIENT_NAME = ENV.fetch('BAILEYS_PROVIDER_DEFAULT_CLIENT_NAME', nil)
DEFAULT_URL = ENV.fetch('BAILEYS_PROVIDER_DEFAULT_URL', nil)
DEFAULT_API_KEY = ENV.fetch('BAILEYS_PROVIDER_DEFAULT_API_KEY', nil)
def self.groups_enabled?
ENV.fetch('BAILEYS_WHATSAPP_GROUPS_ENABLED', 'false') == 'true'
end
def self.status
if DEFAULT_URL.blank? || DEFAULT_API_KEY.blank?
raise ProviderUnavailableError, 'Missing BAILEYS_PROVIDER_DEFAULT_URL or BAILEYS_PROVIDER_DEFAULT_API_KEY setup'
end
response = HTTParty.get(
"#{DEFAULT_URL}/status",
headers: { 'x-api-key' => DEFAULT_API_KEY }
)
unless response.success?
Rails.logger.error response.body
raise ProviderUnavailableError, 'Baileys API is unavailable'
end
response.parsed_response.deep_symbolize_keys
rescue ProviderUnavailableError
raise
rescue StandardError => e
Rails.logger.error e.message
raise ProviderUnavailableError, 'Baileys API is unavailable'
end
def setup_channel_provider
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}",
headers: api_headers,
body: {
clientName: DEFAULT_CLIENT_NAME,
webhookUrl: whatsapp_channel.inbox.callback_webhook_url,
webhookVerifyToken: whatsapp_channel.provider_config['webhook_verify_token'],
# TODO: Remove on Baileys v2, default will be false
includeMedia: false,
groupsEnabled: self.class.groups_enabled?,
autoPresenceSubscribe: whatsapp_channel.provider_config['presence_subscribe'] || false
}.compact.to_json
)
raise ProviderUnavailableError unless process_response(response)
true
end
def disconnect_channel_provider
response = HTTParty.delete(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}",
headers: api_headers
)
raise ProviderUnavailableError unless process_response(response)
true
end
def send_message(recipient_id, message)
@message = message
@recipient_id = recipient_id
if @message.content_attributes[:is_reaction]
@message_content = reaction_message_content
elsif @message.attachments.present?
@message_content = attachment_message_content.merge(reply_context)
elsif @message.outgoing_content.present?
@message_content = { text: @message.outgoing_content }.merge(reply_context)
merge_mention_data
else
@message.update!(is_unsupported: true)
return
end
send_message_request
end
def send_template(phone_number, template_info); end
def sync_templates; end
def allow_group_creation?
true
end
def create_group(subject, participants)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-create",
headers: api_headers,
body: { subject: subject, participants: participants }.to_json
)
raise ProviderUnavailableError unless process_response(response)
response.parsed_response&.deep_symbolize_keys
end
def update_group_subject(group_jid, subject)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-subject",
headers: api_headers,
body: { jid: group_jid, subject: subject }.to_json
)
raise ProviderUnavailableError unless process_response(response)
end
def update_group_description(group_jid, description)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-description",
headers: api_headers,
body: { jid: group_jid, description: description }.to_json
)
raise ProviderUnavailableError unless process_response(response)
end
def update_group_picture(group_jid, image_base64)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/update-profile-picture",
headers: api_headers,
body: { jid: group_jid, image: image_base64 }.to_json
)
raise ProviderUnavailableError unless process_response(response)
end
def update_group_participants(group_jid, participants, action)
Array(participants).each do |participant|
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-participants",
headers: api_headers,
body: { jid: group_jid, participant: participant, action: action }.to_json
)
raise ProviderUnavailableError unless process_response(response)
check_participant_errors(response, action)
end
end
def group_invite_code(group_jid)
response = HTTParty.get(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-invite-code",
headers: api_headers,
query: { jid: group_jid },
format: :json
)
raise ProviderUnavailableError unless process_response(response)
response.parsed_response&.dig('data', 'inviteCode')
end
def revoke_group_invite(group_jid)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-revoke-invite",
headers: api_headers,
body: { jid: group_jid }.to_json
)
raise ProviderUnavailableError unless process_response(response)
response.parsed_response&.dig('data', 'inviteCode')
end
def group_join_requests(group_jid)
response = HTTParty.get(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-request-participants-list",
headers: api_headers,
query: { jid: group_jid },
format: :json
)
return [] if response.code == 403
raise ProviderUnavailableError unless process_response(response)
parsed = response.parsed_response
parsed.is_a?(Array) ? parsed : (parsed&.dig('data') || [])
end
def handle_group_join_requests(group_jid, participants, action)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-request-participants-update",
headers: api_headers,
body: { jid: group_jid, participants: participants, action: action }.to_json
)
raise ProviderUnavailableError unless process_response(response)
end
def group_leave(group_jid)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-leave",
headers: api_headers,
body: { jid: group_jid }.to_json
)
raise ProviderUnavailableError unless process_response(response)
end
PROPERTY_TO_SETTING = {
['announce', true] => 'announcement',
['announce', false] => 'not_announcement',
['restrict', true] => 'locked',
['restrict', false] => 'unlocked'
}.freeze
def group_setting_update(group_jid, property, enabled)
setting = PROPERTY_TO_SETTING[[property, enabled]]
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-setting-update",
headers: api_headers,
body: { jid: group_jid, setting: setting }.to_json
)
raise ProviderUnavailableError unless process_response(response)
end
def group_join_approval_mode(group_jid, mode)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-join-approval-mode",
headers: api_headers,
body: { jid: group_jid, mode: mode }.to_json
)
raise ProviderUnavailableError unless process_response(response)
end
def group_member_add_mode(group_jid, mode)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-member-add-mode",
headers: api_headers,
body: { jid: group_jid, mode: mode }.to_json
)
raise ProviderUnavailableError unless process_response(response)
end
def sync_group(conversation, soft: false)
group_contact = conversation.contact
return true if group_contact.additional_attributes&.dig('group_left')
inbox = conversation.inbox
metadata = group_metadata(group_contact.identifier)
raise ProviderUnavailableError, 'Could not fetch group metadata' if metadata.blank?
update_group_contact_info(group_contact, metadata)
persist_group_settings(group_contact, metadata)
persist_invite_code(group_contact) unless soft
persist_pending_join_requests(group_contact, inbox) unless soft
try_update_group_avatar(group_contact) unless soft
participant_contacts = build_participant_contacts(metadata[:participants], inbox, skip_avatars: soft)
sync_group_members(group_contact, participant_contacts)
persist_sync_status(group_contact)
true
end
def media_url(media_id)
"#{provider_url}/media/#{media_id}"
end
def api_headers
{ 'x-api-key' => api_key, 'Content-Type' => 'application/json' }
end
def validate_provider_config?
response = HTTParty.get(
"#{provider_url}/status/auth",
headers: api_headers
)
process_response(response)
end
def toggle_typing_status(typing_status, recipient_id:, **)
@recipient_id = recipient_id
status_map = {
Events::Types::CONVERSATION_TYPING_ON => 'composing',
Events::Types::CONVERSATION_RECORDING => 'recording',
Events::Types::CONVERSATION_TYPING_OFF => 'paused'
}
response = HTTParty.patch(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/presence",
headers: api_headers,
body: {
toJid: remote_jid,
type: status_map[typing_status]
}.to_json
)
raise ProviderUnavailableError unless process_response(response)
true
end
def presence_subscribe(jids)
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/presence-subscribe",
headers: api_headers,
body: { jids: Array(jids) }.to_json,
timeout: 10
)
raise ProviderUnavailableError unless process_response(response)
response.parsed_response&.dig('data')
end
def update_presence(status)
status_map = {
'online' => 'available',
'offline' => 'unavailable',
'busy' => 'unavailable'
}
response = HTTParty.patch(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/presence",
headers: api_headers,
body: {
type: status_map[status]
}.to_json
)
raise ProviderUnavailableError unless process_response(response)
true
end
def read_messages(messages, recipient_id:, **)
@recipient_id = recipient_id
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/read-messages",
headers: api_headers,
body: {
keys: messages.map { |message| message_key_for(message) }
}.to_json
)
raise ProviderUnavailableError unless process_response(response)
true
end
def unread_message(recipient_id, message)
@recipient_id = recipient_id
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/chat-modify",
headers: api_headers,
body: {
jid: remote_jid,
mod: {
markRead: false,
lastMessages: [{
key: message_key_for(message),
messageTimestamp: message.content_attributes[:external_created_at]
}]
}
}.to_json
)
raise ProviderUnavailableError unless process_response(response)
true
end
def received_messages(recipient_id, messages)
@recipient_id = recipient_id
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/send-receipts",
headers: api_headers,
body: {
keys: messages.map { |message| message_key_for(message) }
}.to_json
)
raise ProviderUnavailableError unless process_response(response)
true
end
def get_profile_pic(jid)
response = HTTParty.get(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/profile-picture-url",
headers: api_headers,
query: { jid: jid },
format: :json
)
return nil unless process_response(response)
response.parsed_response
end
def group_metadata(group_jid)
response = HTTParty.get(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/group-metadata",
headers: api_headers,
query: { jid: group_jid },
format: :json
)
raise ProviderUnavailableError unless process_response(response)
response.parsed_response&.deep_symbolize_keys
end
def on_whatsapp(recipient_id)
@recipient_id = recipient_id
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/on-whatsapp",
headers: api_headers,
body: {
jids: [remote_jid]
}.to_json
)
raise ProviderUnavailableError unless process_response(response)
result = response.parsed_response
result = result.is_a?(Array) ? result : result&.dig('data')
result&.first || { 'jid' => remote_jid, 'exists' => false }
end
def delete_message(recipient_id, message)
@recipient_id = recipient_id
response = HTTParty.delete(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/messages",
headers: api_headers,
body: {
jid: remote_jid,
key: message_key_for(message)
}.to_json
)
raise ProviderUnavailableError unless process_response(response)
true
end
def edit_message(recipient_id, message, new_content)
@recipient_id = recipient_id
response = HTTParty.patch(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/messages",
headers: api_headers,
body: {
jid: remote_jid,
key: message_key_for(message),
messageContent: { text: new_content }
}.to_json
)
raise ProviderUnavailableError unless process_response(response)
true
end
private
def provider_url
whatsapp_channel.provider_config['provider_url'].presence || DEFAULT_URL
end
def api_key
whatsapp_channel.provider_config['api_key'].presence || DEFAULT_API_KEY
end
def reaction_message_content
reply_to = Message.find(@message.in_reply_to)
{
react: {
key: message_key_for(reply_to),
text: @message.outgoing_content
}
}
end
def reply_context
reply_to_external_id = @message.content_attributes[:in_reply_to_external_id]
return {} if reply_to_external_id.blank?
reply_to_message = @message.conversation.messages.find_by(source_id: reply_to_external_id)
return {} unless reply_to_message
{
quotedMessage: {
key: message_key_for(reply_to_message),
message: quoted_message_content(reply_to_message)
}
}
end
def message_key_for(message)
{
id: message.source_id,
remoteJid: remote_jid,
fromMe: message.message_type == 'outgoing',
participant: group_participant_jid(message)
}.compact
end
def group_participant_jid(message)
return unless remote_jid.ends_with?('@g.us')
return if message.message_type == 'outgoing'
message.sender&.identifier
end
def quoted_message_content(message)
if message.attachments.present?
attachment = message.attachments.first
case attachment.file_type
when 'image'
{ imageMessage: { caption: message.content } }
when 'video'
{ videoMessage: { caption: message.content } }
when 'audio'
{ audioMessage: {} }
when 'file'
{ documentMessage: { caption: message.content, fileName: attachment.file.filename.to_s } }
else
{ conversation: message.content.to_s }
end
else
{ conversation: message.content.to_s }
end
end
def attachment_message_content # rubocop:disable Metrics/MethodLength
attachment = @message.attachments.first
buffer = attachment_to_base64(attachment)
content = {
fileName: attachment.file.filename,
caption: @message.outgoing_content
}
case attachment.file_type
when 'image'
content[:image] = buffer
when 'audio'
content[:audio] = buffer
content[:ptt] = attachment.meta&.dig('is_recorded_audio')
when 'file'
content[:document] = buffer
content[:mimetype] = attachment.file.content_type
when 'sticker'
content[:sticker] = buffer
when 'video'
content[:video] = buffer
end
content.compact
end
def send_message_request
response = HTTParty.post(
"#{provider_url}/connections/#{whatsapp_channel.phone_number}/send-message",
headers: api_headers,
body: {
jid: remote_jid,
messageContent: @message_content,
# baileys-api uses this as an idempotency key. Reactions UPDATE a single
# Message row in place across toggle/replace/remove cycles, so reusing
# only `id` would make every follow-up send hit the cached response and
# never reach WhatsApp. Suffixing with updated_at gives each send a fresh
# key while still letting Sidekiq retries of the same attempt dedupe.
chatwootMessageId: "#{@message.id}:#{@message.updated_at.to_f}"
}.to_json,
timeout: 120
)
raise MessageAlreadyProcessingError if response.code == 409
raise ProviderUnavailableError unless process_response(response)
update_external_created_at(response)
response.parsed_response.dig('data', 'key', 'id')
end
def process_response(response)
Rails.logger.error response.body unless response.success?
response.success?
end
def check_participant_errors(response, action)
return unless action.in?(%w[demote remove])
results = response.parsed_response
return unless results.is_a?(Array)
failed = results.find { |r| r['status'].to_s == '406' }
return if failed.blank?
raise GroupParticipantNotAllowedError, 'group_creator_not_modifiable'
end
def merge_mention_data
return if @message.content.blank?
mention_data = Whatsapp::MentionConverterService.extract_mentions_for_whatsapp(@message.content, whatsapp_channel.account)
@message_content.merge!(mention_data) if mention_data.present?
# Replace @DisplayName with @lid/@phone in text so Baileys can match mentions
@message_content[:text] = Whatsapp::MentionConverterService.replace_mentions_in_outgoing_text(
@message.content, @message_content[:text], whatsapp_channel.account
)
end
def remote_jid
return @recipient_id if @recipient_id.ends_with?('@lid')
return @recipient_id if @recipient_id.ends_with?('@g.us')
"#{@recipient_id.delete('+')}@s.whatsapp.net"
end
def update_external_created_at(response)
timestamp = response.parsed_response.dig('data', 'messageTimestamp')
return unless timestamp
external_created_at = baileys_extract_message_timestamp(timestamp)
@message.update!(external_created_at: external_created_at)
end
def build_participant_contacts(participants, inbox, skip_avatars: false)
return [] if participants.blank?
participants.filter_map do |participant|
contact = find_or_create_participant_contact(participant, inbox)
next if contact.blank?
try_update_participant_avatar(contact) unless skip_avatars
{ contact: contact, admin: participant[:admin] }
end
end
def update_group_contact_info(group_contact, metadata)
update_params = {}
update_params[:name] = metadata[:subject] if metadata[:subject].present? && group_contact.name != metadata[:subject]
new_attrs = (group_contact.additional_attributes || {}).merge(
'description' => metadata[:desc].presence,
'owner' => metadata[:owner],
'owner_pn' => metadata[:ownerPn].presence
)
update_params[:additional_attributes] = new_attrs if new_attrs != group_contact.additional_attributes
group_contact.update!(update_params) if update_params.present?
end
def sync_group_members(group_contact, participant_contacts)
return if participant_contacts.blank?
new_contact_ids = participant_contacts.filter_map do |entry|
role = entry[:admin].in?(%w[admin superadmin]) ? :admin : :member
member = GroupMember.find_or_initialize_by(group_contact: group_contact, contact: entry[:contact])
member.assign_attributes(role: role, is_active: true)
member.save! if member.changed?
entry[:contact].id
end
group_contact.group_memberships.active.where.not(contact_id: new_contact_ids).find_each do |member|
member.update!(is_active: false)
end
end
TRACKED_GROUP_SETTINGS = {
announce: 'announce',
restrict: 'restrict',
joinApprovalMode: 'join_approval_mode',
memberAddMode: 'member_add_mode'
}.freeze
def persist_group_settings(group_contact, metadata)
settings = TRACKED_GROUP_SETTINGS.each_with_object({}) do |(api_key, attr_key), hash|
hash[attr_key] = metadata[api_key] if metadata.key?(api_key)
end
return if settings.blank?
new_attrs = (group_contact.additional_attributes || {}).merge(settings)
group_contact.update!(additional_attributes: new_attrs) if new_attrs != group_contact.additional_attributes
end
def persist_sync_status(group_contact)
new_attrs = (group_contact.additional_attributes || {}).merge(
'group_last_synced_at' => Time.current.to_i,
'group_left' => false
)
group_contact.update!(additional_attributes: new_attrs) if new_attrs != group_contact.additional_attributes
end
def persist_invite_code(group_contact)
code = group_invite_code(group_contact.identifier)
return if code.blank?
new_attrs = (group_contact.additional_attributes || {}).merge('invite_code' => code)
group_contact.update!(additional_attributes: new_attrs) if new_attrs != group_contact.additional_attributes
rescue StandardError => e
Rails.logger.error "Failed to fetch invite code for group #{group_contact.identifier}: #{e.message}"
end
def persist_pending_join_requests(group_contact, inbox)
raw_requests = group_join_requests(group_contact.identifier)
requests = raw_requests.filter_map do |req|
contact = find_or_create_participant_contact({ id: req['jid'], phoneNumber: req['phone_number'] }, inbox)
next if contact.blank?
{ 'jid' => req['jid'], 'contact_id' => contact.id, 'request_time' => req['request_time'] }
end
new_attrs = (group_contact.additional_attributes || {}).merge('pending_join_requests' => requests)
group_contact.update!(additional_attributes: new_attrs) if new_attrs != group_contact.additional_attributes
rescue StandardError => e
Rails.logger.error "Failed to fetch pending join requests for group #{group_contact.identifier}: #{e.message}"
end
public
def try_update_group_avatar(group_contact, force: false)
if force
reset_avatar_state(group_contact)
elsif group_contact.avatar.attached?
return
end
response = get_profile_pic(group_contact.identifier)
profile_pic_url = response&.dig('data', 'profilePictureUrl')
::Avatar::AvatarFromUrlJob.perform_later(group_contact, profile_pic_url) if profile_pic_url
rescue StandardError => e
Rails.logger.error "Failed to update avatar for group #{group_contact.identifier}: #{e.message}"
end
private
def reset_avatar_state(group_contact)
group_contact.avatar.purge if group_contact.avatar.attached?
attrs = (group_contact.additional_attributes || {}).except('last_avatar_sync_at', 'avatar_url_hash')
group_contact.update_columns(additional_attributes: attrs) # rubocop:disable Rails/SkipsModelValidations
end
def try_update_participant_avatar(contact)
return if contact.avatar.attached?
phone = contact.phone_number&.delete('+')
return if phone.blank?
profile_pic_url = fetch_profile_picture_url(phone)
::Avatar::AvatarFromUrlJob.perform_later(contact, profile_pic_url) if profile_pic_url
rescue StandardError => e
Rails.logger.error "Failed to update avatar for contact #{contact.id}: #{e.message}"
end
def fetch_profile_picture_url(phone_number)
jid = "#{phone_number}@s.whatsapp.net"
response = get_profile_pic(jid)
response&.dig('data', 'profilePictureUrl')
end
def find_or_create_participant_contact(participant, inbox)
lid = extract_lid_from_participant(participant)
phone = extract_phone_from_participant(participant)
identifier = lid ? "#{lid}@lid" : nil
source_id = lid || phone
return nil if source_id.blank?
Whatsapp::ContactInboxConsolidationService.new(
inbox: inbox, phone: phone, lid: lid, identifier: identifier
).perform
contact_inbox = ::ContactInboxWithContactBuilder.new(
source_id: source_id,
inbox: inbox,
contact_attributes: {
name: phone,
phone_number: ("+#{phone}" if phone),
identifier: identifier
}
).perform
return nil if contact_inbox.blank?
update_participant_contact_info(contact_inbox.contact, phone, identifier)
end
def update_participant_contact_info(contact, phone, identifier)
update_params = {
phone_number: ("+#{phone}" if phone && contact.phone_number.blank?),
identifier: (identifier if identifier && contact.identifier.blank?)
}.compact
contact.update!(update_params) if update_params.present?
contact
end
def extract_lid_from_participant(participant)
return nil if participant[:id].blank?
jid_part, jid_suffix = participant[:id].split('@')
jid_part if jid_suffix == 'lid' && jid_part.match?(/^\d+$/)
end
def extract_phone_from_participant(participant)
return nil if participant[:phoneNumber].blank?
phone = participant[:phoneNumber].split('@').first
phone if phone.match?(/^\d+$/)
end
private_class_method def self.with_error_handling(*method_names)
method_names.each do |method_name|
original_method = instance_method(method_name)
define_method("#{method_name}_without_error_handling") do |*args, **kwargs, &block|
original_method.bind_call(self, *args, **kwargs, &block)
end
define_method(method_name) do |*args, **kwargs, &block|
original_method.bind_call(self, *args, **kwargs, &block)
rescue MessageAlreadyProcessingError
raise
rescue StandardError => e
handle_channel_error
raise e
end
end
end
def handle_channel_error
whatsapp_channel.update_provider_connection!(connection: 'close')
return if @handling_error
@handling_error = true
begin
setup_channel_provider_without_error_handling
rescue StandardError => e
Rails.logger.error "Failed to reconnect channel after error: #{e.message}"
ensure
@handling_error = false
end
end
with_error_handling :setup_channel_provider,
:disconnect_channel_provider,
:send_message,
:toggle_typing_status,
:presence_subscribe,
:update_presence,
:read_messages,
:unread_message,
:received_messages,
:group_metadata,
:sync_group,
:on_whatsapp,
:delete_message,
:edit_message,
:group_leave,
:group_setting_update,
:group_join_approval_mode,
:group_member_add_mode
end