diff --git a/app/controllers/api/v1/accounts/data_imports_controller.rb b/app/controllers/api/v1/accounts/data_imports_controller.rb
index 7f0d28e81..70432bc33 100644
--- a/app/controllers/api/v1/accounts/data_imports_controller.rb
+++ b/app/controllers/api/v1/accounts/data_imports_controller.rb
@@ -4,7 +4,7 @@ class Api::V1::Accounts::DataImportsController < Api::V1::Accounts::BaseControll
DATA_IMPORT_FEATURE = 'data_import'.freeze
before_action :ensure_data_import_feature_enabled
- before_action :set_data_import, only: [:show, :start, :abandon, :error_logs, :skip_logs]
+ before_action :set_data_import, only: [:show, :start, :retry_import, :abandon, :error_logs, :skip_logs]
before_action :check_authorization
def index
@@ -59,6 +59,24 @@ class Api::V1::Accounts::DataImportsController < Api::V1::Accounts::BaseControll
render_show
end
+ def retry_import
+ retry_service = DataImports::Intercom::RetryService.new(account: Current.account, data_import: @data_import)
+ retry_result = retry_service.perform
+ @data_import = retry_service.data_import
+
+ case retry_result
+ when :enqueue
+ DataImports::Intercom::ImportJob.perform_later(@data_import, @data_import.active_intercom_import_run_id)
+ render_show
+ when :not_stalled
+ render json: { message: 'This Intercom import is no longer stalled.' }, status: :unprocessable_entity
+ when :active_import_exists
+ render json: { message: 'Another Intercom import is already in progress.' }, status: :unprocessable_entity
+ when :access_token_missing
+ render json: { message: 'The Intercom access key for this import is unavailable.' }, status: :unprocessable_entity
+ end
+ end
+
def abandon
@data_import.abandon!
render_show
diff --git a/app/javascript/dashboard/api/dataImports.js b/app/javascript/dashboard/api/dataImports.js
index b4c15b98a..e6929d420 100644
--- a/app/javascript/dashboard/api/dataImports.js
+++ b/app/javascript/dashboard/api/dataImports.js
@@ -11,6 +11,10 @@ class DataImportsAPI extends ApiClient {
return axios.post(`${this.url}/${id}/start`);
}
+ retry(id) {
+ return axios.post(`${this.url}/${id}/retry`);
+ }
+
abandon(id) {
return axios.post(`${this.url}/${id}/abandon`);
}
diff --git a/app/javascript/dashboard/i18n/locale/en/settings.json b/app/javascript/dashboard/i18n/locale/en/settings.json
index c013a177f..5852b5769 100644
--- a/app/javascript/dashboard/i18n/locale/en/settings.json
+++ b/app/javascript/dashboard/i18n/locale/en/settings.json
@@ -462,6 +462,7 @@
"STATUS": "Status",
"IMPORTED": "Imported",
"CREATED": "Created",
+ "RETRY": "Retry",
"ABANDON": "Abandon"
},
"DETAIL": {
@@ -508,6 +509,8 @@
},
"ALERTS": {
"IMPORT_STARTED": "Intercom import has started.",
+ "IMPORT_RETRIED": "Intercom import has been queued to resume.",
+ "IMPORT_RETRY_FAILED": "Could not retry the Intercom import.",
"IMPORT_ABANDONED": "Intercom import has been abandoned.",
"IMPORT_FAILED": "Could not start the Intercom import."
}
diff --git a/app/javascript/dashboard/routes/dashboard/settings/data/Show.vue b/app/javascript/dashboard/routes/dashboard/settings/data/Show.vue
index ed89c66a9..89f1d651e 100644
--- a/app/javascript/dashboard/routes/dashboard/settings/data/Show.vue
+++ b/app/javascript/dashboard/routes/dashboard/settings/data/Show.vue
@@ -26,6 +26,7 @@ const dataImport = ref(null);
const isLoading = ref(true);
const isRefreshing = ref(false);
const isPolling = ref(false);
+const isRetrying = ref(false);
const isAbandoning = ref(false);
const isDownloadingErrorLogs = ref(false);
const isDownloadingSkipLogs = ref(false);
@@ -35,6 +36,7 @@ const errorsOpen = ref(true);
const skipLogsOpen = ref(true);
let pollTimer;
let isPageActive = false;
+let importRequestVersion = 0;
const hasActiveImport = computed(() => isActiveImport(dataImport.value));
@@ -50,6 +52,7 @@ const fetchImport = async ({
manual = false,
requestedSkipLogsType = selectedSkipLogsType.value,
} = {}) => {
+ const requestVersion = importRequestVersion;
if (showLoader) {
isLoading.value = true;
} else if (manual) {
@@ -60,6 +63,8 @@ const fetchImport = async ({
const response = await DataImportsAPI.show(route.params.dataImportId, {
skip_logs_type: requestedSkipLogsType || undefined,
});
+ if (requestVersion !== importRequestVersion) return;
+
dataImport.value = response.data;
selectedSkipLogsType.value =
response.data.skip_logs_filters?.selected_source_object_type ||
@@ -105,6 +110,13 @@ const refreshImportInBackground = async () => {
}
};
+const startPolling = () => {
+ stopPolling();
+ if (!isPageActive || !hasActiveImport.value) return;
+
+ pollTimer = window.setInterval(refreshImportInBackground, POLL_INTERVAL_MS);
+};
+
const abandonImport = async () => {
isAbandoning.value = true;
try {
@@ -117,6 +129,22 @@ const abandonImport = async () => {
}
};
+const retryImport = async () => {
+ isRetrying.value = true;
+ importRequestVersion += 1;
+ stopPolling();
+ try {
+ const response = await DataImportsAPI.retry(dataImport.value.id);
+ dataImport.value = response.data;
+ useAlert(t('DATA_IMPORTS.ALERTS.IMPORT_RETRIED'));
+ } catch {
+ useAlert(t('DATA_IMPORTS.ALERTS.IMPORT_RETRY_FAILED'));
+ } finally {
+ isRetrying.value = false;
+ if (hasActiveImport.value) startPolling();
+ }
+};
+
const downloadCsv = (response, filename) => {
const url = window.URL.createObjectURL(
new Blob([response.data], { type: 'text/csv' })
@@ -150,13 +178,6 @@ const downloadSkipLogs = async () => {
}
};
-const startPolling = () => {
- stopPolling();
- if (!isPageActive || !hasActiveImport.value) return;
-
- pollTimer = window.setInterval(refreshImportInBackground, POLL_INTERVAL_MS);
-};
-
const handleVisibilityChange = () => {
if (isPageActive && !document.hidden && hasActiveImport.value) {
refreshImportInBackground();
@@ -197,9 +218,11 @@ onBeforeUnmount(() => {
diff --git a/app/javascript/dashboard/routes/dashboard/settings/data/components/ImportDetailHeader.vue b/app/javascript/dashboard/routes/dashboard/settings/data/components/ImportDetailHeader.vue
index 935d63b87..66324a384 100644
--- a/app/javascript/dashboard/routes/dashboard/settings/data/components/ImportDetailHeader.vue
+++ b/app/javascript/dashboard/routes/dashboard/settings/data/components/ImportDetailHeader.vue
@@ -21,6 +21,10 @@ const props = defineProps({
type: Boolean,
default: false,
},
+ isRetrying: {
+ type: Boolean,
+ default: false,
+ },
isAbandoning: {
type: Boolean,
default: false,
@@ -31,7 +35,7 @@ const props = defineProps({
},
});
-defineEmits(['refresh', 'abandon']);
+defineEmits(['refresh', 'retry', 'abandon']);
const { t } = useI18n();
@@ -90,11 +94,23 @@ const canAbandonImport = computed(() => isAbandonableImport(props.dataImport));
:title="$t('DATA_IMPORTS.MONITOR.REFRESH')"
@click="$emit('refresh')"
/>
+
diff --git a/app/javascript/dashboard/routes/dashboard/settings/data/specs/ImportDetailHeader.spec.js b/app/javascript/dashboard/routes/dashboard/settings/data/specs/ImportDetailHeader.spec.js
new file mode 100644
index 000000000..671f58fae
--- /dev/null
+++ b/app/javascript/dashboard/routes/dashboard/settings/data/specs/ImportDetailHeader.spec.js
@@ -0,0 +1,85 @@
+import { mount } from '@vue/test-utils';
+import ImportDetailHeader from '../components/ImportDetailHeader.vue';
+
+vi.mock('vue-i18n', () => ({
+ useI18n: () => ({ t: key => key }),
+}));
+
+const ButtonStub = {
+ name: 'Button',
+ props: {
+ label: { type: String, default: '' },
+ icon: { type: String, default: '' },
+ isLoading: { type: Boolean, default: false },
+ },
+ emits: ['click'],
+ template: `
+
+ `,
+};
+
+const BaseSettingsHeaderStub = {
+ template: `
+
+ `,
+};
+
+const mountHeader = props =>
+ mount(ImportDetailHeader, {
+ props,
+ global: {
+ stubs: {
+ Button: ButtonStub,
+ BaseSettingsHeader: BaseSettingsHeaderStub,
+ },
+ mocks: {
+ $t: key => key,
+ },
+ },
+ });
+
+describe('ImportDetailHeader', () => {
+ const activeImport = {
+ id: 1,
+ name: 'Intercom import',
+ data_type: 'intercom',
+ source_provider: 'intercom',
+ status: 'processing',
+ stalled: true,
+ };
+
+ it('shows Retry between Refresh and Abandon for stalled imports', async () => {
+ const wrapper = mountHeader({ dataImport: activeImport });
+ const buttons = wrapper.findAll('button');
+
+ expect(buttons.map(button => button.attributes('data-label'))).toEqual([
+ '',
+ 'DATA_IMPORTS.TABLE.RETRY',
+ 'DATA_IMPORTS.TABLE.ABANDON',
+ ]);
+
+ await buttons[1].trigger('click');
+
+ expect(wrapper.emitted('retry')).toHaveLength(1);
+ });
+
+ it('hides Retry when the server does not report the import as stalled', () => {
+ const wrapper = mountHeader({
+ dataImport: { ...activeImport, stalled: false },
+ });
+
+ expect(
+ wrapper.find('[data-label="DATA_IMPORTS.TABLE.RETRY"]').exists()
+ ).toBe(false);
+ });
+});
diff --git a/app/javascript/dashboard/routes/dashboard/settings/data/specs/showActions.spec.js b/app/javascript/dashboard/routes/dashboard/settings/data/specs/showActions.spec.js
new file mode 100644
index 000000000..8552e03de
--- /dev/null
+++ b/app/javascript/dashboard/routes/dashboard/settings/data/specs/showActions.spec.js
@@ -0,0 +1,161 @@
+import { flushPromises, mount } from '@vue/test-utils';
+import { KeepAlive, defineComponent, h, nextTick } from 'vue';
+import { useAlert } from 'dashboard/composables';
+import DataImportsAPI from 'dashboard/api/dataImports';
+import Show from '../Show.vue';
+import { POLL_INTERVAL_MS } from '../importStatus';
+
+vi.mock('dashboard/api/dataImports', () => ({
+ default: {
+ show: vi.fn(),
+ retry: vi.fn(),
+ },
+}));
+
+vi.mock('dashboard/composables', () => ({
+ useAlert: vi.fn(),
+}));
+
+vi.mock('vue-i18n', () => ({
+ useI18n: () => ({ t: key => key }),
+}));
+
+vi.mock('vue-router', async importOriginal => ({
+ ...(await importOriginal()),
+ useRoute: () => ({ params: { dataImportId: 1 } }),
+}));
+
+const SettingsLayoutStub = {
+ template: `
+
+
+
+
+ `,
+};
+
+const ImportDetailHeaderStub = {
+ name: 'ImportDetailHeader',
+ props: {
+ dataImport: {
+ type: Object,
+ default: null,
+ },
+ },
+ emits: ['retry'],
+ template: `
+
+ `,
+};
+
+const deferredRequest = () => {
+ let resolve;
+ const promise = new Promise(resolvePromise => {
+ resolve = resolvePromise;
+ });
+ return { promise, resolve };
+};
+
+const mountShow = () => {
+ const Host = defineComponent({
+ render() {
+ return h(KeepAlive, null, { default: () => h(Show) });
+ },
+ });
+
+ return mount(Host, {
+ global: {
+ stubs: {
+ SettingsLayout: SettingsLayoutStub,
+ ImportDetailHeader: ImportDetailHeaderStub,
+ ImportSummaryTiles: true,
+ ImportProgress: true,
+ ImportErrorsSection: true,
+ ImportSkipLogsSection: true,
+ },
+ mocks: {
+ $t: key => key,
+ },
+ },
+ });
+};
+
+describe('data import detail actions', () => {
+ beforeEach(() => {
+ vi.useFakeTimers();
+ DataImportsAPI.show.mockResolvedValue({
+ data: {
+ id: 1,
+ status: 'processing',
+ stalled: true,
+ skip_logs_filters: {},
+ },
+ });
+ });
+
+ afterEach(() => {
+ vi.useRealTimers();
+ vi.clearAllMocks();
+ });
+
+ it('retries a stalled import and replaces the page state', async () => {
+ DataImportsAPI.retry.mockResolvedValue({
+ data: {
+ id: 1,
+ status: 'pending',
+ stalled: false,
+ skip_logs_filters: {},
+ },
+ });
+ const wrapper = mountShow();
+ await nextTick();
+ await flushPromises();
+
+ await wrapper.find('[data-test="retry"]').trigger('click');
+ await flushPromises();
+
+ expect(DataImportsAPI.retry).toHaveBeenCalledWith(1);
+ expect(useAlert).toHaveBeenCalledWith('DATA_IMPORTS.ALERTS.IMPORT_RETRIED');
+ wrapper.unmount();
+ });
+
+ it('ignores an older poll response after retry succeeds', async () => {
+ const pollRequest = deferredRequest();
+ DataImportsAPI.retry.mockResolvedValue({
+ data: {
+ id: 1,
+ status: 'pending',
+ stalled: false,
+ skip_logs_filters: {},
+ },
+ });
+ const wrapper = mountShow();
+ await nextTick();
+ await flushPromises();
+
+ DataImportsAPI.show.mockReturnValueOnce(pollRequest.promise);
+ await vi.advanceTimersByTimeAsync(POLL_INTERVAL_MS);
+ expect(DataImportsAPI.show).toHaveBeenCalledTimes(2);
+
+ await wrapper.find('[data-test="retry"]').trigger('click');
+ await flushPromises();
+ expect(wrapper.find('[data-test="retry"]').exists()).toBe(false);
+
+ pollRequest.resolve({
+ data: {
+ id: 1,
+ status: 'processing',
+ stalled: true,
+ skip_logs_filters: {},
+ },
+ });
+ await flushPromises();
+
+ expect(wrapper.find('[data-test="retry"]').exists()).toBe(false);
+ wrapper.unmount();
+ });
+});
diff --git a/app/models/data_import.rb b/app/models/data_import.rb
index 5f0ad4737..c78c3ce0b 100644
--- a/app/models/data_import.rb
+++ b/app/models/data_import.rb
@@ -33,6 +33,7 @@
#
class DataImport < ApplicationRecord
ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY = 'active_intercom_import_run_id'.freeze
+ INTERCOM_STALLED_AFTER = 15.minutes
LEGACY_DATA_TYPES = ['contacts'].freeze
INTEGRATION_DATA_TYPES = ['intercom'].freeze
IMPORT_TYPES = %w[contacts conversations].freeze
@@ -71,6 +72,10 @@ class DataImport < ApplicationRecord
failed? || abandoned?
end
+ def stalled?
+ intercom_import? && (pending? || processing?) && updated_at <= INTERCOM_STALLED_AFTER.ago
+ end
+
def abandonable?
intercom_import? && (pending? || processing?)
end
diff --git a/app/policies/data_import_policy.rb b/app/policies/data_import_policy.rb
index 3f908ca5d..58bd57722 100644
--- a/app/policies/data_import_policy.rb
+++ b/app/policies/data_import_policy.rb
@@ -19,6 +19,10 @@ class DataImportPolicy < ApplicationPolicy
show?
end
+ def retry_import?
+ show?
+ end
+
def abandon?
show?
end
diff --git a/app/services/data_imports/intercom/importer.rb b/app/services/data_imports/intercom/importer.rb
index d7d040d89..4d40d358c 100644
--- a/app/services/data_imports/intercom/importer.rb
+++ b/app/services/data_imports/intercom/importer.rb
@@ -1,5 +1,7 @@
# rubocop:disable Metrics/ClassLength, Metrics/AbcSize, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity, Metrics/MethodLength, Rails/SkipsModelValidations
class DataImports::Intercom::Importer
+ class InvalidMessagePayloadError < StandardError; end
+
PageResult = Struct.new(:next_cursor, keyword_init: true) do
def done?
next_cursor.blank?
@@ -7,13 +9,29 @@ class DataImports::Intercom::Importer
end
DEFAULT_IMPORT_TYPES = %w[contacts conversations].freeze
+ CONTACTS_PER_PAGE = 50
+ CONVERSATIONS_PER_PAGE = 10
+ MESSAGES_PER_BATCH = 100
+ 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
+ MESSAGE_MAPPING_UNIQUE_INDEX = :idx_data_import_mappings_on_account_and_source
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
+
+ MessageBatchResult = Struct.new(
+ :imported_entries,
+ :skipped_entries,
+ :current_entries,
+ :previous_entries,
+ :messages,
+ :failed_entries,
+ keyword_init: true
+ )
def initialize(data_import:, run_id: nil)
@data_import = data_import
@@ -22,6 +40,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
@@ -42,28 +61,33 @@ class DataImports::Intercom::Importer
end
def finish!
- return if @data_import.reload.abandoned?
+ @data_import.with_lock do
+ next if @data_import.abandoned? || stale_import_run?
- has_failures = @data_import.import_errors.non_skip_logs.exists? || @data_import.import_errors.failed.exists?
- status = has_failures ? :completed_with_errors : :completed
- @data_import.update!(
- status: status,
- completed_at: Time.current,
- stats: @stats,
- total_records: total_processed_records,
- processed_records: total_successful_records
- )
+ 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,
+ stats: @stats,
+ total_records: total_processed_records,
+ processed_records: total_successful_records
+ )
+ end
end
def fail!(error)
- return if @data_import.reload.abandoned?
+ @data_import.with_lock do
+ next if @data_import.abandoned? || stale_import_run?
- record_run_error(error)
- @data_import.update!(status: :failed, last_error_at: Time.current)
+ record_run_error(error)
+ @data_import.update!(status: :failed, last_error_at: Time.current)
+ end
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?
@@ -72,13 +96,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?
@@ -87,6 +112,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)
@@ -153,8 +179,9 @@ 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)
- 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
end
@@ -164,17 +191,21 @@ 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
+
+ return unless import_conversation_messages(conversation, chatwoot_conversation, contact)
- import_source_message(conversation, chatwoot_conversation, contact)
- 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,38 +215,49 @@ 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 false if import_stopped?
+ return true if @data_import.updated_at > HEARTBEAT_INTERVAL.ago
+
+ @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
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)
@@ -369,66 +411,301 @@ 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?
+
record_truncated_conversation_parts(conversation, parts.size)
- parts.each do |part|
- 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)
+ batch.entries.each_slice(MESSAGES_PER_BATCH) do |entries|
+ return false unless continue_import_with_heartbeat?
+
+ import_message_batch(chatwoot_conversation, contact, batch_builder, entries)
+ return false if @import_stopped
+ end
+ return false if import_stopped?
+
+ true
+ end
+
+ def import_message_batch(conversation, contact, batch_builder, entries)
+ result = bulk_message_batch_result(conversation, contact, batch_builder, entries)
+ return if result.blank?
+
+ result.current_entries.each do |entry|
+ reconcile_bulk_message_entry(conversation, entry) do
+ reconcile_current_run_message_mapping(conversation, entry.mapping, entry.part)
+ end
+ end
+ result.previous_entries.each do |entry|
+ reconcile_bulk_message_entry(conversation, entry) do
+ skip_existing_message_mapping(conversation, entry.mapping, entry.part)
+ end
+ end
+ result.skipped_entries.each do |entry|
+ reconcile_bulk_message_entry(conversation, entry) { record_bulk_skipped_message(conversation, entry) }
+ end
+ Array(result.failed_entries).each { |entry, error| fail_message(conversation, entry.source_id, entry.part, error) }
+ increment_stat('messages', 'imported', result.imported_entries.size)
+ result.messages.each { |message| reindex_message_for_search(message) }
+ end
+
+ def bulk_message_batch_result(conversation, contact, batch_builder, entries)
+ with_query_timeout_retry do
+ bulk_write_message_entries(conversation, contact, batch_builder, entries)
+ end
+ rescue ActiveRecord::ActiveRecordError
+ fallback_message_entries(conversation, contact, batch_builder, entries) unless @import_stopped
+ nil
+ end
+
+ def fallback_message_entries(conversation, contact, batch_builder, entries)
+ entries.each do |entry|
+ break unless continue_import_with_heartbeat?
+
+ message = fallback_message_entry(conversation, contact, batch_builder, entry)
+ reindex_message_for_search(message) if message.is_a?(Message)
+ end
+ end
+
+ def fallback_message_entry(conversation, contact, batch_builder, entry)
+ with_query_timeout_retry do
+ @data_import.with_lock do
+ if inactive_import_run?
+ @import_stopped = true
next
end
- skip_existing_message_mapping(chatwoot_conversation, mapping, part)
+ refreshed_entry = batch_builder.refresh([entry]).entries.first
+ import_message(conversation, contact, refreshed_entry, reindex: false)
+ end
+ end
+ rescue ActiveRecord::ActiveRecordError => e
+ fail_message(conversation, entry.source_id, entry.part, e)
+ end
+
+ def bulk_write_message_entries(conversation, contact, batch_builder, entries)
+ @data_import.with_lock do
+ if inactive_import_run?
+ @import_stopped = true
next
end
- create_message(chatwoot_conversation, contact, part, message_source_id)
- rescue StandardError => e
- fail_message(chatwoot_conversation, message_source_id, part, e)
+ refreshed_entries = batch_builder.refresh(entries).entries
+ persist_message_entries(conversation, contact, refreshed_entries)
end
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 inactive_import_run?
+ @data_import.abandoned? || @data_import.failed? || @data_import.completed? ||
+ @data_import.completed_with_errors? || stale_import_run?
+ end
- attrs = message_attributes(conversation, contact, part, message_source_id, content)
- message = nil
- 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))
+ def persist_message_entries(conversation, contact, entries)
+ grouped_entries = entries.group_by(&:classification)
+ writable_entries = entries.select do |entry|
+ %i[repairable_stale_mapping existing_message new_message].include?(entry.classification)
end
+ content_by_source_id = {}
+ attributes_by_source_id = {}
+ failed_entries = []
+ writable_entries.select! do |entry|
+ validate_message_payload!(entry.part)
+ content = content_for(entry.part)
+ content_by_source_id[entry.source_id] = content
+ if content.present? && entry.message.blank?
+ attributes_by_source_id[entry.source_id] = message_attributes(conversation, contact, entry.part, entry.source_id, content)
+ end
+ true
+ rescue InvalidMessagePayloadError => e
+ failed_entries << [entry, e]
+ false
+ end
+ skipped_entries, imported_entries = writable_entries.partition { |entry| content_by_source_id[entry.source_id].blank? }
+ messages = insert_messages(imported_entries, attributes_by_source_id)
+ upsert_message_mappings(conversation, imported_entries, skipped_entries, messages)
+
+ MessageBatchResult.new(
+ imported_entries: imported_entries,
+ skipped_entries: skipped_entries,
+ current_entries: grouped_entries.fetch(:current_import, []),
+ previous_entries: grouped_entries.fetch(:previous_import, []),
+ messages: messages,
+ failed_entries: failed_entries
+ )
+ end
+
+ def insert_messages(entries, attributes_by_source_id)
+ new_entries = entries.reject(&:message)
+ if new_entries.present?
+ attributes = new_entries.map { |entry| attributes_by_source_id.fetch(entry.source_id) }
+ result = Message.insert_all!(attributes, returning: %w[id source_id])
+ inserted_messages = Message.where(id: result.pluck('id')).index_by(&:source_id)
+ end
+
+ entries.map do |entry|
+ entry.message || inserted_messages.fetch("intercom:#{entry.source_id}")
+ end
+ end
+
+ def validate_message_payload!(part)
+ raise InvalidMessagePayloadError, 'Intercom message payload must be an object' unless part.is_a?(Hash)
+
+ %w[author assigned_to event_details].each do |field|
+ value = part[field]
+ next if value.nil? || value.is_a?(Hash)
+
+ raise InvalidMessagePayloadError, "Intercom message #{field} must be an object"
+ end
+
+ participant = part.dig('event_details', 'participant')
+ unless participant.nil? || participant.is_a?(Hash)
+ raise InvalidMessagePayloadError, 'Intercom message event_details.participant must be an object'
+ end
+
+ %w[created_at updated_at].each do |field|
+ value = part[field]
+ valid_timestamp = value.nil? || value.is_a?(Integer) || value.is_a?(Float) ||
+ (value.is_a?(String) && (value.blank? || value.match?(/\A-?\d+(?:\.\d+)?\z/)))
+ raise InvalidMessagePayloadError, "Intercom message #{field} must be a Unix timestamp" unless valid_timestamp
+
+ timestamp_for(value) if value.present?
+ rescue RangeError
+ raise InvalidMessagePayloadError, "Intercom message #{field} must be a Unix timestamp"
+ end
+
+ %w[body subject].each do |field|
+ value = part[field]
+ next unless value.is_a?(String) && !value.valid_encoding?
+
+ raise InvalidMessagePayloadError, "Intercom message #{field} must use valid encoding"
+ end
+ end
+
+ def upsert_message_mappings(conversation, imported_entries, skipped_entries, messages)
+ now = Time.current
+ mapping_attributes = imported_entries.zip(messages).map do |entry, message|
+ message_mapping_attributes(entry, 'Message', message.id, message_metadata(entry.part), now)
+ end
+ mapping_attributes.concat(skipped_entries.filter_map do |entry|
+ next if entry.mapping
+
+ metadata = message_metadata(entry.part).merge(skipped: true, reason: 'blank_or_unsupported_intercom_part')
+ message_mapping_attributes(entry, 'Conversation', conversation.id, metadata, now)
+ end)
+ return if mapping_attributes.empty?
+
+ DataImportMapping.upsert_all(
+ mapping_attributes,
+ unique_by: MESSAGE_MAPPING_UNIQUE_INDEX,
+ update_only: %i[data_import_id chatwoot_record_type chatwoot_record_id metadata updated_at],
+ record_timestamps: false
+ )
+ end
+
+ def message_mapping_attributes(entry, record_type, record_id, metadata, now)
+ {
+ account_id: @account.id,
+ data_import_id: @data_import.id,
+ source_provider: PROVIDER,
+ source_object_type: 'message',
+ source_object_id: entry.source_id,
+ chatwoot_record_type: record_type,
+ chatwoot_record_id: record_id,
+ metadata: metadata,
+ created_at: entry.mapping&.created_at || now,
+ updated_at: now
+ }
+ end
+
+ def record_bulk_skipped_message(conversation, entry)
+ 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
+ end
+
+ def reconcile_bulk_message_entry(conversation, entry, &)
+ with_query_timeout_retry(&)
+ rescue StandardError => e
+ fail_message(conversation, entry.source_id, entry.part, e)
+ end
+
+ 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)
+ with_query_timeout_retry do
+ @data_import.with_lock do
+ if inactive_import_run?
+ @import_stopped = true
+ next
+ end
+
+ entry = batch_builder.perform([source_entry]).entries.first
+ import_message(conversation, contact, entry)
+ end
+ end
+ rescue StandardError => e
+ fail_message(conversation, source_entry[:source_id], source_entry[:part], e)
+ end
+
+ def import_message(conversation, contact, entry, reindex: true)
+ message = with_query_timeout_retry do
+ Message.transaction(requires_new: true) 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
+ end
+ reindex_message_for_search(message) if reindex && message.is_a?(Message)
+ message
+ rescue StandardError => e
+ fail_message(conversation, entry.source_id, entry.part, e)
+ end
+
+ def create_message(conversation, contact, entry)
+ validate_message_payload!(entry.part)
+ 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
+ unless message
+ result = Message.insert_all!([attrs], returning: %w[id])
+ message = Message.find(result.rows.first.first)
+ end
+ record_message_mapping(entry, message)
increment_stat('messages', 'imported')
- reindex_message_for_search(message)
message
end
@@ -440,13 +717,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!(
@@ -454,15 +730,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'])
@@ -503,8 +794,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)
@@ -629,7 +919,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)
@@ -637,34 +927,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
@@ -676,13 +969,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)
@@ -782,7 +1074,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)
@@ -930,9 +1222,45 @@ class DataImports::Intercom::Importer
@import_types ||= (@data_import.import_types.presence || DEFAULT_IMPORT_TYPES)
end
- def increment_stat(group, key)
+ def increment_stat(group, key, amount = 1)
@stats[group] ||= {}
- @stats[group][key] = @stats[group][key].to_i + 1
+ @stats[group][key] = @stats[group][key].to_i + amount
+ 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)
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..b4cb728cc
--- /dev/null
+++ b/app/services/data_imports/intercom/message_batch_builder.rb
@@ -0,0 +1,144 @@
+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)
+ classify(source_entries)
+ end
+
+ def refresh(entries)
+ classify(entries.map do |entry|
+ { source_id: entry.source_id, part: entry.part, position: entry.position }
+ end)
+ end
+
+ def unprepared_entries
+ ordered_source_entries.map.with_index { |entry, position| entry.merge(position: position) }
+ end
+
+ private
+
+ def classify(source_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.with_index do |source_entry, position|
+ build_entry(source_entry, source_entry.fetch(:position, position), mappings, messages)
+ end)
+ end
+
+ 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/app/services/data_imports/intercom/retry_service.rb b/app/services/data_imports/intercom/retry_service.rb
new file mode 100644
index 000000000..57289d917
--- /dev/null
+++ b/app/services/data_imports/intercom/retry_service.rb
@@ -0,0 +1,28 @@
+class DataImports::Intercom::RetryService
+ attr_reader :data_import
+
+ def initialize(account:, data_import:)
+ @account = account
+ @data_import = data_import
+ end
+
+ def perform
+ @account.with_lock do
+ @data_import.with_lock do
+ next :not_stalled unless @data_import.stalled?
+ next :active_import_exists if another_active_import?
+ next :access_token_missing if @data_import.access_token.blank?
+
+ @data_import.assign_active_intercom_import_run_id
+ @data_import.update!(status: :pending)
+ :enqueue
+ end
+ end
+ end
+
+ private
+
+ def another_active_import?
+ @account.data_imports.active_intercom.where.not(id: @data_import.id).exists?
+ end
+end
diff --git a/app/views/api/v1/accounts/data_imports/_data_import.json.jbuilder b/app/views/api/v1/accounts/data_imports/_data_import.json.jbuilder
index ba56ac4d0..46c224fb1 100644
--- a/app/views/api/v1/accounts/data_imports/_data_import.json.jbuilder
+++ b/app/views/api/v1/accounts/data_imports/_data_import.json.jbuilder
@@ -5,6 +5,7 @@ json.source_type data_import.source_type
json.source_provider data_import.source_provider
json.import_types data_import.import_types
json.status data_import.status
+json.stalled data_import.stalled?
json.total_records data_import.total_records
json.processed_records data_import.processed_records
json.stats data_import.stats
diff --git a/config/routes.rb b/config/routes.rb
index fdef3be23..3655561af 100644
--- a/config/routes.rb
+++ b/config/routes.rb
@@ -234,6 +234,7 @@ Rails.application.routes.draw do
end
member do
post :start
+ post :retry, action: :retry_import
post :abandon
get :error_logs
get :skip_logs
diff --git a/spec/models/data_import_spec.rb b/spec/models/data_import_spec.rb
index fc31ebb1d..244d9f8d4 100644
--- a/spec/models/data_import_spec.rb
+++ b/spec/models/data_import_spec.rb
@@ -64,4 +64,36 @@ RSpec.describe DataImport do
expect(data_import.abandoned_at).to be_nil
end
end
+
+ describe '#stalled?' do
+ let(:account) { create(:account) }
+
+ it 'identifies active Intercom imports without updates for fifteen minutes', :aggregate_failures do
+ freeze_time do
+ processing_import = create(:data_import, :intercom, account: account, status: :processing)
+ pending_import = create(:data_import, :intercom, account: account, status: :pending)
+ processing_import.update!(updated_at: 15.minutes.ago)
+ pending_import.update!(updated_at: 15.minutes.ago)
+
+ expect(processing_import.reload).to be_stalled
+ expect(pending_import.reload).to be_stalled
+ end
+ end
+
+ it 'does not identify recent or terminal Intercom imports as stalled', :aggregate_failures do
+ recent_import = create(:data_import, :intercom, account: account, status: :processing)
+ completed_import = create(:data_import, :intercom, account: account, status: :completed)
+ completed_import.update!(updated_at: 1.hour.ago)
+
+ expect(recent_import).not_to be_stalled
+ expect(completed_import.reload).not_to be_stalled
+ end
+
+ it 'does not identify legacy imports as stalled' do
+ legacy_import = create(:data_import, account: account, status: :processing)
+ legacy_import.update!(updated_at: 1.hour.ago)
+
+ expect(legacy_import.reload).not_to be_stalled
+ end
+ end
end
diff --git a/spec/requests/api/v1/accounts/data_imports_spec.rb b/spec/requests/api/v1/accounts/data_imports_spec.rb
index 82a298ef7..4ded7ff4a 100644
--- a/spec/requests/api/v1/accounts/data_imports_spec.rb
+++ b/spec/requests/api/v1/accounts/data_imports_spec.rb
@@ -209,6 +209,63 @@ RSpec.describe 'Data Imports API', type: :request do
end
end
+ describe 'POST /api/v1/accounts/:account_id/data_imports/:id/retry' do
+ let(:data_import) { create(:data_import, :intercom, account: account, status: :processing, started_at: 2.hours.ago) }
+
+ it 'resumes a stalled import with a new run identifier while preserving progress', :aggregate_failures do
+ started_at = data_import.started_at
+ data_import.update!(
+ cursor: { 'conversations' => { 'starting_after' => 'cursor-1' } },
+ stats: { 'conversations' => { 'imported' => 20 } },
+ source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => 'previous-run' },
+ updated_at: 16.minutes.ago
+ )
+
+ expect do
+ post retry_api_v1_account_data_import_url(account_id: account.id, id: data_import.id),
+ headers: admin.create_new_auth_token,
+ as: :json
+ end.to have_enqueued_job(DataImports::Intercom::ImportJob).with(data_import, a_kind_of(String))
+
+ expect(response).to have_http_status(:ok)
+ expect(response.parsed_body).to include('status' => 'pending', 'stalled' => false)
+ expect(data_import.reload.started_at).to eq(started_at)
+ expect(data_import.cursor.dig('conversations', 'starting_after')).to eq('cursor-1')
+ expect(data_import.stats.dig('conversations', 'imported')).to eq(20)
+ expect(data_import.active_intercom_import_run_id).not_to eq('previous-run')
+ end
+
+ it 'rejects duplicate retries after the import becomes active again' do
+ data_import.update!(updated_at: 16.minutes.ago)
+ post retry_api_v1_account_data_import_url(account_id: account.id, id: data_import.id),
+ headers: admin.create_new_auth_token,
+ as: :json
+
+ expect do
+ post retry_api_v1_account_data_import_url(account_id: account.id, id: data_import.id),
+ headers: admin.create_new_auth_token,
+ as: :json
+ end.not_to have_enqueued_job(DataImports::Intercom::ImportJob)
+
+ expect(response).to have_http_status(:unprocessable_entity)
+ expect(response.parsed_body['message']).to eq('This Intercom import is no longer stalled.')
+ end
+
+ it 'rejects retry while another Intercom import is active' do
+ data_import.update!(updated_at: 16.minutes.ago)
+ create(:data_import, :intercom, account: account, status: :processing)
+
+ expect do
+ post retry_api_v1_account_data_import_url(account_id: account.id, id: data_import.id),
+ headers: admin.create_new_auth_token,
+ as: :json
+ end.not_to have_enqueued_job(DataImports::Intercom::ImportJob)
+
+ expect(response).to have_http_status(:unprocessable_entity)
+ expect(response.parsed_body['message']).to eq('Another Intercom import is already in progress.')
+ end
+ end
+
describe 'POST /api/v1/accounts/:account_id/data_imports/:id/abandon' do
let(:data_import) { create(:data_import, :intercom, account: account) }
@@ -283,6 +340,7 @@ RSpec.describe 'Data Imports API', type: :request do
'id' => data_import.id,
'name' => 'July Intercom migration',
'source_provider' => 'intercom',
+ 'stalled' => false,
'import_errors_count' => 1,
'skip_logs_count' => 1
)
diff --git a/spec/services/data_imports/intercom/importer_spec.rb b/spec/services/data_imports/intercom/importer_spec.rb
index 67c2ec576..1ad041315 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 }
@@ -116,6 +116,300 @@ RSpec.describe DataImports::Intercom::Importer do
expect(DataImportMapping.where(data_import: data_import).count).to eq(5)
end
+ it 'writes a normal conversation with one message insert and one mapping upsert', :aggregate_failures do
+ importer = described_class.new(data_import: data_import)
+ message_batch_sizes = []
+ mapping_batch_sizes = []
+ mapping_upsert_options = []
+ allow(Message).to receive(:insert_all!).and_wrap_original do |method, records, **kwargs|
+ message_batch_sizes << records.size
+ method.call(records, **kwargs)
+ end
+ allow(DataImportMapping).to receive(:upsert_all).and_wrap_original do |method, records, **kwargs|
+ mapping_batch_sizes << records.size
+ mapping_upsert_options << kwargs
+ method.call(records, **kwargs)
+ end
+ expect(importer).not_to receive(:create_message)
+
+ importer.perform
+
+ expect(message_batch_sizes).to eq([3])
+ expect(mapping_batch_sizes).to eq([3])
+ expect(mapping_upsert_options).to contain_exactly(
+ include(unique_by: described_class::MESSAGE_MAPPING_UNIQUE_INDEX, record_timestamps: false)
+ )
+ expect(account.messages.count).to eq(3)
+ expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3)
+ end
+
+ it 'preserves provider order when imported messages share a timestamp' do
+ equal_timestamp_conversation = conversation_payload.deep_dup
+ equal_timestamp_conversation['conversation_parts']['conversation_parts'].each do |part|
+ part['created_at'] = conversation_payload['created_at']
+ part['updated_at'] = conversation_payload['created_at']
+ end
+ allow(client).to receive(:retrieve_conversation).with('conversation_1').and_return(equal_timestamp_conversation)
+
+ described_class.new(data_import: data_import).perform
+
+ conversation = account.conversations.find_by!(identifier: 'intercom:conversation_1')
+ expect(conversation.messages.order(:created_at, :id).pluck(:source_id)).to eq(
+ %w[
+ intercom:conversation:conversation_1:source:source_1
+ intercom:conversation:conversation_1:part:part_1
+ intercom:conversation:conversation_1:part:part_2
+ ]
+ )
+ end
+
+ it 'writes conversations in batches of at most 100 messages', :aggregate_failures do
+ bulk_conversation = conversation_payload.deep_dup
+ template_part = bulk_conversation.dig('conversation_parts', 'conversation_parts').first
+ bulk_conversation['conversation_parts']['conversation_parts'] = Array.new(205) do |index|
+ template_part.merge(
+ 'id' => "part_#{index + 1}",
+ 'created_at' => 1_700_000_100 + index,
+ 'updated_at' => 1_700_000_100 + index
+ )
+ end
+ allow(client).to receive(:retrieve_conversation).with('conversation_1').and_return(bulk_conversation)
+ message_batch_sizes = []
+ mapping_batch_sizes = []
+ allow(Message).to receive(:insert_all!).and_wrap_original do |method, records, **kwargs|
+ message_batch_sizes << records.size
+ method.call(records, **kwargs)
+ end
+ allow(DataImportMapping).to receive(:upsert_all).and_wrap_original do |method, records, **kwargs|
+ mapping_batch_sizes << records.size
+ method.call(records, **kwargs)
+ end
+ importer = described_class.new(data_import: data_import)
+ expect(importer).to receive(:update_conversation_activity).once.and_call_original
+
+ importer.import_conversations_page
+
+ expect(message_batch_sizes).to eq([100, 100, 6])
+ expect(mapping_batch_sizes).to eq([100, 100, 6])
+ expect(account.messages.order(:created_at).pick(:source_id)).to eq('intercom:conversation:conversation_1:source:source_1')
+ expect(account.messages.count).to eq(206)
+ expect(data_import.reload.stats.dig('messages', 'imported')).to eq(206)
+ end
+
+ context 'when a bulk message chunk fails' do
+ it 'retries a query timeout once and keeps the successful bulk result', :aggregate_failures do
+ importer = described_class.new(data_import: data_import)
+ mapping_attempts = 0
+ allow(importer).to receive(:sleep)
+ allow(DataImportMapping).to receive(:upsert_all).and_wrap_original do |method, records, **kwargs|
+ mapping_attempts += 1
+ raise ActiveRecord::QueryCanceled, 'statement timeout' if mapping_attempts == 1
+
+ method.call(records, **kwargs)
+ end
+ expect(importer).not_to receive(:fallback_message_entries)
+
+ importer.perform
+
+ expect(mapping_attempts).to eq(2)
+ expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once
+ expect(account.messages.count).to eq(3)
+ expect(data_import.mappings.where(source_object_type: 'message').count).to eq(3)
+ expect(data_import.import_errors).to be_empty
+ expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3)
+ end
+
+ it 'falls back individually after the query timeout retry is exhausted', :aggregate_failures do
+ importer = described_class.new(data_import: data_import)
+ mapping_attempts = 0
+ reindex_transaction_depths = []
+ transaction_depth_before_import = Message.connection.open_transactions
+ allow(importer).to receive(:sleep)
+ allow(importer).to receive(:fallback_message_entries).and_call_original
+ allow(importer).to receive(:create_message).and_call_original
+ allow(importer).to receive(:reindex_message_for_search).and_wrap_original do |method, message|
+ reindex_transaction_depths << Message.connection.open_transactions
+ method.call(message)
+ end
+ allow(DataImportMapping).to receive(:upsert_all) do
+ mapping_attempts += 1
+ raise ActiveRecord::QueryCanceled, 'statement timeout'
+ end
+
+ importer.perform
+
+ expect(mapping_attempts).to eq(2)
+ expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once
+ expect(importer).to have_received(:fallback_message_entries).once
+ expect(importer).to have_received(:create_message).exactly(3).times
+ expect(account.messages.count).to eq(3)
+ expect(data_import.mappings.where(source_object_type: 'message').count).to eq(3)
+ expect(data_import.import_errors).to be_empty
+ expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3)
+ expect(reindex_transaction_depths).to all(eq(transaction_depth_before_import))
+ end
+
+ it 'falls back immediately for a non-timeout database error', :aggregate_failures do
+ importer = described_class.new(data_import: data_import)
+ mapping_attempts = 0
+ allow(importer).to receive(:sleep)
+ allow(importer).to receive(:fallback_message_entries).and_call_original
+ allow(DataImportMapping).to receive(:upsert_all) do
+ mapping_attempts += 1
+ raise ActiveRecord::StatementInvalid, 'bulk mapping failed'
+ end
+
+ importer.perform
+
+ expect(mapping_attempts).to eq(1)
+ expect(importer).not_to have_received(:sleep)
+ expect(importer).to have_received(:fallback_message_entries).once
+ expect(account.messages.count).to eq(3)
+ expect(data_import.import_errors).to be_empty
+ expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3)
+ end
+
+ it 'validates malformed messages before inserting them through the fallback path', :aggregate_failures do
+ malformed_conversation = conversation_payload.deep_dup
+ malformed_conversation['source']['subject'] = nil
+ malformed_conversation['source']['body'] = nil
+ malformed_conversation.dig('conversation_parts', 'conversation_parts').first['created_at'] = { 'unexpected' => true }
+ allow(client).to receive(:retrieve_conversation).with('conversation_1').and_return(malformed_conversation)
+ importer = described_class.new(data_import: data_import)
+ allow(importer).to receive(:bulk_write_message_entries).and_raise(ActiveRecord::StatementInvalid, 'bulk failed')
+
+ importer.perform
+
+ expect(account.messages.pluck(:source_id)).to eq(['intercom:conversation:conversation_1:part:part_2'])
+ error = data_import.import_errors.find_by!(
+ source_object_type: 'message',
+ source_object_id: 'conversation:conversation_1:part:part_1'
+ )
+ expect(error).to have_attributes(
+ error_code: described_class::InvalidMessagePayloadError.name,
+ message: 'Intercom message created_at must be a Unix timestamp'
+ )
+ expect(data_import.reload).to be_completed_with_errors
+ expect(data_import.stats.dig('messages', 'imported')).to eq(1)
+ expect(data_import.stats.dig('errors', 'count')).to eq(1)
+ end
+
+ it 'refreshes fallback entries when another worker repairs a message', :aggregate_failures do
+ importer = described_class.new(data_import: data_import)
+ allow(importer).to receive(:bulk_write_message_entries).and_raise(ActiveRecord::StatementInvalid, 'bulk failed')
+ allow(importer).to receive(:fallback_message_entries).and_wrap_original do |method, conversation, contact, batch_builder, entries|
+ source_entry = entries.first
+ message = create(
+ :message,
+ account: account,
+ inbox: conversation.inbox,
+ conversation: conversation,
+ source_id: "intercom:#{source_entry.source_id}",
+ created_at: Time.zone.at(source_entry.part['created_at']),
+ updated_at: Time.zone.at(source_entry.part['created_at'])
+ )
+ DataImportMapping.create!(
+ account: account,
+ data_import: data_import,
+ source_provider: 'intercom',
+ source_object_type: 'message',
+ source_object_id: source_entry.source_id,
+ chatwoot_record_type: 'Message',
+ chatwoot_record_id: message.id,
+ metadata: {}
+ )
+ method.call(conversation, contact, batch_builder, entries)
+ end
+
+ importer.perform
+
+ source_id = 'intercom:conversation:conversation_1:source:source_1'
+ expect(account.messages.where(source_id: source_id).count).to eq(1)
+ expect(account.messages.count).to eq(3)
+ expect(data_import.mappings.where(source_object_type: 'message').count).to eq(3)
+ expect(data_import.import_errors).to be_empty
+ expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3)
+ end
+
+ it 'retries a fallback refresh timeout after the lock transaction rolls back', :aggregate_failures do
+ importer = described_class.new(data_import: data_import)
+ refresh_attempts = 0
+ allow(importer).to receive(:sleep)
+ allow(importer).to receive(:bulk_write_message_entries).and_raise(ActiveRecord::StatementInvalid, 'bulk failed')
+ allow(DataImports::Intercom::MessageBatchBuilder).to receive(:new).and_wrap_original do |method, **kwargs|
+ method.call(**kwargs).tap do |batch_builder|
+ allow(batch_builder).to receive(:refresh).and_wrap_original do |refresh, entries|
+ refresh_attempts += 1
+ raise ActiveRecord::QueryCanceled, 'statement timeout' if refresh_attempts == 1
+
+ refresh.call(entries)
+ end
+ end
+ end
+
+ importer.perform
+
+ expect(refresh_attempts).to eq(4)
+ expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once
+ expect(account.messages.count).to eq(3)
+ expect(data_import.mappings.where(source_object_type: 'message').count).to eq(3)
+ expect(data_import.import_errors).to be_empty
+ expect(data_import.reload.stats.dig('messages', 'imported')).to eq(3)
+ end
+ end
+
+ it 'isolates a malformed message payload while importing valid messages', :aggregate_failures do
+ malformed_conversation = conversation_payload.deep_dup
+ malformed_conversation['source']['subject'] = nil
+ malformed_conversation['source']['body'] = nil
+ malformed_conversation.dig('conversation_parts', 'conversation_parts').first['created_at'] = { 'unexpected' => true }
+ allow(client).to receive(:retrieve_conversation).with('conversation_1').and_return(malformed_conversation)
+ importer = described_class.new(data_import: data_import)
+ expect(importer).not_to receive(:fallback_message_entries)
+
+ importer.perform
+
+ expect(account.messages.pluck(:source_id)).to eq(['intercom:conversation:conversation_1:part:part_2'])
+ error = data_import.import_errors.find_by!(
+ source_object_type: 'message',
+ source_object_id: 'conversation:conversation_1:part:part_1'
+ )
+ expect(error).to have_attributes(
+ error_code: described_class::InvalidMessagePayloadError.name,
+ message: 'Intercom message created_at must be a Unix timestamp'
+ )
+ expect(data_import.import_errors.where(source_object_type: 'message').count).to eq(1)
+ expect(data_import.reload).to be_completed_with_errors
+ expect(data_import.stats.dig('messages', 'imported')).to eq(1)
+ expect(data_import.stats.dig('errors', 'count')).to eq(1)
+ end
+
+ it 'isolates an out-of-range message timestamp while importing valid messages', :aggregate_failures do
+ malformed_conversation = conversation_payload.deep_dup
+ malformed_conversation['source']['subject'] = nil
+ malformed_conversation['source']['body'] = nil
+ malformed_conversation.dig('conversation_parts', 'conversation_parts').first['created_at'] = Float::INFINITY
+ allow(client).to receive(:retrieve_conversation).with('conversation_1').and_return(malformed_conversation)
+ importer = described_class.new(data_import: data_import)
+ expect(importer).not_to receive(:fallback_message_entries)
+
+ importer.perform
+
+ expect(account.messages.pluck(:source_id)).to eq(['intercom:conversation:conversation_1:part:part_2'])
+ error = data_import.import_errors.find_by!(
+ source_object_type: 'message',
+ source_object_id: 'conversation:conversation_1:part:part_1'
+ )
+ expect(error).to have_attributes(
+ error_code: described_class::InvalidMessagePayloadError.name,
+ message: 'Intercom message created_at must be a Unix timestamp'
+ )
+ expect(data_import.import_errors.where(source_object_type: 'message').count).to eq(1)
+ expect(data_import.reload).to be_completed_with_errors
+ expect(data_import.stats.dig('messages', 'imported')).to eq(1)
+ expect(data_import.stats.dig('errors', 'count')).to eq(1)
+ end
+
it 'imports historical records without dispatching record events or outbound side effects', :aggregate_failures do
dispatched_events = []
allow(Rails.configuration.dispatcher).to receive(:dispatch) do |event_name, *_args|
@@ -188,14 +482,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
@@ -203,13 +532,19 @@ RSpec.describe DataImports::Intercom::Importer do
allow(ChatwootApp).to receive(:advanced_search_allowed?).and_return(true)
allow(ChatwootApp).to receive(:chatwoot_cloud?).and_return(false)
reindexed_message_ids = []
+ reindex_transaction_depths = []
+ transaction_depth_before_import = Message.connection.open_transactions
original_reindex_for_search = Message.instance_method(:reindex_for_search)
- Message.define_method(:reindex_for_search) { reindexed_message_ids << id }
+ Message.define_method(:reindex_for_search) do
+ reindexed_message_ids << id
+ reindex_transaction_depths << self.class.connection.open_transactions
+ end
Message.__send__(:private, :reindex_for_search)
described_class.new(data_import: data_import).perform
expect(reindexed_message_ids).to match_array(Message.where(account_id: account.id).pluck(:id))
+ expect(reindex_transaction_depths).to all(eq(transaction_depth_before_import))
ensure
Message.define_method(:reindex_for_search, original_reindex_for_search)
Message.__send__(:private, :reindex_for_search)
@@ -267,7 +602,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' } }
)
@@ -285,6 +620,193 @@ 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 message batches' 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(400) 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)
+ empty_result = described_class::MessageBatchResult.new(
+ imported_entries: [],
+ skipped_entries: [],
+ current_entries: [],
+ previous_entries: [],
+ messages: []
+ )
+ allow(importer).to receive(:bulk_write_message_entries) do
+ travel 30.seconds
+ empty_result
+ end
+
+ importer.import_conversations_page
+
+ expect(heartbeat_times).to eq([started_at + 1.minute, started_at + 2.minutes])
+ expect(importer).to have_received(:bulk_write_message_entries).exactly(5).times
+ end
+ end
+
+ it 'stops batches 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 })
+ long_conversation = conversation_payload.deep_dup
+ template_part = long_conversation.dig('conversation_parts', 'conversation_parts').first
+ long_conversation['conversation_parts']['conversation_parts'] = Array.new(100) do |index|
+ template_part.merge('id' => "part_#{index + 1}")
+ end
+ allow(client).to receive(:retrieve_conversation).with('conversation_1').and_return(long_conversation)
+ allow(data_import).to receive(:touch).and_call_original
+ importer = described_class.new(data_import: data_import, run_id: run_id)
+ batch_write_count = 0
+ allow(importer).to receive(:bulk_write_message_entries).and_wrap_original do |method, *args|
+ result = method.call(*args)
+ batch_write_count += 1
+ if batch_write_count == 1
+ DataImport.find(data_import.id).update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => 'new-run' })
+ end
+ result
+ end
+
+ importer.import_conversations_page
+
+ expect(importer).to have_received(:bulk_write_message_entries).once
+ expect(data_import).not_to have_received(:touch)
+ expect(data_import.reload.stats.dig('conversations', 'imported')).to eq(0)
+ expect(account.messages.count).to eq(100)
+ end
+ end
+
+ it 'does not persist stale stats when a newer run takes over during the final batch', :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(:bulk_write_message_entries).and_wrap_original do |method, *args|
+ method.call(*args).tap do
+ 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 'does not persist stale stats when a newer run takes over during the final prefetch fallback entry', :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(:sleep)
+ allow(DataImports::Intercom::MessageBatchBuilder).to receive(:new).and_wrap_original do |method, **kwargs|
+ method.call(**kwargs).tap do |batch_builder|
+ allow(batch_builder).to receive(:perform).and_wrap_original do |perform, *args|
+ raise ActiveRecord::QueryCanceled, 'statement timeout' if args.empty?
+
+ perform.call(*args)
+ end
+ end
+ end
+ allow(importer).to receive(:import_unprepared_message).and_wrap_original do |method, *args|
+ method.call(*args).tap do
+ source_entry = args.last
+ next unless source_entry.dig(:part, 'id') == 'part_2'
+
+ DataImport.find(data_import.id).update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => 'new-run' })
+ end
+ end
+ expect(importer).not_to receive(:update_conversation_activity)
+
+ importer.import_conversations_page
+
+ expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once
+ expect(importer).to have_received(:import_unprepared_message).exactly(3).times
+ expect(account.messages.count).to eq(3)
+ expect(data_import.reload.stats).to include(
+ 'conversations' => include('imported' => 0),
+ 'messages' => include('imported' => 0)
+ )
+ end
+
+ it 'does not write prefetch fallback entries after a newer 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 })
+ importer = described_class.new(data_import: data_import, run_id: run_id)
+ allow(importer).to receive(:sleep)
+ allow(DataImports::Intercom::MessageBatchBuilder).to receive(:new).and_wrap_original do |method, **kwargs|
+ method.call(**kwargs).tap do |batch_builder|
+ allow(batch_builder).to receive(:perform).and_wrap_original do |perform, *args|
+ raise ActiveRecord::QueryCanceled, 'statement timeout' if args.empty?
+
+ perform.call(*args)
+ end
+ end
+ end
+ allow(importer).to receive(:import_unprepared_message).and_wrap_original do |method, *args|
+ DataImport.find(data_import.id).update!(
+ source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => 'new-run' }
+ )
+ method.call(*args)
+ end
+
+ importer.import_conversations_page
+
+ expect(importer).to have_received(:import_unprepared_message).once
+ expect(account.messages).to be_empty
+ expect(data_import.mappings.where(source_object_type: 'message')).to be_empty
+ end
+
+ it 'stops individual fallback entries when a newer 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 })
+ importer = described_class.new(data_import: data_import, run_id: run_id)
+ allow(importer).to receive(:bulk_write_message_entries).and_raise(ActiveRecord::StatementInvalid, 'bulk failed')
+ allow(importer).to receive(:fallback_message_entry).and_wrap_original do |method, *args|
+ method.call(*args).tap do
+ DataImport.find(data_import.id).update!(source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => 'new-run' })
+ end
+ end
+
+ importer.import_conversations_page
+
+ expect(importer).to have_received(:fallback_message_entry).once
+ expect(account.messages.count).to eq(1)
+ 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(:bulk_write_message_entries).and_wrap_original do |method, *args|
+ method.call(*args).tap do
+ 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('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:|
@@ -322,10 +844,11 @@ 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'
+ allow(DataImportMapping).to receive(:upsert_all).and_raise(ActiveRecord::StatementInvalid, 'bulk mapping failed')
+ allow(importer).to receive(:record_message_mapping).and_wrap_original do |method, entry, message|
+ raise StandardError, 'mapping failed' if entry.source_id == 'conversation:conversation_1:source:source_1'
- method.call(object_type, source_id, record, metadata: metadata)
+ method.call(entry, message)
end
importer.import_conversations_page
@@ -338,6 +861,81 @@ 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(DataImportMapping).to receive(:upsert_all).and_raise(ActiveRecord::StatementInvalid, 'bulk unavailable')
+ 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(DataImportMapping).to receive(:upsert_all).and_raise(ActiveRecord::StatementInvalid, 'bulk unavailable')
+ 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
@@ -355,6 +953,25 @@ RSpec.describe DataImports::Intercom::Importer do
expect(data_import.reload).to be_abandoned
expect(data_import.completed_at).to be_nil
end
+
+ it 'does not finish an import after a newer run takes over', :aggregate_failures do
+ run_id = 'intercom-run-1'
+ data_import.update!(
+ status: :processing,
+ source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => run_id }
+ )
+ importer = described_class.new(data_import: data_import, run_id: run_id)
+
+ DataImport.find(data_import.id).update!(
+ status: :pending,
+ source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => 'intercom-run-2' }
+ )
+
+ importer.finish!
+
+ expect(data_import.reload).to be_pending
+ expect(data_import.completed_at).to be_nil
+ end
end
describe '#fail!' do
@@ -373,6 +990,26 @@ RSpec.describe DataImports::Intercom::Importer do
expect(data_import.last_error_at).to be_nil
expect(data_import.import_errors.exists?).to be(false)
end
+
+ it 'does not fail an import after a newer run takes over', :aggregate_failures do
+ run_id = 'intercom-run-1'
+ data_import.update!(
+ status: :processing,
+ source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => run_id }
+ )
+ importer = described_class.new(data_import: data_import, run_id: run_id)
+
+ DataImport.find(data_import.id).update!(
+ status: :pending,
+ source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => 'intercom-run-2' }
+ )
+
+ importer.fail!(StandardError.new('boom'))
+
+ expect(data_import.reload).to be_pending
+ expect(data_import.last_error_at).to be_nil
+ expect(data_import.import_errors.exists?).to be(false)
+ end
end
context 'when the Intercom records were imported by an earlier run' do
@@ -404,6 +1041,60 @@ RSpec.describe DataImports::Intercom::Importer do
'message' => 3
)
expect(next_data_import.import_errors.skip_logs.pluck(:details).map { |details| details['reason'] }.uniq).to eq(['already_imported'])
+ expect(
+ DataImportMapping.where(account: account, source_provider: 'intercom', source_object_type: 'message').distinct.pluck(:data_import_id)
+ ).to eq([data_import.id])
+ 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
@@ -429,6 +1120,34 @@ RSpec.describe DataImports::Intercom::Importer do
expect(next_data_import.import_errors.skip_logs.where(source_object_type: 'message')).to be_empty
message_mappings = DataImportMapping.where(account: account, source_provider: 'intercom', source_object_type: 'message')
expect(message_mappings.filter_map(&:chatwoot_record).count).to eq(3)
+ expect(message_mappings.distinct.pluck(:data_import_id)).to eq([next_data_import.id])
+ 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
@@ -493,7 +1212,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
@@ -741,6 +1464,42 @@ RSpec.describe DataImports::Intercom::Importer do
expect(data_import.reload.stats.dig('messages', 'skipped')).to eq(1)
end
+ it 'retries skip-log reconciliation after a query timeout', :aggregate_failures do
+ importer = described_class.new(data_import: data_import)
+ attempts = 0
+ allow(importer).to receive(:sleep)
+ allow(importer).to receive(:record_skipped_message_log).and_wrap_original do |method, *args|
+ attempts += 1
+ raise ActiveRecord::QueryCanceled, 'statement timeout' if attempts == 1
+
+ method.call(*args)
+ end
+
+ importer.perform
+
+ expect(attempts).to eq(2)
+ expect(importer).to have_received(:sleep).with(be_between(0.2, 0.5)).once
+ expect(data_import.reload).to be_completed
+ expect(data_import.stats.dig('messages', 'skipped')).to eq(1)
+ end
+
+ it 'isolates a persistent skip-log reconciliation failure to the message', :aggregate_failures do
+ importer = described_class.new(data_import: data_import)
+ allow(importer).to receive(:record_skipped_message_log).and_raise(ActiveRecord::StatementInvalid, 'skip log failed')
+
+ importer.perform
+
+ item = data_import.items.find_by!(source_object_type: 'conversation', source_object_id: 'conversation_1')
+ error = data_import.import_errors.find_by!(
+ source_object_type: 'message',
+ source_object_id: 'conversation:conversation_1:part:blank_part'
+ )
+ expect(item).to be_imported
+ expect(error).to have_attributes(error_code: 'ActiveRecord::StatementInvalid', message: 'skip log failed')
+ expect(data_import.reload).to be_completed_with_errors
+ expect(data_import.stats.dig('errors', 'count')).to eq(1)
+ end
+
it 'records the skip log again for a later import run', :aggregate_failures do
described_class.new(data_import: data_import).perform
next_data_import = create(
@@ -876,6 +1635,29 @@ 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(:bulk_write_message_entries).and_wrap_original do |method, *args|
+ method.call(*args).tap do
+ 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
@@ -922,6 +1704,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' => {
@@ -942,7 +1726,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.any? { |record| record[:source_id] == 'intercom:conversation:conversation_1:part:bad_part' }
+ insert_attempts << 'intercom:conversation:conversation_1:part:bad_part'
+ raise ActiveRecord::StatementInvalid, 'bad message'
+ end
method.call(records, **kwargs)
end
@@ -963,6 +1750,9 @@ 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.size).to eq(2)
+ expect(data_import.import_errors.where(source_object_type: 'message').count).to eq(1)
+ expect(account.messages.find_by(source_id: 'intercom:conversation:conversation_1:source:source_1')).to be_present
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..9a2420320
--- /dev/null
+++ b/spec/services/data_imports/intercom/message_batch_builder_spec.rb
@@ -0,0 +1,205 @@
+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 'refreshes classifications while preserving source positions' do
+ batch = builder.perform
+ target_entry = batch.entries.second
+ message = create(
+ :message,
+ account: account,
+ conversation: conversation,
+ inbox: conversation.inbox,
+ source_id: "intercom:#{target_entry.source_id}"
+ )
+ DataImportMapping.create!(
+ account: account,
+ data_import: data_import,
+ source_provider: 'intercom',
+ source_object_type: 'message',
+ source_object_id: target_entry.source_id,
+ chatwoot_record_type: 'Message',
+ chatwoot_record_id: message.id,
+ metadata: {}
+ )
+
+ refreshed_batch = builder.refresh(batch.entries)
+
+ expect(refreshed_batch.entries.map(&:position)).to eq([0, 1, 2])
+ expect(refreshed_batch.entries.second).to have_attributes(classification: :current_import, message: message)
+ 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
diff --git a/spec/services/data_imports/intercom/retry_service_spec.rb b/spec/services/data_imports/intercom/retry_service_spec.rb
new file mode 100644
index 000000000..7c6ccc6ea
--- /dev/null
+++ b/spec/services/data_imports/intercom/retry_service_spec.rb
@@ -0,0 +1,82 @@
+require 'rails_helper'
+
+RSpec.describe DataImports::Intercom::RetryService do
+ let(:account) { create(:account) }
+ let(:data_import) { create(:data_import, :intercom, account: account, status: :processing, started_at: 2.hours.ago) }
+
+ before do
+ account.enable_features!('data_import')
+ end
+
+ it 'prepares a stalled import for another run without clearing progress', :aggregate_failures do
+ original_cursor = { 'contacts' => { 'completed' => true }, 'conversations' => { 'starting_after' => 'cursor-1' } }
+ original_stats = { 'contacts' => { 'imported' => 10 }, 'conversations' => { 'imported' => 5 } }
+ started_at = data_import.started_at
+ error = data_import.import_errors.create!(error_code: 'MessageFailed', message: 'Timed out')
+ data_import.update!(
+ cursor: original_cursor,
+ stats: original_stats,
+ source_metadata: { DataImport::ACTIVE_INTERCOM_IMPORT_RUN_ID_KEY => 'previous-run' },
+ updated_at: 16.minutes.ago
+ )
+
+ result = described_class.new(account: account, data_import: data_import).perform
+
+ expect(result).to eq(:enqueue)
+ expect(data_import.reload).to be_pending
+ expect(data_import.started_at).to eq(started_at)
+ expect(data_import.cursor).to eq(original_cursor)
+ expect(data_import.stats).to eq(original_stats)
+ expect(data_import.import_errors).to contain_exactly(error)
+ expect(data_import.active_intercom_import_run_id).not_to eq('previous-run')
+ end
+
+ it 'does not retry an import that is still receiving updates' do
+ result = described_class.new(account: account, data_import: data_import).perform
+
+ expect(result).to eq(:not_stalled)
+ expect(data_import.reload).to be_processing
+ end
+
+ it 'rechecks the import state after acquiring its row lock' do
+ data_import.update!(updated_at: 16.minutes.ago)
+ allow(data_import).to receive(:with_lock).and_wrap_original do |method, *args, &block|
+ data_import.update!(status: :completed, completed_at: Time.current)
+ method.call(*args, &block)
+ end
+
+ result = described_class.new(account: account, data_import: data_import).perform
+
+ expect(result).to eq(:not_stalled)
+ expect(data_import.reload).to be_completed
+ end
+
+ it 'does not retry while another Intercom import is active' do
+ data_import.update!(updated_at: 16.minutes.ago)
+ create(:data_import, :intercom, account: account, status: :processing)
+
+ result = described_class.new(account: account, data_import: data_import).perform
+
+ expect(result).to eq(:active_import_exists)
+ expect(data_import.reload).to be_processing
+ end
+
+ it 'allows an active legacy import to continue alongside the retry' do
+ data_import.update!(updated_at: 16.minutes.ago)
+ create(:data_import, account: account, status: :processing)
+
+ result = described_class.new(account: account, data_import: data_import).perform
+
+ expect(result).to eq(:enqueue)
+ expect(data_import.reload).to be_pending
+ end
+
+ it 'does not retry when the stored access key is unavailable' do
+ data_import.update!(access_token: nil, updated_at: 16.minutes.ago)
+
+ result = described_class.new(account: account, data_import: data_import).perform
+
+ expect(result).to eq(:access_token_missing)
+ expect(data_import.reload).to be_processing
+ end
+end