diff --git a/cmd/backend/main.go b/cmd/backend/main.go index 632b9a70..6852b053 100644 --- a/cmd/backend/main.go +++ b/cmd/backend/main.go @@ -1146,6 +1146,11 @@ func main() { analyticsRepository := repository.NewAnalyticsRepository(primaryDB) emailAccountErrorRepository := repository.NewEmailAccountErrorRepository(primaryDB) analyticsService = analytics.NewService(analyticsRepository, emailRepostory, campaignRepostory, emailAccountErrorRepository, warmupRepository) + // A mailbox out of cold rotation says so in its drawer; an active one + // needs no notice. + if aware, ok := analyticsService.(analytics.LifecycleAware); ok { + aware.WireLifecycle(repository.NewSendLifecycleRepository(primaryDB)) + } rateLimitRepository := repository.NewRateLimitRepository(primaryDB) rateLimitService = ratelimit.NewService(cache, rateLimitRepository) @@ -1219,6 +1224,10 @@ func main() { if aware, ok := schedulerService.(scheduler.OrgRiskAware); ok { aware.WireOrgRisk(orgRiskRepository) } + // A resting mailbox keeps its warmup traffic but leaves cold rotation. + if aware, ok := schedulerService.(scheduler.LifecycleAware); ok { + aware.WireLifecycle(repository.NewSendLifecycleRepository(primaryDB)) + } campaignService = campaign.NewService(campaignRepostory, taskRepository, emailRepostory, campaignLogRepository, featureGateService, dailyThrottleService, schedulerService, tasksClient, streamingPublisher) // The launch gate refuses a list that is known to be largely // undeliverable, using the same projection preflight reports. diff --git a/cmd/consumer/main.go b/cmd/consumer/main.go index 7472cc33..db4196f8 100644 --- a/cmd/consumer/main.go +++ b/cmd/consumer/main.go @@ -373,6 +373,7 @@ func main() { WarmupEngagementRepo: repository.NewWarmupEngagementRepository(primaryDB.Pool), WarmupService: warmupService, WorkerRepo: workerRepo, + LifecycleRepo: repository.NewSendLifecycleRepository(primaryDB), Publisher: eventsPublisher, StreamingPublisher: streamingPublisher, AdvancedService: advancedService, @@ -429,6 +430,9 @@ func main() { // risk_pool worker when the band changes. Skipped if AssignmentService // or WorkerRepo are nil. go jobsService.StartRiskRebalancer(ctx, 1*time.Hour) + // Same cadence, different question: risk_band picks the worker, the + // lifecycle picks whether the mailbox is in cold rotation at all. + go jobsService.StartLifecycleRebalancer(ctx, 1*time.Hour) // Tracking consumer (opens/clicks): a second subscription on the shared bus // for the tracking topic. It records open/click engagement and fires INSTANT diff --git a/docs/content/docs/guides/mailboxes.mdx b/docs/content/docs/guides/mailboxes.mdx index c0d656b9..d383c8fd 100644 --- a/docs/content/docs/guides/mailboxes.mdx +++ b/docs/content/docs/guides/mailboxes.mdx @@ -130,6 +130,27 @@ Within a tier, Warmbly picks the least-loaded healthy worker with capacity rathe Mailboxes migrate automatically when a trial upgrades, a subscription ends, a workspace moves to or from a dedicated worker, or a risk band changes. No action needed. +## Resting a tired mailbox + +Cold sending used to run at full volume until a mailbox crossed a hard health band, with nothing in between. A mailbox showing early fatigue either kept going or was quarantined. + +A mailbox now has a cold-rotation state, separate from its health: + +| State | Meaning | +|-------|---------| +| `active` | In cold rotation. The default, and where every mailbox starts | +| `resting` | Out of cold rotation to recover. Warmup keeps running, so its reputation stays alive | +| `reserve` | Held back by you. Never entered or left automatically | +| `warming` | Building reputation and not yet taking campaign traffic | + +A mailbox rests when its warmup health reaches `throttled` or worse, and returns on its own once it is healthy again **and** has held steady for three days. One good hour does not put it back at full cold volume. + +It does not rest on the `watch` band. That band is deliberately the one that changes nothing you can feel, and leaving cold rotation is very much something you feel. + + +This is not the same as the risk band that decides which sending worker and IP host a mailbox. A resting mailbox is usually still on a clean worker; it is simply not being offered campaign sends. The mailbox drawer says which state it is in and why. + + ## Health The Accounts list groups mailboxes as **Healthy** (sending normally), **Warming** (ramping through warmup), or **Needs attention** (paused, failing, or not sending), each row showing a live state and score. A drop between refreshes notifies you rather than waiting to be noticed. diff --git a/docs/content/docs/guides/workspace-export-import.mdx b/docs/content/docs/guides/workspace-export-import.mdx index d2f97a52..693caeb7 100644 --- a/docs/content/docs/guides/workspace-export-import.mdx +++ b/docs/content/docs/guides/workspace-export-import.mdx @@ -87,6 +87,7 @@ Some things belong to an instance rather than to a workspace, so they are not ap | Warmup pool membership | Pools are shared across every workspace on an instance, so membership is re-earned rather than asserted by a file | | Domain authentication timings | The verdict travels (public DNS reads the same anywhere), but the destination re-checks before it can stop any sending, so a mailbox is never blocked on an observation the new instance never made | | Risk and review status | A workspace's abuse posture is one platform's verdict about a tenant on its own infrastructure, reached from evidence the destination never saw. An archive can neither carry a restriction nor clear one | +| Cold rotation state | Whether a mailbox is resting or held in reserve is this instance's decision about sending it watched. Every mailbox arrives in normal rotation and earns its way out again | | Cold sending ramps | How far a mailbox had eased into cold volume raises its cap, and the destination never watched it send. Mailboxes re-graduate from their warmup maturity, which costs a few days and errs toward sending less | | Scheduled deletions | A pending deletion from the source must never follow the workspace to its new home | | Failure and delivery counters | A webhook endpoint's failure streak and auto-disable state, and whether a notification's email already went out, describe what happened on the source. They start fresh, so an endpoint is not pre-disabled on the new instance and a notification is not re-sent | diff --git a/internal/app/analytics/service.go b/internal/app/analytics/service.go index 76c9d4c3..41fa49b4 100644 --- a/internal/app/analytics/service.go +++ b/internal/app/analytics/service.go @@ -38,6 +38,9 @@ type analyticsService struct { campaignRepo repository.CampaignRepository emailAccountErrorsRepo repository.EmailAccountErrorRepository warmupRepo repository.WarmupRepository + // lifecycleRepo reads whether the mailbox is in cold rotation. + // Optional/nil-safe. + lifecycleRepo repository.SendLifecycleRepository } func NewService( @@ -221,20 +224,22 @@ func (s *analyticsService) GetAccountStatus(ctx context.Context, orgID, accountI } coldRamp := s.coldRampInfo(ctx, email) + lifecycle := s.sendLifecycleInfo(ctx, email.ID) return &models.EmailAccountStatus{ - ID: email.ID, - Email: email.Email, - Provider: email.Provider, - Status: email.Status, - LastSyncedAt: &email.LastSyncedAt, - Health: health, - Errors: errors, - DailyUsage: *usage, - WarmupStatus: warmupStatus, - WarmupHealth: warmupHealth, - InCampaign: inCampaign, - ColdRamp: coldRamp, + ID: email.ID, + Email: email.Email, + Provider: email.Provider, + Status: email.Status, + LastSyncedAt: &email.LastSyncedAt, + Health: health, + Errors: errors, + DailyUsage: *usage, + WarmupStatus: warmupStatus, + WarmupHealth: warmupHealth, + InCampaign: inCampaign, + ColdRamp: coldRamp, + SendLifecycle: lifecycle, }, nil } @@ -574,3 +579,30 @@ func (s *analyticsService) coldRampInfo(ctx context.Context, email *models.Email Held: warmupramp.ColdHeldUntil(rampStart, state.Placements, now, warmupramp.FreezeWindow) != nil, } } + +// sendLifecycleInfo reports the mailbox's cold-rotation state, but only when +// it is not active: an active mailbox is the normal case and needs no notice. +func (s *analyticsService) sendLifecycleInfo(ctx context.Context, accountID uuid.UUID) *models.SendLifecycleState { + if s.lifecycleRepo == nil { + return nil + } + states, err := s.lifecycleRepo.GetSendLifecycles(ctx, []uuid.UUID{accountID}) + if err != nil { + return nil + } + state, ok := states[accountID] + if !ok || state.State.SendsCold() { + return nil + } + return &state +} + +// WireLifecycle attaches the cold-sending lifecycle. +func (s *analyticsService) WireLifecycle(r repository.SendLifecycleRepository) { + s.lifecycleRepo = r +} + +// LifecycleAware is the optional capability the caller uses to attach it. +type LifecycleAware interface { + WireLifecycle(r repository.SendLifecycleRepository) +} diff --git a/internal/app/consumer/lifecycle_rebalancer.go b/internal/app/consumer/lifecycle_rebalancer.go new file mode 100644 index 00000000..9298913a --- /dev/null +++ b/internal/app/consumer/lifecycle_rebalancer.go @@ -0,0 +1,85 @@ +package jobs + +import ( + "context" + "time" + + "github.com/google/uuid" + "github.com/rs/zerolog/log" + + "github.com/warmbly/warmbly/internal/app/lifecycle" + "github.com/warmbly/warmbly/internal/models" +) + +// StartLifecycleRebalancer moves mailboxes in and out of cold rotation on the +// warmup health signal, on the same cadence as the risk rebalancer. +// +// Separate from that one on purpose: risk_band decides which worker hosts a +// mailbox, this decides whether it is offered cold sends at all. They move on +// the same signal but mean different things, and folding them together would +// make a resting mailbox look like a dirty-IP mailbox. +func (s *JobsService) StartLifecycleRebalancer(ctx context.Context, interval time.Duration) { + if s.LifecycleRepo == nil { + return + } + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + s.rebalanceLifecycles(ctx) + } + } +} + +func (s *JobsService) rebalanceLifecycles(ctx context.Context) { + candidates, err := s.LifecycleRepo.ListLifecycleCandidates(ctx, 500) + if err != nil { + log.Warn().Err(err).Msg("lifecycle rebalancer: list candidates failed") + return + } + + now := time.Now() + rested, resumed := 0, 0 + seen := make([]uuid.UUID, 0, len(candidates)) + for _, c := range candidates { + seen = append(seen, c.EmailAccountID) + d := lifecycle.Decide(c.Current, c.Since, c.HealthState, now) + if d.RestartProbation { + if err := s.LifecycleRepo.RestartProbation(ctx, c.EmailAccountID); err != nil { + log.Warn().Err(err).Str("email_account_id", c.EmailAccountID.String()). + Msg("lifecycle rebalancer: could not restart probation") + } + continue + } + if !d.Changed(c.Current) { + continue + } + // force=false: a mailbox its owner put in reserve is never moved here. + moved, err := s.LifecycleRepo.SetSendLifecycle(ctx, c.EmailAccountID, d.Next, d.Reason, false) + if err != nil { + log.Warn().Err(err).Str("email_account_id", c.EmailAccountID.String()). + Msg("lifecycle rebalancer: could not move mailbox") + continue + } + if !moved { + continue + } + if d.Next == models.SendLifecycleResting { + rested++ + } else { + resumed++ + } + } + // Stamp everything examined, so the next pass moves on rather than + // re-reading the same page forever. + if err := s.LifecycleRepo.MarkLifecycleChecked(ctx, seen); err != nil { + log.Warn().Err(err).Msg("lifecycle rebalancer: could not record the checked window") + } + + if rested > 0 || resumed > 0 { + log.Info().Int("rested", rested).Int("resumed", resumed).Msg("lifecycle rebalancer: moved mailboxes") + } +} diff --git a/internal/app/consumer/service.go b/internal/app/consumer/service.go index f8381e5d..eff7cac4 100644 --- a/internal/app/consumer/service.go +++ b/internal/app/consumer/service.go @@ -35,6 +35,9 @@ type JobsService struct { WarmupEngagementRepo repository.WarmupEngagementRepository WarmupService warmupapp.Service WorkerRepo repository.WorkerRepository + // LifecycleRepo moves mailboxes in and out of cold rotation. Nil disables + // the lifecycle rebalancer entirely. + LifecycleRepo repository.SendLifecycleRepository // Publisher for sending events to workers Publisher events.Publisher diff --git a/internal/app/lifecycle/lifecycle.go b/internal/app/lifecycle/lifecycle.go new file mode 100644 index 00000000..83b4c5f6 --- /dev/null +++ b/internal/app/lifecycle/lifecycle.go @@ -0,0 +1,58 @@ +// Package lifecycle decides when a cold mailbox rests and when it may return. +// Resting removes it from cold rotation while warmup keeps its reputation +// alive; it returns only after a clean probation. +package lifecycle + +import ( + "time" + + "github.com/warmbly/warmbly/internal/models" +) + +// Decision is what the rebalancer should do with one mailbox. +type Decision struct { + // Next is the state to move to, equal to the current one when nothing + // changes. + Next models.SendLifecycle + Reason string + // RestartProbation asks the caller to re-stamp the clock without changing + // state. Probation has to measure HEALTHY time: a mailbox that sat resting + // and unhealthy for three days would otherwise resume on its first healthy + // tick, having served no clean time at all. + RestartProbation bool +} + +// Changed reports whether the mailbox should move. +func (d Decision) Changed(current models.SendLifecycle) bool { return d.Next != current } + +// Decide maps warmup health onto the cold lifecycle. Rests at throttled and +// worse, never at watch: watch is defined to change nothing a customer feels. +func Decide(current models.SendLifecycle, since *time.Time, health models.WarmupHealthState, now time.Time) Decision { + if !current.AutoManaged() { + return Decision{Next: current} + } + + switch health { + case models.WarmupHealthThrottled: + return Decision{Next: models.SendLifecycleResting, + Reason: "warmup health is throttled; resting on warmup traffic to recover"} + case models.WarmupHealthQuarantined, models.WarmupHealthBlocked: + return Decision{Next: models.SendLifecycleResting, + Reason: "warmup health is " + string(health) + "; out of cold rotation until it recovers"} + } + + // Healthy or watch. A resting mailbox returns only after a probation, so + // one good hour cannot bounce it straight back to full cold volume. + if current == models.SendLifecycleResting { + if health != models.WarmupHealthHealthy { + // Still not healthy: the clean streak starts again from here. + return Decision{Next: current, RestartProbation: true} + } + state := models.SendLifecycleState{State: current, Since: since} + if state.ReadyToResume(now) { + return Decision{Next: models.SendLifecycleActive, Reason: "recovered and served its rest"} + } + return Decision{Next: current} + } + return Decision{Next: models.SendLifecycleActive} +} diff --git a/internal/app/lifecycle/lifecycle_test.go b/internal/app/lifecycle/lifecycle_test.go new file mode 100644 index 00000000..58ed905d --- /dev/null +++ b/internal/app/lifecycle/lifecycle_test.go @@ -0,0 +1,109 @@ +package lifecycle + +import ( + "testing" + "time" + + "github.com/warmbly/warmbly/internal/models" +) + +func TestDecideRestsOnRealTrouble(t *testing.T) { + now := time.Now() + for _, h := range []models.WarmupHealthState{ + models.WarmupHealthThrottled, models.WarmupHealthQuarantined, models.WarmupHealthBlocked, + } { + d := Decide(models.SendLifecycleActive, nil, h, now) + if d.Next != models.SendLifecycleResting { + t.Errorf("health %q gave %q, want resting", h, d.Next) + } + if d.Reason == "" { + t.Errorf("health %q rested with no reason", h) + } + } +} + +// Watch is the band defined to change nothing a customer can feel. Leaving +// cold rotation is very much something they feel. +func TestDecideDoesNotRestOnWatch(t *testing.T) { + if d := Decide(models.SendLifecycleActive, nil, models.WarmupHealthWatch, time.Now()); d.Next != models.SendLifecycleActive { + t.Errorf("watch gave %q, want active", d.Next) + } +} + +func TestDecideResumesOnlyAfterProbation(t *testing.T) { + now := time.Date(2026, 3, 10, 12, 0, 0, 0, time.UTC) + + served := now.Add(-models.RestProbation) + if d := Decide(models.SendLifecycleResting, &served, models.WarmupHealthHealthy, now); d.Next != models.SendLifecycleActive { + t.Errorf("a recovered mailbox that served its rest gave %q, want active", d.Next) + } + + fresh := now.Add(-time.Hour) + if d := Decide(models.SendLifecycleResting, &fresh, models.WarmupHealthHealthy, now); d.Next != models.SendLifecycleResting { + t.Errorf("an hour of rest gave %q, want it still resting", d.Next) + } + + if d := Decide(models.SendLifecycleResting, &served, models.WarmupHealthWatch, now); d.Next != models.SendLifecycleResting { + t.Errorf("a still-degraded mailbox gave %q, want resting", d.Next) + } +} + +// Reserve is the owner's decision, in both directions. +func TestDecideNeverTouchesReserveOrWarming(t *testing.T) { + now := time.Now() + for _, h := range []models.WarmupHealthState{ + models.WarmupHealthHealthy, models.WarmupHealthThrottled, models.WarmupHealthBlocked, + } { + if d := Decide(models.SendLifecycleReserve, nil, h, now); d.Next != models.SendLifecycleReserve { + t.Errorf("health %q moved a reserved mailbox to %q", h, d.Next) + } + if d := Decide(models.SendLifecycleWarming, nil, h, now); d.Next != models.SendLifecycleWarming { + t.Errorf("health %q moved a warming mailbox to %q", h, d.Next) + } + } +} + +// A row written before the column existed reads as empty and must be treated +// as active rather than left in limbo. +func TestDecideTreatsAnUnsetStateAsActive(t *testing.T) { + if d := Decide("", nil, models.WarmupHealthHealthy, time.Now()); d.Next != models.SendLifecycleActive { + t.Errorf("an unset lifecycle gave %q, want active", d.Next) + } +} + +// The bug this guards: probation has to measure HEALTHY time. A mailbox that +// sat resting and unhealthy for three days would otherwise resume on its first +// healthy tick, having served no clean time at all. +func TestDecideRestartsProbationWhileStillUnhealthy(t *testing.T) { + now := time.Date(2026, 3, 10, 12, 0, 0, 0, time.UTC) + longAgo := now.Add(-10 * 24 * time.Hour) + + d := Decide(models.SendLifecycleResting, &longAgo, models.WarmupHealthWatch, now) + if d.Next != models.SendLifecycleResting { + t.Errorf("next = %q, want it still resting", d.Next) + } + if !d.RestartProbation { + t.Error("a still-unhealthy mailbox must restart its clean streak, not bank the time") + } + + // Once healthy, the clock runs and is not restarted. + healthy := Decide(models.SendLifecycleResting, &longAgo, models.WarmupHealthHealthy, now) + if healthy.RestartProbation { + t.Error("a healthy mailbox must not have its probation restarted") + } + if healthy.Next != models.SendLifecycleActive { + t.Errorf("next = %q, want active after a served probation", healthy.Next) + } +} + +// A mailbox that is not resting has no probation to restart. +func TestDecideDoesNotRestartProbationForOtherStates(t *testing.T) { + now := time.Now() + for _, state := range []models.SendLifecycle{ + models.SendLifecycleActive, models.SendLifecycleReserve, models.SendLifecycleWarming, "", + } { + if d := Decide(state, nil, models.WarmupHealthWatch, now); d.RestartProbation { + t.Errorf("state %q asked to restart probation", state) + } + } +} diff --git a/internal/app/orgtransfer/spec.go b/internal/app/orgtransfer/spec.go index f86f600e..c2acce10 100644 --- a/internal/app/orgtransfer/spec.go +++ b/internal/app/orgtransfer/spec.go @@ -157,7 +157,14 @@ var Tables = []Table{ // mailbox's cold ceiling, and the destination never watched it send. // Cleared, the mailbox re-graduates from its warmup-maturity band, // which costs a few days and is the safe direction to be wrong in. - ResetOnImport: []string{"worker_id", "auth_checked_at", "auth_failing_since", "cold_ramp_started_at"}, + // send_lifecycle is this instance's operational decision about a + // mailbox it watched send. Importing "resting" would silence a mailbox + // on the destination for a reason nothing there observed; importing + // "active" would assert readiness the destination has not seen. + ResetOnImport: []string{ + "worker_id", "auth_checked_at", "auth_failing_since", "cold_ramp_started_at", + "send_lifecycle", "send_lifecycle_since", "send_lifecycle_reason", + }, }, { Name: "email_accounts_smtp_imap", Group: models.OrgDataGroupCore, diff --git a/internal/infrastructure/db/migrations/000102_send_lifecycle.down.sql b/internal/infrastructure/db/migrations/000102_send_lifecycle.down.sql new file mode 100644 index 00000000..282ab1b6 --- /dev/null +++ b/internal/infrastructure/db/migrations/000102_send_lifecycle.down.sql @@ -0,0 +1,11 @@ +DROP INDEX IF EXISTS public.idx_email_accounts_lifecycle_checked; +DROP INDEX IF EXISTS public.idx_email_accounts_send_lifecycle; + +ALTER TABLE public.email_accounts + DROP CONSTRAINT IF EXISTS email_accounts_send_lifecycle_check; + +ALTER TABLE public.email_accounts + DROP COLUMN IF EXISTS send_lifecycle_checked_at, + DROP COLUMN IF EXISTS send_lifecycle_reason, + DROP COLUMN IF EXISTS send_lifecycle_since, + DROP COLUMN IF EXISTS send_lifecycle; diff --git a/internal/infrastructure/db/migrations/000102_send_lifecycle.up.sql b/internal/infrastructure/db/migrations/000102_send_lifecycle.up.sql new file mode 100644 index 00000000..633fe5e8 --- /dev/null +++ b/internal/infrastructure/db/migrations/000102_send_lifecycle.up.sql @@ -0,0 +1,40 @@ +-- Cold mailbox lifecycle (issue #157). +-- +-- Cold sending ran until a hard health band tripped, with nothing in between: +-- a mailbox showing early fatigue either kept sending at full volume or was +-- quarantined. There was no way to pull one out of cold rotation to recover on +-- warmup traffic alone and put it back once it had. +-- +-- This is a different axis from risk_band, which decides WHICH worker and IP +-- host a mailbox. This decides whether the mailbox is offered to cold sending +-- at all; a resting mailbox is still a clean-band mailbox. +-- +-- Defaults to 'active' so no existing mailbox changes behaviour on deploy. +ALTER TABLE public.email_accounts + ADD COLUMN send_lifecycle text NOT NULL DEFAULT 'active', + -- When the current state was entered, so a rest has a measurable length + -- and promotion can require a probation window rather than a sweep tick. + ADD COLUMN send_lifecycle_since timestamptz, + ADD COLUMN send_lifecycle_reason text, + -- When the rebalancer last looked at this mailbox. Ordering candidates by + -- it guarantees rotation: ordering by send_lifecycle_since alone means + -- every never-moved mailbox sorts equal-first, so on an install with more + -- than one page of them the same page is re-examined forever and the rest + -- are never evaluated at all. + ADD COLUMN send_lifecycle_checked_at timestamptz; + +ALTER TABLE public.email_accounts + ADD CONSTRAINT email_accounts_send_lifecycle_check + CHECK (send_lifecycle IN ('warming', 'active', 'resting', 'reserve')); + +-- Cold sender resolution filters on this every scheduling pass. +CREATE INDEX idx_email_accounts_send_lifecycle + ON public.email_accounts USING btree (send_lifecycle) + WHERE send_lifecycle <> 'active'; + +-- The rebalancer's rotation order. +CREATE INDEX idx_email_accounts_lifecycle_checked + ON public.email_accounts USING btree (send_lifecycle_checked_at NULLS FIRST); + +COMMENT ON COLUMN public.email_accounts.send_lifecycle IS + 'Whether the mailbox is offered to cold sending: warming | active | resting | reserve. Orthogonal to risk_band, which decides which worker hosts it.'; diff --git a/internal/models/analytics.go b/internal/models/analytics.go index fe3ad955..e8f446df 100644 --- a/internal/models/analytics.go +++ b/internal/models/analytics.go @@ -107,6 +107,9 @@ type EmailAccountStatus struct { // When true a low-volume health-check warmup keeps running even if the // user has warmup paused/off. InCampaign bool `json:"in_campaign"` + // SendLifecycle is whether the mailbox is in cold rotation, present only + // when it is NOT: an active mailbox needs no explanation. + SendLifecycle *SendLifecycleState `json:"send_lifecycle,omitempty"` // ColdRamp is the warmup-to-cold graduation ceiling, present only while it // is below the mailbox's own cap. Without it the cap just reads lower than // the number the owner configured. diff --git a/internal/models/send_lifecycle.go b/internal/models/send_lifecycle.go new file mode 100644 index 00000000..8eb3eede --- /dev/null +++ b/internal/models/send_lifecycle.go @@ -0,0 +1,72 @@ +package models + +import "time" + +// SendLifecycle is whether a mailbox is offered to cold sending. +// +// Orthogonal to EmailRiskBand: that decides which worker and IP host a +// mailbox, this decides whether it is in cold rotation at all. A resting +// mailbox is still a clean-band mailbox. +type SendLifecycle string + +const ( + // SendLifecycleWarming is building reputation and not yet in cold rotation. + SendLifecycleWarming SendLifecycle = "warming" + // SendLifecycleActive is in cold rotation. The default. + SendLifecycleActive SendLifecycle = "active" + // SendLifecycleResting was pulled out of cold rotation to recover on + // warmup traffic alone, and returns on its own once it has. + SendLifecycleResting SendLifecycle = "resting" + // SendLifecycleReserve is held back deliberately by its owner. Never + // entered or left automatically. + SendLifecycleReserve SendLifecycle = "reserve" +) + +// SendsCold reports whether cold sender resolution may offer this mailbox. +func (l SendLifecycle) SendsCold() bool { + // An empty value is a row written before the column existed, which is + // active: a mailbox must never stop sending because of a missing default. + return l == SendLifecycleActive || l == "" +} + +// AutoManaged reports whether the rebalancer may move this state. Reserve is +// the owner's decision and is never overridden; warming is derived from the +// mailbox's own warmup settings. +func (l SendLifecycle) AutoManaged() bool { + return l == SendLifecycleActive || l == SendLifecycleResting || l == "" +} + +// Valid reports whether l is a state the database will accept. +func (l SendLifecycle) Valid() bool { + switch l { + case SendLifecycleWarming, SendLifecycleActive, SendLifecycleResting, SendLifecycleReserve: + return true + } + return false +} + +// SendLifecycleState is a mailbox's lifecycle with its history. +type SendLifecycleState struct { + State SendLifecycle `json:"state"` + Since *time.Time `json:"since,omitempty"` + Reason string `json:"reason,omitempty"` +} + +// RestProbation is how long a rested mailbox must stay healthy before it is +// offered cold traffic again. Long enough that a mailbox does not bounce +// between states on one good hour. +const RestProbation = 72 * time.Hour + +// RestingFor reports how long the mailbox has been resting, 0 if it is not. +func (s SendLifecycleState) RestingFor(now time.Time) time.Duration { + if s.State != SendLifecycleResting || s.Since == nil { + return 0 + } + return now.Sub(*s.Since) +} + +// ReadyToResume reports whether a resting mailbox has served its probation. +// Health is the caller's business; this only answers the clock. +func (s SendLifecycleState) ReadyToResume(now time.Time) bool { + return s.State == SendLifecycleResting && s.RestingFor(now) >= RestProbation +} diff --git a/internal/models/send_lifecycle_test.go b/internal/models/send_lifecycle_test.go new file mode 100644 index 00000000..5e163cc3 --- /dev/null +++ b/internal/models/send_lifecycle_test.go @@ -0,0 +1,75 @@ +package models + +import ( + "testing" + "time" +) + +func TestSendsCold(t *testing.T) { + if !SendLifecycleActive.SendsCold() { + t.Error("active must send cold") + } + // A row written before the column existed reads as empty. Treating that as + // "not sending" would silently stop every existing mailbox on deploy. + if !SendLifecycle("").SendsCold() { + t.Error("an unset lifecycle must send cold") + } + for _, l := range []SendLifecycle{SendLifecycleWarming, SendLifecycleResting, SendLifecycleReserve} { + if l.SendsCold() { + t.Errorf("%q must not be offered cold traffic", l) + } + } +} + +func TestAutoManaged(t *testing.T) { + // Reserve is the owner's decision; the rebalancer must never take a + // mailbox out of it. + if SendLifecycleReserve.AutoManaged() { + t.Error("reserve must not be automatically managed") + } + // Warming follows the mailbox's own warmup settings, not the rebalancer. + if SendLifecycleWarming.AutoManaged() { + t.Error("warming must not be automatically managed") + } + for _, l := range []SendLifecycle{SendLifecycleActive, SendLifecycleResting, ""} { + if !l.AutoManaged() { + t.Errorf("%q should be automatically managed", l) + } + } +} + +func TestReadyToResume(t *testing.T) { + now := time.Date(2026, 3, 10, 12, 0, 0, 0, time.UTC) + since := now.Add(-RestProbation) + + rested := SendLifecycleState{State: SendLifecycleResting, Since: &since} + if !rested.ReadyToResume(now) { + t.Error("a mailbox that served its full probation should be ready") + } + + fresh := now.Add(-time.Hour) + if (SendLifecycleState{State: SendLifecycleResting, Since: &fresh}).ReadyToResume(now) { + t.Error("an hour of rest is not a probation") + } + // No timestamp means nothing to measure, so it is not ready: promoting on a + // missing clock would defeat the probation entirely. + if (SendLifecycleState{State: SendLifecycleResting}).ReadyToResume(now) { + t.Error("a resting mailbox with no start time must not resume") + } + if (SendLifecycleState{State: SendLifecycleActive, Since: &since}).ReadyToResume(now) { + t.Error("only a resting mailbox resumes") + } +} + +func TestValid(t *testing.T) { + for _, l := range []SendLifecycle{SendLifecycleWarming, SendLifecycleActive, SendLifecycleResting, SendLifecycleReserve} { + if !l.Valid() { + t.Errorf("%q should be valid", l) + } + } + for _, l := range []SendLifecycle{"", "paused", "nonsense"} { + if l.Valid() { + t.Errorf("%q should not be valid", l) + } + } +} diff --git a/internal/repository/pg_send_lifecycle.go b/internal/repository/pg_send_lifecycle.go new file mode 100644 index 00000000..dec3df18 --- /dev/null +++ b/internal/repository/pg_send_lifecycle.go @@ -0,0 +1,149 @@ +package repository + +import ( + "context" + "time" + + "github.com/google/uuid" + + "github.com/warmbly/warmbly/internal/infrastructure/db" + "github.com/warmbly/warmbly/internal/models" +) + +// SendLifecycleRepository reads and moves a mailbox's cold-sending lifecycle. +type SendLifecycleRepository interface { + // GetSendLifecycles resolves a whole candidate pool in one round trip. The + // campaign scheduler reads this per pass, so it must not be per-account. + GetSendLifecycles(ctx context.Context, accountIDs []uuid.UUID) (map[uuid.UUID]models.SendLifecycleState, error) + // SetSendLifecycle moves one mailbox, stamping when and why. Refuses to + // move a mailbox its owner put in reserve unless force is set, so the + // rebalancer cannot override a deliberate hold. + SetSendLifecycle(ctx context.Context, accountID uuid.UUID, state models.SendLifecycle, reason string, force bool) (bool, error) + // ListLifecycleCandidates returns mailboxes the rebalancer may move, with + // the warmup health that decides where they go, oldest-checked first so + // every mailbox is reached rather than the same page every pass. + ListLifecycleCandidates(ctx context.Context, limit int) ([]LifecycleCandidate, error) + // MarkLifecycleChecked records that the rebalancer looked at these + // mailboxes, which is what rotates the candidate window. + MarkLifecycleChecked(ctx context.Context, accountIDs []uuid.UUID) error + // RestartProbation re-stamps a resting mailbox's clock without changing + // its state, so probation measures healthy time rather than time elapsed. + RestartProbation(ctx context.Context, accountID uuid.UUID) error +} + +// LifecycleCandidate is one mailbox the rebalancer is considering. +type LifecycleCandidate struct { + EmailAccountID uuid.UUID + Current models.SendLifecycle + Since *time.Time + HealthState models.WarmupHealthState +} + +type sendLifecycleRepository struct { + DB *db.DB +} + +func NewSendLifecycleRepository(database *db.DB) SendLifecycleRepository { + return &sendLifecycleRepository{DB: database} +} + +func (r *sendLifecycleRepository) GetSendLifecycles(ctx context.Context, accountIDs []uuid.UUID) (map[uuid.UUID]models.SendLifecycleState, error) { + out := make(map[uuid.UUID]models.SendLifecycleState, len(accountIDs)) + if len(accountIDs) == 0 { + return out, nil + } + rows, err := r.DB.Pool.Query(ctx, ` + SELECT id, send_lifecycle, send_lifecycle_since, send_lifecycle_reason + FROM email_accounts WHERE id = ANY($1::uuid[]) + `, accountIDs) + if err != nil { + return nil, err + } + defer rows.Close() + for rows.Next() { + var id uuid.UUID + var state string + var since *time.Time + var reason *string + if err := rows.Scan(&id, &state, &since, &reason); err != nil { + return nil, err + } + s := models.SendLifecycleState{State: models.SendLifecycle(state), Since: since} + if reason != nil { + s.Reason = *reason + } + out[id] = s + } + return out, rows.Err() +} + +func (r *sendLifecycleRepository) SetSendLifecycle(ctx context.Context, accountID uuid.UUID, state models.SendLifecycle, reason string, force bool) (bool, error) { + // The guard is in SQL so a concurrent owner setting reserve cannot be + // overwritten between a read and a write. + tag, err := r.DB.Pool.Exec(ctx, ` + UPDATE email_accounts + SET send_lifecycle = $2, + send_lifecycle_since = NOW(), + send_lifecycle_reason = NULLIF($3, '') + WHERE id = $1 + AND send_lifecycle <> $2 + AND ($4 OR send_lifecycle <> 'reserve') + `, accountID, string(state), reason, force) + if err != nil { + return false, err + } + return tag.RowsAffected() > 0, nil +} + +func (r *sendLifecycleRepository) ListLifecycleCandidates(ctx context.Context, limit int) ([]LifecycleCandidate, error) { + if limit <= 0 { + limit = 500 + } + // Only mailboxes the rebalancer may move, and only those whose warmup + // health is known: a mailbox in no pool has no signal to act on. + rows, err := r.DB.Pool.Query(ctx, ` + SELECT ea.id, ea.send_lifecycle, ea.send_lifecycle_since, + COALESCE(wpp.health_state, 'healthy') + FROM email_accounts ea + LEFT JOIN warmup_pool_participants wpp ON wpp.email_account_id = ea.id + WHERE ea.status = 'active' + AND ea.send_lifecycle IN ('active', 'resting') + ORDER BY ea.send_lifecycle_checked_at NULLS FIRST, ea.id + LIMIT $1 + `, limit) + if err != nil { + return nil, err + } + defer rows.Close() + + var out []LifecycleCandidate + for rows.Next() { + var c LifecycleCandidate + var state, health string + if err := rows.Scan(&c.EmailAccountID, &state, &c.Since, &health); err != nil { + return nil, err + } + c.Current = models.SendLifecycle(state) + c.HealthState = models.WarmupHealthState(health) + out = append(out, c) + } + return out, rows.Err() +} + +func (r *sendLifecycleRepository) MarkLifecycleChecked(ctx context.Context, accountIDs []uuid.UUID) error { + if len(accountIDs) == 0 { + return nil + } + _, err := r.DB.Pool.Exec(ctx, + `UPDATE email_accounts SET send_lifecycle_checked_at = NOW() WHERE id = ANY($1::uuid[])`, accountIDs) + return err +} + +func (r *sendLifecycleRepository) RestartProbation(ctx context.Context, accountID uuid.UUID) error { + _, err := r.DB.Pool.Exec(ctx, ` + UPDATE email_accounts + SET send_lifecycle_since = NOW() + WHERE id = $1 AND send_lifecycle = 'resting' + `, accountID) + return err +} diff --git a/internal/scheduler/campaign_scheduler.go b/internal/scheduler/campaign_scheduler.go index 60cb6a22..f6370b9e 100644 --- a/internal/scheduler/campaign_scheduler.go +++ b/internal/scheduler/campaign_scheduler.go @@ -265,6 +265,10 @@ func (s *schedulerService) CalculateNextCampaignTime(ctx context.Context, campai } riskMultiplier := riskState.CapMultiplier() + // Which mailboxes are in cold rotation at all. A resting mailbox keeps its + // warmup traffic and its reputation; it just is not offered cold sends. + lifecycles, lifecyclesKnown := s.sendLifecycles(ctx, accounts) + effectiveCap := func(acct models.Email) int { lim := min(acct.CampaignLimit, campaign.DailyLimit) if campaign.RampEnabled { @@ -335,6 +339,7 @@ func (s *schedulerService) CalculateNextCampaignTime(ctx context.Context, campai // reported as the DNS problem it is rather than as a scheduling one. enforceAuth, authGrace := s.domainAuthGate(ctx) authGated := 0 + lifecycleGated := 0 var candidates []AccountCandidate for _, acct := range accounts { @@ -351,6 +356,14 @@ func (s *schedulerService) CalculateNextCampaignTime(ctx context.Context, campai continue } + // Not in cold rotation. Checked here, beside the authentication gate, + // so a resting mailbox costs no capacity query. Applied only when the + // states were actually read. + if lifecyclesKnown && !lifecycles[acct.ID].State.SendsCold() { + lifecycleGated++ + continue + } + sentToday, err := s.taskRepo.CountCampaignEmailsSentToday(ctx, acct.ID) if err != nil { return time.Time{}, nil, uuid.Nil, err @@ -472,6 +485,15 @@ func (s *schedulerService) CalculateNextCampaignTime(ctx context.Context, campai return time.Time{}, nil, uuid.Nil, ErrDomainAuthFailing } + // Every mailbox is out of cold rotation. Say so rather than letting the + // campaign look stalled for no visible reason; they return on their own. + if len(candidates) == 0 && lifecycleGated == len(accounts) { + s.logCampaignDecision(ctx, campaignID, "mailboxes_resting", + "No mailbox is in cold rotation: all are resting or held in reserve", + map[string]interface{}{"resting_mailboxes": lifecycleGated, "pool_size": len(accounts)}) + return s.deferToNextDay(campaign), nil, accounts[0].ID, ErrCampaignDeferred + } + // STEP 8.25: Apply ESP matching to the under-budget candidate set. // strict → only matching mailboxes are eligible; if none, DEFER (never // send cross-provider). diff --git a/internal/scheduler/live_integration_test.go b/internal/scheduler/live_integration_test.go index a1996f82..bc4180aa 100644 --- a/internal/scheduler/live_integration_test.go +++ b/internal/scheduler/live_integration_test.go @@ -1019,3 +1019,72 @@ func TestLiveGraduationDoesNotGateAMailboxThatNeverWarmed(t *testing.T) { t.Errorf("a never-warmed mailbox allows %d, want its full 50 cap", got) } } + +// setLifecycle moves the fixture's mailbox in or out of cold rotation. +func (f *liveFixture) setLifecycle(t *testing.T, state models.SendLifecycle) { + t.Helper() + if _, err := f.pool.Exec(context.Background(), + `UPDATE email_accounts SET send_lifecycle = $2, send_lifecycle_since = NOW() WHERE id = $1`, + f.mailbox, string(state)); err != nil { + t.Fatalf("set lifecycle: %v", err) + } +} + +// Issue #157: a resting mailbox must actually leave cold rotation. The state +// existing in the database proves nothing; this proves the scheduler reads it. +func TestLiveRestingMailboxLeavesColdRotation(t *testing.T) { + handle, pool := liveDB(t) + f := newLiveFixture(t, pool, "UTC") + s := liveScheduler(t, handle, pool) + if aware, ok := s.(LifecycleAware); ok { + aware.WireLifecycle(repository.NewSendLifecycleRepository(handle)) + } + + // Active: the campaign schedules normally. + if _, _, _, err := s.CalculateNextCampaignTime(context.Background(), f.campaign); err != nil && + !errors.Is(err, ErrCampaignDeferred) { + t.Fatalf("an active mailbox could not be scheduled: %v", err) + } + + // Resting: the only mailbox in the pool is out, so the campaign defers + // rather than sending from a mailbox that is meant to be recovering. + f.setLifecycle(t, models.SendLifecycleResting) + _, _, _, err := s.CalculateNextCampaignTime(context.Background(), f.campaign) + if err == nil { + t.Fatal("a resting mailbox was still offered cold traffic") + } + if !errors.Is(err, ErrCampaignDeferred) { + t.Errorf("err = %v, want a deferral: the mailbox returns on its own", err) + } +} + +// Reserve is the owner holding a mailbox back, and it must be honoured the +// same way. +func TestLiveReservedMailboxLeavesColdRotation(t *testing.T) { + handle, pool := liveDB(t) + f := newLiveFixture(t, pool, "UTC") + s := liveScheduler(t, handle, pool) + if aware, ok := s.(LifecycleAware); ok { + aware.WireLifecycle(repository.NewSendLifecycleRepository(handle)) + } + + f.setLifecycle(t, models.SendLifecycleReserve) + if _, _, _, err := s.CalculateNextCampaignTime(context.Background(), f.campaign); err == nil { + t.Fatal("a reserved mailbox was still offered cold traffic") + } +} + +// The default must not change behaviour for any existing mailbox. +func TestLiveActiveIsTheDefault(t *testing.T) { + _, pool := liveDB(t) + f := newLiveFixture(t, pool, "UTC") + + var state string + if err := pool.QueryRow(context.Background(), + `SELECT send_lifecycle FROM email_accounts WHERE id = $1`, f.mailbox).Scan(&state); err != nil { + t.Fatalf("read lifecycle: %v", err) + } + if models.SendLifecycle(state) != models.SendLifecycleActive { + t.Errorf("a new mailbox is %q, want active", state) + } +} diff --git a/internal/scheduler/service.go b/internal/scheduler/service.go index caa819b1..8bbe3978 100644 --- a/internal/scheduler/service.go +++ b/internal/scheduler/service.go @@ -6,6 +6,7 @@ import ( "time" "github.com/google/uuid" + "github.com/rs/zerolog/log" "github.com/warmbly/warmbly/internal/app/behavior" "github.com/warmbly/warmbly/internal/models" "github.com/warmbly/warmbly/internal/repository" @@ -49,6 +50,43 @@ type schedulerService struct { // orgRiskRepo reads the organization's fused abuse posture. Optional/ // nil-safe: without it no organization is ever risk-capped. orgRiskRepo repository.OrgRiskRepository + // lifecycleRepo reads whether a mailbox is in cold rotation at all. + // Optional/nil-safe: without it every mailbox is treated as active. + lifecycleRepo repository.SendLifecycleRepository +} + +// WireLifecycle attaches the cold-sending lifecycle. +func (s *schedulerService) WireLifecycle(r repository.SendLifecycleRepository) { + s.lifecycleRepo = r +} + +// LifecycleAware is the optional capability the caller uses to attach it. +type LifecycleAware interface { + WireLifecycle(r repository.SendLifecycleRepository) +} + +// sendLifecycles resolves the pool's lifecycle states. The second return says +// whether the answer is usable: a nil map and "no states" are indistinguishable +// to the caller otherwise, and an unresolved state reads as active, which would +// quietly put every resting mailbox back into rotation on a transient error. +// +// The gate is skipped rather than closed on error. Failing closed would stop a +// customer's campaigns entirely on one bad query, which is the worse outcome; +// skipping is logged so it is visible rather than silent. +func (s *schedulerService) sendLifecycles(ctx context.Context, accounts []models.Email) (map[uuid.UUID]models.SendLifecycleState, bool) { + if s.lifecycleRepo == nil || len(accounts) == 0 { + return nil, false + } + ids := make([]uuid.UUID, 0, len(accounts)) + for _, a := range accounts { + ids = append(ids, a.ID) + } + states, err := s.lifecycleRepo.GetSendLifecycles(ctx, ids) + if err != nil { + log.Warn().Err(err).Msg("could not read mailbox lifecycles; cold rotation is unfiltered this pass") + return nil, false + } + return states, true } // DomainAuthPolicy resolves whether the sending-domain authentication gate is diff --git a/web/src/components/app/emails/InboxDetails.tsx b/web/src/components/app/emails/InboxDetails.tsx index 5b4958f8..11f9fd7c 100644 --- a/web/src/components/app/emails/InboxDetails.tsx +++ b/web/src/components/app/emails/InboxDetails.tsx @@ -115,6 +115,32 @@ function RampHoldNotice({ hold }: { hold: import("@/lib/api/models/app/analytics ); } +// A mailbox that has quietly stopped receiving campaign sends looks broken. +function LifecycleNotice({ state }: { state: import("@/lib/api/models/app/analytics/AccountStatus").SendLifecycleState }) { + const resting = state.state === "resting"; + const copy = resting + ? "This mailbox is resting: it keeps its warmup traffic to rebuild reputation, but campaigns are not sending from it. It returns on its own once it has recovered and held steady for three days." + : state.state === "reserve" + ? "This mailbox is held in reserve, so campaigns will not send from it until you put it back." + : "This mailbox is still warming up, so campaigns are not sending from it yet."; + return ( +
+
+ +
+

+ Not sending campaigns ({state.state}) +

+

+ {copy} + {state.reason ? ` ${state.reason}.` : ""} +

+
+
+
+ ); +} + // A cold cap below the configured one reads as a bug unless it says why. function ColdRampNotice({ ramp }: { ramp: import("@/lib/api/models/app/analytics/AccountStatus").ColdRampInfo }) { return ( @@ -483,6 +509,7 @@ function OverviewTab({ status, loading, mailbox }: { status?: import("@/lib/api/ {ws?.ramp_hold && } + {status?.send_lifecycle && } {status?.cold_ramp && } {/* Errors */} diff --git a/web/src/lib/api/models/app/analytics/AccountStatus.ts b/web/src/lib/api/models/app/analytics/AccountStatus.ts index f47730e4..5baf5787 100644 --- a/web/src/lib/api/models/app/analytics/AccountStatus.ts +++ b/web/src/lib/api/models/app/analytics/AccountStatus.ts @@ -72,6 +72,8 @@ export default interface AccountStatus { warmup_status?: WarmupStatusInfo; /** Present only while graduation holds the cold cap below the configured one. */ cold_ramp?: ColdRampInfo; + /** Present only when the mailbox is NOT in cold rotation. */ + send_lifecycle?: SendLifecycleState; warmup_health?: WarmupHealthInfo; // True when the mailbox backs a live campaign — a low-volume health-check // warmup keeps running even if the user has warmup paused/off. @@ -85,3 +87,12 @@ export interface ColdRampInfo { days_to_full_cap: number; held: boolean; } + +export type SendLifecycle = "warming" | "active" | "resting" | "reserve"; + +// Why a mailbox is not being offered cold sends. +export interface SendLifecycleState { + state: SendLifecycle; + since?: string; + reason?: string; +}