feat: make reply recording atomic and require exact repair recipient matches

This commit is contained in:
Matthew Meszaros
2026-09-16 21:25:43 -07:00
parent b01984bec2
commit c91e89980f
5 changed files with 94 additions and 18 deletions
+45 -6
View File
@@ -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 <recipient@example.test>"},
ToAddr: []string{"sender@example.test"},
InReplyTo: []string{"<opener@example.test>"},
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()
+9 -3
View File
@@ -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
@@ -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')
@@ -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
+8 -4
View File
@@ -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