diff --git a/docs/content/docs/guides/sequences.mdx b/docs/content/docs/guides/sequences.mdx
index 2ecbfce60..f12592942 100644
--- a/docs/content/docs/guides/sequences.mdx
+++ b/docs/content/docs/guides/sequences.mdx
@@ -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.
+
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
new file mode 100644
index 000000000..324da97e7
--- /dev/null
+++ b/internal/app/advanced/incoming_reply_test.go
@@ -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 "}}, 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 "},
+ ToAddr: []string{"recipient@example.test"},
+ InReplyTo: []string{""},
+ 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 "},
+ 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 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 "},
+ 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 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 "},
+ ToAddr: []string{"someone-else@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 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 "},
+ 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 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 "},
+ 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 for the matching contact", progress.replied)
+ }
+}
diff --git a/internal/app/advanced/service.go b/internal/app/advanced/service.go
index f7bd868cf..f5d5e7c1e 100644
--- a/internal/app/advanced/service.go
+++ b/internal/app/advanced/service.go
@@ -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
diff --git a/internal/app/consumer/event_new_email.go b/internal/app/consumer/event_new_email.go
index 202beb6e4..024e79637 100644
--- a/internal/app/consumer/event_new_email.go
+++ b/internal/app/consumer/event_new_email.go
@@ -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 {
diff --git a/internal/app/consumer/event_new_email_test.go b/internal/app/consumer/event_new_email_test.go
new file mode 100644
index 000000000..e0c63fe9a
--- /dev/null
+++ b/internal/app/consumer/event_new_email_test.go
@@ -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: "",
+ InReplyTo: []string{""},
+ },
+ })
+ if err != nil {
+ t.Fatal(err)
+ }
+ if advancedService.calls != tc.wantCalls {
+ t.Fatalf("ProcessIncomingReply calls = %d, want %d", advancedService.calls, tc.wantCalls)
+ }
+ })
+ }
+}
diff --git a/internal/models/unibox.go b/internal/models/unibox.go
index 5163b2f11..26296f07f 100644
--- a/internal/models/unibox.go
+++ b/internal/models/unibox.go
@@ -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"`