From 7bdf38732f74a17de47ba44e5c07056557d1e506 Mon Sep 17 00:00:00 2001 From: Sony Mathew <2040199+sony-mathew@users.noreply.github.com> Date: Thu, 23 Jul 2026 00:00:29 +0530 Subject: [PATCH] fix(imports): heartbeat active Intercom runs --- .../data_imports/intercom/importer.rb | 40 ++++++++++++------ .../data_imports/intercom/importer_spec.rb | 42 +++++++++++++++++++ 2 files changed, 70 insertions(+), 12 deletions(-) diff --git a/app/services/data_imports/intercom/importer.rb b/app/services/data_imports/intercom/importer.rb index d7d040d89..1ce663865 100644 --- a/app/services/data_imports/intercom/importer.rb +++ b/app/services/data_imports/intercom/importer.rb @@ -7,6 +7,7 @@ class DataImports::Intercom::Importer end DEFAULT_IMPORT_TYPES = %w[contacts conversations].freeze + HEARTBEAT_INTERVAL = 1.minute PROVIDER = 'intercom'.freeze ALREADY_IMPORTED_ERROR_CODE = 'DataImports::Intercom::AlreadyImported'.freeze SKIPPED_MESSAGE_ERROR_CODE = 'DataImports::Intercom::SkippedMessage'.freeze @@ -154,7 +155,8 @@ class DataImports::Intercom::Importer 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) - import_conversation_parts(conversation, mapped_conversation, contact) + return unless import_conversation_parts(conversation, mapped_conversation, contact) + update_conversation_activity(mapped_conversation) return end @@ -167,14 +169,15 @@ class DataImports::Intercom::Importer increment_stat('conversations', 'imported') unless already_handled import_source_message(conversation, chatwoot_conversation, contact) - import_conversation_parts(conversation, chatwoot_conversation, contact) + return unless import_conversation_parts(conversation, chatwoot_conversation, contact) + update_conversation_activity(chatwoot_conversation) rescue StandardError => e raise if e.is_a?(DataImports::Intercom::Client::Error) fail_item(item, e) ensure - persist_stats + persist_stats unless @import_stopped end def import_stopped? @@ -184,6 +187,14 @@ class DataImports::Intercom::Importer @import_stopped = @data_import.abandoned? || @data_import.completed? || @data_import.completed_with_errors? || stale_import_run? end + def continue_import_with_heartbeat? + return true if @data_import.updated_at > HEARTBEAT_INTERVAL.ago + return false if import_stopped? + + @data_import.touch if @data_import.updated_at <= HEARTBEAT_INTERVAL.ago + true + end + def stale_import_run? active_run_id = @data_import.active_intercom_import_run_id @run_id.present? && active_run_id.present? && active_run_id != @run_id @@ -396,21 +407,26 @@ class DataImports::Intercom::Importer record_truncated_conversation_parts(conversation, parts.size) parts.each do |part| + return false unless continue_import_with_heartbeat? + message_source_id = "conversation:#{source_id_for(conversation)}:part:#{part['id']}" - 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) + 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 - skip_existing_message_mapping(chatwoot_conversation, mapping, part) - next + create_message(chatwoot_conversation, contact, part, message_source_id) + rescue StandardError => e + fail_message(chatwoot_conversation, message_source_id, part, e) end - - create_message(chatwoot_conversation, contact, part, message_source_id) - rescue StandardError => e - fail_message(chatwoot_conversation, message_source_id, part, e) end + true end def create_message(conversation, contact, part, message_source_id) diff --git a/spec/services/data_imports/intercom/importer_spec.rb b/spec/services/data_imports/intercom/importer_spec.rb index 67c2ec576..078e44e8f 100644 --- a/spec/services/data_imports/intercom/importer_spec.rb +++ b/spec/services/data_imports/intercom/importer_spec.rb @@ -285,6 +285,48 @@ RSpec.describe DataImports::Intercom::Importer do expect(data_import.reload.cursor.dig('conversations', 'starting_after')).to be_nil end + it 'heartbeats at most once per minute while processing conversation parts' do + freeze_time do + started_at = Time.current + long_conversation = conversation_payload.deep_dup + template_part = long_conversation.dig('conversation_parts', 'conversation_parts').first + long_conversation['conversation_parts']['conversation_parts'] = Array.new(5) do |index| + template_part.merge('id' => "part_#{index + 1}") + end + allow(client).to receive(:retrieve_conversation).with('conversation_1').and_return(long_conversation) + heartbeat_times = [] + allow(data_import).to receive(:touch).and_wrap_original do |method| + heartbeat_times << Time.current + method.call + end + importer = described_class.new(data_import: data_import) + allow(importer).to receive(:create_message) { travel 30.seconds } + + importer.import_conversations_page + + expect(heartbeat_times).to eq([started_at + 1.minute, started_at + 2.minutes]) + end + end + + it 'stops parts without heartbeating or persisting stale stats when a newer run takes over', :aggregate_failures do + freeze_time do + run_id = 'intercom-run-1' + data_import.update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => run_id }) + allow(data_import).to receive(:touch).and_call_original + 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 => 'new-run' }) + travel 1.minute + end + + importer.import_conversations_page + + expect(importer).to have_received(:create_message).once + expect(data_import).not_to have_received(:touch) + expect(data_import.reload.stats.dig('conversations', 'imported')).to eq(0) + 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:|