feat: store campaign reply classifications

This commit is contained in:
Matthew Meszaros
2026-06-07 11:46:38 +02:00
parent 3117284095
commit 384cf868db
3 changed files with 105 additions and 4 deletions
@@ -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;
@@ -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 <> '';
+61 -4
View File
@@ -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 {