diff --git a/app/services/data_imports/intercom/importer.rb b/app/services/data_imports/intercom/importer.rb index 75794a9a9..882bdcd93 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 @@ -75,6 +78,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) @@ -90,6 +94,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) @@ -203,32 +208,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) @@ -388,19 +396,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) @@ -412,25 +408,28 @@ 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 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? @@ -647,7 +646,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,18 +654,19 @@ 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:) @@ -953,6 +953,40 @@ 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? + + @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 + @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 51876ee76..433d542cd 100644 --- a/spec/services/data_imports/intercom/importer_spec.rb +++ b/spec/services/data_imports/intercom/importer_spec.rb @@ -203,8 +203,11 @@ 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 - described_class.new(data_import: data_import).import_conversations_page + importer.import_conversations_page expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3) end @@ -390,6 +393,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 @@ -545,7 +621,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 @@ -974,6 +1054,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' => { @@ -994,7 +1076,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 @@ -1015,6 +1100,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