feat: fully fence inbound reply processing and repair issue 549 state across campaigns and verification

This commit is contained in:
Matthew Meszaros
2026-09-16 22:03:46 -07:00
parent 35b94252e6
commit 5855d7e32d
9 changed files with 297 additions and 207 deletions
+1 -1
View File
@@ -221,7 +221,7 @@ 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.
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.
<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.
@@ -89,6 +89,7 @@ Some things belong to an instance rather than to a workspace, so they are not ap
| Plan limit overrides | A capacity grant is a decision by one platform's operators, not a property the workspace carries |
| Worker assignment | The destination places mailboxes on its own workers |
| Mailbox sync checkpoints | Replaying a checkpoint would make the destination skip everything that arrived between export and import, so it re-syncs from scratch |
| Reply processing checkpoints | Claims that prevent a synchronized inbox message from being handled twice belong to the source process. They reset so the destination can safely classify its imported inbox state |
| Pending warmup verification | Mail waiting for this instance's local or cloud warmup check is not exported. The destination re-syncs it from the provider and applies its own verification |
| Warmup pool membership | Pools are shared across every workspace on an instance, so membership is re-earned rather than asserted by a file |
| Domain authentication timings | The verdict travels (public DNS reads the same anywhere), but the destination re-checks before it can stop any sending, so a mailbox is never blocked on an observation the new instance never made |
+40 -4
View File
@@ -2,6 +2,7 @@ package advanced
import (
"context"
"errors"
"testing"
"github.com/google/uuid"
@@ -78,6 +79,7 @@ type incomingReplyProgressRepo struct {
sourceInbound bool
sourceClaimed bool
replyAccepted bool
completeErr error
advanced *incomingReplyAdvancedRepo
}
@@ -98,14 +100,17 @@ func (r *incomingReplyProgressRepo) RecordReplyClassification(context.Context, u
return nil
}
func (r *incomingReplyProgressRepo) ClaimIncomingReply(context.Context, uuid.UUID, uuid.UUID) (bool, error) {
func (r *incomingReplyProgressRepo) ClaimIncomingReply(context.Context, uuid.UUID, uuid.UUID) (uuid.UUID, error) {
r.claims++
return r.sourceClaimed, nil
if !r.sourceClaimed {
return uuid.Nil, nil
}
return uuid.New(), nil
}
func (r *incomingReplyProgressRepo) CompleteIncomingReply(context.Context, uuid.UUID, uuid.UUID) error {
func (r *incomingReplyProgressRepo) CompleteIncomingReply(context.Context, uuid.UUID, uuid.UUID, uuid.UUID) error {
r.completed++
return nil
return r.completeErr
}
func (r *incomingReplyProgressRepo) RecordEmailReplied(context.Context, uuid.UUID, uuid.UUID, uuid.UUID, uuid.UUID, uuid.UUID) (bool, error) {
@@ -420,3 +425,34 @@ func TestProcessIncomingReplyAcceptsMatchingThreadSender(t *testing.T) {
t.Fatalf("CompleteIncomingReply calls = %d, want 1 for the matching contact", progress.completed)
}
}
func TestProcessIncomingReplyFencesExpiredClaimBeforeSideEffects(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.completeErr = errors.New("claim was replaced")
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("expected the replaced claim to stop processing")
}
if progress.completed != 1 {
t.Fatalf("CompleteIncomingReply calls = %d, want 1", progress.completed)
}
if progress.advanced.marked != 0 {
t.Fatalf("MarkVariantEvent calls = %d, want 0 after ownership was lost", progress.advanced.marked)
}
if progress.advanced.intents != 0 {
t.Fatalf("CreateReplyIntent calls = %d, want 0 after ownership was lost", progress.advanced.intents)
}
}
+32 -26
View File
@@ -1112,8 +1112,8 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
// model call. verdict is what it decided, read after the block; held is
// when an out-of-office hold lifts, for the notification to name.
var verdict replyclassify.Result
var held *time.Time
replyClaimed := false
replyClaimToken := uuid.Nil
replyClaimCompleted := false
var campaignID *uuid.UUID
var sequenceID *uuid.UUID
@@ -1190,6 +1190,15 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
return true
}
claimToken, err := s.campaignProgressRepo.ClaimIncomingReply(ctx, emailAccountID, msg.ID)
if err != nil {
return toErrx(err)
}
if claimToken == uuid.Nil {
return nil
}
replyClaimToken = claimToken
replyResult := replyclassify.ClassifyGated(ctx, replyclassify.Input{
Headers: buildReplyHeaders(msg),
Subject: msg.Subject,
@@ -1201,25 +1210,8 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
// for every reply, so OOO/unsubscribe stay correct even when the gate
// skipped the model.
verdict = replyResult
claimed, err := s.campaignProgressRepo.ClaimIncomingReply(ctx, emailAccountID, msg.ID)
if err != nil {
return toErrx(err)
}
if !claimed {
return nil
}
replyClaimed = true
_ = s.campaignProgressRepo.RecordReplyClassification(ctx, cID, ctID, sID, replyResult.Class, replyResult.Source, replyResult.Confidence)
// Out of office: park the contact's next step until they are back
// rather than writing to an empty desk. An automated reply never
// stamps replied_at, so without this the follow-up goes out on
// schedule and the sequence is over before they read any of it
// (issue #470).
if replyResult.Class == replyclassify.ClassOutOfOffice && settings.ReplyIntent.HoldOnOutOfOffice {
held = s.holdForOutOfOffice(ctx, ctID, settings.ReplyIntent, msg)
}
// OOO trap fix: only a HUMAN reply stamps replied_at. An auto_reply /
// out_of_office must NOT count as a reply, or it would (a) trip
// stop_on_reply and silently halt the sequence, and (b) match the plain
@@ -1233,7 +1225,7 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
return toErrx(err)
}
if !accepted {
if err := s.campaignProgressRepo.CompleteIncomingReply(ctx, emailAccountID, msg.ID); err != nil {
if err := s.campaignProgressRepo.CompleteIncomingReply(ctx, emailAccountID, msg.ID, replyClaimToken); err != nil {
return toErrx(err)
}
return nil
@@ -1246,6 +1238,12 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
}
s.evidence.RecordEvidence(ctx, ctID, models.Step(&cID, &sID), kind, msg.ID.String(), "")
}
// Fence the claim before non-idempotent effects so an expired worker stops here.
if err := s.campaignProgressRepo.CompleteIncomingReply(ctx, emailAccountID, msg.ID, replyClaimToken); err != nil {
return toErrx(err)
}
replyClaimCompleted = true
if !replyclassify.IsAutomated(replyResult.Class) {
_ = s.repo.MarkVariantEvent(ctx, cID, ctID, string(models.DeliverabilityEventReply))
@@ -1303,14 +1301,15 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
BodyText: firstNonEmpty(msg.BodyText, msg.Snippet),
})
}
if !replyClaimed {
claimed, err := s.campaignProgressRepo.ClaimIncomingReply(ctx, emailAccountID, msg.ID)
if replyClaimToken == uuid.Nil && !replyClaimCompleted {
claimToken, err := s.campaignProgressRepo.ClaimIncomingReply(ctx, emailAccountID, msg.ID)
if err != nil {
return toErrx(err)
}
if !claimed {
if claimToken == uuid.Nil {
return nil
}
replyClaimToken = claimToken
}
intent, confidence := classifyReply(text, settings.ReplyIntent)
@@ -1321,6 +1320,16 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
if replyclassify.IsAutomated(verdict.Class) {
intent, confidence = automatedIntent(verdict)
}
if !replyClaimCompleted {
if err := s.campaignProgressRepo.CompleteIncomingReply(ctx, emailAccountID, msg.ID, replyClaimToken); err != nil {
return toErrx(err)
}
}
var held *time.Time
if campaignID != nil && contactID != nil && verdict.Class == replyclassify.ClassOutOfOffice && settings.ReplyIntent.HoldOnOutOfOffice {
held = s.holdForOutOfOffice(ctx, *contactID, settings.ReplyIntent, msg)
}
actionTaken := ""
if settings.ReplyIntent.AutoPauseOnNegative && intent == models.ReplyIntentNegative && campaignID != nil {
@@ -1456,9 +1465,6 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid.
s.notify(uid, account.OrganizationID, cat, title, body, "/app/unibox", map[string]any{"intent": string(intent)})
}
if err := s.campaignProgressRepo.CompleteIncomingReply(ctx, emailAccountID, msg.ID); err != nil {
return toErrx(err)
}
return nil
}
+2 -1
View File
@@ -583,7 +583,8 @@ var Tables = []Table{
},
{
Name: "unibox_emails", Group: models.OrgDataGroupInbox,
Scope: `email_id IN ` + orgMailboxes,
Scope: `email_id IN ` + orgMailboxes,
ResetOnImport: []string{"campaign_reply_claimed_at", "campaign_reply_claim_token", "campaign_reply_processed_at"},
},
{
Name: "unibox_thread_labels", Group: models.OrgDataGroupInbox,
@@ -1,5 +1,5 @@
-- The data repair is one-way because restoring false reply markers would
-- restore issue #549. Only the processing-claim columns are reversible.
-- False reply data cannot be restored; only processing-claim columns are reversible.
ALTER TABLE unibox_emails
DROP COLUMN campaign_reply_processed_at,
DROP COLUMN campaign_reply_claim_token,
DROP COLUMN campaign_reply_claimed_at;
@@ -1,119 +1,188 @@
-- 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.
-- Repair high-confidence issue #549 false replies and add durable reply-processing claims.
ALTER TABLE unibox_emails
ADD COLUMN campaign_reply_claimed_at timestamptz,
ADD COLUMN campaign_reply_claim_token uuid,
ADD COLUMN campaign_reply_processed_at timestamptz;
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 (
CREATE TEMP TABLE issue_549_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,
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 '5 minutes'
OR EXISTS (
SELECT 1
FROM unnest(ue.in_reply_to) AS refs(message_id)
WHERE btrim(refs.message_id, '<> ') = btrim(parent_task.message_id, '<> ')
FROM reply_intents sender_intent
WHERE sender_intent.campaign_id = p.campaign_id
AND sender_intent.task_id = p.dispatch_task_id
AND sender_intent.created_at BETWEEN p.replied_at - INTERVAL '30 seconds'
AND p.replied_at + INTERVAL '30 seconds'
AND lower(sender_intent.contact_email) IN (
lower(source_account.email),
lower(COALESCE(NULLIF(source_account.send_as_email, ''), source_account.email))
)
)
AND EXISTS (
SELECT 1
FROM unnest(ue.to_addr) AS recipients(address)
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')
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'
)
)
)
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(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 NOT EXISTS (
SELECT 1
FROM unibox_emails inbound
JOIN email_accounts inbound_account
ON inbound_account.id = inbound.email_id
AND inbound_account.organization_id = campaign.organization_id
WHERE inbound.folder NOT IN ('sent', 'drafts')
AND inbound.provider_folder NOT IN ('sent', 'drafts')
AND inbound.created_at BETWEEN p.replied_at - INTERVAL '30 seconds'
AND p.replied_at + INTERVAL '30 seconds'
AND EXISTS (
SELECT 1
FROM unnest(inbound.from_addr) AS senders(address)
WHERE lower(btrim(COALESCE(
substring(senders.address FROM '<([^>]*)>'),
substring(senders.address FROM '\(([^()]*)\)\s*$'),
senders.address
))) = lower(btrim(c.email))
)
)
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
USING issue_549_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(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')
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'
)
WITH affected AS (
SELECT DISTINCT contact_id
FROM issue_549_false_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 'unknown'
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 LEAST(100, floor(remaining.positive_score / 2)::integer)
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 'remaining mail evidence confirms the address'
ELSE 'reply evidence corrected; verification pending'
END,
updated_at = NOW()
FROM remaining
WHERE contact.id = remaining.contact_id;
DELETE FROM reply_intents intent
USING false_replies false_reply, email_accounts sending_account
USING issue_549_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'
@@ -124,58 +193,14 @@ WHERE intent.campaign_id = false_reply.campaign_id
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(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')
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 (
WITH 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
FROM issue_549_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
@@ -196,3 +221,5 @@ WHERE EXISTS (
AND progress.contact_id = assignment.contact_id
AND progress.replied_at IS NOT NULL
);
DROP TABLE issue_549_false_replies;
@@ -43,11 +43,11 @@ func TestLiveReplySourceUsesStoredDirectionAtTheWriteBoundary(t *testing.T) {
if inbound {
t.Fatal("Sent-folder source was accepted as inbound")
}
claimed, err := repo.ClaimIncomingReply(ctx, f.other, sentID)
claimToken, err := repo.ClaimIncomingReply(ctx, f.other, sentID)
if err != nil {
t.Fatalf("claim sent source: %v", err)
}
if claimed {
if claimToken != uuid.Nil {
t.Fatal("Sent-folder source was claimed for reply processing")
}
accepted, err := repo.RecordEmailReplied(ctx, f.campaign, f.contact, step, f.other, sentID)
@@ -90,34 +90,38 @@ func TestLiveReplySourceUsesStoredDirectionAtTheWriteBoundary(t *testing.T) {
if !inbound {
t.Fatal("Inbox source was rejected as outbound")
}
claimed, err = repo.ClaimIncomingReply(ctx, f.mailbox, inboxID)
claimToken, err = repo.ClaimIncomingReply(ctx, f.mailbox, inboxID)
if err != nil {
t.Fatalf("claim inbound source: %v", err)
}
if !claimed {
if claimToken == uuid.Nil {
t.Fatal("Inbox source was not claimed for reply processing")
}
claimed, err = repo.ClaimIncomingReply(ctx, f.mailbox, inboxID)
firstClaimToken := claimToken
claimToken, err = repo.ClaimIncomingReply(ctx, f.mailbox, inboxID)
if err != nil {
t.Fatalf("repeat inbound claim: %v", err)
}
if claimed {
if claimToken != uuid.Nil {
t.Fatal("Inbox source was claimed concurrently twice")
}
if _, err := pool.Exec(ctx, `
UPDATE unibox_emails
SET campaign_reply_claimed_at = NOW() - INTERVAL '3 minutes'
SET campaign_reply_claimed_at = NOW() - INTERVAL '11 minutes'
WHERE id = $1
`, inboxID); err != nil {
t.Fatalf("age reply claim: %v", err)
}
claimed, err = repo.ClaimIncomingReply(ctx, f.mailbox, inboxID)
claimToken, err = repo.ClaimIncomingReply(ctx, f.mailbox, inboxID)
if err != nil {
t.Fatalf("reclaim expired lease: %v", err)
}
if !claimed {
if claimToken == uuid.Nil {
t.Fatal("Expired reply-processing lease was not reclaimed")
}
if err := repo.CompleteIncomingReply(ctx, f.mailbox, inboxID, firstClaimToken); err == nil {
t.Fatal("Expired claim completed work owned by its replacement")
}
accepted, err = repo.RecordEmailReplied(ctx, f.campaign, f.contact, step, f.mailbox, inboxID)
if err != nil {
t.Fatalf("record inbound source: %v", err)
@@ -132,14 +136,17 @@ func TestLiveReplySourceUsesStoredDirectionAtTheWriteBoundary(t *testing.T) {
if !accepted {
t.Fatal("Already-recorded reply was not accepted idempotently")
}
if err := repo.CompleteIncomingReply(ctx, f.mailbox, inboxID); err != nil {
if err := repo.CompleteIncomingReply(ctx, f.mailbox, inboxID, claimToken); err != nil {
t.Fatalf("complete inbound reply: %v", err)
}
claimed, err = repo.ClaimIncomingReply(ctx, f.mailbox, inboxID)
if err := repo.CompleteIncomingReply(ctx, f.mailbox, inboxID, claimToken); err == nil {
t.Fatal("Completed claim was accepted twice")
}
claimToken, err = repo.ClaimIncomingReply(ctx, f.mailbox, inboxID)
if err != nil {
t.Fatalf("claim completed reply: %v", err)
}
if claimed {
if claimToken != uuid.Nil {
t.Fatal("Completed reply was claimed again")
}
if err := pool.QueryRow(ctx, `
+25 -13
View File
@@ -186,9 +186,9 @@ type CampaignProgressRepository interface {
// IsInboundReplySource confirms the stored unibox row is not outbound.
IsInboundReplySource(ctx context.Context, emailAccountID, messageID uuid.UUID) (bool, error)
// ClaimIncomingReply leases one stored inbound message for reply processing.
ClaimIncomingReply(ctx context.Context, emailAccountID, messageID uuid.UUID) (bool, error)
ClaimIncomingReply(ctx context.Context, emailAccountID, messageID uuid.UUID) (uuid.UUID, error)
// CompleteIncomingReply prevents a successfully processed message from being retried.
CompleteIncomingReply(ctx context.Context, emailAccountID, messageID uuid.UUID) error
CompleteIncomingReply(ctx context.Context, emailAccountID, messageID, claimToken uuid.UUID) error
// RecordEmailReplied stamps a human reply while its source is still inbound.
RecordEmailReplied(ctx context.Context, campaignID, contactID, sequenceID, emailAccountID, messageID uuid.UUID) (bool, error)
RecordEmailBounced(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error
@@ -810,13 +810,15 @@ func (r *campaignProgressRepository) IsInboundReplySource(ctx context.Context, e
return inbound, err
}
const incomingReplyClaimLease = "2 minutes"
const incomingReplyClaimLease = "10 minutes"
// ClaimIncomingReply leases one inbound message so duplicate events cannot repeat its effects.
func (r *campaignProgressRepository) ClaimIncomingReply(ctx context.Context, emailAccountID, messageID uuid.UUID) (bool, error) {
func (r *campaignProgressRepository) ClaimIncomingReply(ctx context.Context, emailAccountID, messageID uuid.UUID) (uuid.UUID, error) {
claimToken := uuid.New()
result, err := r.db.Exec(ctx, `
UPDATE unibox_emails
SET campaign_reply_claimed_at = NOW()
SET campaign_reply_claimed_at = NOW(),
campaign_reply_claim_token = $3
WHERE id = $1
AND email_id = $2
AND folder NOT IN ('sent', 'drafts')
@@ -826,23 +828,33 @@ func (r *campaignProgressRepository) ClaimIncomingReply(ctx context.Context, ema
campaign_reply_claimed_at IS NULL
OR campaign_reply_claimed_at < NOW() - INTERVAL '`+incomingReplyClaimLease+`'
)
`, messageID, emailAccountID)
`, messageID, emailAccountID, claimToken)
if err != nil {
return false, err
return uuid.Nil, err
}
return result.RowsAffected() == 1, nil
if result.RowsAffected() == 0 {
return uuid.Nil, nil
}
return claimToken, nil
}
// CompleteIncomingReply marks a claimed message as processed.
func (r *campaignProgressRepository) CompleteIncomingReply(ctx context.Context, emailAccountID, messageID uuid.UUID) error {
_, err := r.db.Exec(ctx, `
func (r *campaignProgressRepository) CompleteIncomingReply(ctx context.Context, emailAccountID, messageID, claimToken uuid.UUID) error {
result, err := r.db.Exec(ctx, `
UPDATE unibox_emails
SET campaign_reply_processed_at = NOW()
WHERE id = $1
AND email_id = $2
AND campaign_reply_claimed_at IS NOT NULL
`, messageID, emailAccountID)
return err
AND campaign_reply_claim_token = $3
AND campaign_reply_processed_at IS NULL
`, messageID, emailAccountID, claimToken)
if err != nil {
return err
}
if result.RowsAffected() == 0 {
return fmt.Errorf("incoming reply claim is no longer owned")
}
return nil
}
// RecordEmailReplied records that a contact replied when its source is still inbound.