From d667c1dacd1950c5dd525924ba8ce82d3c5bd8a6 Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Sat, 19 Sep 2026 09:58:52 -0700 Subject: [PATCH] feat: keep the send plan honest at the edges: a behaviour mailbox reopening after the window closes is closed for the day, only a mailbox coming back on its own (spent, closed, spaced out) holds its bound leads as waiting, the min gap is charged from when the mailbox is next open, an end date later today ends the day there, explainCap is stagedCap's last stage so there is one clamp chain, last-send times and warmup health are read for the pool in one query each, the plan cache is keyed on the campaign's status and updated_at so a start or an edit is answered fresh, and the workspace posture row links to /app/deliverability --- internal/app/campaign/handlers.go | 4 +- internal/repository/pg_task.go | 34 ++++++ internal/repository/pg_warmup.go | 49 ++++++++ internal/scheduler/campaign_sender.go | 29 +---- internal/scheduler/send_plan.go | 108 +++++++++++++----- internal/scheduler/send_plan_test.go | 7 ++ .../components/app/campaigns/SendPlanCard.tsx | 2 +- 7 files changed, 180 insertions(+), 53 deletions(-) diff --git a/internal/app/campaign/handlers.go b/internal/app/campaign/handlers.go index e6a54c7c3..69e5c9438 100644 --- a/internal/app/campaign/handlers.go +++ b/internal/app/campaign/handlers.go @@ -1293,7 +1293,9 @@ func (s *campaignService) SendPlan(ctx context.Context, orgID uuid.UUID, campaig if !ok { return nil, errx.New(errx.Internal, "send planning is not available") } - key := campaign.ID.String() + // Keyed on the campaign's own version, so an edit or a start/stop is + // answered fresh while two viewers of an unchanged campaign share a read. + key := campaign.ID.String() + "|" + campaign.Status + "|" + campaign.UpdatedAt.UTC().Format(time.RFC3339Nano) if s.planCache != nil { if plan, ok := s.planCache.get(key); ok { return plan, nil diff --git a/internal/repository/pg_task.go b/internal/repository/pg_task.go index ffeeb4d17..45c610948 100644 --- a/internal/repository/pg_task.go +++ b/internal/repository/pg_task.go @@ -111,6 +111,9 @@ type TaskRepository interface { // Scheduling queries (CRITICAL for "next best time" calculation) CountEmailsSentToday(ctx context.Context, accountID uuid.UUID) (int, error) GetLastEmailTime(ctx context.Context, accountID uuid.UUID) (*time.Time, error) + // GetLastEmailTimes is GetLastEmailTime for a pool in one read (every + // dispatched task type, as the min-gap clock counts them). + GetLastEmailTimes(ctx context.Context, accountIDs []uuid.UUID) (map[uuid.UUID]time.Time, error) // GetLastSendTimes is the batch form, for rotation across a campaign's // whole sender pool. Accounts that have never sent are absent from the map. GetLastSendTimes(ctx context.Context, accountIDs []uuid.UUID, taskType string) (map[uuid.UUID]time.Time, error) @@ -453,6 +456,37 @@ func (r *taskRepository) CountCampaignEmailsSentToday(ctx context.Context, accou return count, err } +// GetLastEmailTimes is the min-gap clock for a whole pool in one query. +func (r *taskRepository) GetLastEmailTimes(ctx context.Context, accountIDs []uuid.UUID) (map[uuid.UUID]time.Time, error) { + out := make(map[uuid.UUID]time.Time, len(accountIDs)) + if len(accountIDs) == 0 { + return out, nil + } + query := ` + SELECT t.email_account_id, MAX(t.completed_at) + FROM tasks t + WHERE t.email_account_id = ANY($1) + AND t.status = 'completed' + AND t.completed_at IS NOT NULL + AND ` + taskDispatchedEmail + ` + GROUP BY t.email_account_id + ` + rows, err := r.db.Query(ctx, query, accountIDs) + if err != nil { + return nil, err + } + defer rows.Close() + for rows.Next() { + var id uuid.UUID + var at time.Time + if err := rows.Scan(&id, &at); err != nil { + return nil, err + } + out[id] = at + } + return out, rows.Err() +} + // CountCampaignEmailsSentTodayByAccounts is the per-mailbox ledger for a pool // in one query, so a plan over a large workspace does not ask once per mailbox. func (r *taskRepository) CountCampaignEmailsSentTodayByAccounts(ctx context.Context, accountIDs []uuid.UUID) (map[uuid.UUID]int, error) { diff --git a/internal/repository/pg_warmup.go b/internal/repository/pg_warmup.go index d41375b27..1d7d5ae69 100644 --- a/internal/repository/pg_warmup.go +++ b/internal/repository/pg_warmup.go @@ -113,6 +113,9 @@ type WarmupRepository interface { // gate cold sends on warmup health without needing a pool type. Returns // ("healthy", nil) when the account is in no pool. GetHealthState(ctx context.Context, accountID uuid.UUID) (models.WarmupHealthState, *time.Time, error) + // GetHealthStates is GetHealthState for a pool in one read. A mailbox in + // no pool is healthy, as the single read reports it. + GetHealthStates(ctx context.Context, accountIDs []uuid.UUID) (map[uuid.UUID]WarmupHealthRead, error) UnblockFromPool(ctx context.Context, accountID uuid.UUID) error IsInPool(ctx context.Context, accountID uuid.UUID, poolType string) (bool, error) GetParticipantHealth(ctx context.Context, accountID uuid.UUID, poolType string) (*models.WarmupParticipantHealth, error) @@ -400,6 +403,52 @@ func (r *warmupRepository) BlockFromPool(ctx context.Context, accountID uuid.UUI return err } +// WarmupHealthRead is one mailbox's worst live standing across its pools. +type WarmupHealthRead struct { + State models.WarmupHealthState + BlockedUntil *time.Time +} + +// GetHealthStates is GetHealthState over a pool: the worst standing per +// mailbox, in one query. +func (r *warmupRepository) GetHealthStates(ctx context.Context, accountIDs []uuid.UUID) (map[uuid.UUID]WarmupHealthRead, error) { + out := make(map[uuid.UUID]WarmupHealthRead, len(accountIDs)) + for _, id := range accountIDs { + out[id] = WarmupHealthRead{State: models.WarmupHealthHealthy} + } + if len(accountIDs) == 0 { + return out, nil + } + query := ` + SELECT DISTINCT ON (email_account_id) email_account_id, health_state, blocked_until + FROM warmup_pool_participants + WHERE email_account_id = ANY($1) + ORDER BY email_account_id, CASE health_state + WHEN 'blocked' THEN 5 + WHEN 'quarantined' THEN 4 + WHEN 'throttled' THEN 3 + WHEN 'watch' THEN 2 + WHEN 'healthy' THEN 1 + ELSE 0 + END DESC + ` + rows, err := r.db.Query(ctx, query, accountIDs) + if err != nil { + return nil, err + } + defer rows.Close() + for rows.Next() { + var id uuid.UUID + var state string + var until *time.Time + if err := rows.Scan(&id, &state, &until); err != nil { + return nil, err + } + out[id] = WarmupHealthRead{State: models.WarmupHealthState(state), BlockedUntil: until} + } + return out, rows.Err() +} + // GetHealthState returns the account's warmup health state and blocked_until without the // caller naming a pool. The ordering keeps the worst state winning. func (r *warmupRepository) GetHealthState(ctx context.Context, accountID uuid.UUID) (models.WarmupHealthState, *time.Time, error) { diff --git a/internal/scheduler/campaign_sender.go b/internal/scheduler/campaign_sender.go index e71f66848..f938ca5d3 100644 --- a/internal/scheduler/campaign_sender.go +++ b/internal/scheduler/campaign_sender.go @@ -123,31 +123,10 @@ const ( // explainCap is effectiveCap with its working shown: the same clamps, each // recorded when it is the one that lowers the cap. func (p *campaignPass) explainCap(acct models.Email) capClamp { - c := p.campaign - out := capClamp{Cap: acct.CampaignLimit, LimitedBy: capByMailbox} - clamp := func(v int, why string) { - if v < out.Cap { - out.Cap, out.LimitedBy = v, why - } - } - clamp(c.DailyLimit, capByCampaign) - if c.RampEnabled { - clamp(campaignRampCeiling(true, c.RampStart, c.RampIncrement, c.RampCeiling, c.RampLevel), capByRamp) - } - // Graduation ceiling: a mailbox at its warmup ceiling must not reach the - // full cold cap the day it joins a campaign. - clamp(coldCeilingFor(p.coldRamp[acct.ID], out.Cap), capByGraduation) - if m := p.risk.CapMultiplier(); m < 1 { - risked := int(float64(out.Cap)*m + 0.5) - // A restricted organization still sends, just far less. Zeroing it here - // would stop the campaign without ever saying why; suspension is the - // band that stops sending, and it does so at the send gate. - if risked < 1 { - risked = 1 - } - clamp(risked, capByRisk) - } - return out + // One chain: stagedCap (send_plan.go) applies the clamps and keeps each + // stage so the plan can show its working; this is its last stage. + stages, by := stagedCap(p, acct) + return capClamp{Cap: stages[4], LimitedBy: by} } // poolBudget is each mailbox's day as the activity feed reports it: the cap diff --git a/internal/scheduler/send_plan.go b/internal/scheduler/send_plan.go index 8d15852b8..9b5968e00 100644 --- a/internal/scheduler/send_plan.go +++ b/internal/scheduler/send_plan.go @@ -76,9 +76,14 @@ func stagedCap(p *campaignPass, acct models.Email) (stages [5]int, limitedBy str } else { stages[2] = cur } + // Graduation ceiling: a mailbox at its warmup ceiling must not reach the + // full cold cap the day it joins a campaign. clamp(3, coldCeilingFor(p.coldRamp[acct.ID], cur), capByGraduation) if m := p.risk.CapMultiplier(); m < 1 { risked := int(float64(cur)*m + 0.5) + // A restricted organization still sends, just far less. Zeroing it + // here would stop the campaign without ever saying why; suspension is + // the band that stops sending, and it does so at the send gate. if risked < 1 { risked = 1 } @@ -97,7 +102,7 @@ func room(capv, sentThis int) int { // planMailbox walks one mailbox through the day. windowSecondsLeft is the // campaign's sending time still ahead today; zero means the window is closed // for the rest of the day and the campaign-level clamp reports it instead. -func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass, acct models.Email, sentThis, sentAll int, now time.Time, windowSecondsLeft int, windowClosesAt time.Time) mailboxDay { +func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass, acct models.Email, sentThis, sentAll int, lastSend *time.Time, now time.Time, windowSecondsLeft int, windowClosesAt time.Time) mailboxDay { // A cap lowered after sends went out counts what went out, so the // waterfall's arithmetic holds on the very day the cap moved. d := mailboxDay{acct: acct, configured: max(acct.CampaignLimit, sentThis), sentThis: sentThis, sentOther: max(0, sentAll-sentThis), minGap: acct.MinWaitTime} @@ -188,6 +193,12 @@ func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass, } if openAt.After(now) { d.reopensAt = openAt + if !openAt.Before(windowClosesAt) { + d.state = models.MailboxPlanHoursClosed + d.gate = mailboxGate{reason: gateHours, paced: true, reopensAt: openAt} + d.byHours = r + return d + } secondsLeft = min(secondsLeft, int(windowClosesAt.Sub(openAt).Seconds())) } next = min(r, s.behaviorDailyCap(ctx, bhv, r, openAt)) @@ -244,10 +255,16 @@ func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass, // Spacing: one send per min gap, from the later of now and the mailbox's // last send plus its gap. Warmup mail sits on the same clock. var earliest time.Time - if last, err := s.taskRepo.GetLastEmailTime(ctx, acct.ID); err == nil && last != nil && d.minGap > 0 { - if at := last.Add(time.Duration(d.minGap) * time.Second); at.After(now) { + if lastSend != nil && d.minGap > 0 { + // Only the part of the gap that falls after the mailbox is open + // again costs sending time; the rest passes while it is closed. + from := now + if d.reopensAt.After(from) { + from = d.reopensAt + } + if at := lastSend.Add(time.Duration(d.minGap) * time.Second); at.After(from) { earliest = at - secondsLeft -= int(at.Sub(now).Seconds()) + secondsLeft -= int(at.Sub(from).Seconds()) } } paceMax := 0 @@ -294,10 +311,17 @@ func dayWindow(campaign *models.Campaign, now time.Time) (models.CampaignSendWin if campaign.EndDate != nil && campaign.EndDate.Before(now) { return out, 0, midnight } + // An end date later today ends the day there. + endsToday := func(closes time.Time) time.Time { + if campaign.EndDate != nil && campaign.EndDate.Before(closes) { + return *campaign.EndDate + } + return closes + } if windows.IsEmpty() { out.SendingDay = true out.OpenNow = out.StartsAt == nil || !campaign.StartDate.After(now) - closes := midnight.AddDate(0, 0, 1) + closes := endsToday(midnight.AddDate(0, 0, 1)) out.ClosesAt = &closes from := now if out.StartsAt != nil { @@ -314,21 +338,35 @@ func dayWindow(campaign *models.Campaign, now time.Time) (models.CampaignSendWin return out, left, closes } nowMin := local.Hour()*60 + local.Minute() + // Minutes of today's windows at or after `from`, clipped to the end date. + endMin := 24 * 60 + if campaign.EndDate != nil { + if e := campaign.EndDate.In(tz); e.Before(midnight.AddDate(0, 0, 1)) { + endMin = e.Hour()*60 + e.Minute() + } + } + minutesFrom := func(from int) int { + total := 0 + for _, iv := range windows[int(local.Weekday())] { + end := min(iv.End, endMin) + if end > from { + total += end - max(iv.Start, from) + } + } + return total + } closes := midnight - left := 0 for _, iv := range windows[int(local.Weekday())] { out.SendingDay = true - end := midnight.Add(time.Duration(iv.End) * time.Minute) + end := midnight.Add(time.Duration(min(iv.End, endMin)) * time.Minute) if end.After(closes) { closes = end } - if nowMin >= iv.Start && nowMin < iv.End { + if nowMin >= iv.Start && nowMin < min(iv.End, endMin) { out.OpenNow = true } - if iv.End > nowMin { - left += (iv.End - max(iv.Start, nowMin)) * 60 - } } + left := minutesFrom(nowMin) * 60 if out.SendingDay { out.ClosesAt = &closes } @@ -336,13 +374,7 @@ func dayWindow(campaign *models.Campaign, now time.Time) (models.CampaignSendWin if out.StartsAt != nil { out.OpenNow = false if out.StartsAt.Before(closes) { - from := max(nowMin, out.StartsAt.In(tz).Hour()*60+out.StartsAt.In(tz).Minute()) - left = 0 - for _, iv := range windows[int(local.Weekday())] { - if iv.End > from { - left += (iv.End - max(iv.Start, from)) * 60 - } - } + left = minutesFrom(max(nowMin, out.StartsAt.In(tz).Hour()*60+out.StartsAt.In(tz).Minute())) * 60 } else { left = 0 } @@ -387,13 +419,21 @@ func (s *schedulerService) PlanCampaignDay(ctx context.Context, campaignID uuid. return nil, err } pass := s.newCampaignPass(ctx, campaign, accounts) - if err := s.prefillSentToday(ctx, pass, accounts); err != nil { + if err := s.prefillPool(ctx, pass, accounts); err != nil { return nil, err } sentBySender, err := s.taskRepo.CountCampaignSendsTodayBySender(ctx, campaignID) if err != nil { return nil, err } + ids := make([]uuid.UUID, 0, len(accounts)) + for _, a := range accounts { + ids = append(ids, a.ID) + } + lastSends, err := s.taskRepo.GetLastEmailTimes(ctx, ids) + if err != nil { + return nil, err + } days := make([]mailboxDay, 0, len(accounts)) for _, acct := range accounts { @@ -401,7 +441,11 @@ func (s *schedulerService) PlanCampaignDay(ctx context.Context, campaignID uuid. if err != nil { return nil, err } - days = append(days, s.planMailbox(ctx, pass, acct, sentBySender[acct.ID], sentAll, now, windowSecondsLeft, closesAt)) + var last *time.Time + if at, ok := lastSends[acct.ID]; ok { + last = &at + } + days = append(days, s.planMailbox(ctx, pass, acct, sentBySender[acct.ID], sentAll, last, now, windowSecondsLeft, closesAt)) } // The pool's waterfall: each clamp's deltas summed over the mailboxes. @@ -512,10 +556,13 @@ func (s *schedulerService) PlanCampaignDay(ctx context.Context, campaignID uuid. } // The leads: mailboxes can only send to a step that is due today, and a - // lead bound to a mailbox with nothing left waits for it. + // lead bound to a mailbox with nothing left waits for it. Only a mailbox + // that is coming back on its own counts (spent, closed, spaced out), as + // in pacedSenders: a lead on one that is failing authentication, resting + // or held by health is moved to another mailbox, not left waiting. unavailable := map[uuid.UUID]bool{} for _, d := range days { - if d.remaining == 0 { + if d.remaining == 0 && (d.gate.paced || d.gate.reason == "") { unavailable[d.acct.ID] = true } } @@ -571,7 +618,7 @@ func (s *schedulerService) PoolCapacityToday(ctx context.Context, campaign *mode // One read for the whole pool. A miss is not an error here: a mailbox // whose sends are unknown counts toward capacity and nothing toward // remaining, which is the conservative side. - _ = s.prefillSentToday(ctx, pass, accounts) + _ = s.prefillPool(ctx, pass, accounts) out := &models.WorkspaceSendCapacity{} for _, acct := range accounts { out.ConfiguredCeiling += acct.CampaignLimit @@ -602,9 +649,11 @@ func (s *schedulerService) PoolCapacityToday(ctx context.Context, campaign *mode return out, nil } -// prefillSentToday loads the pool's sends today in one query into the pass's -// memo, so the per-mailbox reads that follow cost nothing. -func (s *schedulerService) prefillSentToday(ctx context.Context, pass *campaignPass, accounts []models.Email) error { +// prefillPool loads the pool's sends today and warmup health in one query +// each into the pass's memos, so the per-mailbox reads that follow cost +// nothing. An unreadable health batch leaves the memo empty and the +// per-mailbox read decides, as the send path does. +func (s *schedulerService) prefillPool(ctx context.Context, pass *campaignPass, accounts []models.Email) error { ids := make([]uuid.UUID, 0, len(accounts)) for _, a := range accounts { ids = append(ids, a.ID) @@ -616,5 +665,12 @@ func (s *schedulerService) prefillSentToday(ctx context.Context, pass *campaignP for _, id := range ids { pass.sentToday[id] = counts[id] } + if s.warmupRepo != nil { + if states, err := s.warmupRepo.GetHealthStates(ctx, ids); err == nil { + for id, h := range states { + pass.health[id] = healthRead{state: h.State, blockedUntil: h.BlockedUntil, known: true} + } + } + } return nil } diff --git a/internal/scheduler/send_plan_test.go b/internal/scheduler/send_plan_test.go index 2bf5aee3a..fd7274e89 100644 --- a/internal/scheduler/send_plan_test.go +++ b/internal/scheduler/send_plan_test.go @@ -103,6 +103,13 @@ func TestDayWindow(t *testing.T) { t.Fatalf("got %+v secs %d", w, secs) } }) + t.Run("an end date later today ends the day there", func(t *testing.T) { + end := time.Date(2026, 9, 16, 15, 0, 0, 0, loc) + w, secs, closes := dayWindow(&models.Campaign{Timezone: tz, ScheduleWindows: nineToFive, EndDate: &end}, now) + if !w.OpenNow || w.MinutesLeft != 30 || secs != 30*60 || !closes.Equal(end) { + t.Fatalf("got %+v secs %d closes %v", w, secs, closes.In(loc)) + } + }) t.Run("past the end date nothing is left", func(t *testing.T) { end := time.Date(2026, 9, 15, 9, 0, 0, 0, loc) w, secs, _ := dayWindow(&models.Campaign{Timezone: tz, ScheduleWindows: nineToFive, EndDate: &end}, now) diff --git a/web/src/components/app/campaigns/SendPlanCard.tsx b/web/src/components/app/campaigns/SendPlanCard.tsx index f66a9d5f5..dfb5a9aba 100644 --- a/web/src/components/app/campaigns/SendPlanCard.tsx +++ b/web/src/components/app/campaigns/SendPlanCard.tsx @@ -35,7 +35,7 @@ const LIMIT_META: Record = { workspace_risk: { label: "Workspace sending posture", hint: "The workspace is restricted, so every mailbox sends a fraction of its cap.", - to: "/app/settings/deliverability", + to: "/app/deliverability", }, domain_auth: { label: "Domain authentication failing",