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",