From b6efdae2432183bc51c9dedc65cdff28cb265394 Mon Sep 17 00:00:00 2001 From: Sony Mathew Date: Thu, 23 Jul 2026 23:38:01 +0530 Subject: [PATCH] perf: shorten Intercom import jobs (#15051) ## Description Reduces sustained load during large Intercom imports by limiting conversation list pages to 10 while retaining 50-contact pages. Cursor and provider total-count behavior remain unchanged. The one-minute progress heartbeat now lives in parent PR #15050 because stalled-import detection must be safe when that PR is deployed independently. This child PR therefore contains only the smaller-page delta. This is Phase 1, Task 2 of the [Intercom import optimization plan (CW-7615)](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion). This PR is stacked on #15050 and should be reviewed as the two-file delta from `codex/cw-7519-intercom-stalled-retry-15m`. ## Closes - [CW-7615](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion) ## Type of change - [x] Bug fix (non-breaking change which fixes an issue) ## How to test 1. Start an Intercom import containing contacts and conversations across multiple source pages. 2. Confirm contact list requests retain a page size of 50. 3. Confirm conversation list requests use a page size of 10. 4. Confirm page cursors and displayed provider totals continue advancing normally. 5. Confirm the parent PR heartbeat keeps long conversations fresh while these shorter pages are processed. ## Checklist - [x] My code follows the style guidelines of this project - [x] I have performed a self-review of my code - [x] I have added tests that prove the change is effective - [x] New and existing focused tests pass locally with my changes --- .../data_imports/intercom/importer.rb | 307 +++++++++++------- .../intercom/message_batch_builder.rb | 134 ++++++++ .../data_imports/intercom/importer_spec.rb | 281 +++++++++++++++- .../intercom/message_batch_builder_spec.rb | 178 ++++++++++ 4 files changed, 777 insertions(+), 123 deletions(-) create mode 100644 app/services/data_imports/intercom/message_batch_builder.rb create mode 100644 spec/services/data_imports/intercom/message_batch_builder_spec.rb diff --git a/app/services/data_imports/intercom/importer.rb b/app/services/data_imports/intercom/importer.rb index 78a129f05..071e9bdd4 100644 --- a/app/services/data_imports/intercom/importer.rb +++ b/app/services/data_imports/intercom/importer.rb @@ -7,14 +7,17 @@ class DataImports::Intercom::Importer end DEFAULT_IMPORT_TYPES = %w[contacts conversations].freeze + CONTACTS_PER_PAGE = 50 + CONVERSATIONS_PER_PAGE = 10 HEARTBEAT_INTERVAL = 1.minute + QUERY_TIMEOUT_RETRY_LIMIT = 1 + QUERY_TIMEOUT_RETRY_DELAY_RANGE = (0.2..0.5) PROVIDER = 'intercom'.freeze ALREADY_IMPORTED_ERROR_CODE = 'DataImports::Intercom::AlreadyImported'.freeze SKIPPED_MESSAGE_ERROR_CODE = 'DataImports::Intercom::SkippedMessage'.freeze TRUNCATED_PARTS_ERROR_CODE = 'DataImports::Intercom::TruncatedConversationParts'.freeze E164_REGEX = /\A\+[1-9]\d{1,14}\z/ INTERCOM_NUMBER_REGEX = /\A[1-9]\d{1,14}\z/ - REGULAR_MESSAGE_PART_TYPES = %w[comment note source].freeze def initialize(data_import:, run_id: nil) @data_import = data_import @@ -23,6 +26,7 @@ class DataImports::Intercom::Importer @client = DataImports::Intercom::Client.new(access_token: data_import.access_token) @placeholder_inboxes = DataImports::Intercom::PlaceholderInboxBuilder.new(account: @account) @stats = default_stats.deep_merge(data_import.stats || {}) + @dirty_stat_groups = {} end def perform @@ -45,8 +49,9 @@ class DataImports::Intercom::Importer def finish! return if @data_import.reload.abandoned? - has_failures = @data_import.import_errors.non_skip_logs.exists? || @data_import.import_errors.failed.exists? - status = has_failures ? :completed_with_errors : :completed + error_count = @data_import.import_errors.non_skip_logs.count + @data_import.import_errors.failed.count + @stats['errors']['count'] = error_count + status = error_count.positive? ? :completed_with_errors : :completed @data_import.update!( status: status, completed_at: Time.current, @@ -66,7 +71,7 @@ class DataImports::Intercom::Importer end def import_contacts_page(starting_after: cursor_for('contacts')) - response = @client.list_contacts(starting_after: starting_after) + response = @client.list_contacts(starting_after: starting_after, per_page: CONTACTS_PER_PAGE) update_stat_total('contacts', response['total_count']) if response['total_count'].present? Array(response['data'] || response['contacts']).each do |contact| break if import_stopped? @@ -75,13 +80,14 @@ class DataImports::Intercom::Importer end return PageResult.new(next_cursor: nil) if import_stopped? + reconcile_dirty_stats next_cursor = response.dig('pages', 'next', 'starting_after') update_cursor('contacts', next_cursor) PageResult.new(next_cursor: next_cursor) end def import_conversations_page(starting_after: cursor_for('conversations')) - response = @client.list_conversations(starting_after: starting_after) + response = @client.list_conversations(starting_after: starting_after, per_page: CONVERSATIONS_PER_PAGE) update_stat_total('conversations', response['total_count']) if response['total_count'].present? Array(response['data'] || response['conversations']).each do |conversation_summary| break if import_stopped? @@ -90,6 +96,7 @@ class DataImports::Intercom::Importer end return PageResult.new(next_cursor: nil) if import_stopped? + reconcile_dirty_stats next_cursor = response.dig('pages', 'next', 'starting_after') update_cursor('conversations', next_cursor) PageResult.new(next_cursor: next_cursor) @@ -156,8 +163,8 @@ class DataImports::Intercom::Importer mapped_conversation = mapping&.chatwoot_record if mapped_conversation && mapping.data_import_id != @data_import.id skip_already_imported_item(item, mapping, already_handled: already_handled) - import_source_message(conversation, mapped_conversation, contact) - return unless import_conversation_parts(conversation, mapped_conversation, contact) + reconcile_item_stats('conversation') if already_handled + return unless import_conversation_messages(conversation, mapped_conversation, contact) update_conversation_activity(mapped_conversation) return @@ -168,10 +175,13 @@ class DataImports::Intercom::Importer record_mapping('conversation', source_id, chatwoot_conversation, metadata: conversation_metadata(conversation, inbox, source_type)) end item.update!(status: :imported, chatwoot_record_type: 'Conversation', chatwoot_record_id: chatwoot_conversation.id) - increment_stat('conversations', 'imported') unless already_handled + if already_handled + reconcile_item_stats('conversation') + else + increment_stat('conversations', 'imported') + end - import_source_message(conversation, chatwoot_conversation, contact) - return unless import_conversation_parts(conversation, chatwoot_conversation, contact) + return unless import_conversation_messages(conversation, chatwoot_conversation, contact) update_conversation_activity(chatwoot_conversation) rescue StandardError => e @@ -203,32 +213,35 @@ class DataImports::Intercom::Importer end def import_contact(contact_payload, required_for_conversation: false) - source_id = source_id_for(contact_payload) - if source_id.present? && (mapping = find_mapping('contact', source_id)) && (mapped_contact = mapping.chatwoot_record) - return reuse_mapped_contact(contact_payload, source_id, mapping, mapped_contact) - end + item = nil + with_query_timeout_retry do + source_id = source_id_for(contact_payload) + if source_id.present? && (mapping = find_mapping('contact', source_id)) && (mapped_contact = mapping.chatwoot_record) + return reuse_mapped_contact(contact_payload, source_id, mapping, mapped_contact) + end - contact_payload = retrieve_contact_payload(contact_payload) - source_id = source_id_for(contact_payload) - already_handled = item_handled?('contact', source_id) - item = import_item('contact', source_id, contact_payload) - mapping = find_mapping('contact', source_id) + contact_payload = retrieve_contact_payload(contact_payload) + source_id = source_id_for(contact_payload) + already_handled = item_handled?('contact', source_id) + item = import_item('contact', source_id, contact_payload) + mapping = find_mapping('contact', source_id) - mapped_contact = mapping&.chatwoot_record - if mapped_contact && mapping.data_import_id != @data_import.id - skip_already_imported_item(item, mapping, already_handled: already_handled) - return mapped_contact - end + mapped_contact = mapping&.chatwoot_record + if mapped_contact && mapping.data_import_id != @data_import.id + skip_already_imported_item(item, mapping, already_handled: already_handled) + return mapped_contact + end - contact = Contact.transaction do - imported_contact = mapped_contact || find_existing_contact(contact_payload) || create_contact(contact_payload) - update_existing_contact(imported_contact, contact_payload) - record_mapping('contact', source_id, imported_contact, metadata: contact_metadata(contact_payload)) - item.update!(status: :imported, chatwoot_record_type: 'Contact', chatwoot_record_id: imported_contact.id) - imported_contact + contact = Contact.transaction do + imported_contact = mapped_contact || find_existing_contact(contact_payload) || create_contact(contact_payload) + update_existing_contact(imported_contact, contact_payload) + record_mapping('contact', source_id, imported_contact, metadata: contact_metadata(contact_payload)) + item.update!(status: :imported, chatwoot_record_type: 'Contact', chatwoot_record_id: imported_contact.id) + imported_contact + end + increment_stat('contacts', 'imported') unless already_handled + contact end - increment_stat('contacts', 'imported') unless already_handled - contact rescue StandardError => e raise if e.is_a?(DataImports::Intercom::Client::Error) @@ -382,68 +395,85 @@ class DataImports::Intercom::Importer end end - def import_source_message(conversation, chatwoot_conversation, contact) - source = conversation['source'].to_h - return unless source_message_importable?(source) - - message_source_id = "conversation:#{source_id_for(conversation)}:source:#{source['id'].presence || 'initial'}" - source_part = source.merge('part_type' => 'source', 'created_at' => conversation['created_at']) - if (mapping = find_mapping('message', message_source_id)) && message_mapping_handled?(mapping, source_part) - if mapping.data_import_id == @data_import.id - reconcile_current_run_message_mapping(chatwoot_conversation, mapping, source_part) - return - end - - skip_existing_message_mapping(chatwoot_conversation, mapping, source_part) - return - end - - create_message(chatwoot_conversation, contact, source_part, message_source_id) - rescue StandardError => e - fail_message(chatwoot_conversation, message_source_id, source_part, e) - end - - def import_conversation_parts(conversation, chatwoot_conversation, contact) + def import_conversation_messages(conversation, chatwoot_conversation, contact) parts_payload = conversation['conversation_parts'].to_h parts = Array(parts_payload['conversation_parts']) + batch_builder = DataImports::Intercom::MessageBatchBuilder.new( + data_import: @data_import, + conversation: chatwoot_conversation, + source_conversation: conversation + ) + batch = begin + with_query_timeout_retry { batch_builder.perform } + rescue ActiveRecord::QueryCanceled + nil + end + return import_conversation_messages_individually(conversation, chatwoot_conversation, contact, batch_builder, parts.size) if batch.nil? + + batch.source_entries.each { |entry| import_message(chatwoot_conversation, contact, entry) } record_truncated_conversation_parts(conversation, parts.size) - parts.each do |part| + batch.part_entries.each do |entry| return false unless continue_import_with_heartbeat? - message_source_id = "conversation:#{source_id_for(conversation)}:part:#{part['id']}" - begin - if (mapping = find_mapping('message', message_source_id)) && message_mapping_handled?(mapping, part) - if mapping.data_import_id == @data_import.id - reconcile_current_run_message_mapping(chatwoot_conversation, mapping, part) - next - end - - skip_existing_message_mapping(chatwoot_conversation, mapping, part) - next - end - - create_message(chatwoot_conversation, contact, part, message_source_id) - rescue StandardError => e - fail_message(chatwoot_conversation, message_source_id, part, e) - end + import_message(chatwoot_conversation, contact, entry) end + return false if import_stopped? + true end - def create_message(conversation, contact, part, message_source_id) - content = content_for(part) - return record_skipped_message(conversation, message_source_id, part) if content.blank? + def import_conversation_messages_individually(conversation, chatwoot_conversation, contact, batch_builder, parts_count) + source_entries, part_entries = batch_builder.unprepared_entries.partition { |entry| entry[:part]['part_type'] == 'source' } + source_entries.each { |entry| import_unprepared_message(chatwoot_conversation, contact, batch_builder, entry) } + record_truncated_conversation_parts(conversation, parts_count) - attrs = message_attributes(conversation, contact, part, message_source_id, content) - message = nil + part_entries.each do |entry| + return false unless continue_import_with_heartbeat? + + import_unprepared_message(chatwoot_conversation, contact, batch_builder, entry) + end + return false if import_stopped? + + true + end + + def import_unprepared_message(conversation, contact, batch_builder, source_entry) + entry = with_query_timeout_retry { batch_builder.perform([source_entry]).entries.first } + import_message(conversation, contact, entry) + rescue StandardError => e + fail_message(conversation, source_entry[:source_id], source_entry[:part], e) + end + + def import_message(conversation, contact, entry) + with_query_timeout_retry do + case entry.classification + when :current_import + reconcile_current_run_message_mapping(conversation, entry.mapping, entry.part) + when :previous_import + skip_existing_message_mapping(conversation, entry.mapping, entry.part) + when :repairable_stale_mapping, :existing_message, :new_message + create_message(conversation, contact, entry) + else + raise ArgumentError, "Unsupported Intercom message classification: #{entry.classification}" + end + end + rescue StandardError => e + fail_message(conversation, entry.source_id, entry.part, e) + end + + def create_message(conversation, contact, entry) + content = content_for(entry.part) + return record_skipped_message(conversation, entry) if content.blank? + + attrs = message_attributes(conversation, contact, entry.part, entry.source_id, content) + message = entry.message Message.transaction do - message = conversation.messages.find_by(source_id: attrs[:source_id]) unless message result = Message.insert_all!([attrs], returning: %w[id]) message = Message.find(result.rows.first.first) end - record_mapping('message', message_source_id, message, metadata: message_metadata(part)) + record_message_mapping(entry, message) end increment_stat('messages', 'imported') reindex_message_for_search(message) @@ -458,13 +488,12 @@ class DataImports::Intercom::Importer Rails.logger.warn("Intercom import message reindex failed for message #{message.id}: #{e.class} - #{e.message}") end - def record_skipped_message(conversation, message_source_id, part) - mapping = find_mapping('message', message_source_id) - if mapping - already_recorded = skip_log_recorded?('message', message_source_id, SKIPPED_MESSAGE_ERROR_CODE) - record_skipped_message_log(conversation, message_source_id, part) + def record_skipped_message(conversation, entry) + if entry.mapping + already_recorded = skip_log_recorded?('message', entry.source_id, SKIPPED_MESSAGE_ERROR_CODE) + record_skipped_message_log(conversation, entry.source_id, entry.part) increment_stat('messages', 'skipped') unless already_recorded - return mapping.chatwoot_record + return entry.message end DataImportMapping.create!( @@ -472,15 +501,30 @@ class DataImports::Intercom::Importer data_import: @data_import, source_provider: PROVIDER, source_object_type: 'message', - source_object_id: message_source_id, + source_object_id: entry.source_id, chatwoot_record_type: 'Conversation', chatwoot_record_id: conversation.id, - metadata: message_metadata(part).merge(skipped: true, reason: 'blank_or_unsupported_intercom_part') + metadata: message_metadata(entry.part).merge(skipped: true, reason: 'blank_or_unsupported_intercom_part') ) - record_skipped_message_log(conversation, message_source_id, part) + record_skipped_message_log(conversation, entry.source_id, entry.part) increment_stat('messages', 'skipped') end + def record_message_mapping(entry, message) + (entry.mapping || DataImportMapping.new( + account: @account, + source_provider: PROVIDER, + source_object_type: 'message', + source_object_id: entry.source_id + )).tap do |mapping| + mapping.data_import = @data_import + mapping.chatwoot_record_type = 'Message' + mapping.chatwoot_record_id = message.id + mapping.metadata = message_metadata(entry.part) + mapping.save! + end + end + def message_attributes(conversation, contact, part, message_source_id, content) message_type = message_type_for(part) created_at = timestamp_for(part['created_at']) @@ -521,8 +565,7 @@ class DataImports::Intercom::Importer end def activity_part?(part) - part_type = part['part_type'].to_s - part_type.present? && REGULAR_MESSAGE_PART_TYPES.exclude?(part_type) + DataImports::Intercom::MessageBatchBuilder.activity_part?(part) end def message_content(part) @@ -647,7 +690,7 @@ class DataImports::Intercom::Importer ) item = import_item('contact', source_id, contact_payload) unless item&.imported? item.update!(status: :imported, chatwoot_record_type: 'Contact', chatwoot_record_id: mapped_contact.id) - reconcile_item_stats('contact') + mark_stat_group_dirty('contacts') end def reconcile_item_stats(source_object_type) @@ -655,34 +698,37 @@ class DataImports::Intercom::Importer group = stat_group_for(source_object_type) @stats[group]['imported'] = items.imported.count @stats[group]['skipped'] = items.skipped.count - persist_stats end def reconcile_current_run_message_mapping(conversation, mapping, part) record_skipped_message_log(conversation, mapping.source_object_id, part) if mapping.metadata['skipped'] + mark_stat_group_dirty('messages') + end + def reconcile_message_stats mappings = @data_import.mappings.where(source_provider: PROVIDER, source_object_type: 'message') skipped_mappings = mappings.where("metadata ->> 'skipped' = ?", 'true').count message_logs = @data_import.import_errors.where(source_object_type: 'message') @stats['messages']['imported'] = mappings.count - skipped_mappings @stats['messages']['skipped'] = message_logs.where("details ->> 'kind' = ?", 'skipped').count - persist_stats end def skip_already_imported_item(item, mapping, already_handled:) - item.update!( - status: :skipped, - chatwoot_record_type: mapping.chatwoot_record_type, - chatwoot_record_id: mapping.chatwoot_record_id, - last_error_code: ALREADY_IMPORTED_ERROR_CODE, - last_error_message: 'Already imported in a previous import.' - ) - record_already_imported_log( - data_import_item: item, - source_object_type: item.source_object_type, - source_object_id: item.source_object_id, - mapping: mapping - ) + DataImportItem.transaction do + item.update!( + status: :skipped, + chatwoot_record_type: mapping.chatwoot_record_type, + chatwoot_record_id: mapping.chatwoot_record_id, + last_error_code: ALREADY_IMPORTED_ERROR_CODE, + last_error_message: 'Already imported in a previous import.' + ) + record_already_imported_log( + data_import_item: item, + source_object_type: item.source_object_type, + source_object_id: item.source_object_id, + mapping: mapping + ) + end increment_stat(stat_group_for(item.source_object_type), 'skipped') unless already_handled end @@ -694,13 +740,12 @@ class DataImports::Intercom::Importer already_recorded = skip_log_recorded?('message', mapping.source_object_id, ALREADY_IMPORTED_ERROR_CODE) record_already_imported_log(source_object_type: 'message', source_object_id: mapping.source_object_id, mapping: mapping) end - increment_stat('messages', 'skipped') unless already_recorded - end - - def message_mapping_handled?(mapping, part) - return false if mapping.metadata['skipped'] && activity_part?(part) - - mapping.metadata['skipped'] || mapping.chatwoot_record.present? + if already_recorded + message_logs = @data_import.import_errors.where(source_object_type: 'message') + @stats['messages']['skipped'] = message_logs.where("details ->> 'kind' = ?", 'skipped').count + else + increment_stat('messages', 'skipped') + end end def fail_item(item, error) @@ -800,7 +845,7 @@ class DataImports::Intercom::Importer end def source_message_importable?(source) - source['body'].present? || source['subject'].present? || source['attachments'].present? + DataImports::Intercom::MessageBatchBuilder.source_message_importable?(source) end def skipped_message_log_message(part) @@ -953,6 +998,42 @@ class DataImports::Intercom::Importer @stats[group][key] = @stats[group][key].to_i + 1 end + def mark_stat_group_dirty(group) + @dirty_stat_groups[group] = true + end + + def reconcile_dirty_stats + return if @dirty_stat_groups.empty? + + with_query_timeout_retry do + @dirty_stat_groups.each_key do |group| + case group + when 'contacts' + reconcile_item_stats('contact') + when 'messages' + reconcile_message_stats + else + raise ArgumentError, "Unsupported Intercom import stat group: #{group}" + end + end + persist_stats + end + @dirty_stat_groups.clear + end + + def with_query_timeout_retry + retries = 0 + begin + yield + rescue ActiveRecord::QueryCanceled + raise if retries >= QUERY_TIMEOUT_RETRY_LIMIT + + retries += 1 + sleep(rand(QUERY_TIMEOUT_RETRY_DELAY_RANGE)) + retry + end + end + def update_stat_total(group, total) @stats[group] ||= {} @stats[group]['total'] = total.to_i diff --git a/app/services/data_imports/intercom/message_batch_builder.rb b/app/services/data_imports/intercom/message_batch_builder.rb new file mode 100644 index 000000000..852b51c3c --- /dev/null +++ b/app/services/data_imports/intercom/message_batch_builder.rb @@ -0,0 +1,134 @@ +class DataImports::Intercom::MessageBatchBuilder + PROVIDER = 'intercom'.freeze + REGULAR_PART_TYPES = %w[comment note source].freeze + + Entry = Struct.new(:source_id, :part, :position, :mapping, :message, :classification, keyword_init: true) do + def source? + part['part_type'] == 'source' + end + end + + Batch = Struct.new(:items, keyword_init: true) do + def entries + items + end + + def source_entries + items.select(&:source?) + end + + def part_entries + items.reject(&:source?) + end + end + + def self.activity_part?(part) + part_type = part['part_type'].to_s + part_type.present? && REGULAR_PART_TYPES.exclude?(part_type) + end + + def self.source_message_importable?(source) + source['body'].present? || source['subject'].present? || source['attachments'].present? + end + + def initialize(data_import:, conversation:, source_conversation:) + @data_import = data_import + @account = data_import.account + @conversation = conversation + @source_conversation = source_conversation + end + + def perform(source_entries = unprepared_entries) + return Batch.new(items: []) if source_entries.empty? + + mappings = message_mappings(source_entries) + messages = messages_for(source_entries, mappings) + + Batch.new(items: source_entries.map do |source_entry| + build_entry(source_entry, source_entry[:position], mappings, messages) + end) + end + + def unprepared_entries + ordered_source_entries.map.with_index { |entry, position| entry.merge(position: position) } + end + + private + + def ordered_source_entries + entries = [] + source = @source_conversation['source'].to_h + if self.class.source_message_importable?(source) + entries << { + source_id: "conversation:#{source_conversation_id}:source:#{source['id'].presence || 'initial'}", + part: source.merge('part_type' => 'source', 'created_at' => @source_conversation['created_at']) + } + end + + conversation_parts.each do |part| + entries << { source_id: "conversation:#{source_conversation_id}:part:#{part['id']}", part: part } + end + entries + end + + def conversation_parts + Array(@source_conversation.dig('conversation_parts', 'conversation_parts')) + end + + def source_conversation_id + @source_conversation['id'].presence || @source_conversation['external_id'].presence || @source_conversation['email'].presence + end + + def message_mappings(source_entries) + DataImportMapping.where( + account: @account, + source_provider: PROVIDER, + source_object_type: 'message', + source_object_id: source_entries.pluck(:source_id) + ).index_by(&:source_object_id) + end + + def messages_for(source_entries, mappings) + mapped_message_ids = mappings.values.filter_map do |mapping| + mapping.chatwoot_record_id if mapping.chatwoot_record_type == 'Message' + end + chatwoot_source_ids = source_entries.map { |entry| "intercom:#{entry[:source_id]}" } + messages = Message.where(id: mapped_message_ids).or( + Message.where(conversation_id: @conversation.id, source_id: chatwoot_source_ids) + ).to_a + + { + by_id: messages.index_by(&:id), + by_source_id: messages.index_by(&:source_id) + } + end + + def build_entry(source_entry, position, mappings, messages) + source_id = source_entry[:source_id] + mapping = mappings[source_id] + mapped_message = messages[:by_id][mapping.chatwoot_record_id] if mapping&.chatwoot_record_type == 'Message' + existing_message = messages[:by_source_id]["intercom:#{source_id}"] + + Entry.new( + source_id: source_id, + part: source_entry[:part], + position: position, + mapping: mapping, + message: mapped_message || existing_message, + classification: classification_for(mapping, mapped_message, existing_message, source_entry[:part]) + ) + end + + def classification_for(mapping, mapped_message, existing_message, part) + return existing_message.present? ? :existing_message : :new_message if mapping.blank? + return :repairable_stale_mapping unless mapping_handled?(mapping, mapped_message, part) + + mapping.data_import_id == @data_import.id ? :current_import : :previous_import + end + + def mapping_handled?(mapping, mapped_message, part) + return false if mapping.metadata['skipped'] && self.class.activity_part?(part) + + mapping.metadata['skipped'] || mapped_message.present? + end +end diff --git a/spec/services/data_imports/intercom/importer_spec.rb b/spec/services/data_imports/intercom/importer_spec.rb index b1a217701..622a750fd 100644 --- a/spec/services/data_imports/intercom/importer_spec.rb +++ b/spec/services/data_imports/intercom/importer_spec.rb @@ -66,12 +66,12 @@ RSpec.describe DataImports::Intercom::Importer do before do account.enable_features!('data_import') allow(DataImports::Intercom::Client).to receive(:new).with(access_token: 'intercom-token').and_return(client) - allow(client).to receive(:list_contacts).with(starting_after: nil).and_return( + allow(client).to receive(:list_contacts).with(starting_after: nil, per_page: 50).and_return( 'data' => [contact_payload], 'total_count' => 1, 'pages' => { 'next' => nil } ) - allow(client).to receive(:list_conversations).with(starting_after: nil).and_return( + allow(client).to receive(:list_conversations).with(starting_after: nil, per_page: 10).and_return( 'conversations' => [{ 'id' => 'conversation_1' }], 'total_count' => 1, 'pages' => { 'next' => nil } @@ -188,14 +188,49 @@ RSpec.describe DataImports::Intercom::Importer do expect(item.metadata['message_total_contribution']).to eq(3) end + it 'uses smaller conversation pages while retaining the contact page size' do + importer = described_class.new(data_import: data_import) + + importer.import_contacts_page + importer.import_conversations_page + + expect(client).to have_received(:list_contacts).with(starting_after: nil, per_page: 50) + expect(client).to have_received(:list_conversations).with(starting_after: nil, per_page: 10) + end + it 'reconciles imported message stats from same-run mappings on retry' do described_class.new(data_import: data_import).import_conversations_page stats = data_import.reload.stats.deep_dup stats['messages']['imported'] = 0 data_import.update!(stats: stats) + importer = described_class.new(data_import: data_import) + expect(importer).to receive(:reconcile_message_stats).once.ordered.and_call_original + expect(importer).to receive(:update_cursor).with('conversations', nil).once.ordered.and_call_original + importer.import_conversations_page + + expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3) + end + + it 'retries a query timeout during deferred message stats reconciliation', :aggregate_failures do described_class.new(data_import: data_import).import_conversations_page + stats = data_import.reload.stats.deep_dup + stats['messages']['imported'] = 0 + data_import.update!(stats: stats) + importer = described_class.new(data_import: data_import) + reconciliation_attempts = 0 + allow(importer).to receive(:sleep) + allow(importer).to receive(:reconcile_message_stats).and_wrap_original do |method| + reconciliation_attempts += 1 + raise ActiveRecord::QueryCanceled, 'statement timeout' if reconciliation_attempts == 1 + method.call + end + + importer.import_conversations_page + + expect(reconciliation_attempts).to eq(2) + expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3) end @@ -267,7 +302,7 @@ RSpec.describe DataImports::Intercom::Importer do it 'stops an in-flight page when a newer import run takes over', :aggregate_failures do run_id = 'intercom-run-1' data_import.update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => run_id }) - allow(client).to receive(:list_conversations).with(starting_after: nil).and_return( + allow(client).to receive(:list_conversations).with(starting_after: nil, per_page: 10).and_return( 'conversations' => [{ 'id' => 'conversation_1' }, { 'id' => 'conversation_2' }], 'pages' => { 'next' => { 'starting_after' => 'next-conversation-cursor' } } ) @@ -326,6 +361,49 @@ RSpec.describe DataImports::Intercom::Importer do end end + it 'does not persist stale stats when a newer run takes over during the final part', :aggregate_failures do + run_id = 'intercom-run-1' + data_import.update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => run_id }) + importer = described_class.new(data_import: data_import, run_id: run_id) + allow(importer).to receive(:create_message).and_wrap_original do |method, *args| + method.call(*args).tap do + entry = args[2] + next unless entry.part['id'] == 'part_2' + + DataImport.find(data_import.id).update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => 'new-run' }) + end + end + + importer.import_conversations_page + + conversation = account.conversations.find_by!(identifier: 'intercom:conversation_1') + expect(conversation.messages.count).to eq(3) + expect(data_import.reload.stats.dig('conversations', 'imported')).to eq(0) + end + + it 'reconciles conversation stats when a superseded run is retried', :aggregate_failures do + freeze_time do + run_id = 'intercom-run-1' + next_run_id = 'intercom-run-2' + data_import.update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => run_id }) + importer = described_class.new(data_import: data_import, run_id: run_id) + allow(importer).to receive(:create_message) do + data_import.update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => next_run_id }) + travel 1.minute + end + + importer.import_conversations_page + expect(data_import.reload.stats.dig('conversations', 'imported')).to eq(0) + + retry_importer = described_class.new(data_import: data_import, run_id: next_run_id) + retry_importer.import_conversations_page + retry_importer.finish! + + expect(data_import.reload.stats.dig('conversations', 'imported')).to eq(1) + expect(data_import.processed_records).to eq(5) + end + end + it 'rolls back a newly inserted conversation when mapping persistence fails', :aggregate_failures do importer = described_class.new(data_import: data_import) allow(importer).to receive(:record_mapping).and_wrap_original do |method, object_type, source_id, record, metadata:| @@ -363,11 +441,7 @@ RSpec.describe DataImports::Intercom::Importer do it 'rolls back a newly inserted message when mapping persistence fails', :aggregate_failures do importer = described_class.new(data_import: data_import) - allow(importer).to receive(:record_mapping).and_wrap_original do |method, object_type, source_id, record, metadata:| - raise StandardError, 'mapping failed' if object_type == 'message' - - method.call(object_type, source_id, record, metadata: metadata) - end + allow(importer).to receive(:record_message_mapping).and_raise(StandardError, 'mapping failed') importer.import_conversations_page @@ -379,6 +453,79 @@ RSpec.describe DataImports::Intercom::Importer do ) expect(error).to have_attributes(error_code: 'StandardError', message: 'mapping failed') end + + it 'retries a contact query timeout after rolling back the first transaction', :aggregate_failures do + importer = described_class.new(data_import: data_import) + mapping_attempts = 0 + allow(importer).to receive(:sleep) + allow(importer).to receive(:record_mapping).and_wrap_original do |method, object_type, source_id, record, metadata:| + if object_type == 'contact' + mapping_attempts += 1 + raise ActiveRecord::QueryCanceled, 'statement timeout' if mapping_attempts == 1 + end + + method.call(object_type, source_id, record, metadata: metadata) + end + + importer.import_contacts_page + + expect(mapping_attempts).to eq(2) + expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once + expect(account.contacts.where(email: 'customer@example.com').count).to eq(1) + expect(data_import.mappings.where(source_object_type: 'contact', source_object_id: 'contact_1').count).to eq(1) + expect(data_import.import_errors).to be_empty + expect(data_import.reload.stats.dig('contacts', 'imported')).to eq(1) + end + + it 'retries a message query timeout after rolling back the first transaction', :aggregate_failures do + importer = described_class.new(data_import: data_import) + mapping_attempts = 0 + target_source_id = 'conversation:conversation_1:part:part_1' + allow(importer).to receive(:sleep) + allow(importer).to receive(:record_message_mapping).and_wrap_original do |method, entry, message| + if entry.source_id == target_source_id + mapping_attempts += 1 + raise ActiveRecord::QueryCanceled, 'statement timeout' if mapping_attempts == 1 + end + + method.call(entry, message) + end + + importer.import_conversations_page + + conversation = account.conversations.find_by!(identifier: 'intercom:conversation_1') + expect(mapping_attempts).to eq(2) + expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once + expect(conversation.messages.where(source_id: "intercom:#{target_source_id}").count).to eq(1) + expect(data_import.mappings.where(source_object_type: 'message', source_object_id: target_source_id).count).to eq(1) + expect(data_import.import_errors).to be_empty + expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3) + end + + it 'records one message failure only after the query timeout retry is exhausted', :aggregate_failures do + importer = described_class.new(data_import: data_import) + mapping_attempts = 0 + target_source_id = 'conversation:conversation_1:part:part_1' + allow(importer).to receive(:sleep) + allow(importer).to receive(:record_message_mapping).and_wrap_original do |method, entry, message| + if entry.source_id == target_source_id + mapping_attempts += 1 + raise ActiveRecord::QueryCanceled, 'statement timeout' + end + + method.call(entry, message) + end + + importer.import_conversations_page + + conversation = account.conversations.find_by!(identifier: 'intercom:conversation_1') + error = data_import.import_errors.find_by!(source_object_type: 'message', source_object_id: target_source_id) + expect(mapping_attempts).to eq(2) + expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once + expect(conversation.messages.where(source_id: "intercom:#{target_source_id}")).to be_empty + expect(error).to have_attributes(error_code: 'ActiveRecord::QueryCanceled', message: 'statement timeout') + expect(data_import.reload.stats.dig('errors', 'count')).to eq(1) + end end describe '#finish!' do @@ -467,6 +614,57 @@ RSpec.describe DataImports::Intercom::Importer do expect(next_data_import.import_errors.skip_logs.pluck(:details).map { |details| details['reason'] }.uniq).to eq(['already_imported']) end + it 'reconciles skipped stats when a superseded run is retried', :aggregate_failures do + described_class.new(data_import: data_import).perform + run_id = 'intercom-run-1' + next_run_id = 'intercom-run-2' + next_data_import.update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => run_id }) + + importer = described_class.new(data_import: next_data_import, run_id: run_id) + allow(importer).to receive(:skip_existing_message_mapping).and_wrap_original do |method, *args| + method.call(*args) + part = args[2] + next unless part['id'] == 'part_2' + + DataImport.find(next_data_import.id).update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => next_run_id }) + end + + importer.import_conversations_page + expect(next_data_import.reload.stats.dig('messages', 'skipped')).to eq(0) + + retry_importer = described_class.new(data_import: next_data_import, run_id: next_run_id) + retry_importer.import_conversations_page + retry_importer.finish! + + expect(next_data_import.reload.stats.dig('conversations', 'skipped')).to eq(1) + expect(next_data_import.stats.dig('messages', 'skipped')).to eq(3) + expect(next_data_import.total_records).to eq(5) + end + + it 'keeps skipped contact stats when a query timeout is retried after the item update', :aggregate_failures do + described_class.new(data_import: data_import).perform + importer = described_class.new(data_import: next_data_import) + contact_log_attempts = 0 + allow(importer).to receive(:sleep) + allow(importer).to receive(:record_already_imported_log).and_wrap_original do |method, **attributes| + if attributes[:source_object_type] == 'contact' + contact_log_attempts += 1 + raise ActiveRecord::QueryCanceled, 'statement timeout' if contact_log_attempts == 1 + end + + method.call(**attributes) + end + + importer.import_contacts_page + + contact_item = next_data_import.items.find_by!(source_object_type: 'contact', source_object_id: 'contact_1') + expect(contact_log_attempts).to eq(2) + expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once + expect(contact_item).to be_skipped + expect(next_data_import.reload.stats.dig('contacts', 'skipped')).to eq(1) + expect(next_data_import.import_errors.skip_logs.exists?(data_import_item: contact_item)).to be(true) + end + it 'recreates messages when existing message mappings point to deleted records', :aggregate_failures do described_class.new(data_import: data_import).perform conversation = account.conversations.find_by!(identifier: 'intercom:conversation_1') @@ -492,6 +690,33 @@ RSpec.describe DataImports::Intercom::Importer do expect(message_mappings.filter_map(&:chatwoot_record).count).to eq(3) end + it 'repairs a missing mapping without recreating the existing message', :aggregate_failures do + described_class.new(data_import: data_import).import_conversations_page + conversation = account.conversations.find_by!(identifier: 'intercom:conversation_1') + message = conversation.messages.find_by!(source_id: 'intercom:conversation:conversation_1:part:part_1') + DataImportMapping.find_by!( + account: account, + source_provider: 'intercom', + source_object_type: 'message', + source_object_id: 'conversation:conversation_1:part:part_1' + ).destroy! + importer = described_class.new(data_import: data_import) + allow(importer).to receive(:find_mapping).and_call_original + + importer.import_conversations_page + + repaired_mapping = DataImportMapping.find_by!( + account: account, + source_provider: 'intercom', + source_object_type: 'message', + source_object_id: 'conversation:conversation_1:part:part_1' + ) + expect(conversation.messages.where(source_id: message.source_id).count).to eq(1) + expect(repaired_mapping.chatwoot_record).to eq(message) + expect(importer).not_to have_received(:find_mapping).with('message', anything) + expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3) + end + it 'updates conversation activity when a later import adds new messages to the mapped conversation', :aggregate_failures do new_part = { 'id' => 'part_3', @@ -554,7 +779,11 @@ RSpec.describe DataImports::Intercom::Importer do end it 'repairs the item and imported count on retry', :aggregate_failures do - described_class.new(data_import: data_import).import_contacts_page + importer = described_class.new(data_import: data_import) + expect(importer).to receive(:reconcile_item_stats).with('contact').once.ordered.and_call_original + expect(importer).to receive(:update_cursor).with('contacts', nil).once.ordered.and_call_original + + importer.import_contacts_page item = data_import.items.find_by!(source_object_type: 'contact', source_object_id: 'contact_1') expect(item).to be_imported @@ -937,6 +1166,32 @@ RSpec.describe DataImports::Intercom::Importer do expect(data_import.reload).to be_completed_with_errors expect(data_import.stats.dig('errors', 'count')).to eq(1) end + + it 'reconciles error stats when a superseded run is retried', :aggregate_failures do + run_id = 'intercom-run-1' + next_run_id = 'intercom-run-2' + data_import.update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => run_id }) + importer = described_class.new(data_import: data_import, run_id: run_id) + allow(importer).to receive(:import_message).and_wrap_original do |method, *args| + method.call(*args).tap do + entry = args[2] + next unless entry.part['id'] == 'part_2' + + DataImport.find(data_import.id).update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => next_run_id }) + end + end + + importer.import_conversations_page + expect(data_import.reload.stats.dig('errors', 'count')).to eq(0) + + retry_importer = described_class.new(data_import: data_import, run_id: next_run_id) + retry_importer.import_conversations_page + retry_importer.finish! + + expect(data_import.reload).to be_completed_with_errors + expect(data_import.stats.dig('errors', 'count')).to eq(1) + expect(data_import.total_records).to eq(6) + end end context 'when the conversation parts total matches the returned parts' do @@ -983,6 +1238,8 @@ RSpec.describe DataImports::Intercom::Importer do end context 'when a specific Intercom message part fails to persist' do + let(:insert_attempts) { [] } + let(:conversation_payload) do super().deep_merge( 'conversation_parts' => { @@ -1003,7 +1260,10 @@ RSpec.describe DataImports::Intercom::Importer do before do allow(Message).to receive(:insert_all!).and_wrap_original do |method, records, **kwargs| - raise ActiveRecord::StatementInvalid, 'bad message' if records.first[:source_id] == 'intercom:conversation:conversation_1:part:bad_part' + if records.first[:source_id] == 'intercom:conversation:conversation_1:part:bad_part' + insert_attempts << records.first[:source_id] + raise ActiveRecord::StatementInvalid, 'bad message' + end method.call(records, **kwargs) end @@ -1024,6 +1284,7 @@ RSpec.describe DataImports::Intercom::Importer do ) expect(data_import.reload).to be_completed_with_errors expect(data_import.stats.dig('errors', 'count')).to eq(1) + expect(insert_attempts.one?).to be(true) end end end diff --git a/spec/services/data_imports/intercom/message_batch_builder_spec.rb b/spec/services/data_imports/intercom/message_batch_builder_spec.rb new file mode 100644 index 000000000..dc61cbabb --- /dev/null +++ b/spec/services/data_imports/intercom/message_batch_builder_spec.rb @@ -0,0 +1,178 @@ +require 'rails_helper' + +RSpec.describe DataImports::Intercom::MessageBatchBuilder do + let(:account) { create(:account) } + let(:data_import) { create(:data_import, :intercom, account: account) } + let(:conversation) { create(:conversation, account: account) } + let(:source_conversation) do + { + 'id' => 'conversation_1', + 'created_at' => 1_700_000_000, + 'source' => { + 'id' => 'source_1', + 'part_type' => 'conversation', + 'body' => '

Initial message

', + 'created_at' => 1_700_000_000 + }, + 'conversation_parts' => { + 'conversation_parts' => [ + { + 'id' => 'part_1', + 'part_type' => 'comment', + 'body' => '

First reply

', + 'created_at' => 1_700_000_000 + }, + { + 'id' => 'part_2', + 'part_type' => 'note', + 'body' => '

Internal note

', + 'created_at' => 1_700_000_100 + } + ] + } + } + end + let(:builder) do + described_class.new( + data_import: data_import, + conversation: conversation, + source_conversation: source_conversation + ) + end + + it 'preserves source order and prefetches mappings and messages once per batch', :aggregate_failures do + sql_queries = [] + subscriber = lambda do |_name, _start, _finish, _id, payload| + sql_queries << payload unless payload[:name] == 'SCHEMA' + end + batch_builder = builder + + batch = ActiveSupport::Notifications.subscribed(subscriber, 'sql.active_record') { batch_builder.perform } + + expect(batch.entries.map(&:source_id)).to eq( + %w[ + conversation:conversation_1:source:source_1 + conversation:conversation_1:part:part_1 + conversation:conversation_1:part:part_2 + ] + ) + expect(batch.entries.map(&:position)).to eq([0, 1, 2]) + expect(batch.entries.map(&:classification)).to all(eq(:new_message)) + expect(sql_queries.count { |query| query[:name] == 'DataImportMapping Load' }).to eq(1) + message_queries = sql_queries.select { |query| query[:name] == 'Message Load' } + expect(message_queries.size).to eq(1), message_queries.pluck(:sql).join("\n") + end + + it 'returns an empty batch without prefetch queries when the conversation has no messages' do + source_conversation['source'] = nil + source_conversation['conversation_parts']['conversation_parts'] = [] + + expect(DataImportMapping).not_to receive(:where) + expect(Message).not_to receive(:where) + + expect(builder.perform.entries).to be_empty + end + + it 'classifies a live mapping from the current import as already handled' do + message = create( + :message, + account: account, + conversation: conversation, + inbox: conversation.inbox, + source_id: 'intercom:conversation:conversation_1:part:part_1' + ) + DataImportMapping.create!( + account: account, + data_import: data_import, + source_provider: 'intercom', + source_object_type: 'message', + source_object_id: 'conversation:conversation_1:part:part_1', + chatwoot_record_type: 'Message', + chatwoot_record_id: message.id, + metadata: {} + ) + + entry = builder.perform.entries.find { |batch_entry| batch_entry.source_id.end_with?('part:part_1') } + + expect(entry).to have_attributes(classification: :current_import, message: message) + end + + it 'classifies a live mapping from a previous import without changing its owner' do + previous_import = create(:data_import, :intercom, account: account) + message = create( + :message, + account: account, + conversation: conversation, + inbox: conversation.inbox, + source_id: 'intercom:conversation:conversation_1:part:part_1' + ) + mapping = DataImportMapping.create!( + account: account, + data_import: previous_import, + source_provider: 'intercom', + source_object_type: 'message', + source_object_id: 'conversation:conversation_1:part:part_1', + chatwoot_record_type: 'Message', + chatwoot_record_id: message.id, + metadata: {} + ) + + entry = builder.perform.entries.find { |batch_entry| batch_entry.source_id.end_with?('part:part_1') } + + expect(entry).to have_attributes(classification: :previous_import, mapping: mapping, message: message) + expect(mapping.reload.data_import).to eq(previous_import) + end + + it 'classifies a mapping whose message was deleted as repairable' do + mapping = DataImportMapping.create!( + account: account, + data_import: data_import, + source_provider: 'intercom', + source_object_type: 'message', + source_object_id: 'conversation:conversation_1:part:part_1', + chatwoot_record_type: 'Message', + chatwoot_record_id: 0, + metadata: {} + ) + + entry = builder.perform.entries.find { |batch_entry| batch_entry.source_id.end_with?('part:part_1') } + + expect(entry).to have_attributes(classification: :repairable_stale_mapping, mapping: mapping, message: nil) + end + + it 'classifies an existing conversation message without a mapping for repair' do + message = create( + :message, + account: account, + conversation: conversation, + inbox: conversation.inbox, + source_id: 'intercom:conversation:conversation_1:part:part_1' + ) + + entry = builder.perform.entries.find { |batch_entry| batch_entry.source_id.end_with?('part:part_1') } + + expect(entry).to have_attributes(classification: :existing_message, mapping: nil, message: message) + end + + it 'repairs a skipped mapping when the source part is now an activity' do + source_conversation.dig('conversation_parts', 'conversation_parts').first.merge!( + 'part_type' => 'assignment', + 'body' => nil, + 'assigned_to' => { 'name' => 'Support' } + ) + mapping = DataImportMapping.create!( + account: account, + data_import: data_import, + source_provider: 'intercom', + source_object_type: 'message', + source_object_id: 'conversation:conversation_1:part:part_1', + chatwoot_record_type: 'Conversation', + chatwoot_record_id: conversation.id, + metadata: { skipped: true } + ) + + entry = builder.perform.entries.find { |batch_entry| batch_entry.source_id.end_with?('part:part_1') } + + expect(entry).to have_attributes(classification: :repairable_stale_mapping, mapping: mapping, message: nil) + end +end