diff --git a/docs/content/docs/guides/sequences.mdx b/docs/content/docs/guides/sequences.mdx index f12592942..ee896236a 100644 --- a/docs/content/docs/guides/sequences.mdx +++ b/docs/content/docs/guides/sequences.mdx @@ -221,6 +221,8 @@ A campaign-level toggle in the flow toolbar. When on and a contact replies, the Only inbound mail from that campaign contact to the sending mailbox counts. A threaded follow-up synchronized from Sent or Drafts stays outbound even though its `In-Reply-To` header points at an earlier campaign step. +When an upgrade repairs a reply marker that was provably created by a synchronized Sent message, the contact resumes at the next due step and the campaign reply totals update with it. Warmbly leaves ambiguous historical replies untouched. + With it off and more than one step, Warmbly warns you on the canvas: replies that don't match a reply branch would otherwise keep receiving cold follow-ups. Turning it on still runs your reply branches, so it is strictly safer. diff --git a/internal/app/advanced/incoming_reply_test.go b/internal/app/advanced/incoming_reply_test.go index 324da97e7..db97dd36c 100644 --- a/internal/app/advanced/incoming_reply_test.go +++ b/internal/app/advanced/incoming_reply_test.go @@ -66,8 +66,13 @@ func (r incomingReplyContactRepo) GetByID(context.Context, uuid.UUID) (*models.C type incomingReplyProgressRepo struct { repository.CampaignProgressRepository - replied int - latest *repository.CampaignSequencePair + replied int + latest *repository.CampaignSequencePair + sourceInbound bool +} + +func (r *incomingReplyProgressRepo) IsInboundReplySource(context.Context, uuid.UUID, uuid.UUID) (bool, error) { + return r.sourceInbound, nil } func (r *incomingReplyProgressRepo) GetLatestCampaignSequenceForContact(context.Context, uuid.UUID) (*repository.CampaignSequencePair, error) { @@ -82,7 +87,7 @@ func (r *incomingReplyProgressRepo) RecordReplyClassification(context.Context, u return nil } -func (r *incomingReplyProgressRepo) RecordEmailReplied(context.Context, uuid.UUID, uuid.UUID, uuid.UUID) error { +func (r *incomingReplyProgressRepo) RecordEmailReplied(context.Context, uuid.UUID, uuid.UUID, uuid.UUID, uuid.UUID, uuid.UUID) error { r.replied++ return nil } @@ -95,7 +100,7 @@ func (incomingReplyCampaignRepo) GetSequencesRoutingByCampaignID(context.Context func newIncomingReplyService(account *models.Email, senderContact *models.Contact, taskContact uuid.UUID) (*service, *incomingReplyProgressRepo) { taskID, campaignID, sequenceID := uuid.New(), uuid.New(), uuid.New() - progress := &incomingReplyProgressRepo{} + progress := &incomingReplyProgressRepo{sourceInbound: true} taskContactRecord := &models.Contact{ID: taskContact, Email: "task-contact@example.test"} if senderContact != nil && senderContact.ID == taskContact { taskContactRecord = senderContact @@ -195,6 +200,31 @@ func TestProcessIncomingReplyRejectsOutboundFolder(t *testing.T) { } } +func TestProcessIncomingReplyTrustsPersistedDirectionOverEventPayload(t *testing.T) { + orgID, accountID, contactID := uuid.New(), uuid.New(), uuid.New() + account := &models.Email{ID: accountID, OrganizationID: &orgID, Email: "sender@example.test"} + service, progress := newIncomingReplyService(account, &models.Contact{ + ID: contactID, Email: "recipient@example.test", + }, contactID) + progress.sourceInbound = false + + xerr := service.ProcessIncomingReply(context.Background(), accountID, &models.EmailMessageStoreData{ + ID: uuid.New(), + EmailID: accountID, + Folder: models.FolderInbox, + FromAddr: []string{"Recipient "}, + ToAddr: []string{"sender@example.test"}, + InReplyTo: []string{""}, + Subject: "Re: Hello", + }) + if xerr != nil { + t.Fatal(xerr) + } + if progress.replied != 0 { + t.Fatalf("RecordEmailReplied calls = %d, want 0 when the stored source is outbound", progress.replied) + } +} + func TestProcessIncomingReplyRequiresThreadSenderToMatchContact(t *testing.T) { orgID, accountID := uuid.New(), uuid.New() taskContactID, senderContactID := uuid.New(), uuid.New() diff --git a/internal/app/advanced/service.go b/internal/app/advanced/service.go index f5d5e7c1e..b2a9b06e7 100644 --- a/internal/app/advanced/service.go +++ b/internal/app/advanced/service.go @@ -1069,6 +1069,13 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid. if !msg.MayBeInbound() { return nil } + inbound, err := s.campaignProgressRepo.IsInboundReplySource(ctx, emailAccountID, msg.ID) + if err != nil { + return toErrx(err) + } + if !inbound { + return nil + } account, xerr := s.emailRepo.GetByID(ctx, emailAccountID) if xerr != nil { return xerr @@ -1219,7 +1226,9 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid. s.evidence.RecordEvidence(ctx, ctID, models.Step(&cID, &sID), kind, msg.ID.String(), "") } if !replyclassify.IsAutomated(replyResult.Class) { - _ = s.campaignProgressRepo.RecordEmailReplied(ctx, cID, ctID, sID) + if err := s.campaignProgressRepo.RecordEmailReplied(ctx, cID, ctID, sID, emailAccountID, msg.ID); err != nil { + return toErrx(err) + } _ = s.repo.MarkVariantEvent(ctx, cID, ctID, string(models.DeliverabilityEventReply)) // Live org-wide pulse: the team sees the reply land on the diff --git a/internal/infrastructure/db/migrations/000168_repair_sent_followup_replies.down.sql b/internal/infrastructure/db/migrations/000168_repair_sent_followup_replies.down.sql new file mode 100644 index 000000000..8907d3579 --- /dev/null +++ b/internal/infrastructure/db/migrations/000168_repair_sent_followup_replies.down.sql @@ -0,0 +1,2 @@ +-- One-way data repair: restoring reply markers proven to come from outbound +-- Sent-folder messages would restore issue #549, so this migration is a no-op. diff --git a/internal/infrastructure/db/migrations/000168_repair_sent_followup_replies.up.sql b/internal/infrastructure/db/migrations/000168_repair_sent_followup_replies.up.sql new file mode 100644 index 000000000..2cd2b6683 --- /dev/null +++ b/internal/infrastructure/db/migrations/000168_repair_sent_followup_replies.up.sql @@ -0,0 +1,182 @@ +-- Repair reply markers created when an outbound threaded follow-up was +-- consumed as though it were the contact's inbound reply (issue #549). +-- +-- The match is intentionally narrow. The reply stamp must be within two +-- seconds of a Sent-folder message to that contact, the message must answer +-- the exact campaign task whose progress was stamped, its classifier must be +-- inconclusive, and no reply-intent row from the contact may exist nearby. + +WITH false_replies AS ( + SELECT DISTINCT + p.campaign_id, + p.contact_id, + p.sequence_id, + p.dispatch_task_id, + p.replied_at, + ue.id AS source_message_id + FROM campaign_contact_progress p + JOIN campaigns campaign ON campaign.id = p.campaign_id + JOIN contacts c + ON c.id = p.contact_id + AND c.organization_id = campaign.organization_id + JOIN tasks parent_task ON parent_task.id = p.dispatch_task_id + JOIN email_accounts source_account + ON source_account.organization_id = campaign.organization_id + JOIN unibox_emails ue + ON ue.email_id = source_account.id + AND (ue.folder = 'sent' OR ue.provider_folder = 'sent') + AND p.replied_at BETWEEN ue.created_at + AND ue.created_at + INTERVAL '2 seconds' + AND EXISTS ( + SELECT 1 + FROM unnest(ue.in_reply_to) AS refs(message_id) + WHERE btrim(refs.message_id, '<> ') = btrim(parent_task.message_id, '<> ') + ) + AND EXISTS ( + SELECT 1 + FROM unnest(ue.to_addr) AS recipients(address) + WHERE lower(recipients.address) LIKE '%' || lower(c.email) || '%' + ) + WHERE p.replied_at IS NOT NULL + AND p.reply_class IN ('', 'unknown') + AND p.reply_confidence = 0 + AND p.reply_source = '' + AND NOT EXISTS ( + SELECT 1 + FROM reply_intents ri + WHERE ri.campaign_id = p.campaign_id + AND lower(ri.contact_email) = lower(c.email) + AND ri.created_at BETWEEN p.replied_at - INTERVAL '30 seconds' + AND p.replied_at + INTERVAL '30 seconds' + ) +) +DELETE FROM contact_verification_evidence evidence +USING false_replies false_reply +WHERE evidence.contact_id = false_reply.contact_id + AND evidence.kind = 'replied' + AND evidence.ref = false_reply.source_message_id::text; + +WITH false_replies AS ( + SELECT DISTINCT + p.campaign_id, + p.contact_id, + p.sequence_id, + p.dispatch_task_id, + p.replied_at, + ue.email_id AS source_account_id + FROM campaign_contact_progress p + JOIN campaigns campaign ON campaign.id = p.campaign_id + JOIN contacts c + ON c.id = p.contact_id + AND c.organization_id = campaign.organization_id + JOIN tasks parent_task ON parent_task.id = p.dispatch_task_id + JOIN email_accounts source_account + ON source_account.organization_id = campaign.organization_id + JOIN unibox_emails ue + ON ue.email_id = source_account.id + AND (ue.folder = 'sent' OR ue.provider_folder = 'sent') + AND p.replied_at BETWEEN ue.created_at + AND ue.created_at + INTERVAL '2 seconds' + AND EXISTS ( + SELECT 1 + FROM unnest(ue.in_reply_to) AS refs(message_id) + WHERE btrim(refs.message_id, '<> ') = btrim(parent_task.message_id, '<> ') + ) + AND EXISTS ( + SELECT 1 + FROM unnest(ue.to_addr) AS recipients(address) + WHERE lower(recipients.address) LIKE '%' || lower(c.email) || '%' + ) + WHERE p.replied_at IS NOT NULL + AND p.reply_class IN ('', 'unknown') + AND p.reply_confidence = 0 + AND p.reply_source = '' + AND NOT EXISTS ( + SELECT 1 + FROM reply_intents ri + WHERE ri.campaign_id = p.campaign_id + AND lower(ri.contact_email) = lower(c.email) + AND ri.created_at BETWEEN p.replied_at - INTERVAL '30 seconds' + AND p.replied_at + INTERVAL '30 seconds' + ) +) +DELETE FROM reply_intents intent +USING false_replies false_reply, email_accounts sending_account +WHERE intent.campaign_id = false_reply.campaign_id + AND intent.task_id = false_reply.dispatch_task_id + AND intent.created_at BETWEEN false_reply.replied_at - INTERVAL '30 seconds' + AND false_reply.replied_at + INTERVAL '30 seconds' + AND sending_account.id = false_reply.source_account_id + AND lower(intent.contact_email) IN ( + lower(sending_account.email), + lower(COALESCE(NULLIF(sending_account.send_as_email, ''), sending_account.email)) + ); + +WITH false_replies AS ( + SELECT DISTINCT + p.campaign_id, + p.contact_id, + p.sequence_id + FROM campaign_contact_progress p + JOIN campaigns campaign ON campaign.id = p.campaign_id + JOIN contacts c + ON c.id = p.contact_id + AND c.organization_id = campaign.organization_id + JOIN tasks parent_task ON parent_task.id = p.dispatch_task_id + JOIN email_accounts source_account + ON source_account.organization_id = campaign.organization_id + JOIN unibox_emails ue + ON ue.email_id = source_account.id + AND (ue.folder = 'sent' OR ue.provider_folder = 'sent') + AND p.replied_at BETWEEN ue.created_at + AND ue.created_at + INTERVAL '2 seconds' + AND EXISTS ( + SELECT 1 + FROM unnest(ue.in_reply_to) AS refs(message_id) + WHERE btrim(refs.message_id, '<> ') = btrim(parent_task.message_id, '<> ') + ) + AND EXISTS ( + SELECT 1 + FROM unnest(ue.to_addr) AS recipients(address) + WHERE lower(recipients.address) LIKE '%' || lower(c.email) || '%' + ) + WHERE p.replied_at IS NOT NULL + AND p.reply_class IN ('', 'unknown') + AND p.reply_confidence = 0 + AND p.reply_source = '' + AND NOT EXISTS ( + SELECT 1 + FROM reply_intents ri + WHERE ri.campaign_id = p.campaign_id + AND lower(ri.contact_email) = lower(c.email) + AND ri.created_at BETWEEN p.replied_at - INTERVAL '30 seconds' + AND p.replied_at + INTERVAL '30 seconds' + ) +), cleared AS ( + UPDATE campaign_contact_progress progress + SET replied_at = NULL, + reply_class = '', + reply_confidence = 0, + reply_source = '', + instant_fired = array_remove(instant_fired, 'reply') + FROM false_replies false_reply + WHERE progress.campaign_id = false_reply.campaign_id + AND progress.contact_id = false_reply.contact_id + AND progress.sequence_id = false_reply.sequence_id + RETURNING progress.campaign_id, progress.contact_id +) +UPDATE campaign_ab_assignments assignment +SET replied_at = NULL +WHERE EXISTS ( + SELECT 1 + FROM cleared + WHERE cleared.campaign_id = assignment.campaign_id + AND cleared.contact_id = assignment.contact_id + ) + AND NOT EXISTS ( + SELECT 1 + FROM campaign_contact_progress progress + WHERE progress.campaign_id = assignment.campaign_id + AND progress.contact_id = assignment.contact_id + AND progress.replied_at IS NOT NULL + ); diff --git a/internal/repository/campaign_reply_source_live_test.go b/internal/repository/campaign_reply_source_live_test.go new file mode 100644 index 000000000..1e5d98a50 --- /dev/null +++ b/internal/repository/campaign_reply_source_live_test.go @@ -0,0 +1,95 @@ +package repository + +import ( + "context" + "testing" + + "github.com/google/uuid" +) + +// A stored Sent row must never stamp campaign progress as replied (issue #549). +func TestLiveReplySourceUsesStoredDirectionAtTheWriteBoundary(t *testing.T) { + _, pool := liveContactDB(t) + f := newThreadParentFixture(t, pool) + step := f.step(1, "Hello", true) + parentTask := f.send(step, f.mailbox, "", "thread-1", 60) + ctx := context.Background() + + if _, err := pool.Exec(ctx, ` + INSERT INTO campaign_contact_progress (campaign_id, contact_id, sequence_id, dispatch_task_id, sent_at) + VALUES ($1, $2, $3, $4, NOW() - INTERVAL '1 minute') + `, f.campaign, f.contact, step, parentTask); err != nil { + t.Fatalf("insert progress: %v", err) + } + + sentID := uuid.New() + if _, err := pool.Exec(ctx, ` + INSERT INTO unibox_emails (id, user_id, email_id, folder, provider_folder) + VALUES ($1, $2, $3, 'sent', 'sent') + `, sentID, f.owner, f.other); err != nil { + t.Fatalf("insert sent source: %v", err) + } + t.Cleanup(func() { + if _, err := pool.Exec(context.Background(), `DELETE FROM unibox_emails WHERE id = $1`, sentID); err != nil { + t.Errorf("cleanup sent source: %v", err) + } + }) + + repo := NewCampaignProgressRepository(pool) + inbound, err := repo.IsInboundReplySource(ctx, f.other, sentID) + if err != nil { + t.Fatalf("verify sent source: %v", err) + } + if inbound { + t.Fatal("Sent-folder source was accepted as inbound") + } + if err := repo.RecordEmailReplied(ctx, f.campaign, f.contact, step, f.other, sentID); err != nil { + t.Fatalf("record sent source: %v", err) + } + + var replied bool + if err := pool.QueryRow(ctx, ` + SELECT replied_at IS NOT NULL + FROM campaign_contact_progress + WHERE campaign_id = $1 AND contact_id = $2 AND sequence_id = $3 + `, f.campaign, f.contact, step).Scan(&replied); err != nil { + t.Fatalf("read progress: %v", err) + } + if replied { + t.Fatal("Sent-folder source stamped replied_at") + } + + inboxID := uuid.New() + if _, err := pool.Exec(ctx, ` + INSERT INTO unibox_emails (id, user_id, email_id, folder, provider_folder) + VALUES ($1, $2, $3, 'inbox', 'inbox') + `, inboxID, f.owner, f.mailbox); err != nil { + t.Fatalf("insert inbound source: %v", err) + } + t.Cleanup(func() { + if _, err := pool.Exec(context.Background(), `DELETE FROM unibox_emails WHERE id = $1`, inboxID); err != nil { + t.Errorf("cleanup inbound source: %v", err) + } + }) + + inbound, err = repo.IsInboundReplySource(ctx, f.mailbox, inboxID) + if err != nil { + t.Fatalf("verify inbound source: %v", err) + } + if !inbound { + t.Fatal("Inbox source was rejected as outbound") + } + if err := repo.RecordEmailReplied(ctx, f.campaign, f.contact, step, f.mailbox, inboxID); err != nil { + t.Fatalf("record inbound source: %v", err) + } + if err := pool.QueryRow(ctx, ` + SELECT replied_at IS NOT NULL + FROM campaign_contact_progress + WHERE campaign_id = $1 AND contact_id = $2 AND sequence_id = $3 + `, f.campaign, f.contact, step).Scan(&replied); err != nil { + t.Fatalf("read inbound progress: %v", err) + } + if !replied { + t.Fatal("Inbox source did not stamp replied_at") + } +} diff --git a/internal/repository/pg_campaign_progress.go b/internal/repository/pg_campaign_progress.go index 8550a7a81..d636f78ba 100644 --- a/internal/repository/pg_campaign_progress.go +++ b/internal/repository/pg_campaign_progress.go @@ -183,7 +183,9 @@ type CampaignProgressRepository interface { // the reference point for telling an instant machine open or click from a // person's. GetStepSentAt(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) (*time.Time, error) - RecordEmailReplied(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error + // IsInboundReplySource confirms the stored unibox row is not outbound. + IsInboundReplySource(ctx context.Context, emailAccountID, messageID uuid.UUID) (bool, error) + RecordEmailReplied(ctx context.Context, campaignID, contactID, sequenceID, emailAccountID, messageID uuid.UUID) error RecordEmailBounced(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error RecordEmailComplained(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error @@ -787,8 +789,24 @@ func (r *campaignProgressRepository) GetStepSentAt(ctx context.Context, campaign return sentAt, nil } -// RecordEmailReplied records that a contact replied -func (r *campaignProgressRepository) RecordEmailReplied(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error { +// IsInboundReplySource verifies direction from the row stored by the consumer. +func (r *campaignProgressRepository) IsInboundReplySource(ctx context.Context, emailAccountID, messageID uuid.UUID) (bool, error) { + var inbound bool + err := r.db.QueryRow(ctx, ` + SELECT EXISTS ( + SELECT 1 + FROM unibox_emails + WHERE id = $1 + AND email_id = $2 + AND folder NOT IN ('sent', 'drafts') + AND provider_folder NOT IN ('sent', 'drafts') + ) + `, messageID, emailAccountID).Scan(&inbound) + return inbound, err +} + +// RecordEmailReplied records that a contact replied when its source is still inbound. +func (r *campaignProgressRepository) RecordEmailReplied(ctx context.Context, campaignID, contactID, sequenceID, emailAccountID, messageID uuid.UUID) error { query := ` UPDATE campaign_contact_progress SET replied_at = NOW() @@ -796,9 +814,17 @@ func (r *campaignProgressRepository) RecordEmailReplied(ctx context.Context, cam AND contact_id = $2 AND sequence_id = $3 AND replied_at IS NULL + AND EXISTS ( + SELECT 1 + FROM unibox_emails + WHERE id = $4 + AND email_id = $5 + AND folder NOT IN ('sent', 'drafts') + AND provider_folder NOT IN ('sent', 'drafts') + ) ` - _, err := r.db.Exec(ctx, query, campaignID, contactID, sequenceID) + _, err := r.db.Exec(ctx, query, campaignID, contactID, sequenceID, messageID, emailAccountID) return err }