feat: verify reply direction from persisted mailbox state and repair issue 549 false Sent-folder replies

This commit is contained in:
Matthew Meszaros
2026-09-16 21:11:54 -07:00
parent baf38a4b03
commit b01984bec2
7 changed files with 355 additions and 9 deletions
+2
View File
@@ -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.
<Callout type="warn" title="Turn on stop on reply">
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.
</Callout>
+34 -4
View File
@@ -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 <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 != 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()
+10 -1
View File
@@ -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
@@ -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.
@@ -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
);
@@ -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, "<opener@test.local>", "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")
}
}
+30 -4
View File
@@ -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
}