Merge pull request #595 from warmbly/fix/warmup-system-issue-592

feat: hold a pending warmup send when the day's target is cut after it was scheduled
This commit is contained in:
Matthew Meszaros
2026-09-19 06:03:47 +00:00
committed by GitHub
10 changed files with 656 additions and 77 deletions
+3 -1
View File
@@ -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.
@@ -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()
+3 -1
View File
@@ -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
@@ -0,0 +1,2 @@
-- task_status keeps 'skipped_org_suspended': Postgres cannot drop an enum value.
SELECT 1;
@@ -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';
+4
View File
@@ -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)
+143 -71
View File
@@ -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
@@ -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)
}
})
}
}
+79 -3
View File
@@ -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,53 @@ 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 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. 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().
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 +894,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 +930,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)
+392
View File
@@ -0,0 +1,392 @@
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, "<sent-"+uuid.New().String()+"@test.local>")
}
}
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
err error
}
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. A failed
// campaign read surfaces as "not warming", so that sentinel holds the send too.
func TestLiveWarmupUnreadableBudgetHoldsTheSendForRetry(t *testing.T) {
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)
}
})
}
}
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)
}
}
// 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)
}
})
}
}