Compare commits

...
Author SHA1 Message Date
Sony Mathew 4a9f961fa3 Merge branch 'codex/cw-7519-intercom-stalled-retry-15m' into codex/cw-7615-intercom-bulk-message-writes 2026-07-23 23:42:03 +05:30
Sony MathewandGitHub b6efdae243 perf: shorten Intercom import jobs (#15051)
## Description

Reduces sustained load during large Intercom imports by limiting
conversation list pages to 10 while retaining 50-contact pages. Cursor
and provider total-count behavior remain unchanged.

The one-minute progress heartbeat now lives in parent PR #15050 because
stalled-import detection must be safe when that PR is deployed
independently. This child PR therefore contains only the smaller-page
delta.

This is Phase 1, Task 2 of the [Intercom import optimization plan
(CW-7615)](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion).

This PR is stacked on #15050 and should be reviewed as the two-file
delta from `codex/cw-7519-intercom-stalled-retry-15m`.

## Closes

-
[CW-7615](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion)

## Type of change

- [x] Bug fix (non-breaking change which fixes an issue)

## How to test

1. Start an Intercom import containing contacts and conversations across
multiple source pages.
2. Confirm contact list requests retain a page size of 50.
3. Confirm conversation list requests use a page size of 10.
4. Confirm page cursors and displayed provider totals continue advancing
normally.
5. Confirm the parent PR heartbeat keeps long conversations fresh while
these shorter pages are processed.

## Checklist

- [x] My code follows the style guidelines of this project
- [x] I have performed a self-review of my code
- [x] I have added tests that prove the change is effective
- [x] New and existing focused tests pass locally with my changes
2026-07-23 23:38:01 +05:30
Sony Mathew 52a28db6d8 fix(imports): recheck final fallback ownership 2026-07-23 23:34:53 +05:30
Sony Mathew 0f439f2921 Merge branch 'codex/cw-7519-intercom-shorter-jobs-heartbeats' into codex/cw-7615-intercom-bulk-message-writes 2026-07-23 23:18:17 +05:30
Sony MathewandGitHub 9a489d38c7 chore: (refactor) prepare Intercom message batches (#15110)
## Description

Prepares each Intercom conversation's source message and parts as one
deterministic batch before persistence. The builder preserves provider
order, prefetches mappings and live messages once, and classifies every
entry as current-run, previous-run, stale-mapping repair,
existing-message repair, or new.

The importer consumes those prepared entries through the existing
individual transaction and retry path, so activity conversion, private
notes, sender attribution, metadata, timestamps, skip logs, idempotency,
indexing, and heartbeat behavior remain unchanged. This PR intentionally
does not bulk insert messages; batching writes is Phase 2 Task 2.

This is Phase 2, Task 1 (Task 4 overall) of
[CW-7615](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion).

This PR is stacked on #15052 and should be reviewed as the four-file
delta from `codex/cw-7615-intercom-query-timeout-retries`.

## Closes

- Tracking plan:
[CW-7615](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion)

## Type of change

- [x] Bug fix (non-breaking change which fixes an issue)
- [ ] New feature (non-breaking change which adds functionality)
- [ ] Breaking change (fix or feature that would cause existing
functionality not to work as expected)
- [ ] This change requires a documentation update

## How to test

1. Import a conversation containing a source message, comments, notes,
activities, and entries with equal timestamps.
2. Retry the same run and verify messages and counters remain
idempotent.
3. Retry after a previous run and verify existing mappings are reported
as skips without changing their owner.
4. Delete a mapped Message, rerun, and verify the message and mapping
are repaired.
5. Delete only a mapping, rerun, and verify the existing Message is
reused without duplication.
6. Rerun a previously skipped part that now classifies as an activity
and verify it is repaired.
7. Import an empty conversation and verify no message or mapping
prefetch is performed.

## Checklist

- [x] My code follows the style guidelines of this project
- [x] I have performed a self-review of my code
- [x] I have commented on my code, particularly in hard-to-understand
areas
- [ ] I have made corresponding changes to the documentation
- [x] My changes generate no new warnings
- [x] I have added tests that prove my fix is effective or that my
feature works
- [x] New and existing unit tests pass locally with my changes
- [ ] Any dependent changes have been merged and published in downstream
modules
2026-07-23 23:13:15 +05:30
Sony Mathew ed080ac253 fix(imports): validate message timestamp range 2026-07-23 23:07:51 +05:30
Sony MathewandGitHub 4e410b78ff perf: reduce Intercom import query pressure (#15052)
## Description

Reduces database pressure during large Intercom imports by deferring
same-run contact and message statistic recounts until the page boundary,
immediately before the cursor advances. Repeated mappings within a page
now mark their statistic group dirty instead of running aggregate counts
and persisting statistics for every record.

Contact and message database attempts now retry once when PostgreSQL
cancels a query. The failed transaction is allowed to roll back before
the importer waits for a randomized 200-500 ms delay and retries. A
recovered timeout creates no error; a second timeout follows the
existing item/message failure path exactly once. Provider, validation,
and other database errors are not retried.

This is Phase 1, Task 3 of the [Intercom import optimization plan
(CW-7615)](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion).

This PR is stacked on #15051 and should be reviewed as the two-file
delta from `codex/cw-7519-intercom-shorter-jobs-heartbeats`.

## Closes

- Tracking plan:
[CW-7615](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion)

## Type of change

- [x] Bug fix (non-breaking change which fixes an issue)

## How to test

1. Retry a page whose contacts or messages already have same-run
mappings and confirm imported/skipped statistics are reconciled before
the page cursor advances.
2. Raise one `ActiveRecord::QueryCanceled` while persisting a contact
mapping and confirm the transaction rolls back, retries, and creates one
contact and mapping without an error.
3. Repeat for a message mapping and confirm one message and mapping are
created without an error.
4. Raise the timeout on both attempts and confirm one final item/message
error is recorded.
5. Raise an unrelated database error and confirm it is not retried.

## Checklist

- [x] My code follows the style guidelines of this project
- [x] I have performed a self-review of my code
- [x] I have added tests that prove the change is effective
- [x] New and existing focused tests pass locally with my changes
2026-07-23 22:53:52 +05:30
Sony MathewandGitHub 580cbf91c5 fix: preserve Intercom bulk failure isolation (#15116)
This preserves per-message failure isolation when an Intercom bulk
message chunk cannot be committed. Query timeouts are retried once after
the failed transaction rolls back, while other database failures move
directly to an entry-by-entry fallback.

The fallback reacquires the import lock and refreshes each entry before
writing, so a competing worker cannot create duplicate messages. Each
individual write uses a savepoint, and search indexing happens only
after the import lock transaction commits.

This is Phase 2, Task 6 of CW-7615 and is stacked on #15111. The review
delta is one commit across the importer and its focused spec.

## Closes

-
[CW-7615](https://linear.app/chatwoot/issue/CW-7615/optimize-intercom-import-reliability-and-bulk-message-ingestion)

## How to test

1. Run an Intercom import where a bulk mapping write times out once and
confirm the chunk succeeds on retry.
2. Force a persistent timeout or database write failure and confirm
valid messages are imported through the fallback.
3. Force one fallback record to fail and confirm only that message is
recorded as an import error.
4. Simulate another worker repairing a message before fallback and
confirm no duplicate is created.
2026-07-23 22:52:06 +05:30
Sony Mathew d38ef6c9c8 fix(imports): harden bulk message recovery 2026-07-23 15:02:28 +05:30
Sony Mathew ae57547b84 fix(imports): reconcile skipped message stats on retry 2026-07-23 14:47:41 +05:30
Sony Mathew d6a5938916 style(imports): format prefetch fallback 2026-07-23 14:47:26 +05:30
Sony Mathew 5d479d36e1 Merge branch 'codex/cw-7615-intercom-message-batches' into codex/cw-7615-intercom-bulk-message-writes
# Conflicts:
#	app/services/data_imports/intercom/importer.rb
#	app/services/data_imports/intercom/message_batch_builder.rb
#	spec/services/data_imports/intercom/importer_spec.rb
2026-07-23 14:44:36 +05:30
Sony Mathew d238df2375 Merge branch 'codex/cw-7615-intercom-query-timeout-retries' into codex/cw-7615-intercom-message-batches
# Conflicts:
#	app/services/data_imports/intercom/importer.rb
2026-07-23 14:42:31 +05:30
Sony Mathew 21b83c6be8 Merge remote-tracking branch 'origin/codex/cw-7615-intercom-query-timeout-retries' into codex/cw-7615-intercom-query-timeout-retries
# Conflicts:
#	spec/services/data_imports/intercom/importer_spec.rb
2026-07-23 14:40:44 +05:30
Sony Mathew 03f720c0ef Merge branch 'codex/cw-7519-intercom-shorter-jobs-heartbeats' into codex/cw-7615-intercom-query-timeout-retries 2026-07-23 14:39:16 +05:30
Sony Mathew 2ae36d53c3 Merge branch 'codex/cw-7519-intercom-stalled-retry-15m' into codex/cw-7519-intercom-shorter-jobs-heartbeats 2026-07-23 14:38:52 +05:30
Sony Mathew ea783a89ea fix(imports): check run ownership before heartbeat 2026-07-23 14:38:23 +05:30
Sony Mathew 003cdbc417 fix(imports): guard Intercom mapped message type 2026-07-23 00:54:14 +05:30
Sony Mathew bee3eeae42 fix(imports): isolate Intercom prefetch failures 2026-07-23 00:53:31 +05:30
Sony Mathew 4f8215bc5c fix(imports): reload run state before heartbeat checks 2026-07-23 00:51:31 +05:30
Sony Mathew b5d4a5b6af fix(imports): preserve skipped stats on retry 2026-07-23 00:51:23 +05:30
Sony Mathew 25dd72da4e fix(imports): retry deferred stat reconciliation 2026-07-23 00:50:29 +05:30
Sony Mathew 5c9ffe4ce4 fix(imports): reconcile conversation stats on retry 2026-07-23 00:50:08 +05:30
Sony Mathew 0410211e23 Merge branch 'codex/cw-7615-intercom-message-batches' into codex/cw-7615-intercom-bulk-message-writes 2026-07-23 00:30:30 +05:30
Sony Mathew 1fdf128a74 Merge branch 'codex/cw-7615-intercom-query-timeout-retries' into codex/cw-7615-intercom-message-batches 2026-07-23 00:29:58 +05:30
Sony Mathew a3fda2a238 Merge branch 'codex/cw-7519-intercom-shorter-jobs-heartbeats' into codex/cw-7615-intercom-query-timeout-retries 2026-07-23 00:29:29 +05:30
Sony Mathew 045dab192c Merge branch 'codex/cw-7519-intercom-stalled-retry-15m' into codex/cw-7519-intercom-shorter-jobs-heartbeats 2026-07-23 00:29:01 +05:30
Sony Mathew b8945e8d87 fix(imports): guard stale run failures 2026-07-23 00:21:09 +05:30
Sony Mathew c8794e6e2d Merge branch 'codex/cw-7615-intercom-message-batches' into codex/cw-7615-intercom-bulk-message-writes 2026-07-23 00:05:54 +05:30
Sony Mathew fcf8972f25 Merge branch 'codex/cw-7615-intercom-query-timeout-retries' into codex/cw-7615-intercom-message-batches 2026-07-23 00:05:38 +05:30
Sony Mathew d7ee6eec9f Merge branch 'codex/cw-7519-intercom-shorter-jobs-heartbeats' into codex/cw-7615-intercom-query-timeout-retries 2026-07-23 00:05:10 +05:30
Sony Mathew 9156ccdaed Merge branch 'codex/cw-7519-intercom-stalled-retry-15m' into codex/cw-7519-intercom-shorter-jobs-heartbeats
# Conflicts:
#	app/services/data_imports/intercom/importer.rb
2026-07-23 00:03:23 +05:30
Sony Mathew 5b5cdf521a fix(imports): ignore stale detail polls after retry 2026-07-23 00:01:59 +05:30
Sony Mathew 7bdf38732f fix(imports): heartbeat active Intercom runs 2026-07-23 00:00:29 +05:30
Sony Mathew 6a43fd5b5a fix(imports): lock stalled import retries 2026-07-22 23:59:30 +05:30
Sony MathewandGitHub 35167631bb Merge branch 'develop' into codex/cw-7519-intercom-stalled-retry-15m 2026-07-22 23:48:36 +05:30
Shivam MishraandGitHub ddb0535a93 perf: reuse resolved count for reopen rate (#15122)
This improves the Captain overview by loading reporting metrics and FAQ
stats from separate endpoints. Range changes now refresh only the
metrics, while reopen-rate calculation reuses the resolved conversation
count to avoid redundant database queries.

## What changed

- Split Captain overview metrics and FAQ stats into separate APIs.
- Fetch FAQ stats independently from range-based metrics.
- Reuse resolved conversation totals when calculating reopen rate.
- Skip the reopen query when there are no resolved conversations.
2026-07-22 22:03:25 +05:30
Sivin VargheseandGitHub 42cbf7d3b9 fix: stray backslash after hard breaks before formatted list items (#15112) 2026-07-22 20:07:00 +05:30
887897ea98 fix: lock agent quota checks (#15029)
# Pull Request Template

## Description

Locks the agent quota check to the account row while creating account
users. This fixes a race where concurrent agent-create requests could
all observe the same remaining seat before any `account_users` row was
inserted.

The API continues to return the existing `402 Account limit exceeded.
Please purchase more licenses` response when the limit is reached. Bulk
create now preflights the requested email count while holding the
account lock, then creates each agent through the same locked builder
path. The Enterprise custom-role hook now no-ops when create did not
produce an agent.

Fixes:
[CW-7039](https://linear.app/chatwoot/issue/CW-7039/race-condition-in-agent-creation-bypasses-plan-agent-seat-limit)

## Type of change

- [x] Bug fix (non-breaking change which fixes an issue)

## How Has This Been Tested?

- `POSTGRES_DATABASE=chatwoot_test_c20f_agent_quota REDIS_DB=9 bundle
exec rspec spec/builders/agent_builder_spec.rb
spec/enterprise/builders/agent_builder_spec.rb
spec/controllers/api/v1/accounts/agents_controller_spec.rb
spec/enterprise/controllers/api/v1/accounts/agents_controller_spec.rb
spec/enterprise/controllers/enterprise/api/v1/accounts/agents_controller_spec.rb`
- `bundle exec rubocop app/builders/agent_builder.rb
app/controllers/api/v1/accounts/agents_controller.rb
enterprise/app/controllers/enterprise/api/v1/accounts/agents_controller.rb
spec/builders/agent_builder_spec.rb
spec/enterprise/controllers/api/v1/accounts/agents_controller_spec.rb`
- `git diff --check`
- One-off threaded Rails validation with 8 concurrent `AgentBuilder`
calls against an account with one remaining seat: `created: 1`,
`limited: 7`, final `count=2`, `limit=2`.

## Checklist:

- [x] My code follows the style guidelines of this project
- [x] I have performed a self-review of my code
- [x] I have commented on my code, particularly in hard-to-understand
areas
- [ ] I have made corresponding changes to the documentation
- [x] My changes generate no new warnings
- [x] I have added tests that prove my fix is effective or that my
feature works
- [x] New and existing unit tests pass locally with my changes
- [ ] Any dependent changes have been merged and published in downstream
modules

Co-authored-by: Muhsin Keloth <muhsinkeramam@gmail.com>
2026-07-22 18:43:00 +05:30
8aee518149 fix(integrations): restrict Linear/Notion/Shopify hook deletion to admins (#15126)
Non-admin agents could delete an account's Linear, Notion, or Shopify
integration through the dedicated integration endpoints, which — unlike
the generic hooks endpoint — never checked the caller's role. This
restores the intended admin-only boundary for removing an integration.

## Closes
- https://linear.app/chatwoot/issue/CW-7383
- https://linear.app/chatwoot/issue/CW-7384
- https://linear.app/chatwoot/issue/CW-7189

## How to reproduce
As a non-admin **agent**, `DELETE
/api/v1/accounts/:id/integrations/{linear,notion,shopify}` returned
`200` and removed the account-wide integration. After this change it
returns `401` and the integration is preserved; administrators can still
remove it.

## What changed
- Route integration-hook deletion through `HookPolicy` (admin-only) via
a shared `Integrations::BaseController`, matching the generic hooks
controller.

Co-authored-by: Vishnu Narayanan <iamwishnu@gmail.com>
2026-07-22 17:53:25 +05:30
Vishnu NarayananandGitHub 5733b822e3 fix: guard widget email-transcript button against repeat clicks (#15095)
## Description

The widget email-transcript button (`ChatFooter.vue`) had no client-side
guard. Its visibility depends only on whether the contact has an email,
and the click handler fired a request on every click with no in-flight
lock, no disabled state, and no post-send handling. A user could
therefore trigger a large number of duplicate transcript emails from a
single conversation just by clicking repeatedly.

This adds a re-entry guard in the handler, disables the button while a
send is in flight, and applies a short cooldown (15s) after a successful
send. Normal use is unaffected: the button sends once, shows the success
toast, then briefly disables and automatically re-enables so a genuine
later re-request still works. On failure the button stays enabled so the
user can retry immediately. The cooldown timer is cleared on unmount.

Using a timed cooldown (rather than a permanent post-send lock) also
avoids the button getting stuck disabled if a resolved conversation is
reopened and later re-resolved.

This is the client-side complement to the server-side rate limit added
in #15085.

Fixes https://linear.app/chatwoot/issue/CW-7640
2026-07-22 17:36:36 +05:30
Sivin VargheseandGitHub 1e52d23d7a fix: guard agent sort against null names in assignment dropdown (#15125) 2026-07-22 15:27:04 +05:30
Sivin VargheseandGitHub fbb3479263 fix: guard agent sort against null names in assignment dropdown (#15125)
# Pull Request Template

## Description

This PR fixes a crash where opening a conversation threw `TypeError:
Cannot read properties of null (reading 'localeCompare')` and prevented
the agent assignment dropdown from rendering.

Since #14866, agent bots are included in the assignable agents list.
`AgentBot#name` is not presence-validated, so system bots (account-less,
global) can have a `null` name. Those nameless bots flowed into
name-based operations that assumed a string, causing crashes and
warnings across multiple surfaces:

* **Assignment dropdown sort:** `getAgentsByAvailability` called
`a.name.localeCompare(b.name)`, causing a `localeCompare` `TypeError`.
* **Dropdown search:** `MultiselectDropdownItems` called
`option.name.toLowerCase()`, causing a `toLowerCase` `TypeError`.
* **Agent Bots settings:** `Avatar` received `name=null` for a `String`
prop, triggering a Vue prop validation warning.

### What changed

* Keep nameless agent bots in the assignment dropdown and render a `-`
fallback label in `useAgentsList`. These are still valid,
assignable-by-ID records: the assignable agents API includes accessible
bots, and `Conversations::AssignmentService` assigns them by ID.
Preserving them avoids hiding valid assignment targets. Bots are still
included only when `includeAgentBots` is enabled.
* Make the sort in `getAgentsByAvailability` null-safe by coercing
missing names to an empty string (defense in depth).
* Make the search filter in `MultiselectDropdownItems` null-safe
(defense in depth).
* Pass a null-safe `name` prop to `Avatar` in the Agent Bots settings
list to eliminate the Vue prop validation warning.

Fixes
https://linear.app/chatwoot/issue/CW-7670/agent-assignment-dropdown-crashes-with-cannot-read-properties-of-null

## Type of change

- [x] Bug fix (non-breaking change which fixes an issue)

## How Has This Been Tested?


1. Have a system agent bot (`name: null`) that is assignable to an
inbox.
2. Open any conversation in that inbox.
   * The agent assignment dropdown renders without console errors.
   * The nameless bot is listed with a `-` label and can be assigned.
3. Go to **Settings → Agent Bots**.
   * The page renders without the `Avatar` prop validation warning.


## Checklist:

- [x] My code follows the style guidelines of this project
- [x] I have performed a self-review of my code
- [x] I have commented on my code, particularly in hard-to-understand
areas
- [ ] I have made corresponding changes to the documentation
- [x] My changes generate no new warnings
- [x] I have added tests that prove my fix is effective or that my
feature works
- [x] New and existing unit tests pass locally with my changes
- [ ] Any dependent changes have been merged and published in downstream
modules
2026-07-22 15:23:19 +05:30
Sony MathewandGitHub 45278937bd Merge branch 'codex/cw-7615-intercom-query-timeout-retries' into codex/cw-7615-intercom-message-batches 2026-07-21 18:35:27 +05:30
Sony MathewandGitHub 114d7b5613 Merge branch 'codex/cw-7519-intercom-shorter-jobs-heartbeats' into codex/cw-7615-intercom-query-timeout-retries 2026-07-21 18:09:45 +05:30
Sony MathewandGitHub fbee7c2e23 Merge branch 'codex/cw-7519-intercom-stalled-retry-15m' into codex/cw-7519-intercom-shorter-jobs-heartbeats 2026-07-21 18:09:28 +05:30
Sony MathewandGitHub f866adf2c2 Merge branch 'develop' into codex/cw-7519-intercom-stalled-retry-15m 2026-07-21 18:09:06 +05:30
Sony Mathew 92d3536601 perf(imports): bulk write Intercom messages 2026-07-21 18:02:07 +05:30
Sony Mathew 901b234b4d refactor(imports): prepare Intercom message batches 2026-07-21 17:32:07 +05:30
Sony Mathew 1c16c427f4 perf(imports): reduce Intercom query pressure 2026-07-17 15:36:17 +05:30
Sony MathewandGitHub 0b5881889a Merge branch 'develop' into codex/cw-7519-intercom-stalled-retry-15m 2026-07-17 15:27:31 +05:30
Sony Mathew c58ba2f9c2 perf(imports): shorten Intercom jobs and add heartbeats 2026-07-17 14:49:58 +05:30
Sony Mathew 2decd8c544 feat(imports): retry stalled Intercom imports 2026-07-17 14:39:35 +05:30
50 changed files with 2391 additions and 293 deletions
+19 -3
View File
@@ -2,6 +2,14 @@
# It initializes with necessary attributes and provides a perform method
# to create a user and account user in a transaction.
class AgentBuilder
LIMIT_EXCEEDED_MESSAGE = 'Account limit exceeded. Please purchase more licenses'.freeze
class LimitExceededError < StandardError
def initialize
super(AgentBuilder::LIMIT_EXCEEDED_MESSAGE)
end
end
# Initializes an AgentBuilder with necessary attributes.
# @param email [String] the email of the user.
# @param name [String] the name of the user.
@@ -14,15 +22,23 @@ class AgentBuilder
# Creates a user and account user in a transaction.
# @return [User] the created user.
def perform
ActiveRecord::Base.transaction do
@user = find_or_create_user
create_account_user
account.with_lock do
raise LimitExceededError unless can_add_agent?
ActiveRecord::Base.transaction do
@user = find_or_create_user
create_account_user
end
end
@user
end
private
def can_add_agent?
account.usage_limits[:agents] > account.account_users.count
end
# Finds a user by email or creates a new one with a temporary password.
# @return [User] the found or created user.
def find_or_create_user
@@ -1,8 +1,6 @@
class Api::V1::Accounts::AgentsController < Api::V1::Accounts::BaseController
before_action :fetch_agent, except: [:create, :index, :bulk_create]
before_action :check_authorization
before_action :validate_limit, only: [:create]
before_action :validate_limit_for_bulk_create, only: [:bulk_create]
def index
@agents = agents
@@ -20,6 +18,8 @@ class Api::V1::Accounts::AgentsController < Api::V1::Accounts::BaseController
)
@agent = builder.perform
rescue AgentBuilder::LimitExceededError => e
render_payment_required(e.message)
end
def update
@@ -36,25 +36,13 @@ class Api::V1::Accounts::AgentsController < Api::V1::Accounts::BaseController
def bulk_create
emails = params[:emails]
emails.each do |email|
builder = AgentBuilder.new(
email: email,
name: email.split('@').first,
inviter: current_user,
account: Current.account
)
begin
builder.perform
rescue ActiveRecord::RecordInvalid => e
Rails.logger.info "[Agent#bulk_create] ignoring email #{email}, errors: #{e.record.errors}"
end
end
bulk_create_agents(emails)
# This endpoint is used to bulk create agents during onboarding
# onboarding_step key in present in Current account custom attributes, since this is a one time operation
Current.account.custom_attributes.delete('onboarding_step')
Current.account.save!
clear_onboarding_step
head :ok
rescue AgentBuilder::LimitExceededError => e
render_payment_required(e.message)
end
private
@@ -87,22 +75,33 @@ class Api::V1::Accounts::AgentsController < Api::V1::Accounts::BaseController
@agents ||= Current.account.users.order_by_full_name.includes(:account_users, { avatar_attachment: [:blob] })
end
def validate_limit_for_bulk_create
limit_available = params[:emails].count <= available_agent_count
def bulk_create_agents(emails)
Current.account.with_lock do
raise AgentBuilder::LimitExceededError if emails.count > available_agent_count
render_payment_required('Account limit exceeded. Please purchase more licenses') unless limit_available
emails.each { |email| create_agent_from_email(email) }
end
end
def validate_limit
render_payment_required('Account limit exceeded. Please purchase more licenses') unless can_add_agent?
def create_agent_from_email(email)
builder = AgentBuilder.new(
email: email,
name: email.split('@').first,
inviter: current_user,
account: Current.account
)
builder.perform
rescue ActiveRecord::RecordInvalid => e
Rails.logger.info "[Agent#bulk_create] ignoring email #{email}, errors: #{e.record.errors}"
end
def clear_onboarding_step
Current.account.custom_attributes.delete('onboarding_step')
Current.account.save!
end
def available_agent_count
Current.account.usage_limits[:agents] - agents.count
end
def can_add_agent?
available_agent_count.positive?
Current.account.usage_limits[:agents] - Current.account.account_users.count
end
def delete_user_record(agent)
@@ -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
@@ -0,0 +1,9 @@
class Api::V1::Accounts::Integrations::BaseController < Api::V1::Accounts::BaseController
private
# Managing an integration hook (create/update/destroy) is admin-only, enforced via HookPolicy.
# Subclasses opt in per action with `before_action :check_authorization, only: [...]`.
def check_authorization
authorize(:hook)
end
end
@@ -1,4 +1,4 @@
class Api::V1::Accounts::Integrations::HooksController < Api::V1::Accounts::BaseController
class Api::V1::Accounts::Integrations::HooksController < Api::V1::Accounts::Integrations::BaseController
before_action :fetch_hook, except: [:create]
before_action :check_authorization
@@ -35,10 +35,6 @@ class Api::V1::Accounts::Integrations::HooksController < Api::V1::Accounts::Base
@hook = Current.account.hooks.find(params[:id])
end
def check_authorization
authorize(:hook)
end
def permitted_params
params.require(:hook).permit(:app_id, :inbox_id, :status, settings: {})
end
@@ -1,6 +1,7 @@
class Api::V1::Accounts::Integrations::LinearController < Api::V1::Accounts::BaseController
class Api::V1::Accounts::Integrations::LinearController < Api::V1::Accounts::Integrations::BaseController
before_action :fetch_conversation, only: [:create_issue, :link_issue, :unlink_issue, :linked_issues]
before_action :fetch_hook, only: [:destroy]
before_action :check_authorization, only: [:destroy]
def destroy
revoke_linear_token
@@ -1,5 +1,6 @@
class Api::V1::Accounts::Integrations::NotionController < Api::V1::Accounts::BaseController
class Api::V1::Accounts::Integrations::NotionController < Api::V1::Accounts::Integrations::BaseController
before_action :fetch_hook, only: [:destroy]
before_action :check_authorization, only: [:destroy]
def destroy
@hook.destroy!
@@ -1,7 +1,8 @@
class Api::V1::Accounts::Integrations::ShopifyController < Api::V1::Accounts::BaseController
class Api::V1::Accounts::Integrations::ShopifyController < Api::V1::Accounts::Integrations::BaseController
include Shopify::IntegrationHelper
before_action :setup_shopify_context, only: [:orders]
before_action :fetch_hook, except: [:auth]
before_action :check_authorization, only: [:destroy]
before_action :validate_contact, only: [:orders]
def auth
@@ -26,13 +26,20 @@ class CaptainAssistant extends ApiClient {
});
}
getStats({ assistantId, range, signal }) {
getMetrics({ assistantId, range, signal }) {
const requestConfig = {
params: { range, timezone_offset: getTimezoneOffset() },
};
if (signal) requestConfig.signal = signal;
return axios.get(`${this.url}/${assistantId}/stats`, requestConfig);
return axios.get(`${this.url}/${assistantId}/metrics`, requestConfig);
}
getFaqStats({ assistantId, signal }) {
const requestConfig = {};
if (signal) requestConfig.signal = signal;
return axios.get(`${this.url}/${assistantId}/faq_stats`, requestConfig);
}
getSummary({ assistantId, range, stats }) {
@@ -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`);
}
@@ -8,6 +8,8 @@ import {
EditorState,
Selection,
imageResizeView,
toggleMark,
wrapInList,
} from '@chatwoot/prosemirror-schema';
import {
suggestionsPlugin,
@@ -17,8 +19,6 @@ import imagePastePlugin from '@chatwoot/prosemirror-schema/src/plugins/image';
import embedPreviewPlugin from '@chatwoot/prosemirror-schema/src/plugins/embedPreview';
import trailingParagraphPlugin from '@chatwoot/prosemirror-schema/src/plugins/trailingParagraph';
import { embeds as markdownEmbeds } from 'dashboard/helper/markdownEmbeds';
import { toggleMark } from 'prosemirror-commands';
import { wrapInList } from 'prosemirror-schema-list';
import { toggleBlockType } from '@chatwoot/prosemirror-schema/src/menu/common';
import { checkFileSizeLimit } from 'shared/helpers/FileHelper';
import { isEscape } from 'shared/helpers/KeyboardHelpers';
@@ -1,9 +1,9 @@
import { ref } from 'vue';
import { describe, it, expect, vi, beforeEach } from 'vitest';
import { useAgentsList } from '../useAgentsList';
import { useMapGetter } from 'dashboard/composables/store';
import { allAgentsData, formattedAgentsData } from './fixtures/agentFixtures';
import * as agentHelper from 'dashboard/helper/agentHelper';
import { beforeEach, describe, expect, it, vi } from 'vitest';
import { ref } from 'vue';
import { useAgentsList } from '../useAgentsList';
import { allAgentsData, formattedAgentsData } from './fixtures/agentFixtures';
// Mock vue-i18n
vi.mock('vue-i18n', () => ({
@@ -94,6 +94,32 @@ describe('useAgentsList', () => {
expect(agentsList.value.length).toBe(formattedAgentsData.slice(1).length);
});
it('keeps nameless agent bots and applies a fallback label', () => {
const namelessBot = {
id: 91,
name: null,
assignee_type: 'AgentBot',
availability_status: 'offline',
};
mockUseMapGetter({
'inboxAssignableAgents/getAssignableAgents': ref(() => [
...allAgentsData,
namelessBot,
]),
});
const { agentsList } = useAgentsList();
// access the computed to trigger evaluation
expect(agentsList.value).toBeDefined();
const passedAgents =
agentHelper.getAgentsByUpdatedPresence.mock.calls[0][0];
expect(passedAgents).toContainEqual({
...namelessBot,
name: '-',
});
});
it('handles empty assignable agents', () => {
mockUseMapGetter({
'inboxAssignableAgents/getAssignableAgents': ref(() => []),
@@ -1,10 +1,10 @@
import { computed } from 'vue';
import { useMapGetter } from 'dashboard/composables/store';
import { useI18n } from 'vue-i18n';
import {
getAgentsByUpdatedPresence,
getSortedAgentsByAvailability,
} from 'dashboard/helper/agentHelper';
import { computed } from 'vue';
import { useI18n } from 'vue-i18n';
/**
* A composable function that provides a list of agents for assignment.
@@ -53,7 +53,11 @@ export function useAgentsList(
* @type {import('vue').ComputedRef<Array>}
*/
const agentsList = computed(() => {
const agents = assignableAgents.value || [];
const agents = (assignableAgents.value || []).map(agent =>
!agent.name && agent.assignee_type === 'AgentBot'
? { ...agent, name: '-' }
: agent
);
const agentsByUpdatedPresence = getAgentsByUpdatedPresence(
agents,
currentUser.value,
@@ -7,7 +7,7 @@
export const getAgentsByAvailability = (agents, availability) => {
return agents
.filter(agent => agent.availability_status === availability)
.sort((a, b) => a.name.localeCompare(b.name));
.sort((a, b) => (a.name || '').localeCompare(b.name || ''));
};
/**
@@ -1,4 +1,6 @@
import {
InputRule,
inputRules,
MessageMarkdownSerializer,
MessageMarkdownTransformer,
messageSchema,
@@ -9,7 +11,6 @@ import * as Sentry from '@sentry/vue';
import camelcaseKeys from 'camelcase-keys';
import { FORMATTING, MARKDOWN_PATTERNS } from 'dashboard/constants/editor';
import { INBOX_TYPES, TWILIO_CHANNEL_MEDIUM } from 'dashboard/helper/inbox';
import { InputRule, inputRules } from 'prosemirror-inputrules';
/**
* Extract text from markdown, and remove all images, code blocks, links, headers, bold, italic, lists etc.
@@ -26,6 +26,18 @@ describe('agentHelper', () => {
offlineAgentsData
);
});
it('does not throw when an agent has a null name', () => {
const agents = [
{ id: 1, name: null, availability_status: 'offline' },
{ id: 2, name: 'Zoe', availability_status: 'offline' },
];
expect(() => getAgentsByAvailability(agents, 'offline')).not.toThrow();
expect(
getAgentsByAvailability(agents, 'offline').map(agent => agent.id)
).toEqual([1, 2]);
});
});
describe('getSortedAgentsByAvailability', () => {
@@ -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."
}
@@ -26,25 +26,28 @@ const canDrilldown = computed(() => checkPermissions(['administrator']));
const selectedRange = ref('this_month');
const assistantId = computed(() => route.params.assistantId);
const stats = ref(null);
const isFetching = ref(false);
const metricStats = ref(null);
const faqStats = ref(null);
const isFetchingMetrics = ref(false);
// Increments on every fetch so a response (or retry) from a superseded
// range/assistant can't clobber the latest request's state.
let fetchToken = 0;
let abortController = null;
let metricsFetchToken = 0;
let faqStatsFetchToken = 0;
let metricsAbortController = null;
let faqStatsAbortController = null;
const fetchStats = async () => {
fetchToken += 1;
const token = fetchToken;
abortController?.abort();
abortController = new AbortController();
const { signal } = abortController;
stats.value = null;
isFetching.value = true;
const fetchMetrics = async () => {
metricsFetchToken += 1;
const token = metricsFetchToken;
metricsAbortController?.abort();
metricsAbortController = new AbortController();
const { signal } = metricsAbortController;
metricStats.value = null;
isFetchingMetrics.value = true;
const requestStats = () =>
CaptainAssistant.getStats({
const requestMetrics = () =>
CaptainAssistant.getMetrics({
assistantId: assistantId.value,
range: selectedRange.value,
signal,
@@ -52,25 +55,54 @@ const fetchStats = async () => {
let data = null;
try {
({ data } = await requestStats());
({ data } = await requestMetrics());
} catch {
// One silent retry before giving up, unless the request was aborted.
try {
if (token === fetchToken && !signal.aborted)
({ data } = await requestStats());
if (token === metricsFetchToken && !signal.aborted)
({ data } = await requestMetrics());
} catch {
data = null;
}
}
if (token !== fetchToken || signal.aborted) return;
stats.value = data;
isFetching.value = false;
if (token !== metricsFetchToken || signal.aborted) return;
metricStats.value = data;
isFetchingMetrics.value = false;
};
onUnmounted(() => abortController?.abort());
const fetchFaqStats = async () => {
faqStatsFetchToken += 1;
const token = faqStatsFetchToken;
faqStatsAbortController?.abort();
faqStatsAbortController = new AbortController();
const { signal } = faqStatsAbortController;
faqStats.value = null;
watch([selectedRange, assistantId], fetchStats, { immediate: true });
try {
const { data } = await CaptainAssistant.getFaqStats({
assistantId: assistantId.value,
signal,
});
if (token === faqStatsFetchToken && !signal.aborted) faqStats.value = data;
} catch {
if (token === faqStatsFetchToken && !signal.aborted) faqStats.value = null;
}
};
const summaryStats = computed(() => {
if (!metricStats.value || !faqStats.value) return null;
return { ...metricStats.value, knowledge: faqStats.value };
});
onUnmounted(() => {
metricsAbortController?.abort();
faqStatsAbortController?.abort();
});
watch([selectedRange, assistantId], fetchMetrics, { immediate: true });
watch(assistantId, fetchFaqStats, { immediate: true });
// `direction` says whether a rising trend is good ('up'), bad ('down'), or
// neutral, so we can colour the delta independently of its sign.
@@ -90,7 +122,7 @@ const formatDuration = hours =>
hours >= 100 ? `${Math.round(hours / 24)}d` : `${hours}h`;
const metricFor = (statKey, formatValue, direction, trendKind = 'percent') => {
const data = stats.value?.[statKey];
const data = metricStats.value?.[statKey];
if (!data) return { value: '—', trend: '', trendGood: null };
const sign = data.trend > 0 ? '+' : '';
@@ -184,9 +216,9 @@ const closeDrilldown = () => {
<div class="flex flex-col gap-6 pb-8">
<InboxBanner />
<CoverageBanner :knowledge="stats?.knowledge" />
<CoverageBanner :knowledge="faqStats ?? undefined" />
<WelcomeCard :range="selectedRange" :stats="stats" />
<WelcomeCard :range="selectedRange" :stats="summaryStats" />
<div
class="grid grid-cols-1 gap-px overflow-hidden border rounded-xl sm:grid-cols-2 lg:grid-cols-3 bg-n-weak border-n-weak"
@@ -199,13 +231,15 @@ const closeDrilldown = () => {
:trend="metric.trend"
:hint="metric.hint"
:trend-good="metric.trendGood"
:loading="isFetching"
:clickable="canDrilldown && Boolean(metric.metric) && !isFetching"
:loading="isFetchingMetrics"
:clickable="
canDrilldown && Boolean(metric.metric) && !isFetchingMetrics
"
@click="openDrilldown(metric)"
/>
</div>
<KnowledgeCard :knowledge="stats?.knowledge" />
<KnowledgeCard :knowledge="faqStats ?? undefined" />
<QuickLinks />
</div>
@@ -135,7 +135,7 @@ onMounted(() => {
<BaseTableCell class="max-w-0">
<div class="flex items-center gap-4 min-w-0">
<Avatar
:name="bot.name"
:name="bot.name || ''"
:src="bot.thumbnail"
:size="40"
class="flex-shrink-0"
@@ -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(() => {
<ImportDetailHeader
:data-import="dataImport"
:is-refreshing="isRefreshing"
:is-retrying="isRetrying"
:is-abandoning="isAbandoning"
:is-polling="isPolling"
@refresh="fetchImport({ manual: true })"
@retry="retryImport"
@abandon="abandonImport"
/>
</template>
@@ -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')"
/>
<Button
v-if="dataImport?.stalled"
outline
slate
size="sm"
icon="i-lucide-rotate-ccw"
:is-loading="isRetrying"
:disabled="isAbandoning"
:label="$t('DATA_IMPORTS.TABLE.RETRY')"
@click="$emit('retry')"
/>
<Button
v-if="canAbandonImport"
ruby
size="sm"
:is-loading="isAbandoning"
:disabled="isRetrying"
:label="$t('DATA_IMPORTS.TABLE.ABANDON')"
@click="$emit('abandon')"
/>
@@ -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: `
<button
:data-label="label"
:data-icon="icon"
:data-loading="isLoading"
@click="$emit('click')"
>
{{ label }}
</button>
`,
};
const BaseSettingsHeaderStub = {
template: `
<section>
<slot name="title" />
<slot name="description" />
</section>
`,
};
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);
});
});
@@ -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: `
<main>
<slot name="header" />
<slot name="body" />
</main>
`,
};
const ImportDetailHeaderStub = {
name: 'ImportDetailHeader',
props: {
dataImport: {
type: Object,
default: null,
},
},
emits: ['retry'],
template: `
<button
v-if="dataImport?.stalled"
data-test="retry"
@click="$emit('retry')"
/>
`,
};
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();
});
});
@@ -68,6 +68,10 @@ const isAgentBot = computed(
() => props.selectedItem?.assignee_type === 'AgentBot'
);
const selectedItemName = computed(() =>
!props.selectedItem?.name && isAgentBot.value ? '-' : props.selectedItem?.name
);
const selectedThumbnail = computed(
() => props.selectedItem?.thumbnail || props.selectedItem?.avatar_url
);
@@ -95,16 +99,16 @@ const selectedThumbnail = computed(
<h4
v-else
class="items-center overflow-hidden text-sm leading-tight whitespace-nowrap text-ellipsis text-n-slate-12"
:title="selectedItem.name"
:title="selectedItemName"
>
{{ selectedItem.name }}
{{ selectedItemName }}
</h4>
</div>
<Avatar
v-if="hasValue && hasThumbnail && (isAgentBot || !hasIcon)"
:src="selectedThumbnail"
:status="selectedItem.availability_status"
:name="selectedItem.name"
:name="selectedItemName"
:icon-name="isAgentBot ? 'i-lucide-bot' : undefined"
:size="24"
hide-offline-status
@@ -53,7 +53,9 @@ export default {
computed: {
filteredOptions() {
return this.options.filter(option => {
return option.name.toLowerCase().includes(this.search.toLowerCase());
return (option.name || '')
.toLowerCase()
.includes(this.search.toLowerCase());
});
},
noResult() {
+37 -12
View File
@@ -11,6 +11,8 @@ import { IFrameHelper } from '../helpers/utils';
import { CHATWOOT_ON_START_CONVERSATION } from '../constants/sdkEvents';
import { emitter } from 'shared/helpers/mitt';
const TRANSCRIPT_COOLDOWN_MS = 15000;
export default {
components: {
ChatInputWrap,
@@ -24,6 +26,9 @@ export default {
data() {
return {
inReplyTo: null,
isSendingTranscript: false,
transcriptCooldown: false,
transcriptCooldownTimer: null,
};
},
computed: {
@@ -57,6 +62,9 @@ export default {
mounted() {
emitter.on(BUS_EVENTS.TOGGLE_REPLY_TO_MESSAGE, this.toggleReplyTo);
},
beforeUnmount() {
clearTimeout(this.transcriptCooldownTimer);
},
methods: {
...mapActions('conversation', ['sendMessage', 'sendAttachment']),
...mapActions('conversationAttributes', ['getAttributes']),
@@ -90,19 +98,35 @@ export default {
toggleReplyTo(message) {
this.inReplyTo = message;
},
startTranscriptCooldown() {
this.transcriptCooldown = true;
clearTimeout(this.transcriptCooldownTimer);
this.transcriptCooldownTimer = setTimeout(() => {
this.transcriptCooldown = false;
}, TRANSCRIPT_COOLDOWN_MS);
},
async sendTranscript() {
if (this.hasEmail) {
try {
await sendEmailTranscript();
emitter.emit(BUS_EVENTS.SHOW_ALERT, {
message: this.$t('EMAIL_TRANSCRIPT.SEND_EMAIL_SUCCESS'),
type: 'success',
});
} catch (error) {
emitter.$emit(BUS_EVENTS.SHOW_ALERT, {
message: this.$t('EMAIL_TRANSCRIPT.SEND_EMAIL_ERROR'),
});
}
if (
!this.hasEmail ||
this.isSendingTranscript ||
this.transcriptCooldown
) {
return;
}
this.isSendingTranscript = true;
try {
await sendEmailTranscript();
this.startTranscriptCooldown();
emitter.emit(BUS_EVENTS.SHOW_ALERT, {
message: this.$t('EMAIL_TRANSCRIPT.SEND_EMAIL_SUCCESS'),
type: 'success',
});
} catch (error) {
emitter.emit(BUS_EVENTS.SHOW_ALERT, {
message: this.$t('EMAIL_TRANSCRIPT.SEND_EMAIL_ERROR'),
});
} finally {
this.isSendingTranscript = false;
}
},
},
@@ -144,6 +168,7 @@ export default {
v-if="showEmailTranscriptButton"
type="clear"
class="font-normal"
:disabled="isSendingTranscript || transcriptCooldown"
@click="sendTranscript"
>
{{ $t('EMAIL_TRANSCRIPT.BUTTON_TEXT') }}
+5
View File
@@ -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
+4
View File
@@ -19,6 +19,10 @@ class DataImportPolicy < ApplicationPolicy
show?
end
def retry_import?
show?
end
def abandon?
show?
end
+433 -117
View File
@@ -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
@@ -44,8 +63,9 @@ class DataImports::Intercom::Importer
def finish!
return if @data_import.reload.abandoned?
has_failures = @data_import.import_errors.non_skip_logs.exists? || @data_import.import_errors.failed.exists?
status = has_failures ? :completed_with_errors : :completed
error_count = @data_import.import_errors.non_skip_logs.count + @data_import.import_errors.failed.count
@stats['errors']['count'] = error_count
status = error_count.positive? ? :completed_with_errors : :completed
@data_import.update!(
status: status,
completed_at: Time.current,
@@ -56,14 +76,16 @@ class DataImports::Intercom::Importer
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 +94,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 +110,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 +177,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 +189,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 +213,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 +409,291 @@ 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)
entry = with_query_timeout_retry { batch_builder.perform([source_entry]).entries.first }
import_message(conversation, contact, entry)
rescue StandardError => e
fail_message(conversation, source_entry[:source_id], source_entry[:part], e)
end
def import_message(conversation, contact, entry, 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)
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 +705,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 +718,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 +782,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 +907,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 +915,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 +957,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 +1062,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 +1210,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)
@@ -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
@@ -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
@@ -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
+3 -1
View File
@@ -66,7 +66,8 @@ Rails.application.routes.draw do
resources :assistants do
member do
post :playground
get :stats
get :metrics
get :faq_stats
get :summary
get :drilldown
end
@@ -229,6 +230,7 @@ Rails.application.routes.draw do
end
member do
post :start
post :retry, action: :retry_import
post :abandon
get :error_logs
get :skip_logs
@@ -37,6 +37,23 @@ class Captain::AssistantStatsBuilder
build_metrics(current, previous)
end
# Approved/pending FAQ counts and the document total in a single round trip.
def faq_stats
approved, pending, documents = Captain::AssistantResponse.by_assistant(assistant.id).reorder(nil).pick(
Arel.sql("COUNT(*) FILTER (WHERE status = #{Captain::AssistantResponse.statuses['approved']})"),
Arel.sql("COUNT(*) FILTER (WHERE status = #{Captain::AssistantResponse.statuses['pending']})"),
Arel.sql("(SELECT COUNT(*) FROM captain_documents WHERE assistant_id = #{assistant.id.to_i})")
)
total = approved + pending
{
approved: approved,
pending: pending,
documents: documents,
coverage: total.zero? ? 0 : (approved.to_f / total * 100).round
}
end
private
attr_reader :window
@@ -56,8 +73,7 @@ class Captain::AssistantStatsBuilder
handoff_rate: pack(current[:handoff], previous[:handoff], :point),
hours_saved: pack(current[:hours_saved], previous[:hours_saved], :percent),
reopen_rate: pack(current[:reopen], previous[:reopen], :point),
conversation_depth: pack(current[:depth], previous[:depth], :absolute),
knowledge: knowledge
conversation_depth: pack(current[:depth], previous[:depth], :absolute)
}
end
@@ -73,7 +89,7 @@ class Captain::AssistantStatsBuilder
auto_resolution: rate(resolution[:resolved], handled),
handoff: rate(resolution[:handoff], handled),
hours_saved: (public_count * SECONDS_SAVED_PER_REPLY / 3600.0).round,
reopen: reopen_rate(range),
reopen: reopen_rate(range, resolution[:resolved]),
depth: depth_conversations.zero? ? 0 : (public_count.to_f / depth_conversations).round(1)
}
end
@@ -158,7 +174,9 @@ class Captain::AssistantStatsBuilder
# derived from the assistant's handled conversations (not current inbox membership) so a later
# inbox reassignment doesn't drop historical resolves, and covers both the evaluated (inference)
# and time-based (bot) resolve paths so the denominator matches auto_resolution_rate.
def reopen_rate(range)
def reopen_rate(range, resolved_count)
return 0 if resolved_count.zero?
resolved_scope = account.reporting_events
.where(name: RESOLVED_EVENT_NAMES, created_at: range,
conversation_id: handled_scope(range).select(:conversation_id))
@@ -178,24 +196,7 @@ class Captain::AssistantStatsBuilder
'ON resolves.conversation_id = reporting_events.conversation_id ' \
'AND reporting_events.event_end_time >= resolves.event_end_time')
.distinct.count('reporting_events.conversation_id')
rate(reopened, resolved_scope.distinct.count(:conversation_id))
end
# Approved/pending FAQ counts and the document total in a single round trip.
def knowledge
approved, pending, documents = Captain::AssistantResponse.by_assistant(assistant.id).reorder(nil).pick(
Arel.sql("COUNT(*) FILTER (WHERE status = #{Captain::AssistantResponse.statuses['approved']})"),
Arel.sql("COUNT(*) FILTER (WHERE status = #{Captain::AssistantResponse.statuses['pending']})"),
Arel.sql("(SELECT COUNT(*) FROM captain_documents WHERE assistant_id = #{assistant.id.to_i})")
)
total = approved + pending
{
approved: approved,
pending: pending,
documents: documents,
coverage: total.zero? ? 0 : (approved.to_f / total * 100).round
}
rate(reopened, resolved_count)
end
def rate(numerator, denominator)
@@ -1,7 +1,7 @@
class Api::V1::Accounts::Captain::AssistantsController < Api::V1::Accounts::BaseController
before_action -> { check_authorization(Captain::Assistant) }
before_action :set_assistant, only: [:show, :update, :destroy, :playground, :stats, :summary, :drilldown]
before_action :set_assistant, only: [:show, :update, :destroy, :playground, :metrics, :faq_stats, :summary, :drilldown]
def index
@assistants = account_assistants.ordered
@@ -42,10 +42,14 @@ class Api::V1::Accounts::Captain::AssistantsController < Api::V1::Accounts::Base
@tools = assistant.available_agent_tools
end
def stats
def metrics
render json: Captain::AssistantStatsBuilder.new(@assistant, params[:range], params[:timezone_offset]).metrics
end
def faq_stats
render json: Captain::AssistantStatsBuilder.new(@assistant).faq_stats
end
def summary
window = Captain::AssistantStatsWindow.new(params[:range], params[:timezone_offset])
result = cached_or_generated_summary(window, summary_stats)
@@ -1,6 +1,8 @@
module Enterprise::Api::V1::Accounts::AgentsController
def create
super
return if @agent.blank?
associate_agent_with_custom_role
end
@@ -7,7 +7,11 @@ class Captain::AssistantPolicy < ApplicationPolicy
true
end
def stats?
def metrics?
true
end
def faq_stats?
true
end
+1 -4
View File
@@ -34,7 +34,7 @@
"@amplitude/analytics-browser": "^2.11.10",
"@breezystack/lamejs": "^1.2.7",
"@chatwoot/ninja-keys": "1.2.3",
"@chatwoot/prosemirror-schema": "1.3.22",
"@chatwoot/prosemirror-schema": "1.3.23",
"@chatwoot/utils": "^0.0.56",
"@formkit/core": "^1.7.2",
"@formkit/vue": "^1.7.2",
@@ -86,9 +86,6 @@
"mitt": "^3.0.1",
"opus-recorder": "^8.0.5",
"pinia": "^3.0.4",
"prosemirror-commands": "^1.7.1",
"prosemirror-inputrules": "^1.4.0",
"prosemirror-schema-list": "^1.5.1",
"qrcode": "^1.5.4",
"semver": "7.6.3",
"snakecase-keys": "^8.0.1",
+6 -22
View File
@@ -25,8 +25,8 @@ importers:
specifier: 1.2.3
version: 1.2.3
'@chatwoot/prosemirror-schema':
specifier: 1.3.22
version: 1.3.22
specifier: 1.3.23
version: 1.3.23
'@chatwoot/utils':
specifier: ^0.0.56
version: 0.0.56
@@ -180,15 +180,6 @@ importers:
pinia:
specifier: ^3.0.4
version: 3.0.4(typescript@5.6.2)(vue@3.5.12(typescript@5.6.2))
prosemirror-commands:
specifier: ^1.7.1
version: 1.7.1
prosemirror-inputrules:
specifier: ^1.4.0
version: 1.4.0
prosemirror-schema-list:
specifier: ^1.5.1
version: 1.5.1
qrcode:
specifier: ^1.5.4
version: 1.5.4
@@ -461,8 +452,8 @@ packages:
'@chatwoot/ninja-keys@1.2.3':
resolution: {integrity: sha512-xM8d9P5ikDMZm2WbaCTk/TW5HFauylrU3cJ75fq5je6ixKwyhl/0kZbVN/vbbZN4+AUX/OaSIn6IJbtCgIF67g==}
'@chatwoot/prosemirror-schema@1.3.22':
resolution: {integrity: sha512-0r+PT8xhQLCKCpoV9k9XVTTRECs/0Nr37wbcLsRS7yvc7WkF9FY05z2hGCRJReWmTOcmmshHtb042LVP+MyB/w==}
'@chatwoot/prosemirror-schema@1.3.23':
resolution: {integrity: sha512-jGxbWELCdlVI64BJiE1wT84ekJHYDXXKiluQIKT3aKPEjPwMR48umKF3A0yHjKoR7IIxCC9oM77TvXOA0ebLtw==}
'@chatwoot/utils@0.0.56':
resolution: {integrity: sha512-A6dmPLfTSrW4qYNY73btyi4PqpfzcXRSaucscZTQdzNqF6G/QUdgnBmHtho8HeiYby/kSHXaSxLJj+0dx3yEQQ==}
@@ -4001,9 +3992,6 @@ packages:
prosemirror-tables@1.5.0:
resolution: {integrity: sha512-VMx4zlYWm7aBlZ5xtfJHpqa3Xgu3b7srV54fXYnXgsAcIGRqKSrhiK3f89omzzgaAgAtDOV4ImXnLKhVfheVNQ==}
prosemirror-transform@1.10.0:
resolution: {integrity: sha512-9UOgFSgN6Gj2ekQH5CTDJ8Rp/fnKR2IkYfGdzzp5zQMFsS4zDllLVx/+jGcX86YlACpG7UR5fwAXiWzxqWtBTg==}
prosemirror-transform@1.12.0:
resolution: {integrity: sha512-GxboyN4AMIsoHNtz5uf2r2Ru551i5hWeCMD6E2Ib4Eogqoub0NflniaBPVQ4MrGE5yZ8JV9tUHg9qcZTTrcN4w==}
@@ -5136,7 +5124,7 @@ snapshots:
hotkeys-js: 3.8.7
lit: 2.2.6
'@chatwoot/prosemirror-schema@1.3.22':
'@chatwoot/prosemirror-schema@1.3.23':
dependencies:
markdown-it-sup: 2.0.0
prosemirror-commands: 1.7.1
@@ -9035,7 +9023,7 @@ snapshots:
dependencies:
prosemirror-model: 1.22.3
prosemirror-state: 1.4.3
prosemirror-transform: 1.10.0
prosemirror-transform: 1.12.0
prosemirror-state@1.4.3:
dependencies:
@@ -9051,10 +9039,6 @@ snapshots:
prosemirror-transform: 1.12.0
prosemirror-view: 1.34.1
prosemirror-transform@1.10.0:
dependencies:
prosemirror-model: 1.22.3
prosemirror-transform@1.12.0:
dependencies:
prosemirror-model: 1.22.3
+18
View File
@@ -23,6 +23,12 @@ RSpec.describe AgentBuilder, type: :model do
end
describe '#perform' do
it 'locks the account while checking and creating the agent' do
expect(account).to receive(:with_lock).and_call_original
agent_builder.perform
end
context 'when user does not exist' do
it 'creates a new user' do
expect { agent_builder.perform }.to change(User, :count).by(1)
@@ -67,5 +73,17 @@ RSpec.describe AgentBuilder, type: :model do
expect(user.encrypted_password).not_to be_empty
end
end
context 'when the account has reached its agent limit' do
before do
allow(account).to receive(:usage_limits).and_return({ agents: account.account_users.count })
end
it 'raises a limit exceeded error without creating a user' do
expect { agent_builder.perform }.to raise_error(described_class::LimitExceededError, described_class::LIMIT_EXCEEDED_MESSAGE)
expect(User.from_email(email)).to be_nil
end
end
end
end
@@ -13,7 +13,9 @@ RSpec.describe 'Linear Integration API', type: :request do
end
describe 'DELETE /api/v1/accounts/:account_id/integrations/linear' do
it 'deletes the linear integration' do
let(:admin) { create(:user, account: account, role: :administrator) }
it 'deletes the linear integration when the user is an administrator' do
# Stub the HTTP call to Linear's revoke endpoint
allow(HTTParty).to receive(:post).with(
'https://api.linear.app/oauth/revoke',
@@ -21,11 +23,19 @@ RSpec.describe 'Linear Integration API', type: :request do
).and_return(instance_double(HTTParty::Response, success?: true))
delete "/api/v1/accounts/#{account.id}/integrations/linear",
headers: agent.create_new_auth_token,
headers: admin.create_new_auth_token,
as: :json
expect(response).to have_http_status(:ok)
expect(account.hooks.count).to eq(0)
end
it 'returns unauthorized for an agent and keeps the integration' do
delete "/api/v1/accounts/#{account.id}/integrations/linear",
headers: agent.create_new_auth_token,
as: :json
expect(response).to have_http_status(:unauthorized)
expect(account.hooks.count).to eq(1)
end
end
describe 'GET /api/v1/accounts/:account_id/integrations/linear/teams' do
@@ -159,15 +159,17 @@ RSpec.describe 'Shopify Integration API', type: :request do
end
describe 'DELETE /api/v1/accounts/:account_id/integrations/shopify' do
let(:admin) { create(:user, account: account, role: :administrator) }
before do
create(:integrations_hook, :shopify, account: account)
end
context 'when it is an authenticated user' do
context 'when it is an administrator' do
it 'deletes the shopify integration' do
expect do
delete "/api/v1/accounts/#{account.id}/integrations/shopify",
headers: agent.create_new_auth_token,
headers: admin.create_new_auth_token,
as: :json
end.to change { account.hooks.count }.by(-1)
@@ -175,6 +177,18 @@ RSpec.describe 'Shopify Integration API', type: :request do
end
end
context 'when it is an agent' do
it 'returns unauthorized and keeps the integration' do
expect do
delete "/api/v1/accounts/#{account.id}/integrations/shopify",
headers: agent.create_new_auth_token,
as: :json
end.not_to(change { account.hooks.count })
expect(response).to have_http_status(:unauthorized)
end
end
context 'when it is an unauthenticated user' do
it 'returns unauthorized' do
delete "/api/v1/accounts/#{account.id}/integrations/shopify",
@@ -27,7 +27,7 @@ RSpec.describe Captain::AssistantStatsBuilder do
expect(metrics.keys).to contain_exactly(
:conversations_handled, :auto_resolution_rate, :handoff_rate,
:hours_saved, :reopen_rate, :conversation_depth, :knowledge
:hours_saved, :reopen_rate, :conversation_depth
)
expect(metrics[:conversations_handled]).to include(:current, :previous, :trend)
end
@@ -229,7 +229,7 @@ RSpec.describe Captain::AssistantStatsBuilder do
end
end
describe '#metrics knowledge' do
describe '#faq_stats' do
before do
create_list(:captain_assistant_response, 3, assistant: assistant, account: account, status: :approved)
create(:captain_assistant_response, assistant: assistant, account: account, status: :pending)
@@ -237,7 +237,7 @@ RSpec.describe Captain::AssistantStatsBuilder do
end
it 'returns approved, pending, document counts and coverage' do
knowledge = described_class.new(assistant, '30').metrics[:knowledge]
knowledge = described_class.new(assistant).faq_stats
expect(knowledge).to eq(approved: 3, pending: 1, documents: 2, coverage: 75)
end
@@ -245,7 +245,7 @@ RSpec.describe Captain::AssistantStatsBuilder do
it 'reports zero coverage when there are no responses' do
Captain::AssistantResponse.where(assistant: assistant).delete_all
knowledge = described_class.new(assistant, '30').metrics[:knowledge]
knowledge = described_class.new(assistant).faq_stats
expect(knowledge[:coverage]).to eq(0)
end
@@ -21,6 +21,27 @@ RSpec.describe 'Agents API', type: :request do
expect(response).to have_http_status(:payment_required)
expect(response.body).to include('Account limit exceeded. Please purchase more licenses')
end
it 'prevents adding an agent if the last seat is consumed before creation' do
account.update!(limits: { agents: account.account_users.count + 1 })
competing_agent_created = false
allow(AgentBuilder).to receive(:new).and_wrap_original do |method, *args|
unless competing_agent_created
create(:user, account: account, role: :agent)
competing_agent_created = true
end
method.call(*args)
end
post "/api/v1/accounts/#{account.id}/agents", params: params, headers: admin.create_new_auth_token, as: :json
expect(response).to have_http_status(:payment_required)
expect(response.body).to include('Account limit exceeded. Please purchase more licenses')
expect(User.from_email(params[:email])).to be_nil
expect(account.account_users.count).to eq(account.usage_limits[:agents])
end
end
end
@@ -12,7 +12,7 @@ RSpec.describe Captain::AssistantPolicy, type: :policy do
let(:administrator_context) { { user: administrator, account: account, account_user: account.account_users.first } }
let(:agent_context) { { user: agent, account: account, account_user: account.account_users.first } }
permissions :index?, :show?, :playground? do
permissions :index?, :show?, :playground?, :metrics?, :faq_stats? do
context 'when administrator' do
it { expect(assistant_policy).to permit(administrator_context, assistant) }
end
+32
View File
@@ -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
@@ -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
)
@@ -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,275 @@ 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 '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 +457,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 +507,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 +577,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 +595,165 @@ 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 '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 +791,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 +808,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
@@ -373,6 +918,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 +969,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 +1048,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 +1140,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 +1392,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 +1563,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 +1632,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 +1654,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 +1678,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
@@ -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' => '<p>Initial message</p>',
'created_at' => 1_700_000_000
},
'conversation_parts' => {
'conversation_parts' => [
{
'id' => 'part_1',
'part_type' => 'comment',
'body' => '<p>First reply</p>',
'created_at' => 1_700_000_000
},
{
'id' => 'part_2',
'part_type' => 'note',
'body' => '<p>Internal note</p>',
'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
@@ -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