feat: reject outbound mailbox copies and mismatched campaign threads before they can mark contacts replied for stop-on-reply (issue #549)

This commit is contained in:
Matthew Meszaros
2026-09-16 08:24:19 -07:00
parent 4597156b6f
commit 2255e2145d
6 changed files with 446 additions and 8 deletions
+2
View File
@@ -219,6 +219,8 @@ The **On reply** button on a step card creates an action step connected by a rep
A campaign-level toggle in the flow toolbar. When on and a contact replies, the rest of the cold sequence stops for them while the reply branch on the email they answered still runs instantly.
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.
<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>
@@ -0,0 +1,294 @@
package advanced
import (
"context"
"testing"
"github.com/google/uuid"
"github.com/warmbly/warmbly/internal/errx"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
type incomingReplyAdvancedRepo struct {
repository.AdvancedOutreachRepository
}
func (incomingReplyAdvancedRepo) GetOutreachSettings(context.Context, uuid.UUID) (*models.AdvancedOutreachSettings, error) {
settings := models.DefaultAdvancedOutreachSettings()
return &settings, nil
}
func (incomingReplyAdvancedRepo) MarkVariantEvent(context.Context, uuid.UUID, uuid.UUID, string) error {
return nil
}
func (incomingReplyAdvancedRepo) CreateReplyIntent(context.Context, *models.ReplyIntentRecord) error {
return nil
}
type incomingReplyEmailRepo struct {
repository.EmailRepository
account *models.Email
}
func (r incomingReplyEmailRepo) GetByID(context.Context, uuid.UUID) (*models.Email, *errx.Error) {
return r.account, nil
}
type incomingReplyTaskRepo struct {
repository.TaskRepository
task *repository.Task
campaign *repository.CampaignTask
}
func (r incomingReplyTaskRepo) GetTaskByMessageID(context.Context, string) (*repository.Task, error) {
return r.task, nil
}
func (r incomingReplyTaskRepo) GetCampaignTask(context.Context, uuid.UUID) (*repository.CampaignTask, error) {
return r.campaign, nil
}
type incomingReplyContactRepo struct {
repository.ContactRepository
senderContact *models.Contact
taskContact *models.Contact
}
func (r incomingReplyContactRepo) GetByEmailAndOrganization(context.Context, uuid.UUID, string) (*models.Contact, *errx.Error) {
return r.senderContact, nil
}
func (r incomingReplyContactRepo) GetByID(context.Context, uuid.UUID) (*models.Contact, *errx.Error) {
return r.taskContact, nil
}
type incomingReplyProgressRepo struct {
repository.CampaignProgressRepository
replied int
latest *repository.CampaignSequencePair
}
func (r *incomingReplyProgressRepo) GetLatestCampaignSequenceForContact(context.Context, uuid.UUID) (*repository.CampaignSequencePair, error) {
return r.latest, nil
}
func (r *incomingReplyProgressRepo) GetLatestReplyClass(context.Context, uuid.UUID, uuid.UUID) (string, error) {
return "", nil
}
func (r *incomingReplyProgressRepo) RecordReplyClassification(context.Context, uuid.UUID, uuid.UUID, uuid.UUID, string, string, float64) error {
return nil
}
func (r *incomingReplyProgressRepo) RecordEmailReplied(context.Context, uuid.UUID, uuid.UUID, uuid.UUID) error {
r.replied++
return nil
}
type incomingReplyCampaignRepo struct{ repository.CampaignRepository }
func (incomingReplyCampaignRepo) GetSequencesRoutingByCampaignID(context.Context, uuid.UUID) ([]models.Sequence, error) {
return nil, nil
}
func newIncomingReplyService(account *models.Email, senderContact *models.Contact, taskContact uuid.UUID) (*service, *incomingReplyProgressRepo) {
taskID, campaignID, sequenceID := uuid.New(), uuid.New(), uuid.New()
progress := &incomingReplyProgressRepo{}
taskContactRecord := &models.Contact{ID: taskContact, Email: "task-contact@example.test"}
if senderContact != nil && senderContact.ID == taskContact {
taskContactRecord = senderContact
}
return &service{
repo: incomingReplyAdvancedRepo{},
campaignRepo: incomingReplyCampaignRepo{},
emailRepo: incomingReplyEmailRepo{account: account},
taskRepo: incomingReplyTaskRepo{
task: &repository.Task{ID: taskID, TaskType: "campaign", EmailAccountID: account.ID},
campaign: &repository.CampaignTask{TaskID: taskID, CampaignID: &campaignID, ContactID: &taskContact, SequenceID: &sequenceID},
},
contactRepo: incomingReplyContactRepo{
senderContact: senderContact,
taskContact: taskContactRecord,
},
campaignProgressRepo: progress,
}, progress
}
func TestMessageAddressesMailbox(t *testing.T) {
account := &models.Email{
Email: "mailbox@example.test",
SendAsEmail: "alias@example.test",
ReplyTo: "replies@example.test",
}
for _, tc := range []struct {
name string
message *models.EmailMessageStoreData
want bool
}{
{name: "mailbox in to", message: &models.EmailMessageStoreData{ToAddr: []string{"Mailbox <mailbox@example.test>"}}, want: true},
{name: "send alias in cc", message: &models.EmailMessageStoreData{CC: []string{"alias@example.test"}}, want: true},
{name: "reply address in bcc", message: &models.EmailMessageStoreData{BCC: []string{"replies@example.test"}}, want: true},
{name: "different recipient", message: &models.EmailMessageStoreData{ToAddr: []string{"other@example.test"}}},
} {
t.Run(tc.name, func(t *testing.T) {
if got := messageAddressesMailbox(tc.message, account); got != tc.want {
t.Fatalf("messageAddressesMailbox() = %t, want %t", got, tc.want)
}
})
}
}
func TestProcessIncomingReplyRejectsMailboxOwnSentCopy(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, nil, contactID)
xerr := service.ProcessIncomingReply(context.Background(), accountID, &models.EmailMessageStoreData{
EmailID: accountID,
Folder: models.FolderInbox,
FromAddr: []string{"Sender <sender@example.test>"},
ToAddr: []string{"recipient@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 for the mailbox's own outbound copy", progress.replied)
}
}
func TestProcessIncomingReplyRejectsOutboundFolder(t *testing.T) {
for _, tc := range []struct {
name, folder, providerFolder string
}{
{name: "sent", folder: models.FolderSent},
{name: "draft", folder: models.FolderDrafts},
{name: "provider sent", folder: models.FolderInbox, providerFolder: models.FolderSent},
} {
t.Run(tc.name, func(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)
xerr := service.ProcessIncomingReply(context.Background(), accountID, &models.EmailMessageStoreData{
EmailID: accountID,
Folder: tc.folder,
ProviderFolder: tc.providerFolder,
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 for outbound folder", progress.replied)
}
})
}
}
func TestProcessIncomingReplyRequiresThreadSenderToMatchContact(t *testing.T) {
orgID, accountID := uuid.New(), uuid.New()
taskContactID, senderContactID := uuid.New(), uuid.New()
account := &models.Email{ID: accountID, OrganizationID: &orgID, Email: "sender@example.test"}
service, progress := newIncomingReplyService(account, &models.Contact{
ID: senderContactID, Email: "other@example.test",
}, taskContactID)
latestCampaignID, latestSequenceID := uuid.New(), uuid.New()
progress.latest = &repository.CampaignSequencePair{CampaignID: latestCampaignID, SequenceID: latestSequenceID}
xerr := service.ProcessIncomingReply(context.Background(), accountID, &models.EmailMessageStoreData{
EmailID: accountID,
Folder: models.FolderInbox,
FromAddr: []string{"Other person <other@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 threaded sender is a different contact", progress.replied)
}
}
func TestProcessIncomingReplyRequiresRecipientToMatchMailbox(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)
xerr := service.ProcessIncomingReply(context.Background(), accountID, &models.EmailMessageStoreData{
EmailID: accountID,
Folder: models.FolderInbox,
FromAddr: []string{"Recipient <recipient@example.test>"},
ToAddr: []string{"someone-else@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 recipient is not the sending mailbox", progress.replied)
}
}
func TestProcessIncomingReplyRequiresThreadToMatchMailbox(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)
service.taskRepo.(incomingReplyTaskRepo).task.EmailAccountID = uuid.New()
latestCampaignID, latestSequenceID := uuid.New(), uuid.New()
progress.latest = &repository.CampaignSequencePair{CampaignID: latestCampaignID, SequenceID: latestSequenceID}
xerr := service.ProcessIncomingReply(context.Background(), accountID, &models.EmailMessageStoreData{
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 thread belongs to another mailbox", progress.replied)
}
}
func TestProcessIncomingReplyAcceptsMatchingThreadSender(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)
xerr := service.ProcessIncomingReply(context.Background(), accountID, &models.EmailMessageStoreData{
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 for the matching contact", progress.replied)
}
}
+60 -7
View File
@@ -880,6 +880,37 @@ func parseSenderEmail(addrs []string) string {
return strings.ToLower(strings.Trim(primary, "<>"))
}
func messageAddressesMailbox(msg *models.EmailMessageStoreData, account *models.Email) bool {
if msg == nil || account == nil {
return false
}
targets := make(map[string]struct{}, 3)
for _, raw := range []string{account.Email, account.SendFrom(), account.ReplyTo} {
if address := parseSenderEmail([]string{raw}); address != "" {
targets[address] = struct{}{}
}
}
for _, fields := range [][]string{msg.ToAddr, msg.CC, msg.BCC} {
for _, raw := range fields {
addresses, err := mail.ParseAddressList(raw)
if err != nil {
if address := parseSenderEmail([]string{raw}); address != "" {
if _, ok := targets[address]; ok {
return true
}
}
continue
}
for _, address := range addresses {
if _, ok := targets[strings.ToLower(strings.TrimSpace(address.Address))]; ok {
return true
}
}
}
}
return false
}
func cleanMessageID(mid string) string {
return strings.TrimSpace(strings.Trim(mid, "<>"))
}
@@ -1035,6 +1066,9 @@ func classifyReply(text string, cfg models.ReplyIntentSettings) (models.ReplyInt
}
func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.UUID, msg *models.EmailMessageStoreData) *errx.Error {
if !msg.MayBeInbound() {
return nil
}
account, xerr := s.emailRepo.GetByID(ctx, emailAccountID)
if xerr != nil {
return xerr
@@ -1055,6 +1089,12 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
if sender == "" {
return nil
}
if sender == parseSenderEmail([]string{account.Email}) || sender == parseSenderEmail([]string{account.SendFrom()}) {
return nil
}
if !messageAddressesMailbox(msg, account) {
return nil
}
text := strings.TrimSpace(msg.Snippet)
text = strings.TrimSpace(text + "\n" + msg.Subject)
@@ -1071,6 +1111,7 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
var sequenceID *uuid.UUID
var contactID *uuid.UUID
var taskID *uuid.UUID
var referencesCampaignThread bool
// First, try exact message threading via In-Reply-To.
for _, mid := range msg.InReplyTo {
@@ -1082,13 +1123,25 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
if err != nil || task == nil || task.TaskType != "campaign" {
continue
}
taskID = &task.ID
ct, err := s.taskRepo.GetCampaignTask(ctx, task.ID)
if err == nil && ct != nil {
campaignID = ct.CampaignID
contactID = ct.ContactID
sequenceID = ct.SequenceID
referencesCampaignThread = true
if task.EmailAccountID != emailAccountID {
continue
}
ct, err := s.taskRepo.GetCampaignTask(ctx, task.ID)
if err != nil || ct == nil || ct.ContactID == nil {
continue
}
contact, contactErr := s.contactRepo.GetByID(ctx, *ct.ContactID)
if contactErr != nil {
return contactErr
}
if contact == nil || !strings.EqualFold(strings.TrimSpace(contact.Email), sender) {
continue
}
taskID = &task.ID
campaignID = ct.CampaignID
contactID = ct.ContactID
sequenceID = ct.SequenceID
break
}
@@ -1102,7 +1155,7 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
}
}
if campaignID == nil && contactID != nil {
if campaignID == nil && contactID != nil && !referencesCampaignThread {
latest, err := s.campaignProgressRepo.GetLatestCampaignSequenceForContact(ctx, *contactID)
if err == nil && latest != nil {
campaignID = &latest.CampaignID
+1 -1
View File
@@ -97,7 +97,7 @@ func (s *JobsService) ingestNewEmail(ctx context.Context, e *models.JobEventNewE
// (replyclassify) and persists reply_class/confidence/source on the contact's
// campaign progress, gating replied_at so automated replies (auto_reply /
// out_of_office) never count as a human reply for stop_on_reply / branching.
if s.AdvancedService != nil {
if s.AdvancedService != nil && e.Message.MayBeInbound() {
// Logged, not propagated: the ingest must survive it, but a silent
// failure here is indistinguishable from a reply that linked fine.
if xerr := s.AdvancedService.ProcessIncomingReply(ctx, e.Message.EmailID, e.Message); xerr != nil {
@@ -0,0 +1,76 @@
package jobs
import (
"context"
"testing"
"github.com/google/uuid"
"github.com/warmbly/warmbly/internal/app/advanced"
"github.com/warmbly/warmbly/internal/errx"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
type replyRecordingAdvanced struct {
advanced.Service
calls int
}
func (s *replyRecordingAdvanced) ProcessIncomingReply(context.Context, uuid.UUID, *models.EmailMessageStoreData) *errx.Error {
s.calls++
return nil
}
type newEmailInboxRepo struct{ repository.UniboxRepository }
func (newEmailInboxRepo) CreateEntry(context.Context, uuid.UUID, *models.EmailMessageStoreData) error {
return nil
}
type newEmailAccountRepo struct{ repository.EmailRepository }
func (newEmailAccountRepo) GetByID(context.Context, uuid.UUID) (*models.Email, *errx.Error) {
return nil, nil
}
func TestHandleNewEmailDoesNotProcessSentMailAsReply(t *testing.T) {
for _, tc := range []struct {
name, folder, providerFolder string
wantCalls int
}{
{name: "sent follow-up", folder: models.FolderSent},
{name: "provider still reports sent", folder: models.FolderInbox, providerFolder: models.FolderSent},
{name: "draft", folder: models.FolderDrafts},
{name: "inbox", folder: models.FolderInbox, wantCalls: 1},
{name: "archived inbound", folder: models.FolderArchive, wantCalls: 1},
} {
t.Run(tc.name, func(t *testing.T) {
advancedService := &replyRecordingAdvanced{}
service := &JobsService{
UniboxRepository: newEmailInboxRepo{},
EmailRepository: newEmailAccountRepo{},
AdvancedService: advancedService,
}
err := service.HandleNewEmail(context.Background(), &models.JobEventNewEmail{
UserID: uuid.New(),
Message: &models.EmailMessageStoreData{
ID: uuid.New(),
EmailID: uuid.New(),
Folder: tc.folder,
ProviderFolder: tc.providerFolder,
FromAddr: []string{"sender@example.test"},
ToAddr: []string{"recipient@example.test"},
MessageID: "<follow-up@example.test>",
InReplyTo: []string{"<opener@example.test>"},
},
})
if err != nil {
t.Fatal(err)
}
if advancedService.calls != tc.wantCalls {
t.Fatalf("ProcessIncomingReply calls = %d, want %d", advancedService.calls, tc.wantCalls)
}
})
}
}
+13
View File
@@ -270,6 +270,19 @@ func NormalizeFolder(folder string, flags []string) string {
return FolderInbox
}
func outboundOnlyFolder(folder string) bool {
return folder == FolderSent || folder == FolderDrafts
}
// MayBeInbound reports whether neither folder placement proves the message outbound.
func (e *EmailMessageStoreData) MayBeInbound() bool {
if e == nil {
return false
}
return !outboundOnlyFolder(NormalizeFolder(e.Folder, e.Flags)) &&
!outboundOnlyFolder(e.ProviderFolder)
}
type MailSearchResult struct {
Data []EmailMessageStoreDataPreview `json:"data"`
Pagination CPagination `json:"pagination"`