From ae57547b840a0f63540248491d94f90ed34bb5a4 Mon Sep 17 00:00:00 2001 From: Sony Mathew <2040199+sony-mathew@users.noreply.github.com> Date: Thu, 23 Jul 2026 14:47:41 +0530 Subject: [PATCH 1/3] fix(imports): reconcile skipped message stats on retry --- .../data_imports/intercom/importer.rb | 7 +++- .../data_imports/intercom/importer_spec.rb | 34 ++++++++++--------- 2 files changed, 24 insertions(+), 17 deletions(-) diff --git a/app/services/data_imports/intercom/importer.rb b/app/services/data_imports/intercom/importer.rb index 0dbda8ff7..38531741f 100644 --- a/app/services/data_imports/intercom/importer.rb +++ b/app/services/data_imports/intercom/importer.rb @@ -703,7 +703,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 + 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 message_mapping_handled?(mapping, part) diff --git a/spec/services/data_imports/intercom/importer_spec.rb b/spec/services/data_imports/intercom/importer_spec.rb index cf54cd9fa..a0e45856d 100644 --- a/spec/services/data_imports/intercom/importer_spec.rb +++ b/spec/services/data_imports/intercom/importer_spec.rb @@ -520,29 +520,31 @@ 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 conversation stats when a superseded run is retried' do + 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 }) - freeze_time do - 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) - next_data_import.update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => next_run_id }) - travel 1.minute - end + 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' - importer.import_conversations_page - expect(next_data_import.reload.stats.dig('conversations', '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) + 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 'recreates messages when existing message mappings point to deleted records', :aggregate_failures do From 4e410b78ff353d3226c12e932545997a3960316d Mon Sep 17 00:00:00 2001 From: Sony Mathew Date: Thu, 23 Jul 2026 22:53:52 +0530 Subject: [PATCH 2/3] perf: reduce Intercom import query pressure (#15052) ## Description Reduces database pressure during large Intercom imports by deferring same-run contact and message statistic recounts until the page boundary, immediately before the cursor advances. Repeated mappings within a page now mark their statistic group dirty instead of running aggregate counts and persisting statistics for every record. Contact and message database attempts now retry once when PostgreSQL cancels a query. The failed transaction is allowed to roll back before the importer waits for a randomized 200-500 ms delay and retries. A recovered timeout creates no error; a second timeout follows the existing item/message failure path exactly once. Provider, validation, and other database errors are not retried. This is Phase 1, Task 3 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 #15051 and should be reviewed as the two-file delta from `codex/cw-7519-intercom-shorter-jobs-heartbeats`. ## Closes - Tracking plan: [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. Retry a page whose contacts or messages already have same-run mappings and confirm imported/skipped statistics are reconciled before the page cursor advances. 2. Raise one `ActiveRecord::QueryCanceled` while persisting a contact mapping and confirm the transaction rolls back, retries, and creates one contact and mapping without an error. 3. Repeat for a message mapping and confirm one message and mapping are created without an error. 4. Raise the timeout on both attempts and confirm one final item/message error is recorded. 5. Raise an unrelated database error and confirm it is not retried. ## 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 | 170 +++++++++++------- .../data_imports/intercom/importer_spec.rb | 136 +++++++++++++- 2 files changed, 238 insertions(+), 68 deletions(-) diff --git a/app/services/data_imports/intercom/importer.rb b/app/services/data_imports/intercom/importer.rb index 38531741f..049508c66 100644 --- a/app/services/data_imports/intercom/importer.rb +++ b/app/services/data_imports/intercom/importer.rb @@ -10,6 +10,8 @@ class DataImports::Intercom::Importer 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 @@ -25,6 +27,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 @@ -77,6 +80,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('contacts', next_cursor) PageResult.new(next_cursor: next_cursor) @@ -92,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) @@ -210,32 +215,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) @@ -395,19 +403,7 @@ class DataImports::Intercom::Importer 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) + import_message(chatwoot_conversation, contact, source_part, message_source_id) end def import_conversation_parts(conversation, chatwoot_conversation, contact) @@ -419,27 +415,30 @@ class DataImports::Intercom::Importer 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, part, message_source_id) end return false if import_stopped? true end + def import_message(conversation, contact, part, message_source_id) + with_query_timeout_retry do + mapping = find_mapping('message', message_source_id) + if mapping && message_mapping_handled?(mapping, part) + if mapping.data_import_id == @data_import.id + reconcile_current_run_message_mapping(conversation, mapping, part) + else + skip_existing_message_mapping(conversation, mapping, part) + end + else + create_message(conversation, contact, part, message_source_id) + end + end + rescue StandardError => e + fail_message(conversation, message_source_id, part, e) + 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? @@ -656,7 +655,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) @@ -664,34 +663,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 @@ -967,6 +969,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/spec/services/data_imports/intercom/importer_spec.rb b/spec/services/data_imports/intercom/importer_spec.rb index a0e45856d..659a038e5 100644 --- a/spec/services/data_imports/intercom/importer_spec.rb +++ b/spec/services/data_imports/intercom/importer_spec.rb @@ -203,9 +203,34 @@ RSpec.describe DataImports::Intercom::Importer do 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 @@ -432,6 +457,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_mapping).and_wrap_original do |method, object_type, source_id, record, metadata:| + if object_type == 'message' && source_id == target_source_id + 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_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_mapping).and_wrap_original do |method, object_type, source_id, record, metadata:| + if object_type == 'message' && source_id == target_source_id + mapping_attempts += 1 + raise ActiveRecord::QueryCanceled, 'statement timeout' + end + + method.call(object_type, source_id, record, metadata: metadata) + 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 @@ -547,6 +645,30 @@ RSpec.describe DataImports::Intercom::Importer do 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') @@ -634,7 +756,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 @@ -1063,6 +1189,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' => { @@ -1083,7 +1211,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 @@ -1104,6 +1235,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 From 9a489d38c725b99d9e4345dcc62dfb5d9f6df903 Mon Sep 17 00:00:00 2001 From: Sony Mathew Date: Thu, 23 Jul 2026 23:13:15 +0530 Subject: [PATCH 3/3] chore: (refactor) prepare Intercom message batches (#15110) ## Description Prepares each Intercom conversation's source message and parts as one deterministic batch before persistence. The builder preserves provider order, prefetches mappings and live messages once, and classifies every entry as current-run, previous-run, stale-mapping repair, existing-message repair, or new. The importer consumes those prepared entries through the existing individual transaction and retry path, so activity conversion, private notes, sender attribution, metadata, timestamps, skip logs, idempotency, indexing, and heartbeat behavior remain unchanged. This PR intentionally does not bulk insert messages; batching writes is Phase 2 Task 2. This is Phase 2, Task 1 (Task 4 overall) of [CW-7615](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion). This PR is stacked on #15052 and should be reviewed as the four-file delta from `codex/cw-7615-intercom-query-timeout-retries`. ## Closes - Tracking plan: [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) - [ ] New feature (non-breaking change which adds functionality) - [ ] Breaking change (fix or feature that would cause existing functionality not to work as expected) - [ ] This change requires a documentation update ## How to test 1. Import a conversation containing a source message, comments, notes, activities, and entries with equal timestamps. 2. Retry the same run and verify messages and counters remain idempotent. 3. Retry after a previous run and verify existing mappings are reported as skips without changing their owner. 4. Delete a mapped Message, rerun, and verify the message and mapping are repaired. 5. Delete only a mapping, rerun, and verify the existing Message is reused without duplication. 6. Rerun a previously skipped part that now classifies as an activity and verify it is repaired. 7. Import an empty conversation and verify no message or mapping prefetch is performed. ## Checklist - [x] My code follows the style guidelines of this project - [x] I have performed a self-review of my code - [x] I have commented on my code, particularly in hard-to-understand areas - [ ] I have made corresponding changes to the documentation - [x] My changes generate no new warnings - [x] I have added tests that prove my fix is effective or that my feature works - [x] New and existing unit tests pass locally with my changes - [ ] Any dependent changes have been merged and published in downstream modules --- .../data_imports/intercom/importer.rb | 134 +++++++------ .../intercom/message_batch_builder.rb | 134 +++++++++++++ .../data_imports/intercom/importer_spec.rb | 49 +++-- .../intercom/message_batch_builder_spec.rb | 178 ++++++++++++++++++ 4 files changed, 429 insertions(+), 66 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 049508c66..fab61478d 100644 --- a/app/services/data_imports/intercom/importer.rb +++ b/app/services/data_imports/intercom/importer.rb @@ -18,7 +18,6 @@ class DataImports::Intercom::Importer 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 @@ -164,8 +163,7 @@ class DataImports::Intercom::Importer if mapped_conversation && mapping.data_import_id != @data_import.id skip_already_imported_item(item, mapping, already_handled: already_handled) reconcile_item_stats('conversation') if already_handled - import_source_message(conversation, mapped_conversation, contact) - return unless import_conversation_parts(conversation, mapped_conversation, contact) + return unless import_conversation_messages(conversation, mapped_conversation, contact) update_conversation_activity(mapped_conversation) return @@ -182,8 +180,7 @@ class DataImports::Intercom::Importer 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 @@ -397,61 +394,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']) - import_message(chatwoot_conversation, contact, source_part, message_source_id) - 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']}" - import_message(chatwoot_conversation, contact, part, message_source_id) + import_message(chatwoot_conversation, contact, entry) end return false if import_stopped? true end - def import_message(conversation, contact, part, message_source_id) + 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) + + 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 - mapping = find_mapping('message', message_source_id) - if mapping && message_mapping_handled?(mapping, part) - if mapping.data_import_id == @data_import.id - reconcile_current_run_message_mapping(conversation, mapping, part) - else - skip_existing_message_mapping(conversation, mapping, part) - end + 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 - create_message(conversation, contact, part, message_source_id) + raise ArgumentError, "Unsupported Intercom message classification: #{entry.classification}" end end rescue StandardError => e - fail_message(conversation, message_source_id, part, e) + fail_message(conversation, entry.source_id, entry.part, e) 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 create_message(conversation, contact, entry) + content = content_for(entry.part) + return record_skipped_message(conversation, entry) if content.blank? - attrs = message_attributes(conversation, contact, part, message_source_id, content) - message = nil + 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) @@ -466,13 +487,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!( @@ -480,15 +500,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']) @@ -529,8 +564,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) @@ -713,12 +747,6 @@ class DataImports::Intercom::Importer end end - def message_mapping_handled?(mapping, part) - return false if mapping.metadata['skipped'] && activity_part?(part) - - mapping.metadata['skipped'] || mapping.chatwoot_record.present? - end - def fail_item(item, error) increment_stat('errors', 'count') item&.update!(status: :failed, last_error_code: error.class.name, last_error_message: error.message) @@ -816,7 +844,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) 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 659a038e5..02e525a1b 100644 --- a/spec/services/data_imports/intercom/importer_spec.rb +++ b/spec/services/data_imports/intercom/importer_spec.rb @@ -367,8 +367,8 @@ RSpec.describe DataImports::Intercom::Importer do 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 - part = args[2] - next unless part['id'] == 'part_2' + 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 @@ -441,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 @@ -486,13 +482,13 @@ RSpec.describe DataImports::Intercom::Importer do mapping_attempts = 0 target_source_id = 'conversation:conversation_1:part:part_1' 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 == 'message' && source_id == target_source_id + 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(object_type, source_id, record, metadata: metadata) + method.call(entry, message) end importer.import_conversations_page @@ -511,13 +507,13 @@ RSpec.describe DataImports::Intercom::Importer do mapping_attempts = 0 target_source_id = 'conversation:conversation_1:part:part_1' 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 == 'message' && source_id == target_source_id + 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(object_type, source_id, record, metadata: metadata) + method.call(entry, message) end importer.import_conversations_page @@ -694,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', 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