From 384cf868db139df900d004c52f7fa00a763d1b09 Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Sun, 7 Jun 2026 11:46:38 +0200 Subject: [PATCH] feat: store campaign reply classifications --- .../000027_reply_classification.down.sql | 12 ++++ .../000027_reply_classification.up.sql | 32 +++++++++ internal/repository/pg_campaign_progress.go | 65 +++++++++++++++++-- 3 files changed, 105 insertions(+), 4 deletions(-) create mode 100644 internal/infrastructure/db/migrations/000027_reply_classification.down.sql create mode 100644 internal/infrastructure/db/migrations/000027_reply_classification.up.sql diff --git a/internal/infrastructure/db/migrations/000027_reply_classification.down.sql b/internal/infrastructure/db/migrations/000027_reply_classification.down.sql new file mode 100644 index 00000000..11867419 --- /dev/null +++ b/internal/infrastructure/db/migrations/000027_reply_classification.down.sql @@ -0,0 +1,12 @@ +DROP INDEX IF EXISTS idx_campaign_progress_reply_class; + +ALTER TABLE campaign_contact_progress + DROP CONSTRAINT IF EXISTS campaign_contact_progress_reply_source_chk; + +ALTER TABLE campaign_contact_progress + DROP CONSTRAINT IF EXISTS campaign_contact_progress_reply_class_chk; + +ALTER TABLE campaign_contact_progress + DROP COLUMN IF EXISTS reply_source, + DROP COLUMN IF EXISTS reply_confidence, + DROP COLUMN IF EXISTS reply_class; diff --git a/internal/infrastructure/db/migrations/000027_reply_classification.up.sql b/internal/infrastructure/db/migrations/000027_reply_classification.up.sql new file mode 100644 index 00000000..8f0b22e4 --- /dev/null +++ b/internal/infrastructure/db/migrations/000027_reply_classification.up.sql @@ -0,0 +1,32 @@ +-- Reply classification: persist the layered classifier verdict for each inbound +-- reply on the per-(campaign, contact, step) progress row. This is what reply +-- branching ("reply_positive" / "reply_negative" / ...) reads at schedule time, +-- and what lets stop_on_reply / the plain "replied" condition IGNORE automated +-- replies (auto_reply / out_of_office) instead of treating them as a human reply. +-- +-- reply_class enum string: positive | negative | neutral | auto_reply | +-- out_of_office | unsubscribe | unknown ('' / NULL when unset). +-- reply_confidence classifier confidence in [0,1]. +-- reply_source which layer decided: header | lexicon | model | '' (unset). +-- +-- Kept as typed text columns (not jsonb) because they are queried/filtered in +-- SQL by the branch evaluator and the OOO-aware reply checks. +ALTER TABLE campaign_contact_progress + ADD COLUMN IF NOT EXISTS reply_class text NOT NULL DEFAULT '', + ADD COLUMN IF NOT EXISTS reply_confidence real NOT NULL DEFAULT 0, + ADD COLUMN IF NOT EXISTS reply_source text NOT NULL DEFAULT ''; + +-- Discriminator guards so a bad write can't poison branch routing. +ALTER TABLE campaign_contact_progress + ADD CONSTRAINT campaign_contact_progress_reply_class_chk + CHECK (reply_class IN ('', 'positive', 'negative', 'neutral', 'auto_reply', 'out_of_office', 'unsubscribe', 'unknown')); + +ALTER TABLE campaign_contact_progress + ADD CONSTRAINT campaign_contact_progress_reply_source_chk + CHECK (reply_source IN ('', 'header', 'lexicon', 'model')); + +-- Partial index for the reply-branch evaluator, which only ever reads rows that +-- carry a non-empty class. +CREATE INDEX IF NOT EXISTS idx_campaign_progress_reply_class + ON campaign_contact_progress (campaign_id, contact_id) + WHERE reply_class <> ''; diff --git a/internal/repository/pg_campaign_progress.go b/internal/repository/pg_campaign_progress.go index 6ea70843..3e474612 100644 --- a/internal/repository/pg_campaign_progress.go +++ b/internal/repository/pg_campaign_progress.go @@ -22,6 +22,12 @@ type CampaignContactProgress struct { RepliedAt *time.Time BouncedAt *time.Time ComplainedAt *time.Time + // ReplyClass is the layered classifier verdict for the contact's reply + // (positive | negative | neutral | auto_reply | out_of_office | unsubscribe | + // unknown; "" when no reply was classified). Read by the reply_* branch + // conditions. RepliedAt is set ONLY for human replies, so an automated reply + // can carry a ReplyClass here without ever tripping "replied"/stop_on_reply. + ReplyClass string } // CampaignProgress represents overall campaign progress @@ -70,6 +76,18 @@ type CampaignProgressRepository interface { RecordEmailBounced(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error RecordEmailComplained(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error + // RecordReplyClassification stores the layered classifier verdict + // (class/confidence/source) for the contact's reply on the given step. This + // is the data the reply_* branch conditions read. It does NOT stamp + // replied_at — only a HUMAN reply should set replied_at (so automated replies + // never trip stop_on_reply / the "replied" condition). Callers stamp + // replied_at separately via RecordEmailReplied for human replies only. + RecordReplyClassification(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID, class, source string, confidence float64) error + // GetLatestReplyClass returns the most-recent classified reply class for a + // contact in a campaign ("" when none). Convenience getter for the branch + // evaluator / callers that need only the class. + GetLatestReplyClass(ctx context.Context, contactID, campaignID uuid.UUID) (string, error) + // Query methods GetCampaignProgress(ctx context.Context, campaignID uuid.UUID) (*CampaignProgress, error) GetCampaignRollingRates(ctx context.Context, campaignID uuid.UUID, since time.Time) (*CampaignRollingRates, error) @@ -186,6 +204,42 @@ func (r *campaignProgressRepository) RecordEmailComplained(ctx context.Context, return err } +// RecordReplyClassification persists the classifier verdict on the progress row. +// It upserts so the classification lands even if the reply arrives before the +// step's progress row is materialized (rare, but threading can race). It never +// touches replied_at — human-vs-automated gating lives in the caller, which +// stamps replied_at via RecordEmailReplied only for human replies. +func (r *campaignProgressRepository) RecordReplyClassification(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID, class, source string, confidence float64) error { + query := ` + INSERT INTO campaign_contact_progress (campaign_id, contact_id, sequence_id, reply_class, reply_confidence, reply_source) + VALUES ($1, $2, $3, $4, $5, $6) + ON CONFLICT (campaign_id, contact_id, sequence_id) + DO UPDATE SET reply_class = EXCLUDED.reply_class, + reply_confidence = EXCLUDED.reply_confidence, + reply_source = EXCLUDED.reply_source + ` + _, err := r.db.Exec(ctx, query, campaignID, contactID, sequenceID, class, confidence, source) + return err +} + +// GetLatestReplyClass returns the most-recent non-empty reply_class for a +// contact in a campaign, or "" when none has been classified. +func (r *campaignProgressRepository) GetLatestReplyClass(ctx context.Context, contactID, campaignID uuid.UUID) (string, error) { + query := ` + SELECT reply_class + FROM campaign_contact_progress + WHERE contact_id = $1 AND campaign_id = $2 AND reply_class <> '' + ORDER BY COALESCE(sent_at, '-infinity'::timestamptz) DESC + LIMIT 1 + ` + var class string + err := r.db.QueryRow(ctx, query, contactID, campaignID).Scan(&class) + if err == sql.ErrNoRows { + return "", nil + } + return class, err +} + // GetCampaignProgress retrieves overall campaign progress statistics func (r *campaignProgressRepository) GetCampaignProgress(ctx context.Context, campaignID uuid.UUID) (*CampaignProgress, error) { query := ` @@ -262,7 +316,7 @@ func (r *campaignProgressRepository) GetCampaignRollingRates(ctx context.Context // GetContactProgress retrieves progress for a specific contact in a campaign func (r *campaignProgressRepository) GetContactProgress(ctx context.Context, campaignID, contactID uuid.UUID) ([]CampaignContactProgress, error) { query := ` - SELECT campaign_id, contact_id, sequence_id, sent_at, opened_at, clicked_at, replied_at, bounced_at, complained_at + SELECT campaign_id, contact_id, sequence_id, sent_at, opened_at, clicked_at, replied_at, bounced_at, complained_at, COALESCE(reply_class, '') FROM campaign_contact_progress WHERE campaign_id = $1 AND contact_id = $2 ORDER BY sent_at ASC @@ -287,6 +341,7 @@ func (r *campaignProgressRepository) GetContactProgress(ctx context.Context, cam &progress.RepliedAt, &progress.BouncedAt, &progress.ComplainedAt, + &progress.ReplyClass, ) if err != nil { return nil, err @@ -498,12 +553,12 @@ func (r *campaignProgressRepository) FindNextRoutedPair(ctx context.Context, cam query := ` SELECT cl.contact_id, - lp.sequence_id, lp.sent_at, lp.opened_at, lp.clicked_at, lp.replied_at, + lp.sequence_id, lp.sent_at, lp.opened_at, lp.clicked_at, lp.replied_at, COALESCE(lp.reply_class, ''), COALESCE(ss.ids, '{}') AS sent_ids FROM campaign_leads cl JOIN contacts c ON c.id = cl.contact_id LEFT JOIN LATERAL ( - SELECT sequence_id, sent_at, opened_at, clicked_at, replied_at + SELECT sequence_id, sent_at, opened_at, clicked_at, replied_at, reply_class FROM campaign_contact_progress p WHERE p.campaign_id = $1 AND p.contact_id = cl.contact_id AND p.sent_at IS NOT NULL ORDER BY p.sent_at DESC LIMIT 1 @@ -546,8 +601,9 @@ func (r *campaignProgressRepository) FindNextRoutedPair(ctx context.Context, cam var contactID uuid.UUID var lastSeq *uuid.UUID var sentAt, openedAt, clickedAt, repliedAt *time.Time + var replyClass string var sentIDs []uuid.UUID - if serr := rows.Scan(&contactID, &lastSeq, &sentAt, &openedAt, &clickedAt, &repliedAt, &sentIDs); serr != nil { + if serr := rows.Scan(&contactID, &lastSeq, &sentAt, &openedAt, &clickedAt, &repliedAt, &replyClass, &sentIDs); serr != nil { return nil, nil, serr } @@ -563,6 +619,7 @@ func (r *campaignProgressRepository) FindNextRoutedPair(ctx context.Context, cam prog := &CampaignContactProgress{ CampaignID: campaignID, ContactID: contactID, SequenceID: *lastSeq, SentAt: sentAt, OpenedAt: openedAt, ClickedAt: clickedAt, RepliedAt: repliedAt, + ReplyClass: replyClass, } sa := time.Time{} if sentAt != nil {