From c91e89980f265a5f8ebcf69e6dc3caade55cfc2f Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Wed, 16 Sep 2026 21:25:43 -0700 Subject: [PATCH] feat: make reply recording atomic and require exact repair recipient matches --- internal/app/advanced/incoming_reply_test.go | 51 ++++++++++++++++--- internal/app/advanced/service.go | 12 +++-- ...000168_repair_sent_followup_replies.up.sql | 18 +++++-- .../campaign_reply_source_live_test.go | 19 ++++++- internal/repository/pg_campaign_progress.go | 12 +++-- 5 files changed, 94 insertions(+), 18 deletions(-) diff --git a/internal/app/advanced/incoming_reply_test.go b/internal/app/advanced/incoming_reply_test.go index db97dd36c..9c17055c6 100644 --- a/internal/app/advanced/incoming_reply_test.go +++ b/internal/app/advanced/incoming_reply_test.go @@ -12,6 +12,8 @@ import ( type incomingReplyAdvancedRepo struct { repository.AdvancedOutreachRepository + marked int + intents int } func (incomingReplyAdvancedRepo) GetOutreachSettings(context.Context, uuid.UUID) (*models.AdvancedOutreachSettings, error) { @@ -19,11 +21,13 @@ func (incomingReplyAdvancedRepo) GetOutreachSettings(context.Context, uuid.UUID) return &settings, nil } -func (incomingReplyAdvancedRepo) MarkVariantEvent(context.Context, uuid.UUID, uuid.UUID, string) error { +func (r *incomingReplyAdvancedRepo) MarkVariantEvent(context.Context, uuid.UUID, uuid.UUID, string) error { + r.marked++ return nil } -func (incomingReplyAdvancedRepo) CreateReplyIntent(context.Context, *models.ReplyIntentRecord) error { +func (r *incomingReplyAdvancedRepo) CreateReplyIntent(context.Context, *models.ReplyIntentRecord) error { + r.intents++ return nil } @@ -69,6 +73,8 @@ type incomingReplyProgressRepo struct { replied int latest *repository.CampaignSequencePair sourceInbound bool + replyClaimed bool + advanced *incomingReplyAdvancedRepo } func (r *incomingReplyProgressRepo) IsInboundReplySource(context.Context, uuid.UUID, uuid.UUID) (bool, error) { @@ -87,9 +93,9 @@ func (r *incomingReplyProgressRepo) RecordReplyClassification(context.Context, u return nil } -func (r *incomingReplyProgressRepo) RecordEmailReplied(context.Context, uuid.UUID, uuid.UUID, uuid.UUID, uuid.UUID, uuid.UUID) error { +func (r *incomingReplyProgressRepo) RecordEmailReplied(context.Context, uuid.UUID, uuid.UUID, uuid.UUID, uuid.UUID, uuid.UUID) (bool, error) { r.replied++ - return nil + return r.replyClaimed, nil } type incomingReplyCampaignRepo struct{ repository.CampaignRepository } @@ -100,13 +106,15 @@ 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{sourceInbound: true} + progress := &incomingReplyProgressRepo{sourceInbound: true, replyClaimed: true} + advancedRepo := &incomingReplyAdvancedRepo{} + progress.advanced = advancedRepo taskContactRecord := &models.Contact{ID: taskContact, Email: "task-contact@example.test"} if senderContact != nil && senderContact.ID == taskContact { taskContactRecord = senderContact } return &service{ - repo: incomingReplyAdvancedRepo{}, + repo: advancedRepo, campaignRepo: incomingReplyCampaignRepo{}, emailRepo: incomingReplyEmailRepo{account: account}, taskRepo: incomingReplyTaskRepo{ @@ -225,6 +233,37 @@ func TestProcessIncomingReplyTrustsPersistedDirectionOverEventPayload(t *testing } } +func TestProcessIncomingReplyStopsWhenWriteBoundaryRejectsSource(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.replyClaimed = 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 != 1 { + t.Fatalf("RecordEmailReplied calls = %d, want 1 write-boundary claim", progress.replied) + } + if progress.advanced.marked != 0 { + t.Fatalf("MarkVariantEvent calls = %d, want 0 after rejected claim", progress.advanced.marked) + } + if progress.advanced.intents != 0 { + t.Fatalf("CreateReplyIntent calls = %d, want 0 after rejected claim", progress.advanced.intents) + } +} + 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 b2a9b06e7..137303d00 100644 --- a/internal/app/advanced/service.go +++ b/internal/app/advanced/service.go @@ -1218,6 +1218,15 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid. // replied_at IS NOT NULL, so gating the stamp here fixes both at once. // Any reply, human or automatic, proves the mailbox is live; only a // human one counts as engagement. + if !replyclassify.IsAutomated(replyResult.Class) { + claimed, err := s.campaignProgressRepo.RecordEmailReplied(ctx, cID, ctID, sID, emailAccountID, msg.ID) + if err != nil { + return toErrx(err) + } + if !claimed { + return nil + } + } if s.evidence != nil { kind := "replied" if replyclassify.IsAutomated(replyResult.Class) { @@ -1226,9 +1235,6 @@ 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) { - 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.up.sql b/internal/infrastructure/db/migrations/000168_repair_sent_followup_replies.up.sql index 2cd2b6683..3065e3bc6 100644 --- a/internal/infrastructure/db/migrations/000168_repair_sent_followup_replies.up.sql +++ b/internal/infrastructure/db/migrations/000168_repair_sent_followup_replies.up.sql @@ -35,7 +35,11 @@ WITH false_replies AS ( AND EXISTS ( SELECT 1 FROM unnest(ue.to_addr) AS recipients(address) - WHERE lower(recipients.address) LIKE '%' || lower(c.email) || '%' + WHERE lower(btrim(COALESCE( + substring(recipients.address FROM '<([^>]*)>'), + substring(recipients.address FROM '\(([^()]*)\)\s*$'), + recipients.address + ))) = lower(btrim(c.email)) ) WHERE p.replied_at IS NOT NULL AND p.reply_class IN ('', 'unknown') @@ -85,7 +89,11 @@ WITH false_replies AS ( AND EXISTS ( SELECT 1 FROM unnest(ue.to_addr) AS recipients(address) - WHERE lower(recipients.address) LIKE '%' || lower(c.email) || '%' + WHERE lower(btrim(COALESCE( + substring(recipients.address FROM '<([^>]*)>'), + substring(recipients.address FROM '\(([^()]*)\)\s*$'), + recipients.address + ))) = lower(btrim(c.email)) ) WHERE p.replied_at IS NOT NULL AND p.reply_class IN ('', 'unknown') @@ -138,7 +146,11 @@ WITH false_replies AS ( AND EXISTS ( SELECT 1 FROM unnest(ue.to_addr) AS recipients(address) - WHERE lower(recipients.address) LIKE '%' || lower(c.email) || '%' + WHERE lower(btrim(COALESCE( + substring(recipients.address FROM '<([^>]*)>'), + substring(recipients.address FROM '\(([^()]*)\)\s*$'), + recipients.address + ))) = lower(btrim(c.email)) ) WHERE p.replied_at IS NOT NULL AND p.reply_class IN ('', 'unknown') diff --git a/internal/repository/campaign_reply_source_live_test.go b/internal/repository/campaign_reply_source_live_test.go index 1e5d98a50..3f1ede734 100644 --- a/internal/repository/campaign_reply_source_live_test.go +++ b/internal/repository/campaign_reply_source_live_test.go @@ -43,9 +43,13 @@ func TestLiveReplySourceUsesStoredDirectionAtTheWriteBoundary(t *testing.T) { 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 { + claimed, err := repo.RecordEmailReplied(ctx, f.campaign, f.contact, step, f.other, sentID) + if err != nil { t.Fatalf("record sent source: %v", err) } + if claimed { + t.Fatal("Sent-folder source claimed reply progress") + } var replied bool if err := pool.QueryRow(ctx, ` @@ -79,9 +83,20 @@ func TestLiveReplySourceUsesStoredDirectionAtTheWriteBoundary(t *testing.T) { 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 { + claimed, err = repo.RecordEmailReplied(ctx, f.campaign, f.contact, step, f.mailbox, inboxID) + if err != nil { t.Fatalf("record inbound source: %v", err) } + if !claimed { + t.Fatal("Inbox source did not claim reply progress") + } + claimed, err = repo.RecordEmailReplied(ctx, f.campaign, f.contact, step, f.mailbox, inboxID) + if err != nil { + t.Fatalf("repeat inbound claim: %v", err) + } + if claimed { + t.Fatal("Already-recorded reply was claimed twice") + } if err := pool.QueryRow(ctx, ` SELECT replied_at IS NOT NULL FROM campaign_contact_progress diff --git a/internal/repository/pg_campaign_progress.go b/internal/repository/pg_campaign_progress.go index d636f78ba..c26a15a0e 100644 --- a/internal/repository/pg_campaign_progress.go +++ b/internal/repository/pg_campaign_progress.go @@ -185,7 +185,8 @@ type CampaignProgressRepository interface { GetStepSentAt(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) (*time.Time, 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 + // RecordEmailReplied claims the first human reply while its source is inbound. + RecordEmailReplied(ctx context.Context, campaignID, contactID, sequenceID, emailAccountID, messageID uuid.UUID) (bool, error) RecordEmailBounced(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error RecordEmailComplained(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error @@ -806,7 +807,7 @@ func (r *campaignProgressRepository) IsInboundReplySource(ctx context.Context, e } // 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 { +func (r *campaignProgressRepository) RecordEmailReplied(ctx context.Context, campaignID, contactID, sequenceID, emailAccountID, messageID uuid.UUID) (bool, error) { query := ` UPDATE campaign_contact_progress SET replied_at = NOW() @@ -824,8 +825,11 @@ func (r *campaignProgressRepository) RecordEmailReplied(ctx context.Context, cam ) ` - _, err := r.db.Exec(ctx, query, campaignID, contactID, sequenceID, messageID, emailAccountID) - return err + result, err := r.db.Exec(ctx, query, campaignID, contactID, sequenceID, messageID, emailAccountID) + if err != nil { + return false, err + } + return result.RowsAffected() == 1, nil } // RecordEmailBounced records that an email bounced