Files
warmbly/internal/tasks/campaign_task.go
T
Matthew Meszaros dd98187231 fix: shorten recipient unsubscribe links to 22-character, 128-bit stored tickets (#498) (#525)
* feat: shorten every recipient unsubscribe link from a 96-character signed token to a 22-character stored ticket carrying 128 bits from crypto/rand, minted once per recipient per campaign and reused by every step, so the address the text/plain half of a cold email prints in full fits on one line and cannot be guessed, keeping the signed form working for links already in inboxes and as the fallback when the store cannot be written, and answering a failed lookup with a retryable 'try again shortly' instead of telling the recipient their opt-out is invalid (issue #498)

* fix: restore the disabled-signer guard in URLOn, which factoring the URL builder moved behind a token mint that dereferences the signing key, so a nil or origin-less signer returns the empty string every caller reads as 'no link can be minted' instead of panicking (PR #525 review)
2026-09-15 00:42:45 -07:00

1731 lines
70 KiB
Go

package tasks
import (
"context"
"encoding/json"
"errors"
"fmt"
"strconv"
"strings"
"time"
"github.com/google/uuid"
"github.com/rs/zerolog/log"
"github.com/warmbly/warmbly/internal/config"
"github.com/warmbly/warmbly/internal/errx"
"github.com/warmbly/warmbly/internal/infrastructure/pubsub"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/observability/errs"
"github.com/warmbly/warmbly/internal/pkg/mailhtml"
"github.com/warmbly/warmbly/internal/repository"
"github.com/warmbly/warmbly/internal/scheduler"
"github.com/warmbly/warmbly/internal/tasks/proto"
)
func (s *tasksService) HandleCampaignTask(task *proto.ProcessTask) *errx.Error {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
// STEP 1: Parse task ID
taskID, err := uuid.Parse(task.TaskId)
if err != nil {
errs.CaptureException(err)
return errx.New(errx.BadRequest, "invalid task ID")
}
executionKey := "campaign:" + taskID.String()
executionStatus := "failed"
if s.advanced != nil {
duplicate, xerr := s.advanced.StartTaskExecution(ctx, taskID, executionKey, map[string]interface{}{
"task_type": "campaign",
})
if xerr != nil {
return xerr
}
if duplicate {
return nil
}
defer func() {
_ = s.advanced.CompleteTaskExecution(ctx, taskID, executionKey, executionStatus, map[string]interface{}{
"task_type": "campaign",
})
}()
}
// STEP 2: Load task record
taskRecord, err := s.taskRepo.GetTask(ctx, taskID)
if err != nil {
errs.CaptureException(err)
return errx.InternalError()
}
if taskRecord == nil {
return errx.ErrNotFound
}
if taskRecord.Status != "pending" {
log.Info().
Str("task_id", taskID.String()).
Str("status", taskRecord.Status).
Msg("campaign task skipped: task not in pending state")
executionStatus = "completed"
return nil
}
// STEP 3: Mark task as active (with advisory lock)
if err := s.taskRepo.UpdateTaskStatusWithLock(ctx, taskID, "active"); err != nil {
errs.CaptureException(err)
return errx.InternalError()
}
// STEP 4: Load campaign task details
campaignTask, err := s.taskRepo.GetCampaignTask(ctx, taskID)
if err != nil {
errs.CaptureException(err)
return errx.InternalError()
}
if campaignTask == nil {
return errx.ErrNotFound
}
// campaign_tasks nulls its link when the campaign is deleted: the chain
// ends here rather than leaving the row claimed as active.
if campaignTask.CampaignID == nil {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "cancelled")
executionStatus = "completed"
return nil
}
// Get campaign progress for task progress events
campaignProgress, _ := s.campaignProgressRepo.GetCampaignProgress(ctx, *campaignTask.CampaignID)
var totalContacts, processedCount int
if campaignProgress != nil {
totalContacts = campaignProgress.TotalContacts
processedCount = campaignProgress.EmailsSent
}
// STEP 5: Load campaign
campaign, err := s.campaignRepo.GetByID(ctx, *campaignTask.CampaignID)
if err != nil {
// Deleted between the pending check and here: the chain ends.
if errors.Is(err, errx.ErrResourceNotFound) {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "cancelled")
executionStatus = "completed"
return nil
}
errs.CaptureException(err)
return errx.InternalError()
}
// Check if campaign is still active
if campaign.Status != "active" {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "cancelled")
executionStatus = "completed"
return nil // Don't create next task
}
// Publish task started progress event
if s.streamingPublisher != nil {
progress := 0
if totalContacts > 0 {
progress = (processedCount * 100) / totalContacts
}
s.streamingPublisher.PublishTaskProgress(ctx, &pubsub.TaskProgressEvent{
BaseEvent: pubsub.BaseEvent{UserID: campaign.UserID},
OrgID: campaignOrgID(campaign),
CampaignID: campaign.ID.String(),
TaskID: taskID.String(),
Status: "active",
Progress: progress,
TotalContacts: totalContacts,
ProcessedCount: processedCount,
})
}
// STEP 5.4: Tenancy gate. organization_id is what scopes the entitlement
// check below and the recipient suppression check in STEP 7; missing it used
// to skip both and send anyway. Fail closed instead — an orgless campaign is
// stopped and surfaced, never mailed unchecked.
if campaign.OrganizationID == nil {
s.haltOrglessCampaign(ctx, campaign.ID, taskID)
executionStatus = "completed"
return nil
}
orgID := *campaign.OrganizationID
// STEP 5.5: Check if organization can send campaign emails (trial expired, etc.)
if s.featureGate != nil {
canSend, _ := s.featureGate.CanSendCampaignEmail(ctx, orgID)
if !canSend {
// Organization cannot send - pause campaign
s.campaignRepo.UpdateStatus(ctx, campaign.ID, "paused_trial_expired")
s.taskRepo.UpdateTaskStatus(ctx, taskID, "skipped_trial_expired")
executionStatus = "completed"
return nil
}
// Check daily limit
limit, _ := s.featureGate.GetDailyEmailLimit(ctx, orgID)
if limit >= 0 {
sentToday, err := s.campaignProgressRepo.CountEmailsSentTodayByOrganization(ctx, orgID)
if err == nil && sentToday >= limit {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "skipped_daily_limit")
if s.campaignLogRepo != nil {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "daily_limit_reached",
Message: "Campaign paused for today: organization daily email limit reached",
Metadata: map[string]interface{}{
"sent_today": sentToday,
"limit": limit,
},
})
}
// Reschedule to the next day to keep campaign progression alive.
nextDay := time.Now().UTC().Truncate(24 * time.Hour).Add(24 * time.Hour).Add(5 * time.Minute)
_, _, nextAccountID, calcErr := s.scheduler.CalculateNextCampaignTime(ctx, *campaignTask.CampaignID)
if calcErr == nil || errors.Is(calcErr, scheduler.ErrCampaignDeferred) {
if err := s.createCampaignTask(ctx, campaign.ID, nextAccountID, nextDay); err != nil {
log.Warn().Err(err).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Msg("Failed to create next campaign task after daily limit")
}
}
executionStatus = "completed"
return nil
}
}
}
// STEP 6: Calculate next email to send
nextTime, nextPair, accountID, err := s.scheduler.CalculateNextCampaignTime(ctx, *campaignTask.CampaignID)
if err != nil {
if errors.Is(err, scheduler.ErrNoEmailAccounts) {
s.autoPauseCampaign(ctx, *campaignTask.CampaignID, taskID, autoPauseReason(err))
executionStatus = "completed"
return nil
}
if errors.Is(err, scheduler.ErrCampaignDeferred) {
// A valid contact exists but no eligible mailbox right now (ESP-strict
// has no same-provider mailbox, the daily new-lead cap is reached, or
// every mailbox has spent its daily budget or is outside its hours).
// Reschedule at the deferred slot WITHOUT sending and WITHOUT touching
// progress / daily counters / rotation — mirrors the daily-limit path.
// This task completes without a send, and completing it must not
// spend the mailbox's budget either (issue #306): the budget counts
// reserved sends, never bare wake-ups.
// Capped: the next-due moment can be days out, and until this chain
// wakes nothing re-reads the campaign, so leads imported meanwhile
// would sit queued until then.
scheduledNext := scheduler.DeferSlot(nextTime)
if cerr := s.createCampaignTask(ctx, campaign.ID, accountID, scheduledNext); cerr != nil {
log.Warn().Err(cerr).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Msg("Failed to schedule deferred campaign task")
}
s.clearIdle(ctx, campaign)
s.taskRepo.UpdateTaskStatus(ctx, taskID, "completed")
executionStatus = "completed"
return nil
}
// Terminal: all emails sent, OR the campaign passed its end date. Both end
// the campaign at the "completed" status (the status enum has no separate
// "ended"); the reason differs in the activity log.
if errors.Is(err, scheduler.ErrCampaignCompleted) || errors.Is(err, scheduler.ErrCampaignEnded) {
reason := "Campaign completed: all emails sent"
if errors.Is(err, scheduler.ErrCampaignEnded) {
reason = "Campaign ended: reached its end date"
}
// Two counts stand between "nothing was routed" and "the campaign is
// finished", and BOTH have to be known before it can be closed:
// leads verification refused are never routed (issue #264), and a
// lead held with no end reports no next-due moment because there is
// none to wake up for (issue #470). Either would otherwise let
// "all emails sent" close a campaign that still has work.
//
// A count that FAILS is not zero. Closing on a database hiccup
// writes a claim nothing can walk back, so an unreadable count
// leaves the campaign active and the next pass asks again.
undeliverable, cerr := s.campaignProgressRepo.CountUndeliverableLeads(ctx, campaign.ID)
held, herr := s.campaignProgressRepo.CountHeldLeads(ctx, campaign.ID)
if cerr != nil || herr != nil {
log.Warn().AnErr("undeliverable", cerr).AnErr("held", herr).
Str("campaign_id", campaign.ID.String()).
Msg("could not tell a finished campaign from a parked one; leaving it active for the next pass")
if taskID != uuid.Nil {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "completed")
}
executionStatus = "completed"
return nil
}
if undeliverable > 0 {
if errors.Is(err, scheduler.ErrCampaignCompleted) {
s.pauseUndeliverable(ctx, campaign.ID, taskID, undeliverable)
executionStatus = "completed"
return nil
}
reason = fmt.Sprintf("%s (%d lead(s) skipped: address verification refused them)", reason, undeliverable)
}
if held > 0 {
if errors.Is(err, scheduler.ErrCampaignCompleted) {
s.parkHeldLeads(ctx, campaign, taskID, held)
executionStatus = "completed"
return nil
}
reason = fmt.Sprintf("%s (%d lead(s) paused)", reason, held)
}
// A continuous campaign out of leads is waiting, not finished
// (issue #336). Only its end date ends it.
if errors.Is(err, scheduler.ErrCampaignCompleted) && campaign.Continuous {
s.idleCampaign(ctx, campaign, taskID)
executionStatus = "completed"
return nil
}
s.campaignRepo.UpdateStatus(ctx, campaign.ID, "completed")
if s.campaignLogRepo != nil {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "completed",
Message: reason,
})
}
// Broadcast live so the dashboard (and the sidebar campaign counters)
// flip from "sending" to "finished" without a manual refresh.
if s.streamingPublisher != nil {
s.streamingPublisher.PublishCampaignEvent(ctx, &pubsub.CampaignEvent{
BaseEvent: pubsub.BaseEvent{
EventType: pubsub.EventCampaignCompleted,
UserID: campaign.UserID,
},
OrgID: campaignOrgID(campaign),
CampaignID: campaign.ID.String(),
Name: campaign.Name,
Status: "completed",
})
}
s.taskRepo.UpdateTaskStatus(ctx, taskID, "completed")
executionStatus = "completed"
return nil
}
// Benign: the campaign was paused/deleted between ticks. Stop this chain
// cleanly; a resume re-seeds it.
if errors.Is(err, scheduler.ErrCampaignNotActive) {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "cancelled")
executionStatus = "completed"
return nil
}
// Transient / unknown error (a DB blip bubbled up from the scheduler). Do
// NOT silently mark the task completed — that strands the campaign with no
// successor. Record the failure for dashboard review, reset the task to
// pending, and return 5xx so Cloud Tasks retries (with backoff). The
// campaign reconciler is the backstop if retries are ever exhausted.
errs.CaptureException(err)
s.recordSchedulerFailure(ctx, campaign.ID, "scheduler_error", "Could not compute the next step; retrying", err)
// Pulse the dashboard so the failure appears live for the whole team. A
// CAMPAIGN_UPDATED with empty status invalidates the campaign logs query
// without flipping the campaign's status.
if s.streamingPublisher != nil {
s.streamingPublisher.PublishCampaignEvent(ctx, &pubsub.CampaignEvent{
BaseEvent: pubsub.BaseEvent{EventType: pubsub.EventCampaignUpdated, UserID: campaign.UserID},
OrgID: campaignOrgID(campaign),
CampaignID: campaign.ID.String(),
})
}
if rerr := s.taskRepo.UpdateTaskStatus(ctx, taskID, "pending"); rerr != nil {
errs.CaptureException(rerr)
}
executionStatus = "failed"
return errx.InternalError()
}
s.clearIdle(ctx, campaign)
// STEP 7: Load contact and sequence
contact, xerr := s.contactRepo.GetByID(ctx, nextPair.ContactID)
if xerr != nil {
return xerr
}
if s.advanced != nil {
suppressed, reason, sxerr := s.advanced.ShouldSuppressRecipient(ctx, orgID, contact.Email)
if sxerr != nil {
return sxerr
}
if suppressed {
_ = s.taskRepo.UpdateTaskStatusWithLock(ctx, taskID, "skipped_suppressed")
if s.campaignLogRepo != nil {
_ = s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "suppressed",
Message: fmt.Sprintf("Suppressed recipient skipped: %s", contact.Email),
Metadata: map[string]interface{}{
"reason": reason,
},
})
}
_ = s.createCampaignTask(ctx, campaign.ID, accountID, nextTime)
executionStatus = "completed"
return nil
}
}
// Per-lead hold gate, and a backstop rather than the main defence: routing
// already excludes held leads, but a pause (or an out-of-office reply) can
// land between that read and this dispatch, and a hold the send raced is
// the one thing this feature exists to stop. Re-read where it is committed,
// alongside suppression and verification.
hold, herr := s.campaignProgressRepo.GetLeadHold(ctx, campaign.ID, contact.ID)
if errors.Is(herr, repository.ErrLeadNotInCampaign) {
// Removed from the campaign between routing and here. Not an error and
// not a hold: there is simply nobody to send to, so skip it the way the
// other pre-send gates do and let the chain carry on.
_ = s.taskRepo.UpdateTaskStatusWithLock(ctx, taskID, "skipped_suppressed")
_ = s.createCampaignTask(ctx, campaign.ID, accountID, nextTime)
executionStatus = "completed"
return nil
}
if herr != nil {
// Fail closed: an unreadable hold is not an absent one, and sending to
// somebody who asked not to be is the failure this gate exists for.
errs.CaptureException(herr)
s.taskRepo.RecordTaskFailure(ctx, taskID, "Could not read the lead's hold", herr.Error())
if uerr := s.taskRepo.UpdateTaskStatus(ctx, taskID, "pending"); uerr != nil {
errs.CaptureException(uerr)
}
executionStatus = "failed"
return errx.InternalError()
}
if hold != nil {
_ = s.taskRepo.UpdateTaskStatusWithLock(ctx, taskID, "skipped_paused")
if s.campaignLogRepo != nil {
_ = s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "paused_lead",
Message: fmt.Sprintf("Paused lead skipped: %s", contact.Email),
Metadata: map[string]interface{}{"reason": hold.Reason, "source": hold.Source},
})
}
_ = s.createCampaignTask(ctx, campaign.ID, accountID, nextTime)
executionStatus = "completed"
return nil
}
// Pre-send verification gate: drop addresses already known to be invalid
// (bad syntax / no MX / 550 RCPT) before a worker sends and earns a hard
// bounce. 'invalid' is always dropped; 'risky' is dropped only when the
// campaign's "send to risky emails" toggle is off (see the next gate).
// 'unknown'/'valid' always send.
if contact.VerificationStatus == "invalid" {
_ = s.taskRepo.UpdateTaskStatusWithLock(ctx, taskID, "skipped_suppressed")
if s.campaignLogRepo != nil {
_ = s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "suppressed",
Message: fmt.Sprintf("Unverifiable recipient skipped: %s", contact.Email),
Metadata: map[string]interface{}{"reason": contact.VerificationReason},
})
}
_ = s.createCampaignTask(ctx, campaign.ID, accountID, nextTime)
executionStatus = "completed"
return nil
}
// Risky-recipient gate: when "send to risky emails" is off, also drop
// addresses verification flagged 'risky' (catch-all / role / low-quality),
// which raise bounce risk. Enforces the campaign.RiskyEmails toggle that the
// settings UI exposes — without this the toggle is stored but inert.
if !campaign.RiskyEmails && contact.VerificationStatus == "risky" {
_ = s.taskRepo.UpdateTaskStatusWithLock(ctx, taskID, "skipped_suppressed")
if s.campaignLogRepo != nil {
_ = s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "suppressed",
Message: fmt.Sprintf("Risky recipient skipped (send to risky emails is off): %s", contact.Email),
Metadata: map[string]interface{}{"reason": contact.VerificationReason},
})
}
_ = s.createCampaignTask(ctx, campaign.ID, accountID, nextTime)
executionStatus = "completed"
return nil
}
sequence, err := s.campaignRepo.GetSequenceByID(ctx, nextPair.SequenceID)
if err != nil {
errs.CaptureException(err)
return errx.InternalError()
}
// STEP 7.6: Non-email nodes (action / wait). These run a control-plane side
// effect and route onward WITHOUT sending mail — the render/send block below
// is reached only for email nodes. We stamp the node visited so routing
// advances past it next tick, then schedule the next campaign tick (now for
// instant actions and "end", now+wait for a wait node). An "end" node has no
// outgoing connection, so the contact drops out of routing afterwards while
// the campaign keeps processing other contacts.
if sequence.Kind != "email" {
var cfg models.ActionConfig
if len(sequence.Action) > 0 {
_ = json.Unmarshal(sequence.Action, &cfg)
}
aerr := s.executeActionNode(ctx, campaign, contact, sequence.ID, &cfg)
if aerr != nil {
log.Warn().Err(aerr).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Str("action", cfg.Type).Msg("Action node execution failed")
}
resumeAt := nextTime
if cfg.Type == "wait" && cfg.WaitMinutes != nil && *cfg.WaitMinutes > 0 {
resumeAt = time.Now().UTC().Add(time.Duration(*cfg.WaitMinutes) * time.Minute)
}
if rerr := s.campaignProgressRepo.RecordEmailSent(ctx, campaign.ID, contact.ID, sequence.ID); rerr != nil {
log.Warn().Err(rerr).Str("campaign_id", campaign.ID.String()).Msg("Failed to record action node progress")
}
if cerr := s.createCampaignTask(ctx, campaign.ID, accountID, resumeAt); cerr != nil {
log.Warn().Err(cerr).Str("campaign_id", campaign.ID.String()).Msg("Failed to schedule next task after action node")
}
if s.campaignLogRepo != nil {
// Record the real outcome: a failed or skipped action (e.g. the linked
// automation is disabled) must be visible, not logged as if it ran.
evt, msg := "action", fmt.Sprintf("Ran '%s' action for %s", cfg.Type, contact.Email)
if aerr != nil {
evt, msg = "action_skipped", fmt.Sprintf("Action '%s' for %s did not run: %v", cfg.Type, contact.Email, aerr)
}
_ = s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: evt,
Message: msg,
})
}
s.taskRepo.UpdateTaskStatus(ctx, taskID, "completed")
executionStatus = "completed"
return nil
}
// Load this step's attachments (campaign-wide files plus the step's own;
// metadata only — the worker fetches the bytes from object storage by S3
// key at send time).
attachmentRefs := s.campaignAttachmentRefs(ctx, campaign.ID, sequence.ID)
// STEP 7.5: Update campaign task with contact_id and sequence_id for tracking
// This allows the tracking consumer to find the correct contact/sequence when
// processing open/click events from the tracking pixel service
if err := s.taskRepo.UpdateCampaignTaskTracking(ctx, taskID, contact.ID, sequence.ID); err != nil {
// Log but don't fail - tracking can still work via fallback methods
log.Warn().Err(err).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Msg("Failed to update campaign task tracking")
}
// stop_on_reply is enforced inside FindRoutedPairs (STEP 6), and it is now
// ROUTE-AWARE: a contact who replied is only handed back when their next step
// is part of the reply flow (the reply branch's own path). The normal cold
// sequence stops there, so there is no longer a blanket "contact has replied,
// skip" check here — that would also kill the reply branch's follow-up emails.
// STEP 9: Load email account
account, xerr := s.emailRepo.GetByID(ctx, accountID)
if xerr != nil {
return xerr
}
// STEP 9.25: Stamp the mailbox this send actually goes out from. The chain
// seeds a task with the previous tick's pick, and rotation chooses again
// here, so tasks.email_account_id (today's budget, the min-gap clock,
// rotation's own position, bounce and complaint rates, the contact's
// activity feed) has to be corrected before the send (issue #392). Failing
// it fails the task like the reservation: a send charged to another mailbox
// would let this one past its daily cap.
if account.ID != taskRecord.EmailAccountID {
if err := s.taskRepo.UpdateTaskEmailAccount(ctx, taskID, account.ID); err != nil {
errs.CaptureException(err)
s.taskRepo.RecordTaskFailure(ctx, taskID, "Could not record the sending mailbox", err.Error())
s.recordSchedulerFailure(ctx, campaign.ID, "sender_attribution_failed",
fmt.Sprintf("Could not record which mailbox is sending to %s; retrying", contact.Email), err)
if uerr := s.taskRepo.UpdateTaskStatus(ctx, taskID, "pending"); uerr != nil {
errs.CaptureException(uerr)
}
executionStatus = "failed"
return errx.InternalError()
}
taskRecord.EmailAccountID = account.ID
}
// STEP 9.4: The conversation this step joins. A follow-up is a nudge on the
// email the contact already has, not a second cold email, so every step
// after their first is threaded onto the last one they received: the
// parent's Message-ID becomes In-Reply-To/References, which is what the
// RECIPIENT's client threads on, and the parent's provider thread handle
// files it in the same conversation in the SENDER's mailbox (issue #472).
// A contact's first email has no parent and opens the thread.
threadParent := s.threadParent(ctx, campaign.ID, contact.ID, sequence)
// STEP 9.5: Resolve {{form_link:...}} markers to per-recipient form URLs
// BEFORE templating, so the substituted literal survives the naive
// fallback and gets wrapped by click tracking in STEP 11 like any link.
rawSubject, rawBodyHTML, rawBodyPlain := sequence.Subject, sequence.BodyHTML, sequence.BodyPlain
// A reply carries the conversation's subject, so a threading step does not
// have one of its own: it inherits it here, before rendering, so the merge
// fields resolve for THIS contact. Empty means the step writes its own.
threadSubject := s.threadSubject(ctx, campaign.ID, sequence, threadParent)
if threadSubject != "" {
rawSubject = threadSubject
}
s.resolveFormLinks(ctx, orgID, campaign, contact, &rawSubject, &rawBodyHTML, &rawBodyPlain)
// STEP 9.75: The recipient's opt-out. The link (when the instance can
// mint one) backs the List-Unsubscribe header, the link-mode footer and
// any {{.UnsubscribeLink}} the step places by hand; the footer mode comes
// from Settings > Sending unless the campaign overrides it.
optOut := s.resolveOptOut(ctx, orgID, campaign)
var unsubscribeURL string
if s.unsubLinks != nil && s.unsubLinks.Enabled() {
// On the workspace's own verified tracking domain when it has one, so
// the opt-out address sits on the sender's domain like every other link
// in the email rather than naming the platform.
unsubscribeURL = s.mintUnsubscribeLink(ctx, resolveOptOutOrigin(account, campaign), orgID, campaign.ID, contact.ID)
}
extra := map[string]string{UnsubscribeLinkVar: unsubscribeURL}
// STEP 10: Render email template with contact variables, then expand any
// {a|b|c} spintax per-recipient (only real |-groups; literal braces/CSS are
// left intact) so each send varies for deliverability.
subject := expandSpintax(RenderTemplateWith(rawSubject, *contact, extra))
bodyHTML := expandSpintax(RenderTemplateWith(rawBodyHTML, *contact, extra))
bodyPlain := expandSpintax(RenderTemplateWith(rawBodyPlain, *contact, extra))
// If no plain text provided, extract from HTML
if bodyPlain == "" && bodyHTML != "" {
bodyPlain = ExtractPlainTextFromHTML(bodyHTML)
}
if s.advanced != nil {
selection, sxerr := s.advanced.SelectVariant(ctx, orgID, campaign.ID, contact.ID, sequence.ID, subject, bodyHTML, bodyPlain)
if sxerr != nil {
return sxerr
}
if selection != nil {
// A chosen variant is stored template text, so it goes through the
// same render as the step's own copy; the control arm comes back
// already rendered, for which this pass is a no-op.
//
// A threading step is the exception: its subject belongs to the
// conversation, not to the arm, so variants on a follow-up vary
// the body only.
if threadSubject == "" {
subject = expandSpintax(RenderTemplateWith(selection.Subject, *contact, extra))
}
bodyHTML = expandSpintax(RenderTemplateWith(selection.BodyHTML, *contact, extra))
bodyPlain = expandSpintax(RenderTemplateWith(selection.BodyPlain, *contact, extra))
// A variant may carry HTML only; keep the plain-text alternative.
if bodyPlain == "" && bodyHTML != "" {
bodyPlain = ExtractPlainTextFromHTML(bodyHTML)
}
}
}
// STEP 10.5: Resolve per-recipient AI variables (AI blocks in the body that
// generate unique copy for THIS contact). Runs before tracking/signature so
// the generated copy gets open/click tracking like any other body text. A
// zero-cost no-op when the body has no AI blocks. A generation failure fails
// the send (recorded like other send failures) so the task retries with the
// same cached output.
if s.aiProvider != nil && s.aiCredits != nil {
var aerr error
subject, bodyHTML, bodyPlain, aerr = s.resolveAIVariables(ctx, campaign, contact, sequence.ID, subject, bodyHTML, bodyPlain)
if aerr != nil {
s.taskRepo.RecordTaskFailure(ctx, taskID, "AI variable resolution failed", aerr.Error())
if s.campaignLogRepo != nil {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "email_failed",
Message: fmt.Sprintf("Failed to resolve AI variables for %s", contact.Email),
Metadata: map[string]interface{}{
"contact_id": contact.ID.String(),
"sequence_id": sequence.ID.String(),
"error": aerr.Error(),
},
})
}
if rerr := s.taskRepo.UpdateTaskStatus(ctx, taskID, "pending"); rerr != nil {
errs.CaptureException(rerr)
}
executionStatus = "failed"
return errx.InternalError()
}
}
// STEP 10.6: A plain-text campaign ships no HTML part at all. Tracking
// below only rewrites HTML, so dropping it here is what makes the
// setting's "disables tracking" promise true.
if campaign.TextOnly {
if bodyPlain == "" && bodyHTML != "" {
bodyPlain = ExtractPlainTextFromHTML(bodyHTML)
}
bodyHTML = ""
}
// A blank HTML alternative never ships as the part the client prefers.
// Shared with the preview and the test send (see dropBlankHTMLPart).
bodyHTML = dropBlankHTMLPart(bodyHTML, bodyPlain)
// STEP 10.7: A hand-placed {{.UnsubscribeLink}} resolved to the bare signed
// URL; give it an anchor so the recipient reads "Unsubscribe" and not the
// API address (issue #341). After the plain part was derived, so plain text
// keeps the URL it needs, and before tracking, which leaves it alone.
bodyHTML = linkifyUnsubscribeURL(bodyHTML, unsubscribeURL, optOut.LinkText)
// STEP 10.75: Score the copy the recipient will actually receive, after
// merge fields, spintax, A/B and AI blocks have resolved. Advisory: it
// warns once per step and never blocks or delays the send.
s.warnOnWeakContent(ctx, orgID, campaign.ID, sequence.ID, sequence.Position+1,
subject, bodyHTML, bodyPlain, len(attachmentRefs))
// STEP 11: Add tracking on the host resolveTrackingHost picks, and log any
// override that was configured but not verified so a customer whose links
// are not going through their own domain can see why.
trackingDomain, ignored := resolveTrackingHost(config.TrackingHost(), account, campaign)
if s.campaignLogRepo != nil {
for _, ign := range ignored {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "tracking_domain_unverified",
Message: ign.Message,
Metadata: map[string]interface{}{"scope": ign.Scope, "tracking_domain": ign.Domain, "mailbox": account.Email},
})
}
}
if campaign.OpenTracking && bodyHTML != "" {
bodyHTML = AddOpenTrackingPixel(bodyHTML, taskID, trackingDomain)
}
if bodyHTML != "" && (campaign.LinkTracking || campaign.UTMTracking) {
linkOpts := LinkTracking{
TaskID: taskID,
CampaignID: campaign.ID,
TrackingDomain: trackingDomain,
Wrap: campaign.LinkTracking,
UTM: CampaignUTM(campaign),
}
tracked, links := TrackLinks(bodyHTML, linkOpts)
if len(links) == 0 {
bodyHTML = tracked
} else if err := s.trackedLinkRepo.CreateBatch(ctx, links); err != nil {
// Tracking is a nicety: ship the original working links rather
// than tickets that would 404 at the tracking service. UTM tags
// need no ticket, so they still go on.
log.Warn().Err(err).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Msg("Failed to store tracked links; sending untracked")
bodyHTML, _ = TrackLinks(bodyHTML, LinkTracking{TrackingDomain: linkOpts.TrackingDomain, UTM: linkOpts.UTM})
} else {
bodyHTML = tracked
}
}
if campaign.UTMTracking && bodyPlain != "" {
bodyPlain = TagPlainTextLinks(bodyPlain, CampaignUTM(campaign), trackingDomain)
}
// STEP 12: Add signature
if account.SignatureSync {
if bodyHTML != "" {
bodyHTML = AddSignature(bodyHTML, account.SignatureHTML, true)
}
if bodyPlain != "" {
bodyPlain = AddSignature(bodyPlain, account.SignaturePlain, false)
}
}
// STEP 12.5: Opt-out footer, after the signature and after click tracking
// so the link is never rewritten into a tracked ticket.
bodyHTML, bodyPlain = appendOptOut(bodyHTML, bodyPlain, optOut, unsubscribeURL)
// STEP 12.75: Copy any embedded stylesheet onto the elements it matches.
// Outlook.com, Yahoo and Gmail's mobile clients drop <style> blocks, so a
// design written with classes arrives unstyled without this. Last, so the
// signature and footer are inlined with the body; a no-op for a body with
// no <style>, which is every campaign written in the visual editor.
bodyHTML = mailhtml.InlineCSS(bodyHTML)
// STEP 13: Warm the organization DEK so the publisher's encrypt pass (the
// one whose ciphertext is actually sent) fails fast here if KMS is down.
if account.OrganizationID == nil {
errs.CaptureException(fmt.Errorf("email account %s has no organization", account.ID))
return errx.InternalError()
}
if _, err := s.cipherService.Cipher(ctx, *account.OrganizationID); err != nil {
errs.CaptureException(err)
return errx.InternalError()
}
// STEP 14: Generate Message-ID, on the domain the message is actually
// From, which is the send-as alias when the mailbox has one.
messageID := generateMessageID(account.SendFrom())
// STEP 15: Build tracking info (worker receives the already-resolved host).
// A plain-text send carries none: there is no HTML for a pixel or a
// wrapped link to live in.
var tracking *models.TrackingInfo
if !campaign.TextOnly && (campaign.OpenTracking || campaign.LinkTracking) {
tracking = &models.TrackingInfo{
OpenTracking: campaign.OpenTracking,
LinkTracking: campaign.LinkTracking,
TrackingDomain: trackingDomain,
}
}
// STEP 15.5: The List-Unsubscribe header carries the same signed link.
// Off when the campaign disabled it, or when no link could be minted: a
// header pointing nowhere is worse than none.
headerURL := ""
if campaign.UnsubscribeHeader {
headerURL = unsubscribeURL
}
// STEP 15.9: Reserve the send BEFORE it goes on the bus. Once the command is
// published the recipient's copy is committed, so the record of the attempt
// has to exist first: without it, a crash or a failed progress write between
// the dispatch and the sent_at stamp leaves the step looking "never sent"
// and the next tick emails the same person again (issue #169). The
// reservation also counts the send against the day's counters, so a lost
// stamp can never let the daily cap over-send either.
//
// It also binds the lead to this mailbox, in the same transaction, so every
// remaining step of this contact's sequence leaves from the address they
// are about to hear from (issue #401).
reserved, rerr := s.campaignProgressRepo.ReserveSend(ctx, campaign.ID, contact.ID, sequence.ID, taskID, account.ID, nextPair.IsNewLead)
if rerr != nil {
// The attempt could not be made durable, so it must not be made at all.
// Retry the whole task rather than sending something nothing remembers.
errs.CaptureException(rerr)
s.taskRepo.RecordTaskFailure(ctx, taskID, "Could not reserve the send", rerr.Error())
s.recordSchedulerFailure(ctx, campaign.ID, "send_reservation_failed",
fmt.Sprintf("Could not record the send to %s before dispatching it; retrying", contact.Email), rerr)
if uerr := s.taskRepo.UpdateTaskStatus(ctx, taskID, "pending"); uerr != nil {
errs.CaptureException(uerr)
}
executionStatus = "failed"
return errx.InternalError()
}
if !reserved {
// The claim was refused: another tick already has this (contact, step)
// in flight or delivered, or the lead was paused in the moment between
// the gate above and this transaction. End this one instead of sending,
// and keep the chain alive for whoever is next. Which of the two it was
// is worth recording, so a pause that landed on a send does not read as
// a duplicate.
outcome, why := "skipped_duplicate", "campaign send skipped: the step is already in flight or sent"
if raced, rherr := s.campaignProgressRepo.GetLeadHold(ctx, campaign.ID, contact.ID); rherr == nil && raced != nil {
outcome, why = "skipped_paused", "campaign send skipped: the lead was paused as the send was reserved"
}
log.Warn().Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).
Str("contact_id", contact.ID.String()).Str("sequence_id", sequence.ID.String()).
Msg(why)
_ = s.taskRepo.UpdateTaskStatusWithLock(ctx, taskID, outcome)
_ = s.createCampaignTask(ctx, campaign.ID, accountID, nextTime)
executionStatus = "completed"
return nil
}
// STEP 16: Send email to worker via Kafka
emailMsg := EmailMessage{
From: account.Email,
To: []string{contact.Email},
CC: campaign.CC,
BCC: campaign.BCC,
Subject: subject,
BodyHTML: bodyHTML,
BodyPlain: bodyPlain,
MessageID: messageID,
IsWarmup: false,
Tracking: tracking,
UnsubscribeURL: headerURL,
Attachments: attachmentRefs,
}
if threadParent != nil {
emailMsg.InReplyTo = threadParent.MessageID
// The provider handle needs two things the headers do not. It is
// meaningless outside the mailbox that owns it, and Gmail will not
// file a message in a thread whose subject it does not match, so it
// only goes on a message actually carrying the conversation's
// subject. Offering one Gmail would refuse costs a failed send; going
// without it costs the thread in the sender's own mailbox, and the
// recipient still sees a reply.
if threadParent.SenderID == account.ID && threadParent.Subject != "" && threadSubject == threadParent.Subject {
emailMsg.ThreadID = threadParent.ThreadID
}
}
if err := s.emailSender.Send(ctx, taskID, emailMsg, *account); err != nil {
// The send never reached a worker (none assigned, worker offline, bus
// or storage down). Nothing is stamped sent; the task is dead-lettered
// for the retry loop and the chain is re-seeded by the reconciler.
//
// Give the reservation back so the next tick retries the step — but ONLY
// when the command provably never left. A failure of the publish call
// itself is ambiguous (the bus may have taken it), so that reservation
// stands and is resolved by the worker's own result, or by the reclaimer
// if none ever comes.
if !errors.Is(err, ErrSendDispatchUnknown) {
if relErr := s.campaignProgressRepo.ReleaseSend(ctx, campaign.ID, contact.ID, sequence.ID, nextPair.IsNewLead); relErr != nil {
errs.CaptureException(relErr)
log.Error().Err(relErr).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Msg("Failed to release the reservation for a send that never left; the reclaimer will retry it")
}
}
s.taskRepo.RecordTaskFailure(ctx, taskID, "Send failed", err.Error())
if s.campaignLogRepo != nil {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "email_failed",
Message: fmt.Sprintf("Could not hand %s's email to a sending worker, will retry: %s", contact.Email, err.Error()),
Metadata: map[string]interface{}{
"level": "error",
"code": "WORKER_UNAVAILABLE",
"contact_id": contact.ID.String(),
"sequence_id": sequence.ID.String(),
"account_id": account.ID.String(),
"error": err.Error(),
"will_retry": true,
},
})
}
// Publish task failure to Pub/Sub
if s.streamingPublisher != nil {
s.streamingPublisher.PublishTaskStatus(ctx, campaign.UserID, taskID, pubsub.EventTaskFailed, "Failed to send email", map[string]string{
"campaign_id": campaign.ID.String(),
"contact_id": contact.ID.String(),
"error": err.Error(),
})
// Publish detailed task progress event for failure
progress := 0
if totalContacts > 0 {
progress = (processedCount * 100) / totalContacts
}
contactName := contact.FirstName
if contact.LastName != "" {
contactName = contactName + " " + contact.LastName
}
s.streamingPublisher.PublishTaskProgress(ctx, &pubsub.TaskProgressEvent{
BaseEvent: pubsub.BaseEvent{UserID: campaign.UserID},
OrgID: campaignOrgID(campaign),
CampaignID: campaign.ID.String(),
TaskID: taskID.String(),
Status: "failed",
ContactID: contact.ID.String(),
ContactEmail: contact.Email,
ContactName: contactName,
SequenceID: sequence.ID.String(),
SequenceName: sequence.Name,
Progress: progress,
TotalContacts: totalContacts,
ProcessedCount: processedCount,
})
}
if s.advanced != nil {
_ = s.advanced.CaptureTaskDeadLetter(ctx, taskID, "campaign", map[string]interface{}{
"campaign_id": campaign.ID.String(),
"contact_id": contact.ID.String(),
"email": contact.Email,
}, err.Error(), 1)
_ = s.taskRepo.UpdateTaskStatus(ctx, taskID, "dead_lettered")
}
return nil
}
// STEP 16: Store sent email metadata (encrypted) in database
// Note: Full email stored in Cassandra by email sync service
taskRecord.MessageID = messageID
taskRecord.Status = "completed"
// STEP 17: Stamp the step sent. The reservation already made the attempt
// durable, so this is the timing stamp follow-up pacing reads, not the
// duplicate guard — but losing it still parks the lead, so retry it, and
// escalate rather than whispering when every retry fails. The worker's own
// EMAIL_SENT repairs it downstream if it never lands.
// Today's counters were bumped by the reservation, inside the same
// transaction, so the new-lead/day cap can never under-count.
if err := s.stampSendRecorded(ctx, campaign.ID, contact.ID, sequence.ID); err != nil {
errs.CaptureException(err)
log.Error().Err(err).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Str("contact_id", contact.ID.String()).Msg("Failed to record email sent; the worker result will repair the stamp")
if s.campaignLogRepo != nil {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "progress_write_failed",
Message: fmt.Sprintf("The email to %s was sent but recording it failed; it will not be sent again", contact.Email),
Metadata: map[string]interface{}{
"level": "error",
"code": "PROGRESS_WRITE_FAILED",
"contact_id": contact.ID.String(),
"sequence_id": sequence.ID.String(),
"error": err.Error(),
},
})
}
}
// Publish campaign progress summary to Pub/Sub for real-time dashboard updates
if s.streamingPublisher != nil {
if progress, pErr := s.campaignProgressRepo.GetCampaignProgress(ctx, campaign.ID); pErr == nil && progress != nil {
s.streamingPublisher.PublishCampaignProgress(ctx, campaign.UserID, campaign.ID, progress)
}
}
// Log email sent
if s.campaignLogRepo != nil {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "email_sent",
Message: fmt.Sprintf("Email sent to %s", contact.Email),
Metadata: map[string]interface{}{
"contact_id": contact.ID.String(),
"sequence_id": sequence.ID.String(),
"account_id": account.ID.String(),
},
})
}
// This contact's sequence changed address. It only happens when the mailbox
// they had been hearing from stopped being one this campaign can send from
// (disconnected, taken off its sending accounts, resting, or held by warmup
// health), and it is the kind of thing an owner should find in the activity
// log rather than in a confused reply.
if prev := nextPair.AssignedSender; prev != nil && *prev != account.ID && s.campaignLogRepo != nil {
from := "its previous mailbox"
if old, oerr := s.emailRepo.GetByID(ctx, *prev); oerr == nil && old != nil {
from = old.Email
}
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: "sender_reassigned",
Message: fmt.Sprintf("%s now hears from %s: %s can no longer send for this campaign",
contact.Email, account.Email, from),
Metadata: map[string]interface{}{
"contact_id": contact.ID.String(),
"account_id": account.ID.String(),
"prev_account_id": prev.String(),
},
})
}
// STEP 18: Mark task as completed (with advisory lock)
if err := s.taskRepo.UpdateTaskStatusWithLock(ctx, taskID, "completed"); err != nil {
errs.CaptureException(err)
return errx.InternalError()
}
// STEP 18.5: Advance the explicit-sender rotation cursor on a GENUINE send
// only (single atomic UPDATE), so round_robin/least_recently_used cursors
// stay coherent and a send-failure/skip never bumps them. The UPDATE is
// scoped to (campaign_id, email_account_id), so it's a harmless no-op for
// tag/all-resolved mailboxes that have no campaign_senders row.
if err := s.campaignRepo.AdvanceCampaignSender(ctx, campaign.ID, account.ID); err != nil {
log.Warn().Err(err).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Msg("Failed to advance campaign sender cursor")
}
// Publish task completion to Pub/Sub
if s.streamingPublisher != nil {
s.streamingPublisher.PublishTaskStatus(ctx, campaign.UserID, taskID, pubsub.EventTaskCompleted, "Email sent successfully", map[string]string{
"campaign_id": campaign.ID.String(),
"contact_id": contact.ID.String(),
})
// Publish detailed task progress event
newProcessedCount := processedCount + 1
progress := 0
if totalContacts > 0 {
progress = (newProcessedCount * 100) / totalContacts
}
contactName := contact.FirstName
if contact.LastName != "" {
contactName = contactName + " " + contact.LastName
}
// Get sequence index
sequences, _ := s.campaignRepo.GetSequencesByCampaignID(ctx, campaign.ID)
seqIndex := 0
for i, seq := range sequences {
if seq.ID == sequence.ID {
seqIndex = i + 1
break
}
}
// EMAIL_SENT (org-scoped): the whole team sees the send + which
// lead/step fired, live in the campaign view.
s.streamingPublisher.PublishEmailSent(ctx, &pubsub.TaskProgressEvent{
BaseEvent: pubsub.BaseEvent{UserID: campaign.UserID},
OrgID: campaignOrgID(campaign),
CampaignID: campaign.ID.String(),
TaskID: taskID.String(),
Status: "completed",
ContactID: contact.ID.String(),
ContactEmail: contact.Email,
ContactName: contactName,
SequenceID: sequence.ID.String(),
SequenceName: sequence.Name,
SequenceIndex: seqIndex,
Progress: progress,
TotalContacts: totalContacts,
ProcessedCount: newProcessedCount,
})
}
// STEP 19: Publish events to Kafka
s.publishEmailSentEvent(ctx, taskRecord, account, campaign, contact, sequence)
// STEP 20: Create next campaign task. The successor serves whichever lead
// is due next, so it must never be shaped for the contact just emailed:
// send-time optimization used to push it to that contact's next preferred
// hour (tomorrow 09:00 by default after 17:00 UTC), leaving every other
// lead queued for a day. Recipient-time placement belongs in slot
// selection for the selected contact (issue #156), not here.
if err := s.createCampaignTask(ctx, campaign.ID, account.ID, nextTime); err != nil {
// Log but don't fail the current task
log.Warn().Err(err).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Msg("Failed to create next campaign task")
}
executionStatus = "completed"
return nil
}
// threadParent resolves the email this step should be sent as a reply to, or
// nil when it must open a new conversation: the step has reply-in-thread
// turned off, the contact has had nothing from this campaign yet, or the
// previous send left no Message-ID to reference.
//
// A lookup failure is never fatal. Losing the thread costs the recipient a
// tidy conversation; refusing the send costs them the email, so a database
// error here degrades to a new thread and is logged.
func (s *tasksService) threadParent(ctx context.Context, campaignID, contactID uuid.UUID, sequence *Sequence) *repository.ThreadParent {
if sequence == nil || !sequence.ThreadReply {
return nil
}
parent, err := s.campaignProgressRepo.ThreadParentForLead(ctx, campaignID, contactID)
if err != nil {
log.Warn().Err(err).Str("campaign_id", campaignID.String()).Str("contact_id", contactID.String()).
Msg("Could not resolve the thread to reply on; sending as a new conversation")
return nil
}
if parent == nil || parent.MessageID == "" {
return nil
}
return parent
}
// threadSubject is the subject a step inherits from the conversation it is
// replying on, or "" when it writes its own (the switch is off, there is no
// earlier email, or the conversation has no subject yet).
//
// The parent answers it whenever there is one, because that is read off what
// the contact was actually sent. The campaign's own step order is the fallback
// for when there is not, and it only has to run for a step with no subject of
// its own: a previous send the worker never confirmed leaves no parent, and a
// threading step authored in the composer has nothing to fall back on, so
// without this it would ship a blank Subject header. A step that does have a
// subject already has something to send, which keeps this off the first-touch
// path, where there is never a parent and never anything to inherit.
func (s *tasksService) threadSubject(ctx context.Context, campaignID uuid.UUID, sequence *Sequence, parent *repository.ThreadParent) string {
if sequence == nil || !sequence.ThreadReply {
return ""
}
if parent != nil && parent.Subject != "" {
return parent.Subject
}
if strings.TrimSpace(sequence.Subject) != "" {
return ""
}
seqs, err := s.campaignRepo.GetSequencesByCampaignID(ctx, campaignID)
if err != nil {
log.Warn().Err(err).Str("campaign_id", campaignID.String()).
Msg("Could not read the campaign's steps for the conversation subject")
return ""
}
for i := range seqs {
if seqs[i].ID != sequence.ID {
continue
}
// StepSubject returns the step's own when there is nothing to inherit,
// which is not an inherited subject and must not read as one.
if sub := models.StepSubject(seqs, i); sub != sequence.Subject {
return sub
}
}
return ""
}
// autoPauseCampaign pauses a campaign when no active email accounts are available.
// Uses advisory lock to prevent concurrent auto-pause from multiple tasks.
// autoPauseReason turns a no-mailbox scheduling error into the sentence the
// owner reads in the activity log. The narrower sentinels wrap
// ErrNoEmailAccounts, so they must be tested first; the difference between them
// is the difference between "fix your DNS", "widen your sending window", and
// "connect a mailbox".
func autoPauseReason(err error) string {
switch {
case errors.Is(err, scheduler.ErrDomainAuthFailing):
return "Campaign auto-paused: every mailbox is sending from a domain that fails SPF or DMARC authentication"
case errors.Is(err, scheduler.ErrNoEligibleMailbox):
return "Campaign auto-paused: no mailbox can send under its current sending settings (check each mailbox's sending behaviour profile and timezone)"
default:
return "Campaign auto-paused: no active email accounts available"
}
}
// UndeliverablePauseReason is the activity-log line for a campaign parked
// because verification refused every remaining lead.
func UndeliverablePauseReason(n int) string {
return fmt.Sprintf("Campaign paused: %d remaining lead(s) were refused by address verification. Re-verify them or mark them deliverable to continue", n)
}
// pauseUndeliverable parks a campaign whose only remaining leads were refused
// by verification. Resumable: re-verifying or marking them deliverable
// restarts it.
func (s *tasksService) pauseUndeliverable(ctx context.Context, campaignID, taskID uuid.UUID, n int) {
s.campaignRepo.UpdateStatusWithLock(ctx, campaignID, "paused_undeliverable")
if taskID != uuid.Nil {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "completed")
}
if s.campaignLogRepo != nil {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaignID,
EventType: "auto_paused",
Message: UndeliverablePauseReason(n),
Metadata: map[string]interface{}{"undeliverable": n},
})
}
if s.streamingPublisher != nil {
s.streamingPublisher.PublishCampaignEvent(ctx, &pubsub.CampaignEvent{
BaseEvent: pubsub.BaseEvent{EventType: pubsub.EventCampaignCompleted},
CampaignID: campaignID.String(),
Status: "paused_undeliverable",
})
}
}
// CampaignIdleEventType is the activity log entry written when a continuous
// campaign runs out of leads and waits; CampaignIdleMessage is its text.
const (
CampaignIdleEventType = "idle"
CampaignIdleMessage = "Waiting for new leads: every lead has finished the sequence. The campaign stays active and sends to leads as they arrive."
)
// HeldLeadsMessage is the activity-log line for a campaign whose only remaining
// leads are paused.
func HeldLeadsMessage(n int) string {
lead := "lead is"
if n != 1 {
lead = "leads are"
}
return fmt.Sprintf("Waiting: %d %s paused. The campaign stays active and continues when they resume.", n, lead)
}
// CampaignHeldEventType is the activity log entry for a campaign waiting on
// paused leads.
const CampaignHeldEventType = "waiting_on_paused_leads"
// parkHeldLeads leaves a campaign active when the only thing left to send to is
// a lead somebody parked. It is deliberately NOT idleCampaign: that one is for
// a continuous campaign out of leads and its MarkIdle refuses a campaign that
// is not continuous, which would leave this case silent. Nothing here changes
// the status — the campaign IS active, it is waiting — and the log line is
// written once per wait rather than once per pass.
func (s *tasksService) parkHeldLeads(ctx context.Context, campaign *models.Campaign, taskID uuid.UUID, n int) {
if taskID != uuid.Nil {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "completed")
}
if s.campaignLogRepo == nil {
return
}
wrote, err := s.campaignLogRepo.CreateLogOnce(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: CampaignHeldEventType,
Message: HeldLeadsMessage(n),
Metadata: map[string]interface{}{"paused_leads": n},
}, "paused_leads", strconv.Itoa(n), time.Now().Add(-campaignHeldLogWindow))
if err != nil {
log.Warn().Err(err).Str("campaign_id", campaign.ID.String()).Msg("could not record the paused-lead wait")
return
}
if wrote && s.streamingPublisher != nil {
s.streamingPublisher.PublishCampaignEvent(ctx, &pubsub.CampaignEvent{
BaseEvent: pubsub.BaseEvent{EventType: pubsub.EventCampaignIdle, UserID: campaign.UserID},
OrgID: campaignOrgID(campaign),
CampaignID: campaign.ID.String(),
Name: campaign.Name,
Status: campaign.Status,
})
}
}
// campaignHeldLogWindow keeps the wait from filling the activity feed: the
// reconciler re-checks every pass, and the fact does not change between them.
const campaignHeldLogWindow = 6 * time.Hour
// idleCampaign parks a continuous campaign that has nothing left to send. It
// stays active with no chain: a lead add wakes it, and the reconciler re-checks
// it every pass. Logged and broadcast once per wait, not once per pass.
func (s *tasksService) idleCampaign(ctx context.Context, campaign *models.Campaign, taskID uuid.UUID) {
if taskID != uuid.Nil {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "completed")
}
transitioned, err := s.campaignRepo.MarkIdle(ctx, campaign.ID)
if err != nil {
log.Warn().Err(err).Str("campaign_id", campaign.ID.String()).Msg("could not mark the campaign idle")
return
}
if !transitioned {
return
}
if s.campaignLogRepo != nil {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaign.ID,
EventType: CampaignIdleEventType,
Message: CampaignIdleMessage,
})
}
if s.streamingPublisher != nil {
s.streamingPublisher.PublishCampaignEvent(ctx, &pubsub.CampaignEvent{
BaseEvent: pubsub.BaseEvent{EventType: pubsub.EventCampaignIdle, UserID: campaign.UserID},
OrgID: campaignOrgID(campaign),
CampaignID: campaign.ID.String(),
Name: campaign.Name,
Status: "active",
})
}
}
// clearIdle ends an idle wait once the campaign has something to send again.
func (s *tasksService) clearIdle(ctx context.Context, campaign *models.Campaign) {
if campaign.IdleSince == nil {
return
}
if err := s.campaignRepo.ClearIdle(ctx, campaign.ID); err != nil {
log.Warn().Err(err).Str("campaign_id", campaign.ID.String()).Msg("could not clear the campaign's idle mark")
}
}
// autoPauseCampaign parks a campaign that has nothing it can send from. The
// reason is carried through to the activity log because "paused_no_accounts"
// covers several very different fixes (connect a mailbox, widen a sending
// window, repair DNS) and the status alone cannot tell them apart.
//
// It is deliberately LOUD. This runs unattended, hours after anyone touched
// the campaign, and a pause nobody is told about is a campaign that quietly
// stops sending: an error-level line in the operator's log, an error-level
// entry in the activity feed so the dashboard tints it red, and an org-scoped
// realtime pulse so every teammate's campaign list moves to "paused — no
// accounts" without a refresh.
func (s *tasksService) autoPauseCampaign(ctx context.Context, campaignID, taskID uuid.UUID, reason string) {
s.campaignRepo.UpdateStatusWithLock(ctx, campaignID, "paused_no_accounts")
if taskID != uuid.Nil {
s.taskRepo.UpdateTaskStatus(ctx, taskID, "completed")
}
log.Error().Str("campaign_id", campaignID.String()).Str("reason", reason).
Msg("campaign auto-paused: no mailbox can send for it")
// CreateLogOnce, not CreateLog: the reconciler re-checks paused campaigns
// too, and a repeat of the SAME reason must not fill the feed or re-pulse
// the dashboard every pass. Keyed on the reason rather than the code, so a
// pause whose cause changed (DNS, then no mailbox at all) still says so.
// The write also reports whether this pause is news, which is what gates
// the announcement below.
if s.campaignLogRepo == nil {
return
}
entry := &repository.CampaignLogEntry{
CampaignID: campaignID,
EventType: "auto_paused",
Message: reason,
Metadata: map[string]interface{}{
"level": "error",
"code": "no_accounts",
"reason": reason,
},
}
written, err := s.campaignLogRepo.CreateLogOnce(ctx, entry, "reason", reason, time.Now().Add(-time.Hour))
if err != nil || !written || s.streamingPublisher == nil {
return
}
campaign, gerr := s.campaignRepo.GetByID(ctx, campaignID)
if gerr != nil || campaign == nil {
return
}
s.streamingPublisher.PublishCampaignEvent(ctx, &pubsub.CampaignEvent{
BaseEvent: pubsub.BaseEvent{EventType: pubsub.EventCampaignPaused, UserID: campaign.UserID},
OrgID: campaignOrgID(campaign),
CampaignID: campaignID.String(),
Name: campaign.Name,
Status: "paused_no_accounts",
})
}
// haltOrglessCampaign stops a campaign that reached the send path with no
// organization. Suppression and the entitlement gate are both org-scoped, so
// there is no way to honour an unsubscribe, bounce or complaint for this
// campaign — it is parked rather than sent unchecked. Since migration 000092
// made campaigns.organization_id NOT NULL this should be unreachable, so it is
// also reported: reaching it means tenancy was lost somewhere else.
func (s *tasksService) haltOrglessCampaign(ctx context.Context, campaignID, taskID uuid.UUID) {
const reason = "Campaign paused: it has no workspace, so unsubscribes, bounces and complaints cannot be checked before sending. Contact support to reattach it."
errs.CaptureException(fmt.Errorf("campaign %s reached the send path with no organization", campaignID))
log.Error().Str("campaign_id", campaignID.String()).Str("task_id", taskID.String()).
Msg("campaign send blocked: no organization, suppression cannot be enforced")
if err := s.campaignRepo.UpdateStatusWithLock(ctx, campaignID, "paused"); err != nil {
errs.CaptureException(err)
}
s.taskRepo.UpdateTaskStatus(ctx, taskID, "cancelled")
if s.campaignLogRepo != nil {
s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaignID,
EventType: "auto_paused",
Message: reason,
Metadata: map[string]interface{}{"level": "error", "code": "no_organization"},
})
}
}
// executeActionNode runs the control-plane side effect for a non-email node.
// "wait" and "end" have no side effect (their behaviour is timing / routing
// only); the others reuse existing repos/services. Everything here is
// control-plane — the worker is never involved for an action node.
func (s *tasksService) executeActionNode(ctx context.Context, campaign *models.Campaign, contact *models.Contact, sequenceID uuid.UUID, cfg *models.ActionConfig) error {
switch cfg.Type {
case "wait", "end", "":
return nil
case "ai_step":
return s.execSequenceAIAgentStep(ctx, campaign, contact, sequenceID, cfg)
case "add_tag":
if cfg.CategoryID == nil || campaign.OrganizationID == nil {
return nil
}
if _, xerr := s.contactRepo.Update(ctx, campaign.UserID, contact.ID.String(), *campaign.OrganizationID, &models.UpdateContact{
AddCategories: []string{cfg.CategoryID.String()},
}); xerr != nil {
return xerr
}
return nil
case "remove_tag":
if cfg.CategoryID == nil || campaign.OrganizationID == nil {
return nil
}
if _, xerr := s.contactRepo.Update(ctx, campaign.UserID, contact.ID.String(), *campaign.OrganizationID, &models.UpdateContact{
RemoveCategories: []string{cfg.CategoryID.String()},
}); xerr != nil {
return xerr
}
return nil
case "add_to_segment", "remove_from_segment":
if cfg.SegmentID == nil || campaign.OrganizationID == nil || s.segmentRepo == nil {
return nil
}
mode := models.SegmentMemberInclude
if cfg.Type == "remove_from_segment" {
mode = models.SegmentMemberExclude
}
if _, xerr := s.segmentRepo.SetMembers(ctx, *campaign.OrganizationID, *cfg.SegmentID, []uuid.UUID{contact.ID}, mode); xerr != nil {
return xerr
}
return nil
case "label_email":
// Apply unibox labels to the contact's most recent conversation. A no-op
// when the contact has no thread yet (returns "" thread, nil error).
if len(cfg.LabelIDs) == 0 || campaign.OrganizationID == nil {
return nil
}
if _, xerr := s.advanced.LabelLatestThreadForContact(ctx, *campaign.OrganizationID, contact.Email, cfg.LabelIDs); xerr != nil {
return xerr
}
return nil
case "unsubscribe":
if xerr := s.advanced.Unsubscribe(ctx, campaign.ID, contact.ID); xerr != nil {
return xerr
}
return nil
case "create_task":
if s.advanced == nil || campaign.OrganizationID == nil {
return nil
}
owner, perr := uuid.Parse(campaign.UserID)
if perr != nil {
return nil
}
title := strings.TrimSpace(cfg.TaskTitle)
if title == "" {
name := strings.TrimSpace(contact.FirstName + " " + contact.LastName)
if name == "" {
name = contact.Email
}
title = "Follow up: " + name
}
// Per-step assignee; fall back to the campaign owner only when neither a
// user nor a team is chosen (a team-assigned task has no single owner).
assignee := cfg.TaskAssignedTo
if assignee == nil && cfg.TaskAssignedTeamID == nil {
assignee = &owner
}
// Task types are user-managed free text; pass the configured name
// through (empty = untyped).
cid := contact.ID
data := &models.CreateCRMTask{
ContactID: &cid,
Title: title,
Type: cfg.TaskType,
Priority: cfg.TaskPriority,
AssignedTo: assignee,
AssignedTeamID: cfg.TaskAssignedTeamID,
}
if cfg.TaskDueOffsetDays != nil {
due := time.Now().UTC().AddDate(0, 0, *cfg.TaskDueOffsetDays)
data.DueDate = &due
}
if _, xerr := s.advanced.CreateContactTask(ctx, *campaign.OrganizationID, owner, data); xerr != nil {
return xerr
}
return nil
case "create_deal":
if s.advanced == nil || campaign.OrganizationID == nil {
return nil
}
if cfg.DealPipelineID == nil || cfg.DealStageID == nil {
// Misconfigured node (no pipeline/stage chosen): skip rather than fail
// the whole chain.
return nil
}
owner, perr := uuid.Parse(campaign.UserID)
if perr != nil {
return nil
}
// Deal name supports the same {{first_name}}/{{company}} templating other
// campaign copy uses; fall back to a contact-derived name when blank.
name := RenderTemplate(strings.TrimSpace(cfg.DealName), *contact)
if name == "" {
cn := strings.TrimSpace(contact.FirstName + " " + contact.LastName)
if cn == "" {
cn = contact.Email
}
name = "Deal: " + cn
}
currency := strings.TrimSpace(cfg.DealCurrency)
if currency == "" {
currency = "USD"
}
cid := contact.ID
cmpID := campaign.ID
data := &models.CreateDeal{
PipelineID: *cfg.DealPipelineID,
StageID: *cfg.DealStageID,
ContactID: &cid,
Name: name,
Value: cfg.DealValue,
Currency: currency,
CampaignID: &cmpID,
AssignedTo: &owner,
}
if _, xerr := s.advanced.CreateContactDeal(ctx, *campaign.OrganizationID, owner, data); xerr != nil {
return xerr
}
return nil
case "move_deal_stage":
if s.advanced == nil || campaign.OrganizationID == nil {
return nil
}
if cfg.DealPipelineID == nil || cfg.DealStageID == nil {
return nil
}
moved, xerr := s.advanced.MoveContactDealStage(ctx, *campaign.OrganizationID, contact.ID, *cfg.DealPipelineID, *cfg.DealStageID)
if xerr != nil {
return xerr
}
if moved == nil {
// No open deal in the target pipeline: documented no-op. Log it so the
// gap is observable instead of silently doing nothing.
log.Info().
Str("campaign_id", campaign.ID.String()).
Str("contact_id", contact.ID.String()).
Str("pipeline_id", cfg.DealPipelineID.String()).
Msg("move_deal_stage no-op: contact has no open deal in pipeline")
}
return nil
case "run_automation":
if s.automationRunner == nil || campaign.OrganizationID == nil || cfg.AutomationID == nil {
return nil
}
// Seed the automation's event data with the standard contact/campaign
// keys so its action templates ({{.contact_email}} etc.) work out of the
// box; the user-supplied values (rendered per contact) add/override extras.
data := map[string]any{
"campaign_id": campaign.ID.String(),
"campaign_name": campaign.Name,
"contact_id": contact.ID.String(),
"contact_email": contact.Email,
"first_name": contact.FirstName,
"last_name": contact.LastName,
"company": contact.Company,
"phone": contact.Phone,
// Stable per-(campaign,contact,step) key. Campaign tasks are
// at-least-once, so a duplicate delivery would re-launch this step;
// downstream actions (notably a webhook.ping to Zapier/Make) can dedupe
// on this. The scheduler already routes each step once per lead, so a
// real double-run is only a narrow retry window, shared by every action.
"idempotency_key": fmt.Sprintf("campaign:%s:%s:%s", campaign.ID, contact.ID, sequenceID),
}
for _, kv := range cfg.AutomationValues {
key := strings.TrimSpace(kv.Key)
if key == "" {
continue
}
data[key] = RenderTemplate(kv.Value, *contact)
}
return s.automationRunner.RunAutomationByID(ctx, *campaign.OrganizationID, *cfg.AutomationID, data)
case "fire_event":
if s.advanced == nil || campaign.OrganizationID == nil {
return nil
}
s.advanced.FireCampaignEvent(ctx, *campaign.OrganizationID, campaign.ID.String(), cfg.EventName, cfg.EventFields, contact)
return nil
case "switch":
return s.execSequenceSwitchStep(ctx, campaign, contact, sequenceID, cfg)
default:
return nil
}
}
// createCampaignTask creates a new campaign task in GCP Cloud Tasks
func (s *tasksService) createCampaignTask(ctx context.Context, campaignID, accountID uuid.UUID, scheduleTime time.Time) error {
// Create task in database with advisory lock
newTaskID := uuid.New()
newTask := &Task{
ID: newTaskID,
TaskType: "campaign",
EmailAccountID: accountID,
Status: "pending",
ScheduledAt: &scheduleTime,
}
campaignTask := &CampaignTask{
TaskID: newTaskID,
CampaignID: &campaignID,
}
created, err := s.taskRepo.CreateTaskWithLock(ctx, newTask, campaignTask)
if err != nil {
return err
}
if !created {
return nil
}
// Create GCP Cloud Task
processTask := &proto.ProcessTask{
TaskId: newTaskID.String(),
}
cloudTaskName, err := s.tasksClient.CreateTask(ctx, processTask, scheduleTime)
if err != nil {
return err
}
// Update task with cloud task name
if err := s.taskRepo.UpdateTaskScheduledAt(ctx, newTaskID, scheduleTime, cloudTaskName); err != nil {
return err
}
return nil
}
// publishEmailSentEvent publishes email sent event to Kafka
func (s *tasksService) publishEmailSentEvent(
ctx context.Context,
task *Task,
account *Email,
campaign *Campaign,
contact *Contact,
sequence *Sequence,
) {
if s.eventsPublisher == nil {
return
}
if err := s.eventsPublisher.PublishEmailSent(ctx, task, account, campaign, contact, sequence); err != nil {
log.Warn().Err(err).Str("campaign_id", campaign.ID.String()).Str("task_id", task.ID.String()).Msg("Failed to publish email sent event")
}
// Fan an opt-in firehose webhook for the send (campaign.email_sent).
if s.advanced != nil && campaign.OrganizationID != nil {
data := map[string]any{
"campaign_id": campaign.ID.String(),
"contact_id": contact.ID.String(),
"sequence_id": sequence.ID.String(),
}
if contact.Email != "" {
data["contact_email"] = contact.Email
}
if account != nil {
data["from_email"] = account.Email
}
s.advanced.EmitCampaignEvent(ctx, *campaign.OrganizationID, models.WebhookEventCampaignEmailSent, data)
}
}
// campaignOrgID returns the campaign's organization id for org-scoped
// realtime events, or "" for legacy orgless rows.
func campaignOrgID(campaign *Campaign) string {
if campaign == nil || campaign.OrganizationID == nil {
return ""
}
return campaign.OrganizationID.String()
}
// stampSendRecorded writes the sent_at stamp for a send already on the bus,
// retrying a transient database error instead of tolerating it. The email
// cannot be un-sent at this point, so a single failed write used to be enough
// to make routing offer the same step again (issue #169); the reservation now
// prevents that, and these retries keep the lead's pacing correct too.
func (s *tasksService) stampSendRecorded(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID) error {
var err error
for attempt := 1; attempt <= config.CampaignSendStampAttempts; attempt++ {
if err = s.campaignProgressRepo.RecordEmailSent(ctx, campaignID, contactID, sequenceID); err == nil {
return nil
}
if attempt < config.CampaignSendStampAttempts {
select {
case <-ctx.Done():
return err
case <-time.After(time.Duration(attempt) * 250 * time.Millisecond):
}
}
}
return err
}
// recordSchedulerFailure writes a campaign-scoped, reviewable failure to the
// activity log so a stalled or retrying step is VISIBLE in the dashboard
// instead of failing silently. metadata.level="error" tints it red in the
// campaign detail (TaskPreview) and the log feed is already realtime-invalidated
// for every teammate. Best-effort and nil-safe — recording a failure must never
// itself break the tick.
func (s *tasksService) recordSchedulerFailure(ctx context.Context, campaignID uuid.UUID, code, message string, cause error) {
if s.campaignLogRepo == nil {
return
}
meta := map[string]interface{}{
"level": "error",
"code": code,
}
if cause != nil {
meta["error"] = cause.Error()
}
_ = s.campaignLogRepo.CreateLog(ctx, &repository.CampaignLogEntry{
CampaignID: campaignID,
EventType: "scheduler_failed",
Message: message,
Metadata: meta,
})
}
// resolveOptOut is the effective in-body opt-out for a campaign: the
// workspace setting with the campaign's own mode applied. A settings read
// failure falls back to the defaults rather than sending without an opt-out.
func (s *tasksService) resolveOptOut(ctx context.Context, orgID uuid.UUID, campaign *models.Campaign) models.UnsubscribeSettings {
base := models.DefaultAdvancedOutreachSettings().Unsubscribe
if s.advanced != nil {
if settings, xerr := s.advanced.GetOrganizationSettings(ctx, orgID); xerr == nil && settings != nil {
base = settings.Unsubscribe
}
}
return base.Effective(campaign.UnsubscribeMode)
}