diff --git a/app/javascript/dashboard/components-next/call/CallCard.vue b/app/javascript/dashboard/components-next/call/CallCard.vue index fcd7ad232..207129282 100644 --- a/app/javascript/dashboard/components-next/call/CallCard.vue +++ b/app/javascript/dashboard/components-next/call/CallCard.vue @@ -35,7 +35,14 @@ const props = defineProps({ }, }); -defineEmits(['accept', 'reject', 'end', 'toggleMute', 'goToConversation']); +defineEmits([ + 'accept', + 'reject', + 'end', + 'toggleMute', + 'goToConversation', + 'dismiss', +]); const { t } = useI18n(); @@ -115,6 +122,19 @@ const channelIcon = computed(() => { {{ statusLabel }} + + diff --git a/app/javascript/dashboard/components-next/call/FloatingCallWidget.vue b/app/javascript/dashboard/components-next/call/FloatingCallWidget.vue index 39e75ddbf..86d9e979c 100644 --- a/app/javascript/dashboard/components-next/call/FloatingCallWidget.vue +++ b/app/javascript/dashboard/components-next/call/FloatingCallWidget.vue @@ -25,6 +25,7 @@ const { joinCall, endCall: endCallSession, rejectIncomingCall, + dismissCall, formattedCallDuration, } = useCallSession(); @@ -51,6 +52,15 @@ const mainCardState = computed(() => { : VOICE_CALL_DIRECTION.INCOMING; }); +// Stacked cards are always non-active (ringing) calls, so reflect each call's +// real direction. An outbound call must render as OUTGOING — otherwise it shows +// the incoming-only dismiss (✕) control and the agent could drop it locally +// without terminating, leaving the customer ringing with no widget to end it. +const stackedCardState = call => + call?.callDirection === VOICE_CALL_DIRECTION.OUTBOUND + ? VOICE_CALL_DIRECTION.OUTGOING + : VOICE_CALL_DIRECTION.INCOMING; + const toggleMute = () => { isMuted.value = !isMuted.value; setWhatsappCallMuted(isMuted.value); @@ -230,10 +240,11 @@ onBeforeUnmount(stopRingtone); v-for="call in stackedIncomingCalls" :key="call.callSid" :call="call" - :state="VOICE_CALL_DIRECTION.INCOMING" + :state="stackedCardState(call)" :call-info="getCallInfo(call)" @accept="handleJoinCall(call)" @reject="rejectIncomingCall(call.callSid)" + @dismiss="dismissCall(call.callSid)" @go-to-conversation="goToConversation(call)" /> @@ -248,6 +259,7 @@ onBeforeUnmount(stopRingtone); :show-mute="hasActiveCall && isWhatsappActive" @accept="handleJoinCall(primaryIncomingCall)" @reject="rejectIncomingCall(primaryIncomingCall?.callSid)" + @dismiss="dismissCall(primaryIncomingCall?.callSid)" @end="handleEndCall" @toggle-mute="toggleMute" @go-to-conversation="goToConversation(activeCall || primaryIncomingCall)" diff --git a/app/javascript/dashboard/components-next/captain/assistant/DocumentBulkActions.vue b/app/javascript/dashboard/components-next/captain/assistant/DocumentBulkActions.vue index 378860b3e..fa4888576 100644 --- a/app/javascript/dashboard/components-next/captain/assistant/DocumentBulkActions.vue +++ b/app/javascript/dashboard/components-next/captain/assistant/DocumentBulkActions.vue @@ -25,7 +25,7 @@ const store = useStore(); const bulkDeleteDialog = ref(null); const isSyncableDocument = doc => - !doc.pdf_document && doc.status === 'available' && !doc.sync_in_progress; + !doc.pdf_document && doc.status === 'available'; const syncableSelectedIds = computed(() => { if (!props.selectedIds.size) return []; diff --git a/app/javascript/dashboard/components-next/captain/assistant/DocumentCard.vue b/app/javascript/dashboard/components-next/captain/assistant/DocumentCard.vue index c8a011b9a..9d6d574ec 100644 --- a/app/javascript/dashboard/components-next/captain/assistant/DocumentCard.vue +++ b/app/javascript/dashboard/components-next/captain/assistant/DocumentCard.vue @@ -127,7 +127,6 @@ const menuItems = computed(() => { value: 'sync', action: 'sync', icon: 'i-lucide-refresh-cw', - disabled: props.syncInProgress, }); } diff --git a/app/javascript/dashboard/components-next/message/Message.vue b/app/javascript/dashboard/components-next/message/Message.vue index d738bfcc6..72477bbd8 100644 --- a/app/javascript/dashboard/components-next/message/Message.vue +++ b/app/javascript/dashboard/components-next/message/Message.vue @@ -554,11 +554,10 @@ provideMessageContext({
diff --git a/app/javascript/dashboard/components-next/message/bubbles/Base.vue b/app/javascript/dashboard/components-next/message/bubbles/Base.vue index c40d63363..457b583ea 100644 --- a/app/javascript/dashboard/components-next/message/bubbles/Base.vue +++ b/app/javascript/dashboard/components-next/message/bubbles/Base.vue @@ -95,7 +95,7 @@ const replyToPreview = computed(() => { diff --git a/app/javascript/dashboard/components/widgets/WootWriter/Editor.vue b/app/javascript/dashboard/components/widgets/WootWriter/Editor.vue index 7881a0d25..3da526fc0 100644 --- a/app/javascript/dashboard/components/widgets/WootWriter/Editor.vue +++ b/app/javascript/dashboard/components/widgets/WootWriter/Editor.vue @@ -589,11 +589,11 @@ function isCmdPlusEnterToSendEnabled() { useKeyboardEvents({ 'Alt+KeyP': { action: focusEditorInputField, - allowOnFocusedInput: true, + allowOnFocusedInput: false, }, 'Alt+KeyL': { action: focusEditorInputField, - allowOnFocusedInput: true, + allowOnFocusedInput: false, }, }); diff --git a/app/javascript/dashboard/components/widgets/WootWriter/ReplyTopPanel.vue b/app/javascript/dashboard/components/widgets/WootWriter/ReplyTopPanel.vue index a2929f7f8..bf45cdc1c 100644 --- a/app/javascript/dashboard/components/widgets/WootWriter/ReplyTopPanel.vue +++ b/app/javascript/dashboard/components/widgets/WootWriter/ReplyTopPanel.vue @@ -105,11 +105,11 @@ export default { const keyboardEvents = { 'Alt+KeyP': { action: () => handleNoteClick(), - allowOnFocusedInput: true, + allowOnFocusedInput: false, }, 'Alt+KeyL': { action: () => handleReplyClick(), - allowOnFocusedInput: true, + allowOnFocusedInput: false, }, }; useKeyboardEvents(keyboardEvents); diff --git a/app/javascript/dashboard/composables/spec/useKeyboardEvents.spec.js b/app/javascript/dashboard/composables/spec/useKeyboardEvents.spec.js index d20babeaa..995a486b4 100644 --- a/app/javascript/dashboard/composables/spec/useKeyboardEvents.spec.js +++ b/app/javascript/dashboard/composables/spec/useKeyboardEvents.spec.js @@ -1,6 +1,47 @@ +const keyboardEventsMock = vi.hoisted(() => ({ + mountedCallbacks: [], + unmountedCallbacks: [], + registeredKeybindings: {}, +})); + +vi.mock('vue', async importOriginal => ({ + ...(await importOriginal()), + onMounted: vi.fn(callback => + keyboardEventsMock.mountedCallbacks.push(callback) + ), + onUnmounted: vi.fn(callback => + keyboardEventsMock.unmountedCallbacks.push(callback) + ), +})); + +vi.mock('dashboard/composables/useDetectKeyboardLayout', () => ({ + useDetectKeyboardLayout: vi.fn(() => Promise.resolve('QWERTY')), +})); + +vi.mock('tinykeys', () => ({ + createKeybindingsHandler: vi.fn(keybindings => { + keyboardEventsMock.registeredKeybindings = keybindings; + return event => keybindings[event.shortcut]?.(event); + }), +})); + +import { useDetectKeyboardLayout } from 'dashboard/composables/useDetectKeyboardLayout'; import { useKeyboardEvents } from 'dashboard/composables/useKeyboardEvents'; +import { LAYOUT_QWERTZ } from 'shared/helpers/KeyboardHelpers'; + +const mount = async events => { + useKeyboardEvents(events); + await keyboardEventsMock.mountedCallbacks[0](); + return keyboardEventsMock.registeredKeybindings; +}; describe('useKeyboardEvents', () => { + beforeEach(() => { + keyboardEventsMock.mountedCallbacks = []; + keyboardEventsMock.unmountedCallbacks = []; + keyboardEventsMock.registeredKeybindings = {}; + }); + it('should be defined', () => { expect(useKeyboardEvents).toBeDefined(); }); @@ -9,18 +50,46 @@ describe('useKeyboardEvents', () => { expect(useKeyboardEvents).toBeInstanceOf(Function); }); - it('should set up listeners on mount and remove them on unmount', async () => { - const events = { - 'ALT+KeyL': () => {}, - }; + it('registers keybindings on mount and removes the listener on unmount', async () => { + const addEventListener = vi.spyOn(document, 'addEventListener'); - const mountedMock = vi.fn(); - const unmountedMock = vi.fn(); - useKeyboardEvents(events); - mountedMock(); - unmountedMock(); + const keybindings = await mount({ 'ALT+KeyL': () => {} }); + expect(Object.keys(keybindings)).toEqual(['ALT+KeyL']); - expect(mountedMock).toHaveBeenCalled(); - expect(unmountedMock).toHaveBeenCalled(); + const { signal } = addEventListener.mock.calls[0][2]; + expect(signal.aborted).toBe(false); + + keyboardEventsMock.unmountedCallbacks[0](); + expect(signal.aborted).toBe(true); + }); + + it('prefixes Shift on QWERTZ layouts for the affected keys', async () => { + useDetectKeyboardLayout.mockResolvedValueOnce(LAYOUT_QWERTZ); + + const keybindings = await mount({ 'Alt+KeyL': () => {} }); + + expect(Object.keys(keybindings)).toEqual(['Shift+Alt+KeyL']); + }); + + it('ignores shortcuts on focused typeable elements by default', async () => { + const action = vi.fn(); + const keybindings = await mount({ + 'Alt+KeyL': { action, allowOnFocusedInput: false }, + }); + + keybindings['Alt+KeyL']({ target: document.createElement('textarea') }); + + expect(action).not.toHaveBeenCalled(); + }); + + it('runs shortcuts on focused typeable elements when allowed', async () => { + const action = vi.fn(); + const keybindings = await mount({ + 'Alt+KeyL': { action, allowOnFocusedInput: true }, + }); + + keybindings['Alt+KeyL']({ target: document.createElement('textarea') }); + + expect(action).toHaveBeenCalledTimes(1); }); }); diff --git a/app/javascript/dashboard/helper/voice.js b/app/javascript/dashboard/helper/voice.js index 849738f6f..eaa523d75 100644 --- a/app/javascript/dashboard/helper/voice.js +++ b/app/javascript/dashboard/helper/voice.js @@ -1,4 +1,5 @@ import { CONTENT_TYPES } from 'dashboard/components-next/message/constants'; +import { MESSAGE_TYPE } from 'shared/constants/messages'; import { useCallsStore } from 'dashboard/stores/calls'; import types from 'dashboard/store/mutation-types'; @@ -58,6 +59,10 @@ function extractCallerSnapshot(message) { // Snapshot caller info from the message at add-time so the widget can keep // rendering it after the user navigates away from a conversation list that // had the conversation hydrated (and Vuex evicts it from the store). + // Only incoming messages carry the contact as the sender; on outbound calls + // the sender is the initiating agent, so skip the snapshot and let the widget + // fall back to the conversation's contact (conversation.meta.sender). + if (message?.message_type !== MESSAGE_TYPE.INCOMING) return null; const sender = message?.sender; if (!sender) return null; return { diff --git a/app/javascript/dashboard/i18n/locale/en/conversation.json b/app/javascript/dashboard/i18n/locale/en/conversation.json index f0ff6811a..daf49ab59 100644 --- a/app/javascript/dashboard/i18n/locale/en/conversation.json +++ b/app/javascript/dashboard/i18n/locale/en/conversation.json @@ -86,13 +86,16 @@ "MISSED_CALL_INBOUND_SUBTEXT": "No agent picked up", "MISSED_CALL_DECLINED_BY": "Declined by {agentName}", "CALL_ENDED": "Call ended", + "HANDLED_BY": "Handled by {agentName}", "NOT_ANSWERED_YET": "Not answered yet", "CALLING": "Calling…", "THEY_ANSWERED": "They answered", "YOU_ANSWERED": "You answered", "AGENT_ANSWERED": "{agentName} answered", "JOIN_CALL": "Join call", - "CALL_BACK": "Call back" + "CALL_BACK": "Call back", + "TRANSCRIPT_SHOW_MORE": "Show more", + "TRANSCRIPT_SHOW_LESS": "Show less" }, "HEADER": { "RESOLVE_ACTION": "Resolve", @@ -312,6 +315,7 @@ "NOT_ANSWERED_YET": "Not answered yet", "HANDLED_IN_ANOTHER_TAB": "Being handled in another tab", "REJECT_CALL": "Reject", + "DISMISS_CALL": "Dismiss", "JOIN_CALL": "Join call", "END_CALL": "End call", "MUTE": "Mute mic", diff --git a/app/javascript/dashboard/routes/dashboard/settings/inbox/Settings.vue b/app/javascript/dashboard/routes/dashboard/settings/inbox/Settings.vue index 62cd2b7bc..9e5206bb3 100644 --- a/app/javascript/dashboard/routes/dashboard/settings/inbox/Settings.vue +++ b/app/javascript/dashboard/routes/dashboard/settings/inbox/Settings.vue @@ -253,7 +253,6 @@ export default { if ( this.isAWhatsAppCloudChannel && - this.isEmbeddedSignupWhatsApp && this.isFeatureEnabledonAccount( this.accountId, FEATURE_FLAGS.CHANNEL_VOICE diff --git a/app/javascript/dashboard/store/modules/auditlogs.js b/app/javascript/dashboard/store/modules/auditlogs.js index c0970cf18..5b355736e 100644 --- a/app/javascript/dashboard/store/modules/auditlogs.js +++ b/app/javascript/dashboard/store/modules/auditlogs.js @@ -7,7 +7,7 @@ const state = { records: [], meta: { currentPage: 1, - perPage: 15, + perPage: 25, totalEntries: 0, }, uiFlags: { diff --git a/app/javascript/dashboard/stores/calls.js b/app/javascript/dashboard/stores/calls.js index 2a634a5d9..9b54e230c 100644 --- a/app/javascript/dashboard/stores/calls.js +++ b/app/javascript/dashboard/stores/calls.js @@ -54,7 +54,9 @@ export const useCallsStore = defineStore('calls', { return; } - this.calls.push({ + // Prepend so the newest call surfaces as the primary card (incomingCalls[0]) + // and older calls drop into the stack below it. + this.calls.unshift({ ...callData, isActive: false, }); diff --git a/app/jobs/internal/trigger_daily_scheduled_items_job.rb b/app/jobs/internal/trigger_daily_scheduled_items_job.rb index 80b92ea6b..b09193e92 100644 --- a/app/jobs/internal/trigger_daily_scheduled_items_job.rb +++ b/app/jobs/internal/trigger_daily_scheduled_items_job.rb @@ -23,3 +23,5 @@ class Internal::TriggerDailyScheduledItemsJob < ApplicationJob @designated_minute ||= Digest::MD5.hexdigest(ChatwootHub.installation_identifier).hex % 1440 end end + +Internal::TriggerDailyScheduledItemsJob.prepend_mod_with('Internal::TriggerDailyScheduledItemsJob') diff --git a/app/models/channel/whatsapp.rb b/app/models/channel/whatsapp.rb index 2c2205b19..00d3277c9 100644 --- a/app/models/channel/whatsapp.rb +++ b/app/models/channel/whatsapp.rb @@ -41,19 +41,19 @@ class Channel::Whatsapp < ApplicationRecord end # Mirrors Channel::TwilioSms#voice_enabled? so the call subsystem can duck-type across providers. - # Meta's Calling API is only available via the embedded-signup whatsapp_cloud flow — - # 360dialog (default provider) and manual whatsapp_cloud setups can't reach the call APIs. + # Meta's Calling API is available to any whatsapp_cloud inbox (embedded-signup or manual keys); + # only 360dialog (default provider) can't reach the call APIs. def voice_enabled? voice_calling_supported? && provider_config['calling_enabled'].present? && account.feature_enabled?('channel_voice') end - # Whether this inbox can do WhatsApp calling at all. Meta's Calling API is only - # reachable via the embedded-signup whatsapp_cloud flow, so manual whatsapp_cloud - # and 360dialog inboxes can't be toggled on even though calling_enabled would persist. + # Whether this inbox can do WhatsApp calling at all. Meta's Calling API is + # reachable by any whatsapp_cloud inbox, so 360dialog inboxes can't be toggled + # on even though calling_enabled would persist. def voice_calling_supported? - provider == 'whatsapp_cloud' && provider_config['source'] == 'embedded_signup' + provider == 'whatsapp_cloud' end def provider_service @@ -69,7 +69,7 @@ class Channel::Whatsapp < ApplicationRecord # Saved with validate: false to skip validate_provider_config's remote credential # re-check, which could spuriously fail and desync the flag from Meta. def enable_voice_calling! - raise 'WhatsApp calling requires an embedded-signup whatsapp_cloud inbox' unless voice_calling_supported? + raise 'WhatsApp calling requires a whatsapp_cloud inbox' unless voice_calling_supported? raise 'WhatsApp calling requires the channel_voice feature' unless account.feature_enabled?('channel_voice') provider_service.update_calling_status('ENABLED') @@ -82,7 +82,7 @@ class Channel::Whatsapp < ApplicationRecord # `calls` from the webhook subscription (best-effort, so a Meta outage can't # trap admins). Leaves Meta's WABA calling.status untouched. def disable_voice_calling! - raise 'WhatsApp calling requires an embedded-signup whatsapp_cloud inbox' unless voice_calling_supported? + raise 'WhatsApp calling requires a whatsapp_cloud inbox' unless voice_calling_supported? self.provider_config = provider_config.merge('calling_enabled' => false) save!(validate: false) diff --git a/app/services/website_branding_service.rb b/app/services/website_branding_service.rb index 3592267ff..026b51e0e 100644 --- a/app/services/website_branding_service.rb +++ b/app/services/website_branding_service.rb @@ -1,3 +1,5 @@ +require 'resolv' + class WebsiteBrandingService include SocialLinkParser @@ -24,7 +26,8 @@ class WebsiteBrandingService colors: extract_colors(doc), logos: extract_logos(doc), socials: build_socials(links), - email: @email + email: @email, + email_provider: detect_email_provider }) rescue StandardError => e Rails.logger.error "[WebsiteBranding] #{e.message}" @@ -104,6 +107,37 @@ class WebsiteBrandingService rescue URI::InvalidURIError nil end + + GOOGLE_MX_DOMAINS = %w[google.com googlemail.com].freeze + MICROSOFT_MX_DOMAINS = %w[outlook.com].freeze + + # Probes the domain's MX records to infer the mailbox provider, returning + # 'google' or 'microsoft' (matching Channel::Email#provider) or nil when unknown. + def detect_email_provider + hosts = mx_records + return 'google' if mx_hosted_by?(hosts, GOOGLE_MX_DOMAINS) + return 'microsoft' if mx_hosted_by?(hosts, MICROSOFT_MX_DOMAINS) + + nil + end + + # Matches on the registrable domain of the MX host (anchored on a label + # boundary) so lookalikes like "notgoogle.com" don't get misclassified. + def mx_hosted_by?(hosts, provider_domains) + hosts.any? do |host| + provider_domains.any? { |domain| host == domain || host.end_with?(".#{domain}") } + end + end + + def mx_records + Resolv::DNS.open do |resolver| + resolver.timeouts = 5 + resolver.getresources(@domain, Resolv::DNS::Resource::IN::MX).map { |record| record.exchange.to_s.downcase } + end + rescue StandardError => e + Rails.logger.error "[WebsiteBranding] MX probe failed for #{@domain}: #{e.message}" + [] + end end WebsiteBrandingService.prepend_mod_with('WebsiteBrandingService') diff --git a/config/installation_config.yml b/config/installation_config.yml index e374d0948..673c6df4c 100644 --- a/config/installation_config.yml +++ b/config/installation_config.yml @@ -215,6 +215,18 @@ value: locked: false type: code +- name: CAPTAIN_DOCUMENT_AUTO_SYNC_PER_ACCOUNT_BATCH_LIMIT + display_title: 'Captain Document Auto Sync Per Account Batch Limit' + description: 'Maximum syncable Captain documents to enqueue per account in one scheduler run. Defaults to 50.' + value: 50 + locked: false + type: number +- name: CAPTAIN_DOCUMENT_AUTO_SYNC_GLOBAL_BATCH_LIMIT + display_title: 'Captain Document Auto Sync Global Batch Limit' + description: 'Maximum syncable Captain documents to enqueue globally in one scheduler run. Defaults to 1000.' + value: 1000 + locked: false + type: number # End of Captain Config # ------- Context.dev Config ------- # diff --git a/enterprise/app/controllers/api/v1/accounts/audit_logs_controller.rb b/enterprise/app/controllers/api/v1/accounts/audit_logs_controller.rb index 29e02c1dd..7c569b3f1 100644 --- a/enterprise/app/controllers/api/v1/accounts/audit_logs_controller.rb +++ b/enterprise/app/controllers/api/v1/accounts/audit_logs_controller.rb @@ -2,7 +2,7 @@ class Api::V1::Accounts::AuditLogsController < Api::V1::Accounts::EnterpriseAcco before_action :check_admin_authorization? before_action :fetch_audit - RESULTS_PER_PAGE = 15 + RESULTS_PER_PAGE = 25 def show @audit_logs = @audit_logs.page(params[:page]).per(RESULTS_PER_PAGE) diff --git a/enterprise/app/controllers/api/v1/accounts/captain/bulk_actions_controller.rb b/enterprise/app/controllers/api/v1/accounts/captain/bulk_actions_controller.rb index b9a6bbcc5..7e2817f69 100644 --- a/enterprise/app/controllers/api/v1/accounts/captain/bulk_actions_controller.rb +++ b/enterprise/app/controllers/api/v1/accounts/captain/bulk_actions_controller.rb @@ -75,7 +75,6 @@ class Api::V1::Accounts::Captain::BulkActionsController < Api::V1::Accounts::Bas Current.account.captain_documents.where(id: params[:ids]).find_each(batch_size: 100) do |document| next unless document.syncable? next unless document.available? - next if document.sync_in_progress? document.update!( sync_status: :syncing, diff --git a/enterprise/app/controllers/api/v1/accounts/captain/documents_controller.rb b/enterprise/app/controllers/api/v1/accounts/captain/documents_controller.rb index 23f410499..273c082b1 100644 --- a/enterprise/app/controllers/api/v1/accounts/captain/documents_controller.rb +++ b/enterprise/app/controllers/api/v1/accounts/captain/documents_controller.rb @@ -37,7 +37,6 @@ class Api::V1::Accounts::Captain::DocumentsController < Api::V1::Accounts::BaseC def sync return render_could_not_create_error(I18n.t('captain.documents.sync_not_supported_for_pdf')) unless @document.syncable? return render_could_not_create_error(I18n.t('captain.documents.sync_only_available_documents')) unless @document.available? - return render_could_not_create_error(I18n.t('captain.documents.sync_already_in_progress')) if @document.sync_in_progress? @document.update!( sync_status: :syncing, diff --git a/enterprise/app/controllers/api/v1/accounts/whatsapp_calls_controller.rb b/enterprise/app/controllers/api/v1/accounts/whatsapp_calls_controller.rb index b54940586..0301d428b 100644 --- a/enterprise/app/controllers/api/v1/accounts/whatsapp_calls_controller.rb +++ b/enterprise/app/controllers/api/v1/accounts/whatsapp_calls_controller.rb @@ -35,8 +35,14 @@ class Api::V1::Accounts::WhatsappCallsController < Api::V1::Accounts::BaseContro def initiate @call = create_outbound_call - @message = Voice::CallMessageBuilder.new(@call).perform! - @call.update!(message_id: @message.id) + # Link the call to its message in one transaction so the message.created + # broadcast (an after_create_commit hook) fires only once call.message_id is + # set. Otherwise the live ringing bubble receives a message with no `call` + # payload (no direction/agent) and renders "Calling…" instead of "Handled by …". + ActiveRecord::Base.transaction do + @message = Voice::CallMessageBuilder.new(@call).perform! + @call.update!(message_id: @message.id) + end end private diff --git a/enterprise/app/jobs/captain/documents/perform_sync_job.rb b/enterprise/app/jobs/captain/documents/perform_sync_job.rb index eaef7d64d..2a33db3f3 100644 --- a/enterprise/app/jobs/captain/documents/perform_sync_job.rb +++ b/enterprise/app/jobs/captain/documents/perform_sync_job.rb @@ -48,9 +48,7 @@ class Captain::Documents::PerformSyncJob < MutexApplicationJob return if document.pdf_document? with_lock(lock_key(document), LOCK_TIMEOUT) do - mark_sync_started(document) - result = Captain::Documents::SyncService.new(document.reload).perform - log_sync_outcome(document, result: result, duration_ms: duration_ms_since(start_time)) + perform_sync(document, start_time) end rescue LockAcquisitionError log_sync_outcome(document, result: :already_syncing) @@ -64,6 +62,12 @@ class Captain::Documents::PerformSyncJob < MutexApplicationJob private + def perform_sync(document, start_time) + mark_sync_started(document) + result = Captain::Documents::SyncService.new(document.reload).perform + log_sync_outcome(document, result: result, duration_ms: duration_ms_since(start_time)) + end + def log_sync_outcome(document, **fields) payload = { document_id: document.id, diff --git a/enterprise/app/jobs/captain/documents/schedule_syncs_job.rb b/enterprise/app/jobs/captain/documents/schedule_syncs_job.rb index 393a307db..82693c0bb 100644 --- a/enterprise/app/jobs/captain/documents/schedule_syncs_job.rb +++ b/enterprise/app/jobs/captain/documents/schedule_syncs_job.rb @@ -1,21 +1,27 @@ class Captain::Documents::ScheduleSyncsJob < ApplicationJob queue_as :scheduled_jobs - PER_ACCOUNT_HOURLY_CAP = 50 - GLOBAL_HOURLY_CAP = 1000 - DUE_DOCUMENT_BATCH_SIZE = PER_ACCOUNT_HOURLY_CAP * 2 # Inspite of skipping, we should at least reach the hourly cap + DEFAULT_PER_ACCOUNT_BATCH_LIMIT = 50 + DEFAULT_GLOBAL_BATCH_LIMIT = 1000 SYNC_STALE_TIMEOUT = Captain::Document::SYNC_STALE_TIMEOUT + DAILY_SYNC_JITTER = 4.hours + WEEKLY_SYNC_JITTER = 1.day + MONTHLY_SYNC_JITTER = 4.days - def perform - @remaining_global_capacity = GLOBAL_HOURLY_CAP + def perform(plan_name = nil) + @per_account_batch_limit = configured_sync_limit('CAPTAIN_DOCUMENT_AUTO_SYNC_PER_ACCOUNT_BATCH_LIMIT', DEFAULT_PER_ACCOUNT_BATCH_LIMIT) + @global_batch_limit = configured_sync_limit('CAPTAIN_DOCUMENT_AUTO_SYNC_GLOBAL_BATCH_LIMIT', DEFAULT_GLOBAL_BATCH_LIMIT) + @remaining_global_capacity = @global_batch_limit + @plan_name = plan_name.to_s.downcase.presence sync_intervals = Enterprise::Account.captain_document_sync_intervals - stats = { accounts_scanned: 0, accounts_enabled: 0, accounts_scheduled: 0, documents_enqueued: 0, documents_skipped: 0 } + stats = { accounts_scanned: 0, accounts_enabled: 0, accounts_scheduled: 0, documents_enqueued: 0 } Account.joins(:captain_documents).distinct.find_each(batch_size: 100) do |account| break if @remaining_global_capacity <= 0 stats[:accounts_scanned] += 1 next unless account.feature_enabled?('captain_document_auto_sync') + next unless account_in_selected_plan?(account) stats[:accounts_enabled] += 1 interval = account.captain_document_sync_interval(sync_intervals) @@ -24,7 +30,6 @@ class Captain::Documents::ScheduleSyncsJob < ApplicationJob stats[:accounts_scheduled] += 1 result = enqueue_due_documents(account, interval) stats[:documents_enqueued] += result[:enqueued] - stats[:documents_skipped] += result[:skipped] end log_scheduler_summary(stats) @@ -33,92 +38,85 @@ class Captain::Documents::ScheduleSyncsJob < ApplicationJob private def enqueue_due_documents(account, interval) - per_account_limit = [PER_ACCOUNT_HOURLY_CAP, @remaining_global_capacity].min - result = { enqueued: 0, skipped: 0 } - skipped_document_ids = [] + per_account_limit = [@per_account_batch_limit, @remaining_global_capacity].min + result = { enqueued: 0 } - while result[:enqueued] < per_account_limit - - documents = due_documents(account, interval, skipped_document_ids).limit(DUE_DOCUMENT_BATCH_SIZE).to_a - break if documents.empty? - - documents.each do |document| - break if result[:enqueued] >= per_account_limit - - process_due_document(document, result, skipped_document_ids) - end + due_documents(account, interval).limit(per_account_limit).to_a.each do |document| + process_due_document(document, interval, result) end result end - def process_due_document(document, result, skipped_document_ids) - return unless document.syncable? + def process_due_document(document, interval, result) + sync_execution_delay = sync_jitter(interval) - # Reserve the sync slot before enqueueing so later scheduler runs skip this document while the job is queued. - unless reserve_sync_slot(document) - result[:skipped] += 1 - skipped_document_ids << document.id - return - end - - Captain::Documents::PerformSyncJob.perform_later(document) + Captain::Documents::PerformSyncJob.set(queue: :purgable, wait: sync_execution_delay).perform_later(document) @remaining_global_capacity -= 1 result[:enqueued] += 1 end - def due_documents(account, interval, skipped_document_ids) + def due_documents(account, interval) syncing = Captain::Document.sync_statuses[:syncing] synced = Captain::Document.sync_statuses[:synced] failed = Captain::Document.sync_statuses[:failed] + stale_cutoff = SYNC_STALE_TIMEOUT.ago + # The scheduler runs at predictable plan windows. Use a wider due window so + # jittered executions do not miss the next window just because they finished later. + sync_due_before = due_window(interval).ago documents = account.captain_documents.syncable.where(status: :available).where( '(sync_status = ? AND last_synced_at < ?) OR (sync_status = ? AND last_sync_attempted_at < ?) OR ' \ '(sync_status = ? AND last_sync_attempted_at < ?)', - synced, interval.ago, failed, interval.ago, syncing, SYNC_STALE_TIMEOUT.ago + synced, sync_due_before, failed, sync_due_before, syncing, stale_cutoff ) - documents = documents.where.not(id: skipped_document_ids) if skipped_document_ids.present? documents.order(Arel.sql('last_sync_attempted_at ASC NULLS FIRST'), :id) end - def reserve_sync_slot(document) - mark_sync_started(document) - true - rescue ActiveRecord::RecordInvalid => e - log_document_skip(document, e) - false + def configured_sync_limit(config_key, default) + configured_value = InstallationConfig.find_by(name: config_key)&.value + limit = configured_value.to_s.to_i + limit.positive? ? limit : default end - def log_document_skip(document, error) - payload = { - event: 'document_skipped', - document_id: document.id, - account_id: document.account_id, - assistant_id: document.assistant_id, - error_class: error.class.name, - error_message: error.message, - validation_errors: document.errors.full_messages - } + def account_in_selected_plan?(account) + return true if @plan_name.blank? - Rails.logger.warn("[Captain::Documents::ScheduleSyncsJob] #{payload.to_json}") + account_sync_plan(account) == @plan_name + end + + def account_sync_plan(account) + plan = account.custom_attributes['plan_name'] + plan = 'enterprise' if plan.blank? && ChatwootApp.self_hosted_enterprise? + plan.to_s.downcase.presence + end + + def sync_jitter(interval) + jitter_window = if interval <= 1.day + DAILY_SYNC_JITTER + elsif interval <= 1.week + WEEKLY_SYNC_JITTER + else + MONTHLY_SYNC_JITTER + end + + rand(0..jitter_window.to_i).seconds + end + + def due_window(interval) + (interval.to_i / 2).seconds end def log_scheduler_summary(stats) payload = { event: 'completed', + plan_name: @plan_name, global_cap_hit: @remaining_global_capacity <= 0, + per_account_batch_limit: @per_account_batch_limit, + global_batch_limit: @global_batch_limit, remaining_global_capacity: @remaining_global_capacity }.merge(stats) Rails.logger.info("[Captain::Documents::ScheduleSyncsJob] #{payload.to_json}") end - - def mark_sync_started(document) - document.update!( - sync_status: :syncing, - sync_step: nil, - last_sync_error_code: nil, - last_sync_attempted_at: Time.current - ) - end end diff --git a/enterprise/app/jobs/enterprise/internal/trigger_daily_scheduled_items_job.rb b/enterprise/app/jobs/enterprise/internal/trigger_daily_scheduled_items_job.rb new file mode 100644 index 000000000..aa22310fb --- /dev/null +++ b/enterprise/app/jobs/enterprise/internal/trigger_daily_scheduled_items_job.rb @@ -0,0 +1,19 @@ +module Enterprise::Internal::TriggerDailyScheduledItemsJob + def perform + super + + Captain::Documents::ScheduleSyncsJob.perform_later('enterprise') + Captain::Documents::ScheduleSyncsJob.perform_later('business') if business_auto_sync_due? + Captain::Documents::ScheduleSyncsJob.perform_later('startups') if startup_auto_sync_due? + end + + private + + def business_auto_sync_due? + Time.current.utc.sunday? + end + + def startup_auto_sync_due? + Time.current.utc.day == 1 + end +end diff --git a/enterprise/app/jobs/enterprise/internal/trigger_hourly_scheduled_items_job.rb b/enterprise/app/jobs/enterprise/internal/trigger_hourly_scheduled_items_job.rb deleted file mode 100644 index 9d3baa2d6..000000000 --- a/enterprise/app/jobs/enterprise/internal/trigger_hourly_scheduled_items_job.rb +++ /dev/null @@ -1,7 +0,0 @@ -module Enterprise::Internal::TriggerHourlyScheduledItemsJob - def perform - super - - Captain::Documents::ScheduleSyncsJob.perform_later - end -end diff --git a/enterprise/app/jobs/enterprise/webhooks/whatsapp_events_job.rb b/enterprise/app/jobs/enterprise/webhooks/whatsapp_events_job.rb index b1e2bc053..db411f068 100644 --- a/enterprise/app/jobs/enterprise/webhooks/whatsapp_events_job.rb +++ b/enterprise/app/jobs/enterprise/webhooks/whatsapp_events_job.rb @@ -31,10 +31,11 @@ module Enterprise::Webhooks::WhatsappEventsJob # and timer/recorder kick off before the contact actually answers. def handle_call_events(channel, params) value = params.dig(:entry, 0, :changes, 0, :value) || {} + contacts = value[:contacts] Array(value[:calls]).each do |call_payload| with_call_lock(channel, call_payload[:id]) do - Whatsapp::IncomingCallService.new(inbox: channel.inbox, params: { calls: [call_payload] }).perform + Whatsapp::IncomingCallService.new(inbox: channel.inbox, params: { calls: [call_payload], contacts: contacts }).perform end end diff --git a/enterprise/app/services/enterprise/billing/reconcile_plan_features_service.rb b/enterprise/app/services/enterprise/billing/reconcile_plan_features_service.rb index d543adad0..b9c3f6856 100644 --- a/enterprise/app/services/enterprise/billing/reconcile_plan_features_service.rb +++ b/enterprise/app/services/enterprise/billing/reconcile_plan_features_service.rb @@ -16,6 +16,7 @@ class Enterprise::Billing::ReconcilePlanFeaturesService advanced_search_indexing advanced_search linear_integration + channel_voice ].freeze BUSINESS_PLAN_FEATURES = %w[sla custom_roles csat_review_notes conversation_required_attributes advanced_assignment custom_tools].freeze diff --git a/enterprise/app/services/enterprise/website_branding_service.rb b/enterprise/app/services/enterprise/website_branding_service.rb index 553317c43..c6f2fbdaa 100644 --- a/enterprise/app/services/enterprise/website_branding_service.rb +++ b/enterprise/app/services/enterprise/website_branding_service.rb @@ -58,6 +58,7 @@ module Enterprise::WebsiteBrandingService socials: brand['socials'] || [], links: brand['links'], email: @email, + email_provider: detect_email_provider, industries: brand.dig('industries', 'eic') || [], stock: brand['stock'], is_nsfw: brand['is_nsfw'] || false diff --git a/enterprise/app/services/messages/audio_transcription_service.rb b/enterprise/app/services/messages/audio_transcription_service.rb index c502b717a..748bf1efa 100644 --- a/enterprise/app/services/messages/audio_transcription_service.rb +++ b/enterprise/app/services/messages/audio_transcription_service.rb @@ -1,11 +1,12 @@ class Messages::AudioTranscriptionService< Llm::LegacyBaseOpenAiService include Integrations::LlmInstrumentation - WHISPER_MODEL = 'whisper-1'.freeze - # Whisper's hard limit is 25 MB *decimal* (25_000_000), not binary (25.megabytes - # = 26_214_400) — using the binary form leaks the 25.0–26.2 MB range to the API - # as 413s. Long audio (~70+ min Opus) keeps the attachment but skips transcription. - WHISPER_BYTE_LIMIT = 25_000_000 + TRANSCRIPTION_MODEL = 'gpt-4o-mini-transcribe'.freeze + # OpenAI's transcription endpoint hard limit is 25 MB *decimal* (25_000_000), not + # binary (25.megabytes = 26_214_400) — using the binary form leaks the 25.0–26.2 MB + # range to the API as 413s. Long audio (~70+ min Opus) keeps the attachment but skips + # transcription. + TRANSCRIPTION_BYTE_LIMIT = 25_000_000 attr_reader :attachment, :message, :account @@ -42,7 +43,7 @@ class Messages::AudioTranscriptionService< Llm::LegacyBaseOpenAiService blob = attachment.file&.blob return false unless blob - blob.byte_size > WHISPER_BYTE_LIMIT + blob.byte_size > TRANSCRIPTION_BYTE_LIMIT end def fetch_audio_file @@ -75,12 +76,12 @@ class Messages::AudioTranscriptionService< Llm::LegacyBaseOpenAiService transcribed_text = nil File.open(temp_file_path, 'rb') do |file| - # temperature: 0.0 minimises Whisper's hallucinations on silence / - # near-silent audio; non-zero values trigger spiraling repeats like - # "Oh, dear. Oh, dear. Oh, dear." — well-documented Whisper behaviour. + # temperature: 0.0 minimises hallucinations on silence / near-silent + # audio; non-zero values trigger spiraling repeats — well-documented + # behaviour across OpenAI transcription models. response = @client.audio.transcribe( parameters: { - model: WHISPER_MODEL, + model: TRANSCRIPTION_MODEL, file: file, temperature: 0.0 } @@ -97,7 +98,7 @@ class Messages::AudioTranscriptionService< Llm::LegacyBaseOpenAiService def instrumentation_params(file_path) { span_name: 'llm.messages.audio_transcription', - model: WHISPER_MODEL, + model: TRANSCRIPTION_MODEL, account_id: account&.id, feature_name: 'audio_transcription', file_path: file_path diff --git a/enterprise/app/services/voice/inbound_call_builder.rb b/enterprise/app/services/voice/inbound_call_builder.rb index b6fb089b7..eef70e76b 100644 --- a/enterprise/app/services/voice/inbound_call_builder.rb +++ b/enterprise/app/services/voice/inbound_call_builder.rb @@ -59,9 +59,17 @@ class Voice::InboundCallBuilder end def ensure_contact! - account.contacts.find_or_create_by!(phone_number: from_number) do |record| - record.name = from_number if record.name.blank? + contact = account.contacts.find_or_create_by!(phone_number: from_number) do |record| + record.name = contact_name.presence || from_number end + contact.update!(name: contact_name) if contact_name.present? && contact.name == from_number + contact + end + + # WhatsApp inbound calls carry the caller's profile name in extra_meta; Twilio + # calls don't, so contact naming falls back to the phone number. + def contact_name + extra_meta['contact_name'].presence end # WhatsApp ContactInbox.source_id must be digits-only (the wa_id); Twilio accepts the +. diff --git a/enterprise/app/services/whatsapp/incoming_call_service.rb b/enterprise/app/services/whatsapp/incoming_call_service.rb index 99f6350ff..afca508a8 100644 --- a/enterprise/app/services/whatsapp/incoming_call_service.rb +++ b/enterprise/app/services/whatsapp/incoming_call_service.rb @@ -73,15 +73,27 @@ class Whatsapp::IncomingCallService def create_inbound_call(payload) sdp_offer = payload.dig(:session, :sdp) + extra_meta = { 'sdp_offer' => sdp_offer, 'ice_servers' => Call.default_ice_servers } + name = caller_profile_name(payload) + extra_meta['contact_name'] = name if name.present? + call = Voice::InboundCallBuilder.perform!( inbox: inbox, from_number: "+#{payload[:from]}", call_sid: payload[:id], - provider: :whatsapp, - extra_meta: { 'sdp_offer' => sdp_offer, 'ice_servers' => Call.default_ice_servers } + provider: :whatsapp, extra_meta: extra_meta ) update_conversation(call) broadcast_incoming(call, sdp_offer) end + # Match strictly on wa_id (== calls[].from): in a batched payload missing this + # call's contact entry, borrowing another caller's name would corrupt this + # contact, so fall back to the phone number (nil here) instead of contacts.first. + def caller_profile_name(payload) + contacts = Array(params[:contacts]).map(&:with_indifferent_access) + match = contacts.find { |c| c[:wa_id].to_s == payload[:from].to_s } + match&.dig(:profile, :name).presence + end + # `connect` is the WebRTC tunnel-ready signal, not the pickup signal. Apply # Meta's SDP answer so the handshake completes during ringing; the call # stays in `ringing` until status=ACCEPTED arrives. Don't gate on @@ -159,11 +171,13 @@ class Whatsapp::IncomingCallService ) end - # Ring the assignee if any, else online inbox agents, else the whole account. + # Ring the assignee if any, else online inbox agents, else fall back to the + # inbox's own agents and account admins. Never the whole-account stream, which + # would ring online agents from unrelated inboxes. def broadcast_incoming(call, sdp_offer) contact = call.contact token = call.conversation.assignee&.pubsub_token - streams = token ? [token] : (online_agent_streams.presence || account_streams) + streams = token ? [token] : (online_agent_streams.presence || fallback_agent_streams) broadcast(call, 'voice_call.incoming', streams: streams, direction: call.direction_label, inbox_id: call.inbox_id, @@ -175,6 +189,11 @@ class Whatsapp::IncomingCallService inbox.available_agents.pluck('users.pubsub_token').compact end + def fallback_agent_streams + user_ids = inbox.member_ids | inbox.account.administrators.ids + User.where(id: user_ids).pluck(:pubsub_token).compact + end + def broadcast(call, event, streams: account_streams, **extra) payload = { event: event, data: base_payload(call).merge(extra) } streams.each { |s| ActionCable.server.broadcast(s, payload) } diff --git a/spec/enterprise/controllers/api/v1/accounts/audit_logs_controller_spec.rb b/spec/enterprise/controllers/api/v1/accounts/audit_logs_controller_spec.rb index 8cf28f93a..7c6c69ec6 100644 --- a/spec/enterprise/controllers/api/v1/accounts/audit_logs_controller_spec.rb +++ b/spec/enterprise/controllers/api/v1/accounts/audit_logs_controller_spec.rb @@ -51,10 +51,13 @@ RSpec.describe 'Enterprise Audit API', type: :request do expect(json_response['audit_logs'][1]['action']).to eql('create') expect(json_response['audit_logs'][1]['audited_changes']['name']).to eql(inbox.name) expect(json_response['audit_logs'][1]['associated_id']).to eql(account.id) - expect(json_response['current_page']).to be(1) # contains audit log for account user as well # contains audit logs for account update(enable audit logs) - expect(json_response['total_entries']).to be(3) + expect(json_response.slice('current_page', 'per_page', 'total_entries')).to eql( + 'current_page' => 1, + 'per_page' => 25, + 'total_entries' => 3 + ) end end end diff --git a/spec/enterprise/controllers/api/v1/accounts/captain/bulk_actions_controller_spec.rb b/spec/enterprise/controllers/api/v1/accounts/captain/bulk_actions_controller_spec.rb index 744fa4ea8..31a0bc36f 100644 --- a/spec/enterprise/controllers/api/v1/accounts/captain/bulk_actions_controller_spec.rb +++ b/spec/enterprise/controllers/api/v1/accounts/captain/bulk_actions_controller_spec.rb @@ -147,7 +147,7 @@ RSpec.describe 'Api::V1::Accounts::Captain::BulkActions', type: :request do params: sync_params, headers: admin.create_new_auth_token, as: :json - end.to have_enqueued_job(Captain::Documents::PerformSyncJob).exactly(documents.size).times + end.to have_enqueued_job(Captain::Documents::PerformSyncJob).on_queue('low').exactly(documents.size).times documents.each do |document| expect(document.reload).to have_attributes( @@ -190,7 +190,7 @@ RSpec.describe 'Api::V1::Accounts::Captain::BulkActions', type: :request do expect(response).to have_http_status(:ok) end - it 'skips documents that already have a sync in progress' do + it 'queues documents that already have a sync in progress' do syncing_document = create(:captain_document, assistant: assistant, account: account, status: :available) syncing_document.update!(sync_status: :syncing, last_sync_attempted_at: 1.minute.ago) @@ -199,9 +199,10 @@ RSpec.describe 'Api::V1::Accounts::Captain::BulkActions', type: :request do params: sync_params.merge(ids: [syncing_document.id]), headers: admin.create_new_auth_token, as: :json - end.not_to have_enqueued_job(Captain::Documents::PerformSyncJob) + end.to have_enqueued_job(Captain::Documents::PerformSyncJob).with(syncing_document).on_queue('low') expect(response).to have_http_status(:ok) + expect(json_response).to eq({ ids: [syncing_document.id], count: 1 }) end it 'queues stale syncing documents again' do diff --git a/spec/enterprise/controllers/api/v1/accounts/captain/documents_controller_spec.rb b/spec/enterprise/controllers/api/v1/accounts/captain/documents_controller_spec.rb index cb7731652..77cb25f49 100644 --- a/spec/enterprise/controllers/api/v1/accounts/captain/documents_controller_spec.rb +++ b/spec/enterprise/controllers/api/v1/accounts/captain/documents_controller_spec.rb @@ -243,7 +243,7 @@ RSpec.describe 'Api::V1::Accounts::Captain::Documents', type: :request do before do create_list(:captain_document, 5, assistant: assistant, account: account) - create(:installation_config, name: 'CAPTAIN_CLOUD_PLAN_LIMITS', value: captain_limits.to_json) + InstallationConfig.find_or_initialize_by(name: 'CAPTAIN_CLOUD_PLAN_LIMITS').update!(value: captain_limits.to_json) post "/api/v1/accounts/#{account.id}/captain/documents", params: valid_attributes, headers: admin.create_new_auth_token @@ -281,7 +281,7 @@ RSpec.describe 'Api::V1::Accounts::Captain::Documents', type: :request do expect do post "/api/v1/accounts/#{account.id}/captain/documents/#{document.id}/sync", headers: admin.create_new_auth_token, as: :json - end.to have_enqueued_job(Captain::Documents::PerformSyncJob).with(document) + end.to have_enqueued_job(Captain::Documents::PerformSyncJob).with(document).on_queue('low') expect(document.reload).to have_attributes( sync_status: 'syncing', @@ -292,15 +292,15 @@ RSpec.describe 'Api::V1::Accounts::Captain::Documents', type: :request do expect(response).to have_http_status(:accepted) end - it 'rejects documents that already have a sync in progress' do + it 'queues documents that already have a sync in progress' do document.update!(sync_status: :syncing, last_sync_attempted_at: 1.minute.ago) expect do post "/api/v1/accounts/#{account.id}/captain/documents/#{document.id}/sync", headers: admin.create_new_auth_token, as: :json - end.not_to have_enqueued_job(Captain::Documents::PerformSyncJob) + end.to have_enqueued_job(Captain::Documents::PerformSyncJob).with(document).on_queue('low') - expect(response).to have_http_status(:unprocessable_entity) + expect(response).to have_http_status(:accepted) end it 'queues stale syncing documents again' do diff --git a/spec/enterprise/jobs/captain/documents/perform_sync_job_spec.rb b/spec/enterprise/jobs/captain/documents/perform_sync_job_spec.rb new file mode 100644 index 000000000..dfec1c9d7 --- /dev/null +++ b/spec/enterprise/jobs/captain/documents/perform_sync_job_spec.rb @@ -0,0 +1,60 @@ +require 'rails_helper' + +RSpec.describe Captain::Documents::PerformSyncJob, type: :job do + let(:account) { create(:account) } + let(:assistant) { create(:captain_assistant, account: account) } + let(:document) { create(:captain_document, assistant: assistant, account: account, status: :available) } + + def stub_lock(job) + allow(job).to receive(:with_lock).and_yield + end + + def stub_page_fetch(content: 'Updated content') + fetch_result = Captain::Documents::SinglePageFetcher::Result.new( + success: true, + title: 'Updated title', + content: content + ) + fetcher = instance_double(Captain::Documents::SinglePageFetcher, fetch: fetch_result) + allow(Captain::Documents::SinglePageFetcher).to receive(:new).and_return(fetcher) + end + + def stub_page_fetch_failure + fetcher = instance_double(Captain::Documents::SinglePageFetcher) + allow(fetcher).to receive(:fetch).and_raise(StandardError, 'boom') + allow(Captain::Documents::SinglePageFetcher).to receive(:new).and_return(fetcher) + end + + it 'syncs the document content' do + travel_to Time.zone.local(2026, 5, 18, 10, 0, 0) do + job = described_class.new + stub_lock(job) + stub_page_fetch + + job.perform(document) + + expect(document.reload).to have_attributes( + sync_status: 'synced', + last_sync_attempted_at: Time.current, + last_synced_at: Time.current, + content: 'Updated content' + ) + end + end + + it 'marks unexpected failures as failed' do + travel_to Time.zone.local(2026, 5, 18, 10, 0, 0) do + job = described_class.new + stub_lock(job) + stub_page_fetch_failure + + expect { job.perform(document) }.to raise_error(StandardError, 'boom') + + expect(document.reload).to have_attributes( + sync_status: 'failed', + last_sync_error_code: 'sync_error', + last_sync_attempted_at: Time.current + ) + end + end +end diff --git a/spec/enterprise/jobs/captain/documents/schedule_syncs_job_spec.rb b/spec/enterprise/jobs/captain/documents/schedule_syncs_job_spec.rb index 76604a0b5..3a241d0a1 100644 --- a/spec/enterprise/jobs/captain/documents/schedule_syncs_job_spec.rb +++ b/spec/enterprise/jobs/captain/documents/schedule_syncs_job_spec.rb @@ -5,11 +5,28 @@ RSpec.describe Captain::Documents::ScheduleSyncsJob, type: :job do let(:assistant) { create(:captain_assistant, account: account) } before do - create(:installation_config, name: 'CAPTAIN_DOCUMENT_AUTO_SYNC_INTERVALS', value: { business: 24, hacker: nil }.to_json) + set_installation_config('CAPTAIN_DOCUMENT_AUTO_SYNC_INTERVALS', { business: 168, enterprise: 24, startups: 720, hacker: nil }.to_json) + set_installation_config('CAPTAIN_DOCUMENT_AUTO_SYNC_PER_ACCOUNT_BATCH_LIMIT', 50) + set_installation_config('CAPTAIN_DOCUMENT_AUTO_SYNC_GLOBAL_BATCH_LIMIT', 1000) account.enable_features!('captain_document_auto_sync') clear_enqueued_jobs end + def set_installation_config(name, value) + InstallationConfig.find_or_initialize_by(name: name).tap do |config| + config.value = value + config.save! + end + end + + def update_sync_limit(name, value) + InstallationConfig.find_by!(name: name).update!(value: value) + end + + def sync_job_for(document) + have_enqueued_job(Captain::Documents::PerformSyncJob).with(document) + end + context 'when the account has not enabled auto-sync' do before { account.disable_features!('captain_document_auto_sync') } @@ -32,6 +49,25 @@ RSpec.describe Captain::Documents::ScheduleSyncsJob, type: :job do end end + context 'when a plan name is passed' do + it 'queues due documents only for that plan' do + enterprise_account = create(:account, custom_attributes: { plan_name: 'Enterprise' }) + enterprise_account.enable_features!('captain_document_auto_sync') + enterprise_assistant = create(:captain_assistant, account: enterprise_account) + business_document = create(:captain_document, assistant: assistant, account: account, status: :available) + enterprise_document = create(:captain_document, assistant: enterprise_assistant, account: enterprise_account, status: :available) + + business_document.update!(sync_status: :synced, last_synced_at: 3.days.ago, last_sync_attempted_at: 3.days.ago) + enterprise_document.update!(sync_status: :synced, last_synced_at: 3.days.ago, last_sync_attempted_at: 3.days.ago) + clear_enqueued_jobs + + described_class.new.perform('enterprise') + + expect(Captain::Documents::PerformSyncJob).to have_been_enqueued.with(enterprise_document) + expect(Captain::Documents::PerformSyncJob).not_to have_been_enqueued.with(business_document) + end + end + context 'when an available document has backfilled sync metadata' do it 'leaves it alone when last synced within the plan cadence' do create( @@ -54,53 +90,34 @@ RSpec.describe Captain::Documents::ScheduleSyncsJob, type: :job do account: account, status: :available, sync_status: :synced, - last_synced_at: 3.days.ago + last_synced_at: 8.days.ago ) clear_enqueued_jobs expect { described_class.new.perform } - .to have_enqueued_job(Captain::Documents::PerformSyncJob).with(document) + .to sync_job_for(document).on_queue('purgable') end - it 'marks the due document as syncing before queueing' do + it 'delays only the queued sync job' do travel_to Time.zone.local(2026, 4, 27, 10, 0, 0) do + job = described_class.new + allow(job).to receive(:rand).and_return(30.minutes.to_i) document = create( :captain_document, assistant: assistant, account: account, status: :available, sync_status: :synced, - last_synced_at: 3.days.ago + last_synced_at: 8.days.ago ) clear_enqueued_jobs - described_class.new.perform + job.perform - expect(document.reload).to have_attributes( - sync_status: 'syncing', - last_sync_attempted_at: Time.current - ) + expect(Captain::Documents::PerformSyncJob) + .to have_been_enqueued.with(document).at(30.minutes.from_now) end end - - it 'does not queue the same document again while the reserved sync is fresh' do - document = create( - :captain_document, - assistant: assistant, - account: account, - status: :available, - sync_status: :synced, - last_synced_at: 2.days.ago - ) - clear_enqueued_jobs - - expect { described_class.new.perform } - .to have_enqueued_job(Captain::Documents::PerformSyncJob).with(document) - - clear_enqueued_jobs - - expect { described_class.new.perform }.not_to have_enqueued_job(Captain::Documents::PerformSyncJob) - end end context 'when an available document was synced within the plan cadence' do @@ -116,78 +133,59 @@ RSpec.describe Captain::Documents::ScheduleSyncsJob, type: :job do context 'when an available document was last synced before the plan cadence' do it 'queues a sync for that document' do document = create(:captain_document, assistant: assistant, account: account, status: :available) - document.update!(sync_status: :synced, last_synced_at: 2.days.ago, last_sync_attempted_at: 2.days.ago) + document.update!(sync_status: :synced, last_synced_at: 8.days.ago, last_sync_attempted_at: 8.days.ago) clear_enqueued_jobs expect { described_class.new.perform } - .to have_enqueued_job(Captain::Documents::PerformSyncJob).with(document) + .to sync_job_for(document) + end + end + + context 'when jitter spreads queued sync execution' do + it 'uses a widened due window so jittered syncs do not skip the next plan run' do + travel_to Time.zone.local(2026, 4, 27, 10, 0, 0) do + job = described_class.new + interval = 1.week + due_window = (interval.to_i / 2).seconds + allow(job).to receive(:rand).and_return(2.hours.to_i) + document = create(:captain_document, assistant: assistant, account: account, status: :available) + + document.update!(sync_status: :synced, last_synced_at: (due_window - 1.minute).ago) + clear_enqueued_jobs + + expect { job.perform }.not_to have_enqueued_job(Captain::Documents::PerformSyncJob) + + document.update!(sync_status: :synced, last_synced_at: (due_window + 1.minute).ago) + clear_enqueued_jobs + + expect { job.perform } + .to have_enqueued_job(Captain::Documents::PerformSyncJob) + .with(document) + .on_queue('purgable') + .at(2.hours.from_now) + end end - it 'skips invalid legacy documents without counting them against the account cap' do - stub_const("#{described_class}::PER_ACCOUNT_HOURLY_CAP", 1) - create( - :captain_document, - assistant: assistant, - account: account, - status: :in_progress, - content: nil, - external_link: 'https://example.com' - ) - invalid_document = build( - :captain_document, - assistant: assistant, - account: account, - status: :available, - sync_status: :synced, - last_synced_at: 2.days.ago, - last_sync_attempted_at: 2.days.ago, - external_link: 'https://example.com/' - ) - invalid_document.save!(validate: false) - valid_document = create(:captain_document, assistant: assistant, account: account, status: :available) - valid_document.update!(sync_status: :synced, last_synced_at: 2.days.ago, last_sync_attempted_at: 2.days.ago) - clear_enqueued_jobs + it 'uses a random delay inside the cadence window' do + travel_to Time.zone.local(2026, 4, 27, 10, 0, 0) do + document = create(:captain_document, assistant: assistant, account: account, status: :available) + document.update!(sync_status: :synced, last_synced_at: 8.days.ago) + job = described_class.new + sync_execution_delay = 12_345.seconds - expect { described_class.new.perform }.not_to raise_error - expect(Captain::Documents::PerformSyncJob).not_to have_been_enqueued.with(invalid_document) - expect(Captain::Documents::PerformSyncJob).to have_been_enqueued.with(valid_document) - end + clear_enqueued_jobs + allow(job).to receive(:rand).with(0..described_class::WEEKLY_SYNC_JITTER.to_i).and_return(sync_execution_delay.to_i) - it 'keeps paging due documents when invalid documents fill the first batch' do - stub_const("#{described_class}::PER_ACCOUNT_HOURLY_CAP", 1) - stub_const("#{described_class}::DUE_DOCUMENT_BATCH_SIZE", 1) - create( - :captain_document, - assistant: assistant, - account: account, - status: :in_progress, - content: nil, - external_link: 'https://example.com' - ) - invalid_document = build( - :captain_document, - assistant: assistant, - account: account, - status: :available, - sync_status: :synced, - last_synced_at: 2.days.ago, - last_sync_attempted_at: 3.days.ago, - external_link: 'https://example.com/' - ) - invalid_document.save!(validate: false) - valid_document = create(:captain_document, assistant: assistant, account: account, status: :available) - valid_document.update!(sync_status: :synced, last_synced_at: 2.days.ago, last_sync_attempted_at: 2.days.ago) - clear_enqueued_jobs - - described_class.new.perform - - expect(Captain::Documents::PerformSyncJob).to have_been_enqueued.with(valid_document) + expect { job.perform } + .to sync_job_for(document) + .at(sync_execution_delay.from_now) + end end end context 'when more documents are due than the account cap allows' do before do - stub_const("#{described_class}::PER_ACCOUNT_HOURLY_CAP", 2) + update_sync_limit('CAPTAIN_DOCUMENT_AUTO_SYNC_PER_ACCOUNT_BATCH_LIMIT', 2) end it 'queues backfilled and oldest-attempted documents first' do @@ -195,26 +193,47 @@ RSpec.describe Captain::Documents::ScheduleSyncsJob, type: :job do oldest_document = create(:captain_document, assistant: assistant, account: account, status: :available) backfilled_document = create(:captain_document, assistant: assistant, account: account, status: :available) - newest_document.update!(sync_status: :synced, last_synced_at: 2.days.ago, last_sync_attempted_at: 2.days.ago) - oldest_document.update!(sync_status: :synced, last_synced_at: 3.days.ago, last_sync_attempted_at: 3.days.ago) - backfilled_document.update!(sync_status: :synced, last_synced_at: 4.days.ago, last_sync_attempted_at: nil) + newest_document.update!(sync_status: :synced, last_synced_at: 8.days.ago, last_sync_attempted_at: 8.days.ago) + oldest_document.update!(sync_status: :synced, last_synced_at: 9.days.ago, last_sync_attempted_at: 9.days.ago) + backfilled_document.update!(sync_status: :synced, last_synced_at: 10.days.ago, last_sync_attempted_at: nil) clear_enqueued_jobs expect { described_class.new.perform } - .to have_enqueued_job(Captain::Documents::PerformSyncJob).with(backfilled_document) - .and have_enqueued_job(Captain::Documents::PerformSyncJob).with(oldest_document) + .to sync_job_for(backfilled_document) + .and sync_job_for(oldest_document) expect(Captain::Documents::PerformSyncJob).not_to have_been_enqueued.with(newest_document) end end + context 'when sync caps are configured' do + it 'uses installation config caps for per-account and global limits' do + update_sync_limit('CAPTAIN_DOCUMENT_AUTO_SYNC_PER_ACCOUNT_BATCH_LIMIT', 2) + update_sync_limit('CAPTAIN_DOCUMENT_AUTO_SYNC_GLOBAL_BATCH_LIMIT', 3) + + second_account = create(:account, custom_attributes: { plan_name: 'business' }) + second_account.enable_features!('captain_document_auto_sync') + second_assistant = create(:captain_assistant, account: second_account) + + first_account_documents = create_list(:captain_document, 3, assistant: assistant, account: account, status: :available) + second_account_documents = create_list(:captain_document, 3, assistant: second_assistant, account: second_account, status: :available) + (first_account_documents + second_account_documents).each do |document| + document.update!(sync_status: :synced, last_synced_at: 8.days.ago, last_sync_attempted_at: 8.days.ago) + end + clear_enqueued_jobs + + expect { described_class.new.perform } + .to have_enqueued_job(Captain::Documents::PerformSyncJob).exactly(3).times + end + end + context 'when an available document failed before the plan cadence' do it 'queues a sync for that document' do document = create(:captain_document, assistant: assistant, account: account, status: :available) - document.update!(sync_status: :failed, last_sync_attempted_at: 2.days.ago) + document.update!(sync_status: :failed, last_sync_attempted_at: 8.days.ago) clear_enqueued_jobs expect { described_class.new.perform } - .to have_enqueued_job(Captain::Documents::PerformSyncJob).with(document) + .to sync_job_for(document) end end @@ -228,7 +247,7 @@ RSpec.describe Captain::Documents::ScheduleSyncsJob, type: :job do clear_enqueued_jobs expect { described_class.new.perform } - .to have_enqueued_job(Captain::Documents::PerformSyncJob).with(document) + .to sync_job_for(document) end end diff --git a/spec/enterprise/jobs/enterprise/internal/trigger_daily_scheduled_items_job_spec.rb b/spec/enterprise/jobs/enterprise/internal/trigger_daily_scheduled_items_job_spec.rb new file mode 100644 index 000000000..3afc20328 --- /dev/null +++ b/spec/enterprise/jobs/enterprise/internal/trigger_daily_scheduled_items_job_spec.rb @@ -0,0 +1,42 @@ +require 'rails_helper' + +RSpec.describe Internal::TriggerDailyScheduledItemsJob do + before do + allow(ChatwootHub).to receive(:installation_identifier).and_return('test-installation-id') + allow(Captain::Documents::ScheduleSyncsJob).to receive(:perform_later) + end + + it 'enqueues enterprise Captain document auto-sync every day' do + travel_to Time.zone.parse('2026-05-26 00:00:00 UTC') do + described_class.perform_now + end + + expect(Captain::Documents::ScheduleSyncsJob).to have_received(:perform_later).with('enterprise') + end + + it 'enqueues business Captain document auto-sync weekly' do + travel_to Time.zone.parse('2026-05-24 00:00:00 UTC') do + described_class.perform_now + end + + expect(Captain::Documents::ScheduleSyncsJob).to have_received(:perform_later).with('business') + end + + it 'enqueues startup Captain document auto-sync monthly' do + travel_to Time.zone.parse('2026-06-01 00:00:00 UTC') do + described_class.perform_now + end + + expect(Captain::Documents::ScheduleSyncsJob).to have_received(:perform_later).with('startups') + end + + it 'does not enqueue business or startup Captain document auto-sync before their plan window' do + travel_to Time.zone.parse('2026-05-25 00:00:00 UTC') do + described_class.perform_now + end + + expect(Captain::Documents::ScheduleSyncsJob).to have_received(:perform_later).with('enterprise') + expect(Captain::Documents::ScheduleSyncsJob).not_to have_received(:perform_later).with('business') + expect(Captain::Documents::ScheduleSyncsJob).not_to have_received(:perform_later).with('startups') + end +end diff --git a/spec/enterprise/services/messages/audio_transcription_service_spec.rb b/spec/enterprise/services/messages/audio_transcription_service_spec.rb index 212c0bf01..32752c2b2 100644 --- a/spec/enterprise/services/messages/audio_transcription_service_spec.rb +++ b/spec/enterprise/services/messages/audio_transcription_service_spec.rb @@ -72,7 +72,7 @@ RSpec.describe Messages::AudioTranscriptionService, type: :service do content_type: 'audio/mpeg' ) allow(service).to receive(:can_transcribe?).and_return(true) - allow(attachment.file.blob).to receive(:byte_size).and_return(described_class::WHISPER_BYTE_LIMIT + 1) + allow(attachment.file.blob).to receive(:byte_size).and_return(described_class::TRANSCRIPTION_BYTE_LIMIT + 1) end it 'returns an error without calling Whisper' do diff --git a/spec/enterprise/services/whatsapp/incoming_call_service_spec.rb b/spec/enterprise/services/whatsapp/incoming_call_service_spec.rb index 53c7974f4..53b8ec52f 100644 --- a/spec/enterprise/services/whatsapp/incoming_call_service_spec.rb +++ b/spec/enterprise/services/whatsapp/incoming_call_service_spec.rb @@ -32,6 +32,9 @@ describe Whatsapp::IncomingCallService do describe 'inbound connect' do let(:sdp_offer) { "v=0\r\n...sdp..." } + let!(:agent) { create(:user, account: account) } + + before { create(:inbox_member, inbox: inbox, user: agent) } it 'creates the Call + Conversation + voice_call message and broadcasts voice_call.incoming' do allow(ActionCable.server).to receive(:broadcast) @@ -44,10 +47,16 @@ describe Whatsapp::IncomingCallService do expect(call).to have_attributes(provider: 'whatsapp', direction: 'incoming', status: 'ringing', provider_call_id: provider_call_id) expect(call.meta['sdp_offer']).to eq(sdp_offer) + # No agent is online, so the call falls back to the inbox's agents (and + # account admins) — never the whole-account stream. expect(ActionCable.server).to have_received(:broadcast).with( - "account_#{account.id}", + agent.pubsub_token, hash_including(event: 'voice_call.incoming', data: hash_including(sdp_offer: sdp_offer)) ) + expect(ActionCable.server).not_to have_received(:broadcast).with( + "account_#{account.id}", + hash_including(event: 'voice_call.incoming') + ) end end diff --git a/spec/models/channel/whatsapp_spec.rb b/spec/models/channel/whatsapp_spec.rb index 9afc21f7c..95bc83932 100644 --- a/spec/models/channel/whatsapp_spec.rb +++ b/spec/models/channel/whatsapp_spec.rb @@ -223,12 +223,12 @@ RSpec.describe Channel::Whatsapp do expect(channel.voice_enabled?).to be true end - it 'returns false for whatsapp_cloud channels without embedded_signup source' do + it 'returns true for manual whatsapp_cloud channels with calling_enabled' do channel = create(:channel_whatsapp, account: account, provider: 'whatsapp_cloud', validate_provider_config: false, sync_templates: false) channel.update!(provider_config: channel.provider_config.merge('source' => 'manual', 'calling_enabled' => true)) - expect(channel.voice_enabled?).to be false + expect(channel.voice_enabled?).to be true end it 'returns false for default-provider channels (360dialog) even with calling_enabled' do diff --git a/swagger/paths/application/audit_logs/index.yml b/swagger/paths/application/audit_logs/index.yml index 8fcd22256..df7966c57 100644 --- a/swagger/paths/application/audit_logs/index.yml +++ b/swagger/paths/application/audit_logs/index.yml @@ -24,7 +24,7 @@ responses: per_page: type: integer description: Number of items per page - example: 15 + example: 25 total_entries: type: integer description: Total number of audit log entries @@ -49,4 +49,4 @@ responses: content: application/json: schema: - $ref: '#/components/schemas/bad_request_error' \ No newline at end of file + $ref: '#/components/schemas/bad_request_error' diff --git a/swagger/swagger.json b/swagger/swagger.json index 859453d89..8867bf52a 100644 --- a/swagger/swagger.json +++ b/swagger/swagger.json @@ -1649,7 +1649,7 @@ "per_page": { "type": "integer", "description": "Number of items per page", - "example": 15 + "example": 25 }, "total_entries": { "type": "integer", diff --git a/swagger/tag_groups/application_swagger.json b/swagger/tag_groups/application_swagger.json index 3eb055649..98be4a321 100644 --- a/swagger/tag_groups/application_swagger.json +++ b/swagger/tag_groups/application_swagger.json @@ -192,7 +192,7 @@ "per_page": { "type": "integer", "description": "Number of items per page", - "example": 15 + "example": 25 }, "total_entries": { "type": "integer",