From d46cfad5978099c6c753157d9d3e6e2a69908e48 Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Fri, 18 Sep 2026 22:39:52 -0700 Subject: [PATCH 1/2] feat: check the warmup daily target again at the moment a send executes rather than only when the next one is placed, so a spam placement, health band or partner loss that cuts the target while a send is pending holds it as skipped_daily_limit and parks the chain at the next opening with a reply-back's aim intact, fail closed when that count cannot be read, share one target resolver between the placer and the send-time gate, add the skipped_org_suspended task status the suspended-workspace hold has written since #233 without a migration so its write stops failing and leaving the task pending for the dispatcher to re-fire, and anchor the ramp live fixture in UTC off the day boundary so its assertions no longer depend on the host timezone --- docs/content/docs/guides/warmup.mdx | 4 +- .../app/analytics/warmup_ramp_live_test.go | 7 +- internal/app/consumer/warmup_reply_back.go | 4 +- .../000181_task_status_org_suspended.down.sql | 2 + .../000181_task_status_org_suspended.up.sql | 4 + internal/scheduler/service.go | 4 + internal/scheduler/warmup_scheduler.go | 214 ++++++++---- internal/scheduler/warmup_scheduler_test.go | 20 ++ internal/tasks/email_task.go | 81 ++++- internal/tasks/warmup_cut_cap_live_test.go | 329 ++++++++++++++++++ 10 files changed, 592 insertions(+), 77 deletions(-) create mode 100644 internal/infrastructure/db/migrations/000181_task_status_org_suspended.down.sql create mode 100644 internal/infrastructure/db/migrations/000181_task_status_org_suspended.up.sql create mode 100644 internal/tasks/warmup_cut_cap_live_test.go diff --git a/docs/content/docs/guides/warmup.mdx b/docs/content/docs/guides/warmup.mdx index 5f0b9be41..44e2ca422 100644 --- a/docs/content/docs/guides/warmup.mdx +++ b/docs/content/docs/guides/warmup.mdx @@ -80,6 +80,8 @@ At defaults a mailbox sends 10 on day one, 11 the next, and levels off at 40 aft Four things shape the real daily number: sends are spread across your warmup hours with jitter, the target is capped by how many eligible partners exist, a mailbox whose health drops gets reduced volume and wider spacing until it recovers, and a recent spam placement holds the ramp where it is. +The target is checked twice: when the next send is placed, and again as it goes out. A signal that lowers the day's number while a send is already waiting (a spam placement, a health band, partners leaving the pool) holds that send rather than letting it go out over the new number, and the mailbox picks up again at its next opening. A cut can still land below what the mailbox has already sent that day, and the drawer shows that honestly; what it cannot do is add to the excess. + ### Holding the ramp on an early signal If any warmup email lands in a recipient's spam folder, that mailbox stops climbing immediately: @@ -104,7 +106,7 @@ Real inboxes reply, so a configurable share of the time a mailbox answers an exi Replies thread properly with a real `Re:` subject and `In-Reply-To` header, and candidates must be between 45 minutes and seven days old, so nothing is answered instantly or revived indefinitely. -Receiving warmup mail can also prompt an answer directly. When a verified warmup email arrives, the recipient sometimes points its next scheduled send back at whoever wrote, 25 minutes to 5 hours later and inside its own warmup hours. That is a re-pointing, not extra work: each mailbox has one warmup send queued at a time, so a reply-back moves that send earlier and aims it, and can never push a send the mailbox had already planned sooner. The chance is the recipient's own reply rate, drawn once when the reply is scheduled rather than again when it sends, and it stops before a thread reaches its message cap so replies cannot answer replies indefinitely. +Receiving warmup mail can also prompt an answer directly. When a verified warmup email arrives, the recipient sometimes points its next scheduled send back at whoever wrote, 25 minutes to 5 hours later and inside its own warmup hours. That is a re-pointing, not extra work: each mailbox has one warmup send queued at a time, so a reply-back moves that send earlier and aims it, and can never push a send the mailbox had already planned sooner. If the send it moved was parked for tomorrow because today's target is spent, the day's cap still holds: the answer waits for the next opening and keeps its aim. The chance is the recipient's own reply rate, drawn once when the reply is scheduled rather than again when it sends, and it stops before a thread reaches its message cap so replies cannot answer replies indefinitely. Timing imitates people throughout: sends come in bursts and lulls rather than a fixed rhythm, never land on round clock marks, and opens happen on a natural delay during the recipient's waking hours. No mailbox in the pool reads mail at 3am or reacts within seconds. diff --git a/internal/app/analytics/warmup_ramp_live_test.go b/internal/app/analytics/warmup_ramp_live_test.go index f02f5434c..ccb80c8f3 100644 --- a/internal/app/analytics/warmup_ramp_live_test.go +++ b/internal/app/analytics/warmup_ramp_live_test.go @@ -60,7 +60,12 @@ func newRampFixture(t *testing.T, daysWarming, base, increase, max int) *rampFix VALUES ($1, $2, $3, $4, 'Ramp', '', '', 'smtp_imap', 'active', 50, 600, 'UTC', $5, $6, $7, $8)`, f.mailbox, f.user, f.org, "ramp-"+f.mailbox.String()[:8]+"@test.local", - time.Now().Add(-time.Duration(daysWarming)*24*time.Hour), base, increase, max) + // The column is timestamp without time zone and the app anchors it + // with the database's now(); pgx writes a local time's wall clock, so + // on a machine east of UTC the anchor lands hours late and the ramp + // reads a day short. A few hours into the day, not on its boundary, + // so a freeze shorter than a day cannot cross it either. + time.Now().UTC().Add(-time.Duration(daysWarming)*24*time.Hour-6*time.Hour), base, increase, max) t.Cleanup(func() { c := context.Background() diff --git a/internal/app/consumer/warmup_reply_back.go b/internal/app/consumer/warmup_reply_back.go index aa40b3bf4..cd3c011be 100644 --- a/internal/app/consumer/warmup_reply_back.go +++ b/internal/app/consumer/warmup_reply_back.go @@ -20,7 +20,9 @@ const ( // scheduleWarmupReplyBack occasionally points the RECIPIENT's next warmup send // back at the sender. It re-points an already-pending task and only ever pulls -// it earlier, so budgets and health gating are untouched. +// it earlier; health gating is untouched, and a send pulled into a day that is +// already spent is held to the next opening by the send-time budget check, +// aim intact. func (s *JobsService) scheduleWarmupReplyBack(ctx context.Context, token *models.WarmupToken, recipientAccountID uuid.UUID) { if s.TaskRepo == nil || s.EmailRepository == nil || token == nil { return diff --git a/internal/infrastructure/db/migrations/000181_task_status_org_suspended.down.sql b/internal/infrastructure/db/migrations/000181_task_status_org_suspended.down.sql new file mode 100644 index 000000000..87e78c9bc --- /dev/null +++ b/internal/infrastructure/db/migrations/000181_task_status_org_suspended.down.sql @@ -0,0 +1,2 @@ +-- task_status keeps 'skipped_org_suspended': Postgres cannot drop an enum value. +SELECT 1; diff --git a/internal/infrastructure/db/migrations/000181_task_status_org_suspended.up.sql b/internal/infrastructure/db/migrations/000181_task_status_org_suspended.up.sql new file mode 100644 index 000000000..5cc24b11f --- /dev/null +++ b/internal/infrastructure/db/migrations/000181_task_status_org_suspended.up.sql @@ -0,0 +1,4 @@ +-- A warmup send held for a suspended workspace has written this status since +-- the org risk posture shipped, but nothing added it, so every write failed and +-- left the task pending for the dispatcher to fire again. +ALTER TYPE public.task_status ADD VALUE IF NOT EXISTS 'skipped_org_suspended'; diff --git a/internal/scheduler/service.go b/internal/scheduler/service.go index 8bbe39785..11a80733d 100644 --- a/internal/scheduler/service.go +++ b/internal/scheduler/service.go @@ -16,6 +16,10 @@ import ( type SchedulerService interface { // Warmup scheduling CalculateNextWarmupTime(ctx context.Context, accountID uuid.UUID) (time.Time, error) + // WarmupDailyBudget is today's warmup target and what has gone out against + // it, read again at send time so a target cut after the send was placed + // still holds. + WarmupDailyBudget(ctx context.Context, accountID uuid.UUID) (WarmupBudget, error) // Campaign scheduling CalculateNextCampaignTime(ctx context.Context, campaignID uuid.UUID) (time.Time, *repository.ContactSequencePair, uuid.UUID, error) diff --git a/internal/scheduler/warmup_scheduler.go b/internal/scheduler/warmup_scheduler.go index 572d27f64..57caaea38 100644 --- a/internal/scheduler/warmup_scheduler.go +++ b/internal/scheduler/warmup_scheduler.go @@ -188,74 +188,20 @@ func (s *schedulerService) CalculateNextWarmupTime(ctx context.Context, accountI } } - // STEP 2: One shared resolve, so this target and the one the mailbox - // drawer reports cannot drift apart. - healthState := s.resolveHealthState(ctx, accountID) - rampAnchor := time.Now() - if account.Warmup != nil { - rampAnchor = *account.Warmup + // STEP 2: Today's budget. The send-time gate resolves it through the same + // code, so a send is never placed on one number and checked against another. + day, err := s.resolveWarmupBudget(ctx, account, activelyWarming, inCampaign) + if err != nil { + return time.Time{}, err } - plan := warmupramp.Resolve(ctx, s.warmupRepo, warmupramp.Input{ - AccountID: accountID, - WarmupStart: rampAnchor, - ActivelyWarming: activelyWarming, - Base: account.WarmupBase, - Increase: account.WarmupIncrease, - Max: account.WarmupMax, - InCampaign: inCampaign, - Health: healthState, - Now: time.Now(), - }) - targetVolume := plan.Target - if plan.Cut() { - log.Info(). - Str("email_account_id", accountID.String()). - Int("placements_48h", plan.Placements). - Int("sends_48h", plan.Sends). - Int("target", targetVolume). - Msg("warmup volume cut on an early placement signal; ramp held") - } - - // Vary the day's target so a mailbox doesn't send an identical count every - // day. Deterministic per (account, local day) so it's stable across the - // day's reschedules. Actively-warming mailboxes keep a floor of WarmupBase. - if activelyWarming && targetVolume > 0 { - factor := dailyVolumeFactor(accountID, time.Now().In(loadLocation(account.Timezone))) - varied := int(float64(targetVolume)*factor + 0.5) - if varied < account.WarmupBase { - varied = account.WarmupBase - } - if varied < 1 { - varied = 1 - } - if varied < targetVolume { - targetVolume = varied - } - } - - // STEP 2.1: Cap per-mailbox volume to actual recipient capacity. The - // sender should not send multiple warmup messages to the same recipient - // in a single day just to hit an arbitrary target; that creates obvious - // pool loops when membership is small. Recipient-only participants count - // here, so operators can add inbound capacity without making those - // mailboxes warmup senders. - if s.warmupRepo != nil { - poolType := s.warmupPoolTypeForAccount(ctx, account) - // The set the selector draws from, so the cap never exceeds what a send can reach. - candidates, err := s.warmupRepo.WarmupPartnerCandidates(ctx, poolType, accountID) - if err == nil { - eligibleRecipients := len(candidates) - if eligibleRecipients <= 0 { - return recipientRecheckTime(), nil - } - if targetVolume > eligibleRecipients { - targetVolume = eligibleRecipients - } - } + if day.NoPartners { + return recipientRecheckTime(), nil } + targetVolume := day.Target + emailsSentToday := day.Sent // Resolve owns the band's volume half; only its spacing half applies here. - adj := adjustmentFor(healthState) + adj := adjustmentFor(day.health) // Spacing: a drawn gap from the profile when one is enabled, otherwise the // mailbox's fixed min gap. The health-state multiplier still applies on top, @@ -265,14 +211,8 @@ func (s *schedulerService) CalculateNextWarmupTime(ctx context.Context, accountI minWaitSeconds = int(float64(minWaitSeconds)*adj.minWaitMultiplier + 0.5) } - // STEP 3: Count emails already sent today - emailsSentToday, err := s.taskRepo.CountWarmupEmailsSentToday(ctx, accountID) - if err != nil { - return time.Time{}, err - } - // STEP 4: Check if we've hit today's limit - if emailsSentToday >= targetVolume { + if day.Reached() { // Move to tomorrow's first slot return s.snapWarmupToBehavior(bhv, calculateFirstSlotTomorrowAt(account.Timezone, warmupStart)), nil } @@ -370,6 +310,138 @@ func (s *schedulerService) CalculateNextWarmupTime(ctx context.Context, accountI return s.snapWarmupToBehavior(bhv, candidateTime), nil } +// WarmupBudget is today's warmup volume for one mailbox and what has already +// gone out against it, resolved by the scheduler when it places a send and +// again when the send executes. The mailbox drawer does not read it: it shows +// the plan before daily variation and the recipient cap. +type WarmupBudget struct { + // Target is the day's volume after ramp, early cut, health band, daily + // variation and recipient capacity. + Target int + // Sent is the completed warmup sends counted against today. + Sent int + // NoPartners is set when the pool offers this mailbox nobody to write to. + // Target and Sent are not resolved then, and the partner draw rather than + // the cap is what stops the send, so Reached deliberately reports false. + NoPartners bool +} + +// Reached reports whether today has no volume left. +func (b WarmupBudget) Reached() bool { + return !b.NoPartners && b.Sent >= b.Target +} + +// warmupDay is the budget plus the health band the scheduler still needs for +// spacing. +type warmupDay struct { + WarmupBudget + health models.WarmupHealthState +} + +// WarmupDailyBudget resolves today's budget for the send-time check. The +// scheduler counts today's sends only when it places the NEXT send, so a +// signal that cuts the target between placing a send and executing it (an +// early placement, a health band, a partner leaving the pool) used to leave +// that send going out over the cut number. Reading the budget again at +// execution closes that window. +func (s *schedulerService) WarmupDailyBudget(ctx context.Context, accountID uuid.UUID) (WarmupBudget, error) { + account, xerr := s.emailRepo.GetByID(ctx, accountID) + if xerr != nil { + return WarmupBudget{}, xerr + } + activelyWarming := account.IsWarmingActive() + inCampaign := s.accountInActiveCampaign(ctx, accountID) + if !activelyWarming && !inCampaign { + return WarmupBudget{}, ErrWarmupNotEnabled + } + day, err := s.resolveWarmupBudget(ctx, account, activelyWarming, inCampaign) + if err != nil { + return WarmupBudget{}, err + } + return day.WarmupBudget, nil +} + +// resolveWarmupBudget is the one place today's target is computed: the shared +// ramp policy, the per-day variation, the recipient cap, then the count of +// what has already been sent against it. +func (s *schedulerService) resolveWarmupBudget(ctx context.Context, account *models.Email, activelyWarming, inCampaign bool) (warmupDay, error) { + accountID := account.ID + healthState := s.resolveHealthState(ctx, accountID) + rampAnchor := time.Now() + if account.Warmup != nil { + rampAnchor = *account.Warmup + } + plan := warmupramp.Resolve(ctx, s.warmupRepo, warmupramp.Input{ + AccountID: accountID, + WarmupStart: rampAnchor, + ActivelyWarming: activelyWarming, + Base: account.WarmupBase, + Increase: account.WarmupIncrease, + Max: account.WarmupMax, + InCampaign: inCampaign, + Health: healthState, + Now: time.Now(), + }) + targetVolume := plan.Target + if plan.Cut() { + log.Info(). + Str("email_account_id", accountID.String()). + Int("placements_48h", plan.Placements). + Int("sends_48h", plan.Sends). + Int("target", targetVolume). + Msg("warmup volume cut on an early placement signal; ramp held") + } + + // Vary the day's target so a mailbox doesn't send an identical count every + // day. Deterministic per (account, local day) so it's stable across the + // day's reschedules. Actively-warming mailboxes keep a floor of WarmupBase. + if activelyWarming && targetVolume > 0 { + factor := dailyVolumeFactor(accountID, time.Now().In(loadLocation(account.Timezone))) + varied := int(float64(targetVolume)*factor + 0.5) + if varied < account.WarmupBase { + varied = account.WarmupBase + } + if varied < 1 { + varied = 1 + } + if varied < targetVolume { + targetVolume = varied + } + } + + day := warmupDay{health: healthState} + + // Cap per-mailbox volume to actual recipient capacity. The sender should + // not send multiple warmup messages to the same recipient in a single day + // just to hit an arbitrary target; that creates obvious pool loops when + // membership is small. Recipient-only participants count here, so + // operators can add inbound capacity without making those mailboxes + // warmup senders. + if s.warmupRepo != nil { + poolType := s.warmupPoolTypeForAccount(ctx, account) + // The set the selector draws from, so the cap never exceeds what a send can reach. + candidates, err := s.warmupRepo.WarmupPartnerCandidates(ctx, poolType, accountID) + if err == nil { + eligibleRecipients := len(candidates) + if eligibleRecipients <= 0 { + day.NoPartners = true + return day, nil + } + if targetVolume > eligibleRecipients { + targetVolume = eligibleRecipients + } + } + } + day.Target = targetVolume + + sent, err := s.taskRepo.CountWarmupEmailsSentToday(ctx, accountID) + if err != nil { + return day, err + } + day.Sent = sent + return day, nil +} + // snapWarmupToBehavior moves a warmup candidate onto the mailbox's rolled // workday and randomises its sub-minute component. A profile with no reachable // window leaves the candidate untouched: warmup should degrade to its own diff --git a/internal/scheduler/warmup_scheduler_test.go b/internal/scheduler/warmup_scheduler_test.go index e170a193d..8cc3d82b8 100644 --- a/internal/scheduler/warmup_scheduler_test.go +++ b/internal/scheduler/warmup_scheduler_test.go @@ -62,3 +62,23 @@ func TestWarmupRampTarget(t *testing.T) { }) } } + +func TestWarmupBudgetReached(t *testing.T) { + tests := []struct { + name string + budget WarmupBudget + want bool + }{ + {"under target", WarmupBudget{Target: 10, Sent: 8}, false}, + {"at target", WarmupBudget{Target: 8, Sent: 8}, true}, + {"over target after a cut", WarmupBudget{Target: 8, Sent: 9}, true}, + {"nobody to write to is the partner draw's call, not the cap's", WarmupBudget{NoPartners: true}, false}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := tt.budget.Reached(); got != tt.want { + t.Errorf("Reached() = %v, want %v", got, tt.want) + } + }) + } +} diff --git a/internal/tasks/email_task.go b/internal/tasks/email_task.go index 56fd9dca3..cf4a680fa 100644 --- a/internal/tasks/email_task.go +++ b/internal/tasks/email_task.go @@ -14,6 +14,7 @@ import ( "github.com/warmbly/warmbly/internal/models" "github.com/warmbly/warmbly/internal/observability/errs" "github.com/warmbly/warmbly/internal/repository" + "github.com/warmbly/warmbly/internal/scheduler" "github.com/warmbly/warmbly/internal/tasks/proto" ) @@ -185,7 +186,12 @@ func (s *tasksService) HandleEmailTask(task *proto.ProcessTask) *errx.Error { // same domains, so leaving it running would keep spending the reputation // the suspension exists to protect. if s.orgBlocksSending(ctx, account.OrganizationID) { - _ = s.taskRepo.UpdateTaskStatus(ctx, taskID, "skipped_org_suspended") + // Discarding this error hid a status the enum did not have for months, + // with the task left pending for the dispatcher to fire again. + if err := s.taskRepo.UpdateTaskStatus(ctx, taskID, "skipped_org_suspended"); err != nil { + errs.CaptureException(err) + return errx.InternalError() + } executionStatus = "completed" return nil } @@ -207,6 +213,52 @@ func (s *tasksService) HandleEmailTask(task *proto.ProcessTask) *errx.Error { } } + // STEP 3.8: Today's target, read again now rather than trusted from when + // this send was placed. The scheduler counts the day only when it places + // the NEXT send, so a placement recorded in between cut the target for the + // drawer and the scheduler but not for the send already waiting, and the + // mailbox ended the day one over its own cut number (#592). A reply-back + // pulled forward from tomorrow lands here too. + // + // The aim is read first: it is what makes the successor still a reply, and + // the read is cheap next to being wrong about it. + var aim *uuid.UUID + if warmupTask, aimErr := s.taskRepo.GetWarmupTask(ctx, taskID); aimErr != nil { + log.Warn().Err(aimErr).Str("task_id", taskID.String()).Msg("warmup task aim unreadable; a held reply-back would go out as a fresh message") + } else if warmupTask != nil { + aim = warmupTask.TargetAccountID + } + + budget, budgetErr := s.scheduler.WarmupDailyBudget(ctx, account.ID) + switch { + case errors.Is(budgetErr, scheduler.ErrWarmupNotEnabled): + // The mailbox stopped warming between this handler's own state check + // and now (a campaign ended mid-tick). Benign: the send was already + // cleared, and the chain winds down on the next pass. + case budgetErr != nil: + // Not knowing how many have gone out today is exactly when a send must + // not go out: failing open here would reopen #592 whenever the database + // is struggling. The task is still pending, so this retries. + errs.CaptureException(budgetErr) + return errx.InternalError() + case budget.Reached(): + log.Info(). + Str("task_id", taskID.String()). + Str("email_account_id", account.ID.String()). + Int("sent_today", budget.Sent). + Int("target", budget.Target). + Msg("warmup send skipped: today's target is already reached") + // Acknowledged only once the task is marked, or the row stays pending + // and blocks the successor this chain needs. + if err := s.taskRepo.UpdateTaskStatus(ctx, taskID, "skipped_daily_limit"); err != nil { + errs.CaptureException(err) + return errx.InternalError() + } + s.rescheduleWarmupAfterCap(ctx, account.ID, aim) + executionStatus = "completed" + return nil + } + // STEP 4: Mark task as active (with advisory lock) if err := s.taskRepo.UpdateTaskStatusWithLock(ctx, taskID, "active"); err != nil { errs.CaptureException(err) @@ -841,8 +893,30 @@ func (s *tasksService) EnsureWarmupScheduled(ctx context.Context, accountID uuid return s.createWarmupTask(ctx, accountID, nextTime) } -// createWarmupTask creates a new warmup task in GCP Cloud Tasks +// rescheduleWarmupAfterCap parks the chain at the scheduler's next slot, which +// is tomorrow's opening once today is spent. A send that was aimed at one +// partner (a reply-back) keeps its aim, so the answer goes out first thing +// rather than being lost to the cap. A failure here is logged rather than +// returned: the task is already marked, and the reconciler re-seeds a mailbox +// that ends up with no pending task. +func (s *tasksService) rescheduleWarmupAfterCap(ctx context.Context, accountID uuid.UUID, aim *uuid.UUID) { + nextTime, err := s.scheduler.CalculateNextWarmupTime(ctx, accountID) + if err != nil { + nextTime = warmupPartnerRecheckTime() + } + if err := s.createWarmupTaskAimedAt(ctx, accountID, nextTime, aim); err != nil { + log.Warn().Err(err).Str("email_account_id", accountID.String()).Msg("Failed to reschedule warmup task after the daily target") + } +} + +// createWarmupTask creates the mailbox's next warmup wakeup. func (s *tasksService) createWarmupTask(ctx context.Context, accountID uuid.UUID, scheduleTime time.Time) error { + return s.createWarmupTaskAimedAt(ctx, accountID, scheduleTime, nil) +} + +// createWarmupTaskAimedAt is createWarmupTask with the send pointed at one +// partner, the way a reply-back points it. +func (s *tasksService) createWarmupTaskAimedAt(ctx context.Context, accountID uuid.UUID, scheduleTime time.Time, target *uuid.UUID) error { // Create task in database newTaskID := uuid.New() newTask := &Task{ @@ -855,7 +929,8 @@ func (s *tasksService) createWarmupTask(ctx context.Context, accountID uuid.UUID // Create warmup task entry warmupTask := &WarmupTask{ - TaskID: newTaskID, + TaskID: newTaskID, + TargetAccountID: target, } created, err := s.taskRepo.CreateWarmupTaskWithLock(ctx, newTask, warmupTask) diff --git a/internal/tasks/warmup_cut_cap_live_test.go b/internal/tasks/warmup_cut_cap_live_test.go new file mode 100644 index 000000000..cd40817d9 --- /dev/null +++ b/internal/tasks/warmup_cut_cap_live_test.go @@ -0,0 +1,329 @@ +package tasks + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgxpool" + + warmupapp "github.com/warmbly/warmbly/internal/app/warmup" + "github.com/warmbly/warmbly/internal/models" + "github.com/warmbly/warmbly/internal/pkg/encrypt" + "github.com/warmbly/warmbly/internal/repository" + "github.com/warmbly/warmbly/internal/scheduler" + "github.com/warmbly/warmbly/internal/tasks/proto" +) + +// The day's target is enforced at the moment a warmup send executes, not only +// when the next one is placed (#592). Run with WARMBLY_TEST_DB on a scratch +// database; the premium pool must be empty so the recipient cap is known. + +type capFixture struct { + pool *pgxpool.Pool + svc *tasksService + sender *recordingSender + user uuid.UUID + org uuid.UUID + mailbox uuid.UUID + partners []uuid.UUID +} + +const capFixturePartners = 10 + +func newCapFixture(t *testing.T) *capFixture { + t.Helper() + handle := liveCampaignDB(t) + pool := handle.Pool + requireEmptyPool(t, pool, models.WarmupPoolPremiumID) + requireEmptyPool(t, pool, models.WarmupPoolFreeID) + + f := &capFixture{pool: pool, sender: &recordingSender{}, user: uuid.New(), org: uuid.New(), mailbox: uuid.New()} + t.Cleanup(func() { + c := context.Background() + for _, step := range []struct { + sql string + arg any + }{ + {`DELETE FROM warmup_tokens WHERE sender_account_id = $1`, f.mailbox}, + {`DELETE FROM warmup_spam_reports WHERE reported_account_id = $1`, f.mailbox}, + {`DELETE FROM warmup_statistics WHERE email_account_id = $1`, f.mailbox}, + {`DELETE FROM warmup_tasks WHERE task_id IN (SELECT id FROM tasks WHERE email_account_id = $1)`, f.mailbox}, + {`DELETE FROM task_failures WHERE task_id IN (SELECT id FROM tasks WHERE email_account_id = $1)`, f.mailbox}, + {`DELETE FROM tasks WHERE email_account_id = $1`, f.mailbox}, + {`DELETE FROM warmup_pool_participants WHERE email_account_id IN (SELECT id FROM email_accounts WHERE organization_id = $1)`, f.org}, + {`DELETE FROM email_accounts WHERE organization_id = $1`, f.org}, + {`DELETE FROM organizations WHERE id = $1`, f.org}, + {`DELETE FROM users WHERE id = $1`, f.user}, + } { + if _, err := pool.Exec(c, step.sql, step.arg); err != nil { + t.Errorf("cleanup %q: %v", step.sql, err) + } + } + }) + + f.exec(t, `INSERT INTO users (id, email, first_name, last_name) VALUES ($1, $2, 'Cap', 'Test')`, + f.user, "cap-"+f.user.String()[:8]+"@test.local") + f.exec(t, `INSERT INTO organizations (id, name, slug, owner_user_id) VALUES ($1, 'Cap Test', $2, $3)`, + f.org, "cap-"+f.org.String()[:8], f.user) + // Day one of a base-10 ramp, anchored with the database clock the way the + // app anchors it. An always-open window keeps the clock out of the result. + f.exec(t, `INSERT INTO email_accounts (id, user_id, organization_id, email, name, signature_plain, signature_html, + provider, status, campaign_limit, min_wait_time, timezone, warmup, warmup_base, warmup_increase, + warmup_max, warmup_reply_rate, warmup_pool_type, warmup_start_time, warmup_end_time) + VALUES ($1, $2, $3, $4, 'Cap', '', '', 'smtp_imap', 'active', 50, 0, 'UTC', now(), 10, 1, 40, 0, + 'premium', '00:00', '23:59')`, + f.mailbox, f.user, f.org, "cap-"+f.mailbox.String()[:8]+"@test.local") + // Enough recipients that the recipient cap sits above the ramp target. + for i := 0; i < capFixturePartners; i++ { + id := uuid.New() + f.partners = append(f.partners, id) + f.exec(t, `INSERT INTO email_accounts (id, user_id, organization_id, email, name, signature_plain, signature_html, + provider, status, campaign_limit, min_wait_time, timezone, warmup_pool_type) + VALUES ($1, $2, $3, $4, 'Cap', '', '', 'smtp_imap', 'active', 50, 600, 'UTC', 'premium')`, + id, f.user, f.org, "cap-"+id.String()[:8]+"@partner.test") + f.exec(t, `INSERT INTO warmup_pool_participants (pool_id, email_account_id, participant_role, health_state) + VALUES ($1, $2, 'recipient_only', 'healthy')`, models.WarmupPoolPremiumID, id) + } + + enc, err := encrypt.NewEncrypter([]byte("0123456789abcdef0123456789abcdef")) + if err != nil { + t.Fatalf("encrypter: %v", err) + } + taskRepo := repository.NewTaskRepository(pool) + warmupRepo := repository.NewWarmupRepository(pool) + emailRepo := repository.NewEmailRepostory(handle, enc) + campaignRepo := repository.NewCampaignRepostory(handle) + f.svc = &tasksService{ + tasksClient: noopTaskScheduler{}, + scheduler: scheduler.NewSchedulerService(taskRepo, warmupRepo, nil, emailRepo, campaignRepo, nil, nil), + emailSender: f.sender, + warmupHealth: warmupapp.NewService(warmupRepo), + taskRepo: taskRepo, + warmupRepo: warmupRepo, + emailRepo: emailRepo, + campaignRepo: campaignRepo, + } + return f +} + +func (f *capFixture) exec(t *testing.T, sql string, args ...any) { + t.Helper() + if _, err := f.pool.Exec(context.Background(), sql, args...); err != nil { + t.Fatalf("fixture %q: %v", sql[:min(60, len(sql))], err) + } +} + +// sentToday records n warmup sends already completed against today. +func (f *capFixture) sentToday(t *testing.T, n int) { + t.Helper() + for i := 0; i < n; i++ { + f.exec(t, `INSERT INTO tasks (id, task_type, email_account_id, status, message_id, scheduled_at, completed_at) + VALUES ($1, 'warmup', $2, 'completed', $3, now(), now())`, + uuid.New(), f.mailbox, "") + } +} + +func (f *capFixture) placement(t *testing.T) { + t.Helper() + f.exec(t, `INSERT INTO warmup_spam_reports (id, reporter_account_id, reported_account_id, message_id, report_type, created_at) + VALUES (gen_random_uuid(), $1, $1, $2, 'spam_placement', now())`, + f.mailbox, "msg-"+uuid.New().String()) +} + +// pending places the mailbox's one warmup wakeup at the given time and +// returns its id. +func (f *capFixture) pending(t *testing.T, at time.Time) uuid.UUID { + t.Helper() + if err := f.svc.createWarmupTask(context.Background(), f.mailbox, at); err != nil { + t.Fatalf("create pending warmup task: %v", err) + } + id, _ := f.pendingTask(t) + return id +} + +// pendingTask is the mailbox's pending wakeup with its aim, or uuid.Nil. +func (f *capFixture) pendingTask(t *testing.T) (uuid.UUID, *uuid.UUID) { + t.Helper() + rows, err := f.pool.Query(context.Background(), ` + SELECT t.id, wt.target_account_id + FROM tasks t LEFT JOIN warmup_tasks wt ON wt.task_id = t.id + WHERE t.email_account_id = $1 AND t.task_type = 'warmup' AND t.status = 'pending'`, f.mailbox) + if err != nil { + t.Fatalf("pending tasks: %v", err) + } + defer rows.Close() + var id uuid.UUID + var target *uuid.UUID + n := 0 + for rows.Next() { + if err := rows.Scan(&id, &target); err != nil { + t.Fatalf("scan: %v", err) + } + n++ + } + if n > 1 { + t.Fatalf("%d pending warmup tasks; the chain must hold exactly one", n) + } + return id, target +} + +func (f *capFixture) status(t *testing.T, id uuid.UUID) string { + t.Helper() + var status string + if err := f.pool.QueryRow(context.Background(), `SELECT status FROM tasks WHERE id = $1`, id).Scan(&status); err != nil { + t.Fatalf("task status: %v", err) + } + return status +} + +func (f *capFixture) run(t *testing.T, id uuid.UUID) { + t.Helper() + if xerr := f.svc.HandleEmailTask(&proto.ProcessTask{TaskId: id.String()}); xerr != nil { + t.Fatalf("HandleEmailTask: %v", xerr.Message) + } +} + +func (f *capFixture) budget(t *testing.T) scheduler.WarmupBudget { + t.Helper() + b, err := f.svc.scheduler.WarmupDailyBudget(context.Background(), f.mailbox) + if err != nil { + t.Fatalf("budget: %v", err) + } + return b +} + +// budgetFailingScheduler is the real scheduler with the send-time budget read +// broken, which is what a database blip looks like at that moment. +type budgetFailingScheduler struct { + scheduler.SchedulerService +} + +func (budgetFailingScheduler) WarmupDailyBudget(context.Context, uuid.UUID) (scheduler.WarmupBudget, error) { + return scheduler.WarmupBudget{}, errors.New("budget read failed") +} + +// Not knowing the day's count is exactly when a send must not go out: failing +// open there would reopen the bug whenever the database is struggling. +func TestLiveWarmupUnreadableBudgetHoldsTheSendForRetry(t *testing.T) { + f := newCapFixture(t) + f.sentToday(t, 8) + task := f.pending(t, time.Now()) + f.svc.scheduler = budgetFailingScheduler{f.svc.scheduler} + + if xerr := f.svc.HandleEmailTask(&proto.ProcessTask{TaskId: task.String()}); xerr == nil { + t.Fatal("an unreadable budget was reported as success; the task would be acknowledged and never retried") + } + if f.sender.sent != 0 { + t.Fatalf("%d send(s) dispatched without knowing today's count", f.sender.sent) + } + if got := f.status(t, task); got != "pending" { + t.Fatalf("task status = %q, want pending so the retry picks it up", got) + } +} + +func TestLiveWarmupPendingSendRespectsTargetCutAfterScheduling(t *testing.T) { + f := newCapFixture(t) + f.sentToday(t, 8) + if b := f.budget(t); b.Target != 10 || b.Sent != 8 || b.Reached() { + t.Fatalf("before the placement: budget %+v, want target 10 with 8 sent", b) + } + task := f.pending(t, time.Now()) + + // The placement lands while the send is waiting; the day is now 8. + f.placement(t) + if b := f.budget(t); b.Target != 8 || !b.Reached() { + t.Fatalf("after the placement: budget %+v, want the cut target of 8, reached", b) + } + + f.run(t, task) + + if f.sender.sent != 0 { + t.Fatalf("target was cut to 8 but the pending send still went out: %d send(s) dispatched", f.sender.sent) + } + if got := f.status(t, task); got != "skipped_daily_limit" { + t.Fatalf("task status = %q, want skipped_daily_limit", got) + } + next, _ := f.pendingTask(t) + if next == uuid.Nil { + t.Fatal("the chain was not rescheduled; the mailbox would never warm again") + } + var at time.Time + if err := f.pool.QueryRow(context.Background(), `SELECT scheduled_at FROM tasks WHERE id = $1`, next).Scan(&at); err != nil { + t.Fatal(err) + } + if at.Before(time.Now().Add(time.Hour)) { + t.Fatalf("rescheduled for %s; a spent day parks the chain at the next opening, not now", at) + } +} + +func TestLiveWarmupSendUnderTargetStillGoesOut(t *testing.T) { + f := newCapFixture(t) + f.sentToday(t, 8) + task := f.pending(t, time.Now()) + + f.run(t, task) + + if f.sender.sent != 1 { + t.Fatalf("%d send(s) dispatched, want 1: the gate must only hold a spent day", f.sender.sent) + } + if got := f.status(t, task); got != "completed" { + t.Fatalf("task status = %q, want completed", got) + } + if b := f.budget(t); b.Sent != 9 { + t.Fatalf("sent today = %d after the send, want 9", b.Sent) + } +} + +// A reply-back re-points the pending send and pulls it earlier, including +// out of tomorrow into a day that is already spent. The cap holds it, and the +// aim survives so the answer goes out at the next opening instead of being +// dropped. +func TestLiveWarmupReplyBackPulledIntoASpentDayKeepsItsAim(t *testing.T) { + f := newCapFixture(t) + f.sentToday(t, 10) + task := f.pending(t, time.Now().Add(6*time.Hour)) + writer := f.partners[0] + moved, err := f.svc.taskRepo.DirectPendingWarmupTask(context.Background(), f.mailbox, writer, time.Now()) + if err != nil || !moved { + t.Fatalf("direct pending task: moved=%v err=%v", moved, err) + } + + f.run(t, task) + + if f.sender.sent != 0 { + t.Fatalf("%d send(s) dispatched over a spent day", f.sender.sent) + } + if got := f.status(t, task); got != "skipped_daily_limit" { + t.Fatalf("task status = %q, want skipped_daily_limit", got) + } + next, target := f.pendingTask(t) + if next == uuid.Nil { + t.Fatal("no successor task") + } + if target == nil || *target != writer { + t.Fatalf("successor aimed at %v, want the reply-back's writer %s", target, writer) + } +} + +// A suspended workspace's send is held under a status the enum carries, so the +// row leaves pending. It used to write a value the enum lacked, the write +// failed silently, and the dispatcher fired the task again every tick. +func TestLiveWarmupSuspendedWorkspaceMarksTheTask(t *testing.T) { + f := newCapFixture(t) + handle := liveCampaignDB(t) + f.exec(t, `UPDATE organizations SET risk_state = 'suspended' WHERE id = $1`, f.org) + f.svc.orgRiskRepo = repository.NewOrgRiskRepository(handle) + task := f.pending(t, time.Now()) + + f.run(t, task) + + if f.sender.sent != 0 { + t.Fatalf("%d send(s) dispatched from a suspended workspace", f.sender.sent) + } + if got := f.status(t, task); got != "skipped_org_suspended" { + t.Fatalf("task status = %q, want skipped_org_suspended", got) + } +} From 79bd4f62e3cdc0d3bd6cf7ccadf2113873c89997 Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Fri, 18 Sep 2026 22:46:08 -0700 Subject: [PATCH 2/2] feat: hold a warmup send whose send-time budget reports the mailbox not warming instead of letting it through, because a failed campaign read produces that sentinel and let a health-check mailbox send uncapped, and cover both held-status writes with a test that a hold whose write fails is retried rather than acknowledged --- internal/tasks/email_task.go | 13 +-- internal/tasks/warmup_cut_cap_live_test.go | 93 ++++++++++++++++++---- 2 files changed, 85 insertions(+), 21 deletions(-) diff --git a/internal/tasks/email_task.go b/internal/tasks/email_task.go index cf4a680fa..a98e2ff0d 100644 --- a/internal/tasks/email_task.go +++ b/internal/tasks/email_task.go @@ -231,15 +231,16 @@ func (s *tasksService) HandleEmailTask(task *proto.ProcessTask) *errx.Error { budget, budgetErr := s.scheduler.WarmupDailyBudget(ctx, account.ID) switch { - case errors.Is(budgetErr, scheduler.ErrWarmupNotEnabled): - // The mailbox stopped warming between this handler's own state check - // and now (a campaign ended mid-tick). Benign: the send was already - // cleared, and the chain winds down on the next pass. case budgetErr != nil: // Not knowing how many have gone out today is exactly when a send must // not go out: failing open here would reopen #592 whenever the database - // is struggling. The task is still pending, so this retries. - errs.CaptureException(budgetErr) + // is struggling. The task is still pending, so this retries. That covers + // ErrWarmupNotEnabled too, which a failed campaign read can produce: if + // the mailbox really stopped warming, the retry's own check above winds + // the chain down. + if !errors.Is(budgetErr, scheduler.ErrWarmupNotEnabled) { + errs.CaptureException(budgetErr) + } return errx.InternalError() case budget.Reached(): log.Info(). diff --git a/internal/tasks/warmup_cut_cap_live_test.go b/internal/tasks/warmup_cut_cap_live_test.go index cd40817d9..980bab699 100644 --- a/internal/tasks/warmup_cut_cap_live_test.go +++ b/internal/tasks/warmup_cut_cap_live_test.go @@ -199,28 +199,40 @@ func (f *capFixture) budget(t *testing.T) scheduler.WarmupBudget { // broken, which is what a database blip looks like at that moment. type budgetFailingScheduler struct { scheduler.SchedulerService + err error } -func (budgetFailingScheduler) WarmupDailyBudget(context.Context, uuid.UUID) (scheduler.WarmupBudget, error) { - return scheduler.WarmupBudget{}, errors.New("budget read failed") +func (s budgetFailingScheduler) WarmupDailyBudget(context.Context, uuid.UUID) (scheduler.WarmupBudget, error) { + return scheduler.WarmupBudget{}, s.err } // Not knowing the day's count is exactly when a send must not go out: failing -// open there would reopen the bug whenever the database is struggling. +// open there would reopen the bug whenever the database is struggling. A failed +// campaign read surfaces as "not warming", so that sentinel holds the send too. func TestLiveWarmupUnreadableBudgetHoldsTheSendForRetry(t *testing.T) { - f := newCapFixture(t) - f.sentToday(t, 8) - task := f.pending(t, time.Now()) - f.svc.scheduler = budgetFailingScheduler{f.svc.scheduler} + for _, tc := range []struct { + name string + err error + }{ + {"read failed", errors.New("budget read failed")}, + {"campaign read failed and reported not warming", scheduler.ErrWarmupNotEnabled}, + } { + t.Run(tc.name, func(t *testing.T) { + f := newCapFixture(t) + f.sentToday(t, 8) + task := f.pending(t, time.Now()) + f.svc.scheduler = budgetFailingScheduler{SchedulerService: f.svc.scheduler, err: tc.err} - if xerr := f.svc.HandleEmailTask(&proto.ProcessTask{TaskId: task.String()}); xerr == nil { - t.Fatal("an unreadable budget was reported as success; the task would be acknowledged and never retried") - } - if f.sender.sent != 0 { - t.Fatalf("%d send(s) dispatched without knowing today's count", f.sender.sent) - } - if got := f.status(t, task); got != "pending" { - t.Fatalf("task status = %q, want pending so the retry picks it up", got) + if xerr := f.svc.HandleEmailTask(&proto.ProcessTask{TaskId: task.String()}); xerr == nil { + t.Fatal("an unreadable budget was reported as success; the task would be acknowledged and never retried") + } + if f.sender.sent != 0 { + t.Fatalf("%d send(s) dispatched without knowing today's count", f.sender.sent) + } + if got := f.status(t, task); got != "pending" { + t.Fatalf("task status = %q, want pending so the retry picks it up", got) + } + }) } } @@ -327,3 +339,54 @@ func TestLiveWarmupSuspendedWorkspaceMarksTheTask(t *testing.T) { t.Fatalf("task status = %q, want skipped_org_suspended", got) } } + +// statusFailingRepo is the real repository with one status write refused, +// which is what a database blip looks like at that write. +type statusFailingRepo struct { + repository.TaskRepository + refuse string +} + +func (r statusFailingRepo) UpdateTaskStatus(ctx context.Context, taskID uuid.UUID, status string) error { + if status == r.refuse { + return errors.New("status write failed") + } + return r.TaskRepository.UpdateTaskStatus(ctx, taskID, status) +} + +// A hold whose status write fails must not be reported as handled: the row +// stays pending, and acknowledging it would leave it blocking the successor +// until the overdue sweep. Discarding this error is what hid a status the enum +// did not carry for months. +func TestLiveWarmupHoldIsNotAcknowledgedUntilTheTaskIsMarked(t *testing.T) { + for _, tc := range []struct { + name string + status string + arrange func(t *testing.T, f *capFixture) + }{ + {"daily target reached", "skipped_daily_limit", func(t *testing.T, f *capFixture) { + f.sentToday(t, 10) + }}, + {"workspace suspended", "skipped_org_suspended", func(t *testing.T, f *capFixture) { + f.exec(t, `UPDATE organizations SET risk_state = 'suspended' WHERE id = $1`, f.org) + f.svc.orgRiskRepo = repository.NewOrgRiskRepository(liveCampaignDB(t)) + }}, + } { + t.Run(tc.name, func(t *testing.T) { + f := newCapFixture(t) + tc.arrange(t, f) + task := f.pending(t, time.Now()) + f.svc.taskRepo = statusFailingRepo{TaskRepository: f.svc.taskRepo, refuse: tc.status} + + if xerr := f.svc.HandleEmailTask(&proto.ProcessTask{TaskId: task.String()}); xerr == nil { + t.Fatal("the hold was reported as handled although its status write failed") + } + if f.sender.sent != 0 { + t.Fatalf("%d send(s) dispatched", f.sender.sent) + } + if got := f.status(t, task); got != "pending" { + t.Fatalf("task status = %q, want pending so the retry picks it up", got) + } + }) + } +}