feat(search): index contact labels
This commit is contained in:
@@ -40,7 +40,7 @@ Hermes task landing checklist:
|
||||
| Source plan task family | Tracker owner | Current state | Follow-up trigger |
|
||||
| --- | --- | --- | --- |
|
||||
| Search config, engine interface, Meilisearch client, index settings, DB fallback | Phase 1, B6 | Review. Meilisearch-first engine, bootstrap/settings, no-live tests, release-mode DB fallback rejection, and `cmd/reindex_search` guard are represented here. | Reopen only from live Meilisearch gate failure, frontend search payload drift, or `reference/chatwoot` search behavior not covered by B6. |
|
||||
| Search indexing for conversations, messages, contacts, companies, articles | P5.3, B4, B6 | Review. Service-layer indexing hooks, document builders, CRM search through Meilisearch, and entity/global payload shapes are tracked in the Commit Ledger. | Reopen when a model mutation path writes searchable data without indexing or when live smoke shows stale search results. |
|
||||
| Search indexing for conversations, messages, contacts, companies, articles | P5.3, B4, B6 | Review. Service-layer indexing hooks, document builders, CRM search through Meilisearch, contact label documents, and entity/global payload shapes are tracked in the Commit Ledger. | Reopen when a model mutation path writes searchable data without indexing or when live smoke shows stale search results. |
|
||||
| Automation rule CRUD, listener triggers, condition/action execution, external webhook/email actions | B9, P5.4 | Review. Rule payloads, trigger coverage, execution logs, retryable external actions, durable team email, and enterprise `add_sla` action are tracked. | Reopen from fresh reference evidence for unsupported action params, trigger events, or failed B12 automation smoke. |
|
||||
| Macro CRUD and execute side effects | B9, P5.5 | Review. Macro serializers, permissions, action params, and display-ID conversation execution are tracked. | Reopen from macro frontend smoke failures or new reference action semantics. |
|
||||
| CSAT account reports, public submit/update, resolved-conversation sends, downloads, channel templates | B8, P5.6 | Review. CSAT response/report payloads, public lock behavior, resolve-triggered survey send, CSV download, and queued channel templates are tracked. | Reopen from CSAT live smoke failures, provider-template status drift, or new reference survey settings. |
|
||||
@@ -49,11 +49,11 @@ Hermes task landing checklist:
|
||||
|
||||
## Current Baseline
|
||||
|
||||
- Current tracking checkpoint: 2026-06-07 P3.103 contact bulk-action parity, prepared as `feat(crm): queue contact bulk actions`.
|
||||
- Latest implementation checkpoint: this checkpoint, prepared as `feat(crm): queue contact bulk actions`.
|
||||
- Latest documentation/tooling checkpoint: this tracker update records Chatwoot contact bulk action enqueueing, label add/remove side effects, delete side effects, and raw empty `200 OK` account bulk-action responses.
|
||||
- Current tracking checkpoint: 2026-06-07 P5.3c contact label search indexing, prepared as `feat(search): index contact labels`.
|
||||
- Latest implementation checkpoint: this checkpoint, prepared as `feat(search): index contact labels`.
|
||||
- Latest documentation/tooling checkpoint: this tracker update records Meilisearch contact label document fields, contact-label mutation indexing hooks, and contact bulk-action search index fan-out.
|
||||
- Plan landing status: complete for the current known Hermes plans and user-confirmed scope. Future work should update this file directly instead of opening a parallel tracker.
|
||||
- Worktree status at this implementation checkpoint: reused dashboard contact bulk actions from `ContactsIndex.vue` and `api/bulkActions.js` now enqueue Chatwoot-shaped contact jobs on the `medium` queue, normalize lower-case `type`, allow empty `ids` like the reference controller/job path, return empty `200 OK`, add/remove contact labels account-scoped, and soft-delete only selected contacts in the current account. This retains P3.102 enterprise account route parity, P6.1 webhook placeholder burn-down, P3.101 WhatsApp call route-parameter parity, P3.100 dashboard app route-parameter parity, and prior checkpoints. Live API/browser/enterprise smoke still needs the full PostgreSQL/Redis/Meilisearch/GoChat/Vite/Chrome stack.
|
||||
- Worktree status at this implementation checkpoint: contact search documents now carry label arrays for Meilisearch filter parity, durable contact indexing reloads current label assignments before replay, `UpdateLabels` triggers contact reindexing, and contact bulk-action label/delete jobs enqueue search index/delete follow-up jobs. This retains P3.103 contact bulk-action parity, P3.102 enterprise account route parity, P6.1 webhook placeholder burn-down, and prior checkpoints. Live API/browser/enterprise smoke still needs the full PostgreSQL/Redis/Meilisearch/GoChat/Vite/Chrome stack.
|
||||
- Next executable implementation checkpoint: continue Phase 2/3 drift audit for the next reused-frontend mismatch, or run B12 live smoke when the full PostgreSQL/Redis/Meilisearch/GoChat/Vite/Chrome stack is available. Re-run Phase 6 placeholder audit after future route/smoke changes.
|
||||
- `go test ./...` passes when run outside the restricted socket sandbox for the latest implementation baseline; the latest docs/tooling checkpoint verified `scripts/parity_frontend_smoke.sh --check` with workspace-local temp/cache dirs after `/tmp` was full.
|
||||
- Route dump succeeds with `972` registered routes after enterprise account route tracking.
|
||||
@@ -156,6 +156,7 @@ This table is the shortest authoritative handoff view. If an older lower section
|
||||
|
||||
| Priority | Workstream | Current state | Next checkpoint | Commit close rule |
|
||||
| --- | --- | --- | --- | --- |
|
||||
| 0 | P5.3c contact label search indexing | Implemented for Meilisearch contact-label parity: contact search documents now include `labels`, durable `search:index` contact replay loads current `contact_labels`/`tags`, direct contact label replacement reindexes the contact, and contact bulk-action label/delete workers enqueue search index/delete follow-ups after account-scoped side effects. | Keep in Review; reopen from live Meilisearch gate, contact label search smoke, or a fresh mutation path that changes contact labels without indexing. | Focused service/search tests passed; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. No route artifacts change. |
|
||||
| 0 | P3.103 contact bulk actions | Implemented for reused CRM contact index actions: `POST /api/v1/accounts/:account_id/bulk_actions` now normalizes Chatwoot `type`, accepts Contact payloads without requiring `action_name`, enqueues durable `contact:bulk_action` jobs on the `medium` queue, returns empty `200 OK`, applies label add/remove account-scoped through `contact_labels`/`tags`, soft-deletes only current-account selected contacts, and treats unknown contact operations as no-op success like `Contacts::BulkActionService`. | Keep in Review; reopen from B12 contacts smoke or fresh reference evidence for exact Pundit authorization, label serializer, deleted-association cleanup, or notification side effects beyond the inspected controller/service/frontend contract. | Focused BulkActionHandler tests passed; focused conversation maintenance contact bulk-action worker tests passed; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. No route artifacts change. |
|
||||
| 0 | P3.102 enterprise account limits and billing routes | Implemented for reused enterprise account frontend calls: `GET/POST /enterprise/api/v1/accounts/:account_id/{limits,checkout,subscription,toggle_deletion,topup_checkout}` are now registered and tracked from `routes.rb:523-527`; `limits` returns Chatwoot-shaped usage data for agents, Captain documents/responses, and default-plan conversation/non-web-inbox counts; `toggle_deletion` mutates `marked_for_deletion_at` and `marked_for_deletion_reason`; `subscription` persists the `is_creating_customer` guard when no Stripe customer exists; checkout/top-up return explicit billing-provider errors instead of missing routes. | Keep in Review; reopen from B12 enterprise account smoke or fresh reference evidence for actual Stripe session creation, cloud-env gating, plan-config defaults, or account deletion notification/cancellation jobs beyond the local persisted boundary. | Focused EnterpriseAccountHandler tests passed; route dump/parity artifacts regenerated to `TOTAL: 972` and `435 exact, 0 method-compatible, 9 parameter-compatible, 0 missing out of 444`; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. |
|
||||
| 0 | P6.1 webhook placeholder fallback burn-down | Implemented for Phase 6 placeholder cleanup: `chatwootParityStub` and the unused webhook placeholder helper are removed from product code. Public webhook nil-handler guards now return explicit `503 { error: "webhook provider unavailable", message: "webhook handler is not configured" }` responses instead of `501 not implemented` placeholder bodies, while wired provider handlers still own real Telegram, WhatsApp, TikTok, LINE, Twilio, Twitter, Instagram, and Shopify behavior. | Keep in Review; reopen from placeholder audit or webhook smoke if a frontend/provider-reachable route returns placeholder/not-implemented content or if a provider handler is missing from normal bootstrap. | Focused router tests cover route boot and nil-handler fallback body/status; placeholder audit refreshed; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. No route artifacts change. |
|
||||
@@ -392,6 +393,7 @@ This ledger records the committed parity checkpoints that future slices should b
|
||||
|
||||
| Commit | Scope | Verification summary | Follow-up state |
|
||||
| --- | --- | --- | --- |
|
||||
| `feat(search): index contact labels` | Advances P5.3/B6 for mandatory Meilisearch contact label parity. Contact documents now carry label arrays, durable contact index replay preloads current label assignments, `ContactService.UpdateLabels` reindexes contacts after replacement, and contact bulk-action add/remove/delete jobs enqueue search index/delete follow-up jobs through the durable indexer. | Focused service/search tests passed; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. No route artifacts change. | Move P5.3c to Review; continue Phase 2/3 drift audit, Phase 6 placeholder audit, or B12 live smoke. |
|
||||
| `feat(crm): queue contact bulk actions` | Advances P3.103/P5.12 with Chatwoot contact bulk-action parity. GoChat now accepts reused frontend Contact bulk payloads for label add/remove and delete, normalizes `type` like the Rails controller, enqueues `contact:bulk_action` on the `medium` queue when a WorkerPool is configured, returns empty `200 OK`, and replays account-scoped label/delete side effects through durable workers. | Focused BulkActionHandler tests and conversation maintenance worker contact bulk-action tests passed; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. No route artifacts change. | Move P3.103 to Review; continue Phase 2/3 drift audit, Phase 6 placeholder audit, or B12 live smoke. |
|
||||
| `fix(routes): align dashboard app ids` | Advances P3.100 with exact Chatwoot dashboard app member route parameter parity. GoChat now registers `GET/PATCH/PUT/DELETE /api/v1/accounts/:account_id/dashboard_apps/:id`, keeps handler compatibility with legacy local `:dashboard_app_id`, and regenerates route parity artifacts. | Focused DashboardAppHandler and router tests passed; route dump/parity regenerated to `TOTAL: 967` and `425 exact, 0 method-compatible, 14 parameter-compatible, 0 missing out of 439`; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. | Move P3.100 to Review; continue reducing remaining parameter-compatible rows or run B12 live smoke. |
|
||||
| `feat(conversations): align destroy job` | Advances P3.99 with Chatwoot conversation destroy parity. GoChat now returns empty `200 OK` for `DELETE /conversations/:conversation_id`, wires a durable low-priority `conversation:delete_object` job into the WorkerPool, and lets that job perform the existing soft-delete/event/search cleanup while retaining a synchronous fallback without workers. | Focused ConversationService delete/job tests and ConversationHandler delete response tests passed; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. No route artifacts change. | Move P3.99 to Review; continue Phase 2/3 drift audit, Phase 6 placeholder audit, or B12 live smoke. |
|
||||
@@ -2115,7 +2117,7 @@ Tracking table:
|
||||
| --- | --- | --- | --- | --- | --- |
|
||||
| P5.1 | Implement durable worker core and job model. | `reference/chatwoot/app/jobs/application_job.rb`, `mutex_application_job.rb` | `internal/worker/worker.go`, `internal/model/background_job.go`, `migrations/000026_add_background_jobs.*.sql` | Job table, worker persistence API, enqueue API, worker loop, retry/backoff, scheduled jobs, mutex/idempotency keys, dead-letter state, and restart pickup tests exist. | Review by `feat(worker): add durable background jobs` |
|
||||
| P5.2 | Route async dispatcher events through durable jobs. | `event_dispatcher_job.rb`, Chatwoot async dispatcher listeners | `internal/dispatch/dispatcher.go`, `internal/channel/dispatcher.go` | Heavy listeners can enqueue durable jobs without changing sync listener behavior; tests cover sync vs async routing and replay. | Review by `feat(dispatch): queue async events durably` |
|
||||
| P5.3 | Move Meilisearch indexing and reindex fan-out into retryable jobs. | Meilisearch plan plus Chatwoot callbacks/jobs that index searchable records | search services, contact/company/conversation indexing hooks | Create/update/delete indexing survives handler success, retries on Meilisearch failure, and optional live Meilisearch gate remains green. | Review by `feat(search): queue index updates durably` |
|
||||
| P5.3 | Move Meilisearch indexing and reindex fan-out into retryable jobs. | Meilisearch plan plus Chatwoot callbacks/jobs that index searchable records | search services, contact/company/conversation indexing hooks | Create/update/delete/label indexing survives handler success, retries on Meilisearch failure, and optional live Meilisearch gate remains green. | Review by `feat(search): queue index updates durably` and `feat(search): index contact labels` |
|
||||
| P5.4 | Queue automation webhook and transcript delivery. | `webhook_job.rb`, automation action execution services | `internal/automation/action_delivery.go`, `internal/automation/action_service.go` | Existing timeout/retry fakeable delivery is invoked by durable jobs; logs preserve attempt metadata and idempotency. | Review by `feat(automation): queue external action deliveries` |
|
||||
| P5.5 | Queue delayed automation actions and macro execution. | `trigger_scheduled_items_job.rb`, `macros_execution_job.rb` | automation rule listener, macro service | Delayed actions execute after schedule time if reference params exist; macro execute supports multi-conversation job fan-out, and repeated workers do not duplicate side effects. | Review by `feat(automation): queue macro and csat jobs` and `feat(automation): close delayed action parity`; current reference exposes no delayed automation action params |
|
||||
| P5.6 | Queue CSAT survey sends and channel-specific templates. | CSAT listener/services, WhatsApp/Twilio template services/jobs | `internal/automation/csat_survey_listener.go`, `internal/service/csat_template_service.go`, channel send services | Resolve-triggered CSAT send is durable; WhatsApp/Twilio template delivery and failure states are fakeable and observable. | Review by `feat(automation): queue macro and csat jobs` and `feat(csat): queue channel templates` |
|
||||
@@ -2164,10 +2166,12 @@ P5.3 current checkpoint:
|
||||
|
||||
- `DurableSearchIndexer` now wraps the Meilisearch-backed `SearchService` write side. Conversation, message, contact, company, and article create/update/delete hooks enqueue `search:index` jobs on the `search` queue when a `WorkerPool` is configured.
|
||||
- Worker replay reloads each account-scoped record before indexing so Meilisearch receives current database state instead of stale request-time payloads.
|
||||
- Contact replay now also reloads current `contact_labels`/`tags` so Meilisearch contact documents carry the same label filter field used by reused CRM contact/search flows.
|
||||
- Contact label replacement and contact bulk-action label/delete jobs now reindex or delete contact documents after the label side effect, keeping Meilisearch label filters fresh.
|
||||
- Missing records during an index replay are converted into delegate delete calls, which keeps delayed create/update jobs from resurrecting documents after a database delete.
|
||||
- Search read paths remain Meilisearch-first through `searchService`; contact and company services still use the live `SearchService` reader instead of the durable wrapper.
|
||||
- No-worker construction still falls back to synchronous indexing for focused tests and development paths that do not start the durable worker.
|
||||
- Remaining P5.3 work is limited to explicit reindex fan-out/CLI scheduling and optional live Meilisearch environment gates; normal service-layer writes are now durable.
|
||||
- Remaining P5.3 work is limited to explicit reindex fan-out/CLI scheduling and optional live Meilisearch environment gates; normal service-layer writes and contact-label writes are now durable.
|
||||
|
||||
P5.3 verification:
|
||||
|
||||
@@ -2791,3 +2795,4 @@ Verification milestone gates:
|
||||
- 2026-06-07: P6.1 webhook placeholder fallback checkpoint prepared as `fix(webhooks): replace parity stubs`; refreshed the Phase 6 placeholder audit and burned down the remaining `chatwootParityStub` nil-handler fallbacks in `internal/router/router.go`. Public webhook routes now return explicit `503 webhook provider unavailable` JSON when a provider handler is not configured instead of `501 not implemented` placeholder bodies, and the unused placeholder helper is removed. Focused router tests cover boot and nil-handler fallback behavior; placeholder audit shows no `chatwootParityStub` product-code matches; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. No route artifacts change.
|
||||
- 2026-06-07: P3.102 enterprise account limits checkpoint prepared as `feat(enterprise): align account limits API`; audited Chatwoot enterprise `AccountsController#limits/#toggle_deletion/#subscription/#checkout/#topup_checkout`, `BillingHelper`, `Enterprise::Account::PlanUsageAndLimits`, reused dashboard `api/enterprise/account.js`, and `routes.rb:523-527`. GoChat now registers the enterprise account route family under `/enterprise/api/v1/accounts/:account_id`, returns Chatwoot-shaped account limit payloads, persists scheduled deletion custom attributes, records subscription customer-creation guards, and exposes explicit local billing-provider errors for checkout/top-up paths. Focused EnterpriseAccountHandler tests passed; route dump/parity regenerated to `TOTAL: 972` and `435 exact, 0 method-compatible, 9 parameter-compatible, 0 missing out of 444`; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed.
|
||||
- 2026-06-07: P3.103 contact bulk-action checkpoint prepared as `feat(crm): queue contact bulk actions`; audited Chatwoot `BulkActionsController`, `Contacts::BulkActionJob`, `Contacts::BulkActionService`, reused dashboard `api/bulkActions.js`, and `ContactsIndex.vue` label/delete callers. GoChat account bulk actions now accept Contact payloads for label add/remove and delete, normalize lower-case type values, enqueue durable `contact:bulk_action` jobs on the `medium` queue, return empty `200 OK`, replay account-scoped contact label add/remove through tags/contact_labels, soft-delete only selected current-account contacts, and preserve no-op success for unknown contact bulk payloads. Focused BulkActionHandler and conversation maintenance worker contact bulk-action tests passed; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. No route artifacts change.
|
||||
- 2026-06-07: P5.3c contact label search-index checkpoint prepared as `feat(search): index contact labels`; audited Meilisearch contact document/filter behavior, contact label endpoints, and the new contact bulk-action mutation path. GoChat contact search documents now include `labels`, durable contact index replay preloads current contact labels from `contact_labels`/`tags`, direct contact label replacement reindexes the contact, and contact bulk-action label/delete jobs enqueue search index/delete follow-ups through the durable search indexer. Focused service/search tests passed; full `go test ./...` passed outside the restricted socket sandbox; `git diff --check` passed. No route artifacts change.
|
||||
|
||||
@@ -679,6 +679,7 @@ func Bootstrap(env string) (*App, error) {
|
||||
}
|
||||
searchService := search.NewSearchServiceWithEngine(searchEngine, searchRepo)
|
||||
searchIndexer := service.NewDurableSearchIndexer(db, workerPool, searchService)
|
||||
service.RegisterContactBulkActionSearchIndexer(workerPool, db, searchIndexer)
|
||||
conversationService.SetSearchIndexer(searchIndexer)
|
||||
messageService.SetSearchIndexer(searchIndexer)
|
||||
contactService.SetSearchIndexer(searchIndexer)
|
||||
|
||||
@@ -29,6 +29,7 @@ type Contact struct {
|
||||
SourceID string `gorm:"size:255" json:"source_id,omitempty"`
|
||||
CompanyID *uint `json:"company_id,omitempty"`
|
||||
LastActivityAt *int64 `gorm:"index" json:"last_activity_at,omitempty"`
|
||||
Labels []string `gorm:"-" json:"labels,omitempty"`
|
||||
}
|
||||
|
||||
func (Contact) TableName() string { return "contacts" }
|
||||
|
||||
@@ -372,6 +372,7 @@ func ContactDocument(contact model.Contact) SearchDocument {
|
||||
ContactSource: contact.ContactType,
|
||||
ContactType: contact.ContactType,
|
||||
ContactHasDetails: content != "",
|
||||
Labels: contact.Labels,
|
||||
CreatedAtTS: timestamp(contact.CreatedAt),
|
||||
UpdatedAtTS: timestamp(contact.UpdatedAt),
|
||||
LastActivityAtTS: timestampPtr(contact.LastActivityAt),
|
||||
|
||||
@@ -140,6 +140,9 @@ func TestContactDocumentSetsResolvedScopeFields(t *testing.T) {
|
||||
assert.Equal(t, "lead", doc.ContactSource)
|
||||
assert.True(t, doc.ContactHasDetails)
|
||||
|
||||
labelled := ContactDocument(model.Contact{Base: model.Base{ID: 7}, AccountID: 2, Name: "Labelled", Labels: []string{"vip", "trial"}})
|
||||
assert.Equal(t, []string{"vip", "trial"}, labelled.Labels)
|
||||
|
||||
anonymous := ContactDocument(model.Contact{Base: model.Base{ID: 6}, AccountID: 2, Name: "Anonymous"})
|
||||
assert.False(t, anonymous.ContactHasDetails)
|
||||
}
|
||||
|
||||
@@ -62,11 +62,26 @@ func (s *ContactService) SetWorkerPool(wp *worker.WorkerPool) {
|
||||
|
||||
func (s *ContactService) indexContact(ctx context.Context, contact *model.Contact) {
|
||||
if s.searchIndexer != nil {
|
||||
contact = s.contactWithLabels(ctx, contact)
|
||||
logSearchIndexError("contact", contact.ID, s.searchIndexer.IndexContact(ctx, contact))
|
||||
s.indexContactConversations(ctx, contact)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *ContactService) contactWithLabels(ctx context.Context, contact *model.Contact) *model.Contact {
|
||||
if contact == nil || contact.ID == 0 || s == nil || s.repo == nil {
|
||||
return contact
|
||||
}
|
||||
labels, err := s.GetLabels(ctx, contact.AccountID, contact.ID)
|
||||
if err != nil {
|
||||
applogger.L().Warnf("search index contact labels load failed for contact %d: %v", contact.ID, err)
|
||||
return contact
|
||||
}
|
||||
copy := *contact
|
||||
copy.Labels = labels
|
||||
return ©
|
||||
}
|
||||
|
||||
func (s *ContactService) indexContactConversations(ctx context.Context, contact *model.Contact) {
|
||||
if s == nil || s.repo == nil || s.searchIndexer == nil || contact == nil || contact.ID == 0 {
|
||||
return
|
||||
@@ -1441,6 +1456,12 @@ func (s *ContactService) UpdateLabels(ctx context.Context, accountID, contactID
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
contact, err := s.repo.FindByAccountAndID(ctx, accountID, contactID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
contact.Labels = normalized
|
||||
s.indexContact(ctx, contact)
|
||||
return normalized, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -93,6 +93,14 @@ func RegisterConversationMaintenanceJobs(wp *worker.WorkerPool, db *gorm.DB) {
|
||||
registerConversationMaintenanceJobsWithNow(wp, db, time.Now)
|
||||
}
|
||||
|
||||
func RegisterContactBulkActionSearchIndexer(wp *worker.WorkerPool, db *gorm.DB, indexer SearchIndexer) {
|
||||
if wp == nil || db == nil {
|
||||
return
|
||||
}
|
||||
runner := &conversationMaintenanceRunner{wp: wp, db: db, now: time.Now, searchIndexer: indexer}
|
||||
wp.Register(TaskTypeContactBulkAction, runner.performContactBulkAction)
|
||||
}
|
||||
|
||||
func registerConversationMaintenanceJobsWithNow(wp *worker.WorkerPool, db *gorm.DB, now func() time.Time) {
|
||||
if wp == nil || db == nil {
|
||||
return
|
||||
@@ -170,9 +178,10 @@ func scheduledItemsIdempotencyKey(scheduledAt time.Time) string {
|
||||
}
|
||||
|
||||
type conversationMaintenanceRunner struct {
|
||||
wp *worker.WorkerPool
|
||||
db *gorm.DB
|
||||
now func() time.Time
|
||||
wp *worker.WorkerPool
|
||||
db *gorm.DB
|
||||
now func() time.Time
|
||||
searchIndexer SearchIndexer
|
||||
}
|
||||
|
||||
func (r *conversationMaintenanceRunner) performScheduledTriggerItems(ctx context.Context, job *model.BackgroundJob) error {
|
||||
@@ -407,39 +416,56 @@ func (r *conversationMaintenanceRunner) performContactBulkAction(ctx context.Con
|
||||
|
||||
switch {
|
||||
case payload.Params.ActionName == "delete":
|
||||
return r.db.WithContext(ctx).
|
||||
Where("account_id = ? AND id IN ?", payload.AccountID, payload.Params.IDs).
|
||||
Delete(&model.Contact{}).Error
|
||||
contactIDs, err := scopedContactIDs(ctx, r.db, payload.AccountID, payload.Params.IDs)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(contactIDs) == 0 {
|
||||
return nil
|
||||
}
|
||||
if err := r.db.WithContext(ctx).
|
||||
Where("account_id = ? AND id IN ?", payload.AccountID, contactIDs).
|
||||
Delete(&model.Contact{}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
r.deleteContactSearchIndexes(ctx, payload.AccountID, contactIDs)
|
||||
return nil
|
||||
case len(payload.Params.Labels.Add) > 0:
|
||||
return bulkAddContactLabels(ctx, r.db, payload.AccountID, payload.Params.IDs, payload.Params.Labels.Add)
|
||||
contactIDs, err := bulkAddContactLabels(ctx, r.db, payload.AccountID, payload.Params.IDs, payload.Params.Labels.Add)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return r.indexContactSearchDocuments(ctx, payload.AccountID, contactIDs)
|
||||
case len(payload.Params.Labels.Remove) > 0:
|
||||
return bulkRemoveContactLabels(ctx, r.db, payload.AccountID, payload.Params.IDs, payload.Params.Labels.Remove)
|
||||
contactIDs, err := bulkRemoveContactLabels(ctx, r.db, payload.AccountID, payload.Params.IDs, payload.Params.Labels.Remove)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return r.indexContactSearchDocuments(ctx, payload.AccountID, contactIDs)
|
||||
default:
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
func bulkAddContactLabels(ctx context.Context, db *gorm.DB, accountID uint, contactIDs []uint, labels []string) error {
|
||||
func bulkAddContactLabels(ctx context.Context, db *gorm.DB, accountID uint, contactIDs []uint, labels []string) ([]uint, error) {
|
||||
labels = normalizeContactServiceLabels(labels)
|
||||
if len(contactIDs) == 0 || len(labels) == 0 {
|
||||
return nil
|
||||
return nil, nil
|
||||
}
|
||||
var scopedContactIDs []uint
|
||||
if err := db.WithContext(ctx).Model(&model.Contact{}).
|
||||
Where("account_id = ? AND id IN ?", accountID, contactIDs).
|
||||
Pluck("id", &scopedContactIDs).Error; err != nil {
|
||||
return err
|
||||
scopedIDs, err := scopedContactIDs(ctx, db, accountID, contactIDs)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(scopedContactIDs) == 0 {
|
||||
return nil
|
||||
if len(scopedIDs) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
err = db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
for _, label := range labels {
|
||||
tag := model.Tag{AccountID: accountID, Name: label}
|
||||
if err := tx.Where("account_id = ? AND name = ?", accountID, label).FirstOrCreate(&tag).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
for _, contactID := range scopedContactIDs {
|
||||
for _, contactID := range scopedIDs {
|
||||
contactLabel := model.ContactLabel{AccountID: accountID, ContactID: contactID, TagID: tag.ID}
|
||||
if err := tx.Where("account_id = ? AND contact_id = ? AND tag_id = ?", accountID, contactID, tag.ID).FirstOrCreate(&contactLabel).Error; err != nil {
|
||||
return err
|
||||
@@ -448,25 +474,75 @@ func bulkAddContactLabels(ctx context.Context, db *gorm.DB, accountID uint, cont
|
||||
}
|
||||
return nil
|
||||
})
|
||||
return scopedIDs, err
|
||||
}
|
||||
|
||||
func bulkRemoveContactLabels(ctx context.Context, db *gorm.DB, accountID uint, contactIDs []uint, labels []string) error {
|
||||
func bulkRemoveContactLabels(ctx context.Context, db *gorm.DB, accountID uint, contactIDs []uint, labels []string) ([]uint, error) {
|
||||
labels = normalizeContactServiceLabels(labels)
|
||||
if len(contactIDs) == 0 || len(labels) == 0 {
|
||||
return nil
|
||||
return nil, nil
|
||||
}
|
||||
scopedIDs, err := scopedContactIDs(ctx, db, accountID, contactIDs)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(scopedIDs) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
var tagIDs []uint
|
||||
if err := db.WithContext(ctx).Model(&model.Tag{}).
|
||||
Where("account_id = ? AND name IN ?", accountID, labels).
|
||||
Pluck("id", &tagIDs).Error; err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
if len(tagIDs) == 0 {
|
||||
return scopedIDs, nil
|
||||
}
|
||||
err = db.WithContext(ctx).
|
||||
Where("account_id = ? AND contact_id IN ? AND tag_id IN ?", accountID, scopedIDs, tagIDs).
|
||||
Delete(&model.ContactLabel{}).Error
|
||||
return scopedIDs, err
|
||||
}
|
||||
|
||||
func scopedContactIDs(ctx context.Context, db *gorm.DB, accountID uint, contactIDs []uint) ([]uint, error) {
|
||||
if len(contactIDs) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
var ids []uint
|
||||
err := db.WithContext(ctx).Model(&model.Contact{}).
|
||||
Where("account_id = ? AND id IN ?", accountID, contactIDs).
|
||||
Pluck("id", &ids).Error
|
||||
return ids, err
|
||||
}
|
||||
|
||||
func (r *conversationMaintenanceRunner) indexContactSearchDocuments(ctx context.Context, accountID uint, contactIDs []uint) error {
|
||||
if r.searchIndexer == nil || len(contactIDs) == 0 {
|
||||
return nil
|
||||
}
|
||||
return db.WithContext(ctx).
|
||||
Where("account_id = ? AND contact_id IN ? AND tag_id IN ?", accountID, contactIDs, tagIDs).
|
||||
Delete(&model.ContactLabel{}).Error
|
||||
labelsByContactID, err := contactLabelsForSearch(ctx, r.db, accountID, contactIDs)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var contacts []model.Contact
|
||||
if err := r.db.WithContext(ctx).Where("account_id = ? AND id IN ?", accountID, contactIDs).Find(&contacts).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
for i := range contacts {
|
||||
contacts[i].Labels = labelsByContactID[contacts[i].ID]
|
||||
if err := r.searchIndexer.IndexContact(ctx, &contacts[i]); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *conversationMaintenanceRunner) deleteContactSearchIndexes(ctx context.Context, accountID uint, contactIDs []uint) {
|
||||
if r.searchIndexer == nil {
|
||||
return
|
||||
}
|
||||
for _, contactID := range contactIDs {
|
||||
logSearchIndexError("contact", contactID, r.searchIndexer.DeleteContact(ctx, accountID, contactID))
|
||||
}
|
||||
}
|
||||
|
||||
func parseBulkActionTime(value string) (time.Time, bool) {
|
||||
|
||||
@@ -387,6 +387,49 @@ func TestConversationMaintenanceJobsContactBulkAction(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestConversationMaintenanceJobsContactBulkActionQueuesSearchIndex(t *testing.T) {
|
||||
now := time.Date(2026, 6, 6, 0, 30, 0, 0, time.UTC)
|
||||
db := setupServiceTestDB(t)
|
||||
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return now }))
|
||||
registerConversationMaintenanceJobsWithNow(wp, db, func() time.Time { return now })
|
||||
delegate := &recordingDurableSearchIndexer{}
|
||||
searchIndexer := NewDurableSearchIndexer(db, wp, delegate)
|
||||
RegisterContactBulkActionSearchIndexer(wp, db, searchIndexer)
|
||||
|
||||
account := createTestAccount(t, db)
|
||||
contact := createTestContact(t, db, account.ID)
|
||||
|
||||
_, err := EnqueueContactBulkAction(context.Background(), wp, account.ID, 42, ContactBulkActionParams{
|
||||
Type: "Contact",
|
||||
IDs: []uint{contact.ID},
|
||||
Labels: ConversationBulkActionLabels{Add: []string{"vip"}},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue contact bulk action: %v", err)
|
||||
}
|
||||
processRequiredJob(t, wp, "contact bulk label add")
|
||||
processRequiredJob(t, wp, "contact search index")
|
||||
|
||||
if got := delegate.indexedContactLabels[contact.ID]; len(got) != 1 || got[0] != "vip" {
|
||||
t.Fatalf("expected contact search reindex with vip label, got %#v", got)
|
||||
}
|
||||
|
||||
_, err = EnqueueContactBulkAction(context.Background(), wp, account.ID, 42, ContactBulkActionParams{
|
||||
Type: "Contact",
|
||||
ActionName: "delete",
|
||||
IDs: []uint{contact.ID},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue contact delete: %v", err)
|
||||
}
|
||||
processRequiredJob(t, wp, "contact bulk delete")
|
||||
processRequiredJob(t, wp, "contact search delete")
|
||||
|
||||
if len(delegate.deletedContacts) != 1 || delegate.deletedContacts[0] != contact.ID {
|
||||
t.Fatalf("expected contact search delete, got %#v", delegate.deletedContacts)
|
||||
}
|
||||
}
|
||||
|
||||
func createTestOneoffCampaign(t *testing.T, db *gorm.DB, accountID, inboxID, contactID uint, scheduledAt time.Time) *campaign.Campaign {
|
||||
t.Helper()
|
||||
c := &campaign.Campaign{
|
||||
|
||||
@@ -14,6 +14,7 @@ import (
|
||||
type mockServiceSearchIndexer struct {
|
||||
indexed []string
|
||||
deleted []string
|
||||
indexedContactLabels map[uint][]string
|
||||
indexedConversationIDs []uint
|
||||
indexedConversationNames []string
|
||||
}
|
||||
@@ -46,6 +47,12 @@ func (m *mockServiceSearchIndexer) DeleteMessage(ctx context.Context, accountID
|
||||
|
||||
func (m *mockServiceSearchIndexer) IndexContact(ctx context.Context, contact *model.Contact) error {
|
||||
m.indexed = append(m.indexed, "contact")
|
||||
if contact != nil {
|
||||
if m.indexedContactLabels == nil {
|
||||
m.indexedContactLabels = map[uint][]string{}
|
||||
}
|
||||
m.indexedContactLabels[contact.ID] = append([]string(nil), contact.Labels...)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -155,6 +162,25 @@ func TestContactService_SearchIndexHooksReindexContactConversations(t *testing.T
|
||||
assert.Equal(t, []string{"Ada Lovelace"}, indexer.indexedConversationNames)
|
||||
}
|
||||
|
||||
func TestContactService_UpdateLabelsSearchIndexHookIncludesLabels(t *testing.T) {
|
||||
db := setupServiceTestDB(t)
|
||||
repo := NewContactService(repository.NewContactRepo(db), nil, nil)
|
||||
indexer := &mockServiceSearchIndexer{}
|
||||
repo.SetSearchIndexer(indexer)
|
||||
|
||||
account := createTestAccount(t, db)
|
||||
contact, err := repo.Create(context.Background(), account.ID, CreateContactRequest{Name: "Ada", Email: "ada@example.com"})
|
||||
require.NoError(t, err)
|
||||
indexer.indexed = nil
|
||||
|
||||
labels, err := repo.UpdateLabels(context.Background(), account.ID, contact.ID, []string{"vip", "trial"})
|
||||
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, []string{"vip", "trial"}, labels)
|
||||
assert.Equal(t, []string{"contact"}, indexer.indexed)
|
||||
assert.Equal(t, []string{"vip", "trial"}, indexer.indexedContactLabels[contact.ID])
|
||||
}
|
||||
|
||||
func TestCompanyService_SearchIndexHooks(t *testing.T) {
|
||||
db, _, _, _, svc := setupCompanyServiceTest(t)
|
||||
indexer := &mockServiceSearchIndexer{}
|
||||
|
||||
@@ -77,7 +77,14 @@ func (i *DurableSearchIndexer) IndexContact(ctx context.Context, contact *model.
|
||||
return nil
|
||||
}
|
||||
return i.enqueueOrIndex(ctx, "contact", contact.AccountID, contact.ID, func() error {
|
||||
return i.delegate.IndexContact(ctx, contact)
|
||||
enriched, err := i.loadContact(ctx, contact.AccountID, contact.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if enriched == nil {
|
||||
return nil
|
||||
}
|
||||
return i.delegate.IndexContact(ctx, enriched)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -175,11 +182,14 @@ func (i *DurableSearchIndexer) performIndex(ctx context.Context, payload searchI
|
||||
}
|
||||
return i.delegate.IndexMessage(ctx, item)
|
||||
case "contact":
|
||||
var item model.Contact
|
||||
if err := i.load(ctx, payload, &item); err != nil {
|
||||
item, err := i.loadContact(ctx, payload.AccountID, payload.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return i.delegate.IndexContact(ctx, &item)
|
||||
if item == nil {
|
||||
return nil
|
||||
}
|
||||
return i.delegate.IndexContact(ctx, item)
|
||||
case "company":
|
||||
var item model.Company
|
||||
if err := i.load(ctx, payload, &item); err != nil {
|
||||
@@ -200,6 +210,49 @@ func (i *DurableSearchIndexer) performIndex(ctx context.Context, payload searchI
|
||||
}
|
||||
}
|
||||
|
||||
func (i *DurableSearchIndexer) loadContact(ctx context.Context, accountID, id uint) (*model.Contact, error) {
|
||||
if i.db == nil {
|
||||
return nil, nil
|
||||
}
|
||||
var item model.Contact
|
||||
err := i.db.WithContext(ctx).Where("id = ? AND account_id = ?", id, accountID).First(&item).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil, i.performDelete(ctx, searchIndexJob{Operation: "delete", Entity: "contact", AccountID: accountID, ID: id})
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
labels, err := contactLabelsForSearch(ctx, i.db, accountID, []uint{id})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
item.Labels = labels[id]
|
||||
return &item, nil
|
||||
}
|
||||
|
||||
func contactLabelsForSearch(ctx context.Context, db *gorm.DB, accountID uint, contactIDs []uint) (map[uint][]string, error) {
|
||||
result := map[uint][]string{}
|
||||
if db == nil || len(contactIDs) == 0 {
|
||||
return result, nil
|
||||
}
|
||||
var rows []struct {
|
||||
ContactID uint
|
||||
Name string
|
||||
}
|
||||
if err := db.WithContext(ctx).Table("contact_labels").
|
||||
Select("contact_labels.contact_id, tags.name").
|
||||
Joins("JOIN tags ON tags.id = contact_labels.tag_id AND tags.account_id = contact_labels.account_id").
|
||||
Where("contact_labels.account_id = ? AND contact_labels.contact_id IN ?", accountID, contactIDs).
|
||||
Order("contact_labels.created_at ASC, tags.name ASC").
|
||||
Scan(&rows).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, row := range rows {
|
||||
result[row.ContactID] = append(result[row.ContactID], row.Name)
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (i *DurableSearchIndexer) loadMessage(ctx context.Context, accountID, id uint) (*model.Message, error) {
|
||||
if i.db == nil {
|
||||
return nil, nil
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
type recordingDurableSearchIndexer struct {
|
||||
indexedConversations []model.Conversation
|
||||
indexedContacts []uint
|
||||
indexedContactLabels map[uint][]string
|
||||
indexedArticles []model.Article
|
||||
deletedContacts []uint
|
||||
err error
|
||||
@@ -39,6 +40,10 @@ func (r *recordingDurableSearchIndexer) DeleteMessage(ctx context.Context, accou
|
||||
|
||||
func (r *recordingDurableSearchIndexer) IndexContact(ctx context.Context, contact *model.Contact) error {
|
||||
r.indexedContacts = append(r.indexedContacts, contact.ID)
|
||||
if r.indexedContactLabels == nil {
|
||||
r.indexedContactLabels = map[uint][]string{}
|
||||
}
|
||||
r.indexedContactLabels[contact.ID] = append([]string(nil), contact.Labels...)
|
||||
return r.err
|
||||
}
|
||||
|
||||
@@ -99,6 +104,33 @@ func TestDurableSearchIndexerQueuesAndReplaysContactIndex(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestDurableSearchIndexerPreloadsContactLabels(t *testing.T) {
|
||||
db := setupServiceTestDB(t)
|
||||
account := createTestAccount(t, db)
|
||||
contact := &model.Contact{AccountID: account.ID, Name: "Label Me", Email: "label@example.com"}
|
||||
if err := db.Create(contact).Error; err != nil {
|
||||
t.Fatalf("create contact: %v", err)
|
||||
}
|
||||
tag := model.Tag{AccountID: account.ID, Name: "vip"}
|
||||
if err := db.Create(&tag).Error; err != nil {
|
||||
t.Fatalf("create tag: %v", err)
|
||||
}
|
||||
if err := db.Create(&model.ContactLabel{AccountID: account.ID, ContactID: contact.ID, TagID: tag.ID}).Error; err != nil {
|
||||
t.Fatalf("create contact label: %v", err)
|
||||
}
|
||||
delegate := &recordingDurableSearchIndexer{}
|
||||
indexer := NewDurableSearchIndexer(db, nil, delegate)
|
||||
|
||||
if err := indexer.IndexContact(context.Background(), contact); err != nil {
|
||||
t.Fatalf("index contact: %v", err)
|
||||
}
|
||||
|
||||
labels := delegate.indexedContactLabels[contact.ID]
|
||||
if len(labels) != 1 || labels[0] != "vip" {
|
||||
t.Fatalf("expected preloaded contact labels, got %#v", labels)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDurableSearchIndexerPreloadsConversationContact(t *testing.T) {
|
||||
db := setupServiceTestDB(t)
|
||||
account := createTestAccount(t, db)
|
||||
|
||||
Reference in New Issue
Block a user