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