mirror of
https://github.com/warmbly/warmbly.git
synced 2026-10-04 08:02:01 +00:00
Merge pull request #571 from warmbly/fix/reply-tracking-regression
fix: restore reply tracking for cross-mailbox campaign threads
This commit is contained in:
@@ -221,7 +221,9 @@ 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 also recalculates address verification from the contact's remaining delivery, open, click, reply, and bounce evidence. It leaves ambiguous historical replies untouched.
|
||||
When a campaign rotates between sender accounts, a real inbound reply still links to the campaign if the mail client references an earlier step from another mailbox in the same workspace.
|
||||
|
||||
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. An upgrade also restores provable human replies that were received by one campaign mailbox but referenced a thread started by another mailbox in the workspace. Warmbly recalculates address verification from the corrected delivery, open, click, reply, and bounce evidence. It 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.
|
||||
|
||||
@@ -79,6 +79,7 @@ type incomingReplyProgressRepo struct {
|
||||
sourceInbound bool
|
||||
sourceClaimed bool
|
||||
replyAccepted bool
|
||||
receivingSent bool
|
||||
completeErr error
|
||||
advanced *incomingReplyAdvancedRepo
|
||||
}
|
||||
@@ -87,6 +88,10 @@ func (r *incomingReplyProgressRepo) IsInboundReplySource(context.Context, uuid.U
|
||||
return r.sourceInbound, nil
|
||||
}
|
||||
|
||||
func (r *incomingReplyProgressRepo) CampaignContactSentFromAccount(context.Context, uuid.UUID, uuid.UUID, uuid.UUID) (bool, error) {
|
||||
return r.receivingSent, nil
|
||||
}
|
||||
|
||||
func (r *incomingReplyProgressRepo) GetLatestCampaignSequenceForContact(context.Context, uuid.UUID) (*repository.CampaignSequencePair, error) {
|
||||
return r.latest, nil
|
||||
}
|
||||
@@ -118,7 +123,14 @@ func (r *incomingReplyProgressRepo) RecordEmailReplied(context.Context, uuid.UUI
|
||||
return r.replyAccepted, nil
|
||||
}
|
||||
|
||||
type incomingReplyCampaignRepo struct{ repository.CampaignRepository }
|
||||
type incomingReplyCampaignRepo struct {
|
||||
repository.CampaignRepository
|
||||
campaign *models.Campaign
|
||||
}
|
||||
|
||||
func (r incomingReplyCampaignRepo) GetByID(context.Context, uuid.UUID) (*models.Campaign, error) {
|
||||
return r.campaign, nil
|
||||
}
|
||||
|
||||
func (incomingReplyCampaignRepo) GetSequencesRoutingByCampaignID(context.Context, uuid.UUID) ([]models.Sequence, error) {
|
||||
return nil, nil
|
||||
@@ -126,7 +138,12 @@ 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, sourceClaimed: true, replyAccepted: true}
|
||||
progress := &incomingReplyProgressRepo{
|
||||
sourceInbound: true,
|
||||
sourceClaimed: true,
|
||||
replyAccepted: true,
|
||||
receivingSent: true,
|
||||
}
|
||||
advancedRepo := &incomingReplyAdvancedRepo{}
|
||||
progress.advanced = advancedRepo
|
||||
taskContactRecord := &models.Contact{ID: taskContact, Email: "task-contact@example.test"}
|
||||
@@ -134,9 +151,12 @@ func newIncomingReplyService(account *models.Email, senderContact *models.Contac
|
||||
taskContactRecord = senderContact
|
||||
}
|
||||
return &service{
|
||||
repo: advancedRepo,
|
||||
campaignRepo: incomingReplyCampaignRepo{},
|
||||
emailRepo: incomingReplyEmailRepo{account: account},
|
||||
repo: advancedRepo,
|
||||
campaignRepo: incomingReplyCampaignRepo{campaign: &models.Campaign{
|
||||
ID: campaignID,
|
||||
OrganizationID: account.OrganizationID,
|
||||
}},
|
||||
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},
|
||||
@@ -374,7 +394,7 @@ func TestProcessIncomingReplyRequiresRecipientToMatchMailbox(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestProcessIncomingReplyRequiresThreadToMatchMailbox(t *testing.T) {
|
||||
func TestProcessIncomingReplyAcceptsCrossMailboxThreadWhenStoredSourceIsInbound(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{
|
||||
@@ -384,6 +404,37 @@ func TestProcessIncomingReplyRequiresThreadToMatchMailbox(t *testing.T) {
|
||||
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 != 1 {
|
||||
t.Fatalf("RecordEmailReplied calls = %d, want 1 for an inbound cross-mailbox thread", progress.replied)
|
||||
}
|
||||
if progress.completed != 1 {
|
||||
t.Fatalf("CompleteIncomingReply calls = %d, want 1 for an inbound cross-mailbox thread", progress.completed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProcessIncomingReplyRejectsCrossOrganizationThread(t *testing.T) {
|
||||
orgID, otherOrgID, accountID, contactID := uuid.New(), 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()
|
||||
service.campaignRepo = incomingReplyCampaignRepo{campaign: &models.Campaign{
|
||||
ID: uuid.New(),
|
||||
OrganizationID: &otherOrgID,
|
||||
}}
|
||||
|
||||
xerr := service.ProcessIncomingReply(context.Background(), accountID, &models.EmailMessageStoreData{
|
||||
EmailID: accountID,
|
||||
Folder: models.FolderInbox,
|
||||
@@ -396,7 +447,32 @@ func TestProcessIncomingReplyRequiresThreadToMatchMailbox(t *testing.T) {
|
||||
t.Fatal(xerr)
|
||||
}
|
||||
if progress.replied != 0 {
|
||||
t.Fatalf("RecordEmailReplied calls = %d, want 0 when the thread belongs to another mailbox", progress.replied)
|
||||
t.Fatalf("RecordEmailReplied calls = %d, want 0 for a cross-organization thread", progress.replied)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProcessIncomingReplyRejectsCrossMailboxThreadAtUnrelatedWorkspaceMailbox(t *testing.T) {
|
||||
orgID, accountID, contactID := uuid.New(), uuid.New(), uuid.New()
|
||||
account := &models.Email{ID: accountID, OrganizationID: &orgID, Email: "unrelated@example.test"}
|
||||
service, progress := newIncomingReplyService(account, &models.Contact{
|
||||
ID: contactID, Email: "recipient@example.test",
|
||||
}, contactID)
|
||||
service.taskRepo.(incomingReplyTaskRepo).task.EmailAccountID = uuid.New()
|
||||
progress.receivingSent = false
|
||||
|
||||
xerr := service.ProcessIncomingReply(context.Background(), accountID, &models.EmailMessageStoreData{
|
||||
EmailID: accountID,
|
||||
Folder: models.FolderInbox,
|
||||
FromAddr: []string{"Recipient <recipient@example.test>"},
|
||||
ToAddr: []string{"unrelated@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 an unrelated workspace mailbox", progress.replied)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1132,12 +1132,25 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
|
||||
continue
|
||||
}
|
||||
referencesCampaignThread = true
|
||||
if task.EmailAccountID != emailAccountID {
|
||||
ct, err := s.taskRepo.GetCampaignTask(ctx, task.ID)
|
||||
if err != nil || ct == nil || ct.CampaignID == nil || ct.ContactID == nil {
|
||||
continue
|
||||
}
|
||||
ct, err := s.taskRepo.GetCampaignTask(ctx, task.ID)
|
||||
if err != nil || ct == nil || ct.ContactID == nil {
|
||||
continue
|
||||
if task.EmailAccountID != emailAccountID {
|
||||
campaign, err := s.campaignRepo.GetByID(ctx, *ct.CampaignID)
|
||||
if err != nil || campaign == nil || campaign.OrganizationID == nil ||
|
||||
*campaign.OrganizationID != *account.OrganizationID {
|
||||
continue
|
||||
}
|
||||
sentFromReceivingAccount, err := s.campaignProgressRepo.CampaignContactSentFromAccount(
|
||||
ctx, *ct.CampaignID, *ct.ContactID, emailAccountID,
|
||||
)
|
||||
if err != nil {
|
||||
return toErrx(err)
|
||||
}
|
||||
if !sentFromReceivingAccount {
|
||||
continue
|
||||
}
|
||||
}
|
||||
contact, contactErr := s.contactRepo.GetByID(ctx, *ct.ContactID)
|
||||
if contactErr != nil {
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
-- Historical reply attribution cannot be safely reversed.
|
||||
@@ -0,0 +1,251 @@
|
||||
-- Repair real replies skipped when their thread root came from another workspace mailbox.
|
||||
|
||||
CREATE TEMP TABLE issue_549_missed_replies AS
|
||||
WITH task_matches AS (
|
||||
SELECT
|
||||
ue.id AS source_message_id,
|
||||
ue.campaign_reply_processed_at AS replied_at,
|
||||
ue.subject,
|
||||
campaign_task.task_id,
|
||||
campaign_task.campaign_id,
|
||||
campaign_task.contact_id,
|
||||
campaign_task.sequence_id,
|
||||
campaign.organization_id,
|
||||
contact.email AS contact_email,
|
||||
count(*) OVER (PARTITION BY ue.id) AS task_match_count
|
||||
FROM unibox_emails ue
|
||||
JOIN email_accounts receiving_account
|
||||
ON receiving_account.id = ue.email_id
|
||||
JOIN tasks parent_task
|
||||
ON parent_task.task_type = 'campaign'
|
||||
AND parent_task.email_account_id <> ue.email_id
|
||||
AND EXISTS (
|
||||
SELECT 1
|
||||
FROM unnest(ue.in_reply_to) AS refs(message_id)
|
||||
WHERE btrim(refs.message_id, '<> ') = btrim(parent_task.message_id, '<> ')
|
||||
)
|
||||
JOIN campaign_tasks campaign_task
|
||||
ON campaign_task.task_id = parent_task.id
|
||||
AND campaign_task.campaign_id IS NOT NULL
|
||||
AND campaign_task.contact_id IS NOT NULL
|
||||
AND campaign_task.sequence_id IS NOT NULL
|
||||
JOIN campaigns campaign
|
||||
ON campaign.id = campaign_task.campaign_id
|
||||
AND campaign.organization_id = receiving_account.organization_id
|
||||
JOIN email_accounts sending_account
|
||||
ON sending_account.id = parent_task.email_account_id
|
||||
AND sending_account.organization_id = campaign.organization_id
|
||||
JOIN contacts contact
|
||||
ON contact.id = campaign_task.contact_id
|
||||
AND contact.organization_id = campaign.organization_id
|
||||
AND EXISTS (
|
||||
SELECT 1
|
||||
FROM unnest(ue.from_addr) AS senders(address)
|
||||
WHERE lower(btrim(COALESCE(
|
||||
substring(senders.address FROM '<([^>]*)>'),
|
||||
substring(senders.address FROM '\(([^()]*)\)\s*$'),
|
||||
senders.address
|
||||
))) = lower(btrim(contact.email))
|
||||
)
|
||||
JOIN campaign_contact_progress progress
|
||||
ON progress.campaign_id = campaign_task.campaign_id
|
||||
AND progress.contact_id = campaign_task.contact_id
|
||||
AND progress.sequence_id = campaign_task.sequence_id
|
||||
AND progress.replied_at IS NULL
|
||||
WHERE ue.campaign_reply_processed_at IS NOT NULL
|
||||
AND ue.folder NOT IN ('sent', 'drafts')
|
||||
AND ue.provider_folder NOT IN ('sent', 'drafts')
|
||||
AND EXISTS (
|
||||
SELECT 1
|
||||
FROM campaign_tasks receiving_campaign_task
|
||||
JOIN tasks receiving_task ON receiving_task.id = receiving_campaign_task.task_id
|
||||
WHERE receiving_campaign_task.campaign_id = campaign_task.campaign_id
|
||||
AND receiving_campaign_task.contact_id = campaign_task.contact_id
|
||||
AND receiving_task.task_type = 'campaign'
|
||||
AND receiving_task.status = 'completed'
|
||||
AND receiving_task.email_account_id = ue.email_id
|
||||
)
|
||||
AND EXISTS (
|
||||
SELECT 1
|
||||
FROM unnest(ue.to_addr || ue.cc || ue.bcc) AS recipients(address)
|
||||
WHERE lower(btrim(COALESCE(
|
||||
substring(recipients.address FROM '<([^>]*)>'),
|
||||
substring(recipients.address FROM '\(([^()]*)\)\s*$'),
|
||||
recipients.address
|
||||
))) IN (
|
||||
lower(receiving_account.email),
|
||||
lower(COALESCE(NULLIF(receiving_account.send_as_email, ''), receiving_account.email)),
|
||||
lower(COALESCE(NULLIF(receiving_account.reply_to, ''), receiving_account.email))
|
||||
)
|
||||
)
|
||||
), intent_matches AS (
|
||||
SELECT
|
||||
task_matches.*,
|
||||
intent.id AS intent_id,
|
||||
intent.intent,
|
||||
intent.metadata,
|
||||
count(*) OVER (PARTITION BY task_matches.source_message_id) AS intent_match_count,
|
||||
count(*) OVER (PARTITION BY intent.id) AS source_match_count
|
||||
FROM task_matches
|
||||
JOIN reply_intents intent
|
||||
ON intent.organization_id = task_matches.organization_id
|
||||
AND intent.campaign_id IS NULL
|
||||
AND intent.task_id IS NULL
|
||||
AND lower(btrim(intent.contact_email)) = lower(btrim(task_matches.contact_email))
|
||||
AND COALESCE(intent.metadata->>'subject', '') = task_matches.subject
|
||||
AND intent.created_at BETWEEN task_matches.replied_at - INTERVAL '30 seconds'
|
||||
AND task_matches.replied_at + INTERVAL '30 seconds'
|
||||
WHERE task_matches.task_match_count = 1
|
||||
)
|
||||
SELECT
|
||||
source_message_id,
|
||||
replied_at,
|
||||
task_id,
|
||||
campaign_id,
|
||||
contact_id,
|
||||
sequence_id,
|
||||
intent_id,
|
||||
metadata
|
||||
FROM intent_matches
|
||||
WHERE intent_match_count = 1
|
||||
AND source_match_count = 1
|
||||
AND intent NOT IN ('out_of_office', 'automated')
|
||||
AND COALESCE(metadata->>'reply_class', '') NOT IN ('auto_reply', 'out_of_office');
|
||||
|
||||
UPDATE campaign_contact_progress progress
|
||||
SET replied_at = repair.replied_at,
|
||||
reply_class = CASE
|
||||
WHEN COALESCE(repair.metadata->>'reply_class', '') IN (
|
||||
'positive', 'negative', 'neutral', 'unsubscribe', 'unknown'
|
||||
) THEN repair.metadata->>'reply_class'
|
||||
ELSE 'unknown'
|
||||
END,
|
||||
reply_confidence = 0,
|
||||
reply_source = CASE
|
||||
WHEN COALESCE(repair.metadata->>'classified_by', '') IN ('header', 'lexicon', 'model')
|
||||
THEN repair.metadata->>'classified_by'
|
||||
ELSE ''
|
||||
END
|
||||
FROM issue_549_missed_replies repair
|
||||
WHERE progress.campaign_id = repair.campaign_id
|
||||
AND progress.contact_id = repair.contact_id
|
||||
AND progress.sequence_id = repair.sequence_id
|
||||
AND progress.replied_at IS NULL;
|
||||
|
||||
UPDATE reply_intents intent
|
||||
SET campaign_id = repair.campaign_id,
|
||||
task_id = repair.task_id
|
||||
FROM issue_549_missed_replies repair
|
||||
WHERE intent.id = repair.intent_id
|
||||
AND intent.campaign_id IS NULL
|
||||
AND intent.task_id IS NULL;
|
||||
|
||||
UPDATE campaign_ab_assignments assignment
|
||||
SET replied_at = COALESCE(assignment.replied_at, repair.replied_at)
|
||||
FROM issue_549_missed_replies repair
|
||||
WHERE assignment.campaign_id = repair.campaign_id
|
||||
AND assignment.contact_id = repair.contact_id;
|
||||
|
||||
INSERT INTO contact_verification_evidence (contact_id, kind, ref, observed_at)
|
||||
SELECT repair.contact_id, 'replied', repair.source_message_id::text, repair.replied_at
|
||||
FROM issue_549_missed_replies repair
|
||||
JOIN contacts contact ON contact.id = repair.contact_id
|
||||
JOIN campaign_contact_progress progress
|
||||
ON progress.campaign_id = repair.campaign_id
|
||||
AND progress.contact_id = repair.contact_id
|
||||
AND progress.sequence_id = repair.sequence_id
|
||||
WHERE contact.verification_evidence_reset_at IS NULL
|
||||
OR COALESCE(progress.dispatched_at, progress.sent_at) > contact.verification_evidence_reset_at
|
||||
ON CONFLICT (contact_id, kind, ref) DO NOTHING;
|
||||
|
||||
WITH affected AS (
|
||||
SELECT DISTINCT contact_id
|
||||
FROM issue_549_missed_replies
|
||||
), newest_evidence AS (
|
||||
SELECT
|
||||
evidence.contact_id,
|
||||
evidence.kind,
|
||||
evidence.observed_at,
|
||||
row_number() OVER (
|
||||
PARTITION BY evidence.contact_id
|
||||
ORDER BY evidence.observed_at DESC
|
||||
) AS overall_rank
|
||||
FROM contact_verification_evidence evidence
|
||||
JOIN affected ON affected.contact_id = evidence.contact_id
|
||||
), ranked_evidence AS (
|
||||
SELECT
|
||||
contact_id,
|
||||
kind,
|
||||
observed_at,
|
||||
row_number() OVER (
|
||||
PARTITION BY contact_id, kind
|
||||
ORDER BY observed_at DESC
|
||||
) AS kind_rank
|
||||
FROM newest_evidence
|
||||
WHERE overall_rank <= 50
|
||||
), remaining AS (
|
||||
SELECT
|
||||
affected.contact_id,
|
||||
COALESCE(sum(
|
||||
CASE
|
||||
WHEN evidence.kind = 'replied' AND evidence.kind_rank <= 3
|
||||
THEN 45 * power(0.5, GREATEST(0, extract(epoch FROM (NOW() - evidence.observed_at))) / extract(epoch FROM INTERVAL '365 days'))
|
||||
WHEN evidence.kind = 'auto_replied' AND evidence.kind_rank <= 2
|
||||
THEN 30 * power(0.5, GREATEST(0, extract(epoch FROM (NOW() - evidence.observed_at))) / extract(epoch FROM INTERVAL '180 days'))
|
||||
WHEN evidence.kind = 'clicked' AND evidence.kind_rank <= 3
|
||||
THEN 35 * power(0.5, GREATEST(0, extract(epoch FROM (NOW() - evidence.observed_at))) / extract(epoch FROM INTERVAL '270 days'))
|
||||
WHEN evidence.kind = 'opened' AND evidence.kind_rank <= 3
|
||||
THEN 25 * power(0.5, GREATEST(0, extract(epoch FROM (NOW() - evidence.observed_at))) / extract(epoch FROM INTERVAL '180 days'))
|
||||
WHEN evidence.kind = 'delivered' AND evidence.kind_rank <= 4
|
||||
THEN 14 * power(0.5, GREATEST(0, extract(epoch FROM (NOW() - evidence.observed_at))) / extract(epoch FROM INTERVAL '180 days'))
|
||||
ELSE 0
|
||||
END
|
||||
), 0) AS positive_score,
|
||||
COALESCE(sum(
|
||||
CASE
|
||||
WHEN evidence.kind = 'bounced_recipient' AND evidence.kind_rank <= 2
|
||||
THEN 70 * power(0.5, GREATEST(0, extract(epoch FROM (NOW() - evidence.observed_at))) / extract(epoch FROM INTERVAL '365 days'))
|
||||
ELSE 0
|
||||
END
|
||||
), 0) AS negative_score,
|
||||
max(evidence.observed_at) FILTER (
|
||||
WHERE evidence.kind IN ('delivered', 'opened', 'clicked', 'replied', 'auto_replied')
|
||||
) AS last_positive_at,
|
||||
max(evidence.observed_at) FILTER (
|
||||
WHERE evidence.kind = 'bounced_recipient'
|
||||
) AS last_negative_at
|
||||
FROM affected
|
||||
LEFT JOIN ranked_evidence evidence
|
||||
ON evidence.contact_id = affected.contact_id
|
||||
GROUP BY affected.contact_id
|
||||
)
|
||||
UPDATE contacts contact
|
||||
SET verification_evidence_at = remaining.last_positive_at,
|
||||
verification_status = CASE
|
||||
WHEN contact.verification_source = 'manual' THEN contact.verification_status
|
||||
WHEN remaining.last_negative_at > COALESCE(remaining.last_positive_at, '-infinity'::timestamptz) THEN 'invalid'
|
||||
WHEN remaining.positive_score >= 20 THEN 'valid'
|
||||
ELSE contact.verification_status
|
||||
END,
|
||||
verification_confidence = CASE
|
||||
WHEN contact.verification_source = 'manual'
|
||||
THEN LEAST(100, 95 + floor(remaining.positive_score / 4)::integer)
|
||||
WHEN remaining.last_negative_at > COALESCE(remaining.last_positive_at, '-infinity'::timestamptz)
|
||||
THEN LEAST(100, 60 + floor(remaining.negative_score / 2)::integer)
|
||||
WHEN remaining.positive_score >= 20
|
||||
THEN LEAST(100, 70 + floor(remaining.positive_score / 2)::integer)
|
||||
ELSE contact.verification_confidence
|
||||
END,
|
||||
verification_reason = CASE
|
||||
WHEN contact.verification_source = 'manual' THEN contact.verification_reason
|
||||
WHEN remaining.last_negative_at > COALESCE(remaining.last_positive_at, '-infinity'::timestamptz)
|
||||
THEN 'recipient address bounced after the last positive mail evidence'
|
||||
WHEN remaining.positive_score >= 20
|
||||
THEN 'real mail evidence confirms the address'
|
||||
ELSE contact.verification_reason
|
||||
END,
|
||||
updated_at = NOW()
|
||||
FROM remaining
|
||||
WHERE contact.id = remaining.contact_id;
|
||||
|
||||
DROP TABLE issue_549_missed_replies;
|
||||
@@ -160,3 +160,52 @@ func TestLiveReplySourceUsesStoredDirectionAtTheWriteBoundary(t *testing.T) {
|
||||
t.Fatal("Inbox source did not stamp replied_at")
|
||||
}
|
||||
}
|
||||
|
||||
// A cross-provider reply is eligible only when its receiving mailbox sent to the lead.
|
||||
func TestLiveCrossProviderReplyRequiresReceivingMailboxUsedForLead(t *testing.T) {
|
||||
_, pool := liveContactDB(t)
|
||||
f := newThreadParentFixture(t, pool)
|
||||
ctx := context.Background()
|
||||
root := f.step(1, "Hello", true)
|
||||
rootTask := f.send(root, f.mailbox, "<opener@test.local>", "thread-1", 120)
|
||||
|
||||
if _, err := pool.Exec(ctx, `
|
||||
UPDATE email_accounts
|
||||
SET provider = CASE id WHEN $1 THEN 'outlook'::email_provider ELSE 'smtp_imap'::email_provider END
|
||||
WHERE id IN ($1, $2)
|
||||
`, f.mailbox, f.other); err != nil {
|
||||
t.Fatalf("set cross-provider fixture: %v", err)
|
||||
}
|
||||
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 '2 minutes')
|
||||
`, f.campaign, f.contact, root, rootTask); err != nil {
|
||||
t.Fatalf("insert root progress: %v", err)
|
||||
}
|
||||
|
||||
repo := NewCampaignProgressRepository(pool)
|
||||
used, err := repo.CampaignContactSentFromAccount(ctx, f.campaign, f.contact, f.other)
|
||||
if err != nil {
|
||||
t.Fatalf("check unused SMTP/IMAP mailbox: %v", err)
|
||||
}
|
||||
if used {
|
||||
t.Fatal("unused workspace mailbox was accepted for a cross-mailbox reply")
|
||||
}
|
||||
|
||||
followup := f.step(2, "Re: Hello", true)
|
||||
followupTask := f.send(followup, f.other, "<followup@test.local>", "thread-1", 60)
|
||||
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, followup, followupTask); err != nil {
|
||||
t.Fatalf("insert cross-provider progress: %v", err)
|
||||
}
|
||||
|
||||
used, err = repo.CampaignContactSentFromAccount(ctx, f.campaign, f.contact, f.other)
|
||||
if err != nil {
|
||||
t.Fatalf("check used SMTP/IMAP mailbox: %v", err)
|
||||
}
|
||||
if !used {
|
||||
t.Fatal("mailbox that sent a later campaign step was rejected for a cross-provider reply")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -185,6 +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)
|
||||
// CampaignContactSentFromAccount confirms the mailbox sent this contact a campaign step.
|
||||
CampaignContactSentFromAccount(ctx context.Context, campaignID, contactID, emailAccountID uuid.UUID) (bool, error)
|
||||
// ClaimIncomingReply leases one stored inbound message for reply processing.
|
||||
ClaimIncomingReply(ctx context.Context, emailAccountID, messageID uuid.UUID) (uuid.UUID, error)
|
||||
// CompleteIncomingReply prevents a successfully processed message from being retried.
|
||||
@@ -810,6 +812,24 @@ func (r *campaignProgressRepository) IsInboundReplySource(ctx context.Context, e
|
||||
return inbound, err
|
||||
}
|
||||
|
||||
// CampaignContactSentFromAccount proves a cross-mailbox reply arrived at a sender used for this lead.
|
||||
func (r *campaignProgressRepository) CampaignContactSentFromAccount(ctx context.Context, campaignID, contactID, emailAccountID uuid.UUID) (bool, error) {
|
||||
var sent bool
|
||||
err := r.db.QueryRow(ctx, `
|
||||
SELECT EXISTS (
|
||||
SELECT 1
|
||||
FROM campaign_tasks campaign_task
|
||||
JOIN tasks task ON task.id = campaign_task.task_id
|
||||
WHERE campaign_task.campaign_id = $1
|
||||
AND campaign_task.contact_id = $2
|
||||
AND task.task_type = 'campaign'
|
||||
AND task.status = 'completed'
|
||||
AND task.email_account_id = $3
|
||||
)
|
||||
`, campaignID, contactID, emailAccountID).Scan(&sent)
|
||||
return sent, err
|
||||
}
|
||||
|
||||
const incomingReplyClaimLease = "10 minutes"
|
||||
|
||||
// ClaimIncomingReply leases one inbound message so duplicate events cannot repeat its effects.
|
||||
|
||||
Reference in New Issue
Block a user