feat: make the send plan count a lead bound to a spent mailbox as waiting for it (leads.waiting_on_sender), charge a behaviour profile's spent budget and hourly ceiling to the plan rather than to hours or spacing, fold a foreign-timezone mailbox's 8pm close into its pacing, report the UTC budget day and keep the waterfall adding up when a cap was lowered after sends, read the pool's sends today in one query and cache the plan and workspace capacity for a few seconds, share one cold-ramp notice builder between the drawer and the plan, let the wizard estimate survive a counter miss, and stop a malformed plan payload from taking the campaign overview down

This commit is contained in:
Matthew Meszaros
2026-09-19 09:45:03 -07:00
parent 150dc7df6e
commit 89daae3f29
16 changed files with 337 additions and 97 deletions
@@ -904,7 +904,7 @@ What the campaign sends today and every limit that decided it, worked out on the
### Response
The arithmetic adds up: `configured_ceiling` minus every `limits[].emails` minus `sent_today` is `expected_remaining`, and `projected_today` is `sent_today` plus `expected_remaining`. `limits` lists only the clamps that removed something, in the order the scheduler applies them; `bottleneck` names the one that decides the number (`""` when nothing binds, `budget_spent` when every mailbox has used its day). `day` and the times are in the campaign's timezone. `organization` is present only when the workspace's plan has a daily allowance.
The arithmetic adds up: `configured_ceiling` minus every `limits[].emails` minus `sent_today` is `expected_remaining`, and `projected_today` is `sent_today` plus `expected_remaining`. `limits` lists only the clamps that removed something, in the order the scheduler applies them; `bottleneck` names the one that decides the number (`""` when nothing binds, `budget_spent` when every mailbox has used its day). `day` is the budget day in UTC (every daily counter resets at midnight UTC, whatever the campaign's timezone); the window's times are in the campaign's timezone. `leads.waiting_on_sender` counts due steps whose own mailbox has nothing left today, since each contact keeps the address they first heard from. `organization` is present only when the workspace's plan has a daily allowance.
```json
{
@@ -934,6 +934,7 @@ The arithmetic adds up: `configured_ceiling` minus every `limits[].emails` minus
"waiting_on_step": 200,
"waiting_on_condition": 0,
"held": 3,
"waiting_on_sender": 0,
"new_leads_started_today": 8,
"max_new_leads_per_day": 0
},
+1 -1
View File
@@ -110,7 +110,7 @@ The number is the sends already out plus what is still expected, and the list un
| New leads per day | Only this many contacts may receive their first email today; follow-ups keep going |
| Not enough leads due | The mailboxes could send more, but no more steps are due today |
Under the list, the card says how many leads are due now, due later today, waiting on a step's delay, inside an undecided branch window or held, and when the campaign next works through its queue. Expand the mailbox rows to see each mailbox's cap for this campaign today and which limit set it, what it has sent (for this campaign, and for others), what it is still expected to send, and whether it is sending, closed for the day, out of budget, resting, held or failing authentication. Each limit links to the page where it is changed.
Under the list, the card says how many leads are due now, due later today, waiting on a step's delay, inside an undecided branch window, waiting for their own mailbox (each contact keeps the address they first heard from, so a follow-up whose mailbox is spent waits rather than switching) or held, and when the campaign next works through its queue. Daily budgets reset at midnight UTC whatever the campaign's timezone, so a campaign far from UTC sees its "sent today" reset partway through its own day. Expand the mailbox rows to see each mailbox's cap for this campaign today and which limit set it, what it has sent (for this campaign, and for others), what it is still expected to send, and whether it is sending, closed for the day, out of budget, resting, held or failing authentication. Each limit links to the page where it is changed.
The sidebar's **sent today** meter reads its denominator the same way: what the workspace's mailboxes can send today under these limits, not their caps added up. The two agree by construction, because both are read through the scheduler.
+14 -2
View File
@@ -23163,6 +23163,7 @@
"waiting_on_step",
"waiting_on_condition",
"held",
"waiting_on_sender",
"new_leads_started_today",
"max_new_leads_per_day"
],
@@ -23191,6 +23192,10 @@
"type": "integer",
"description": "Leads paused with no end."
},
"waiting_on_sender": {
"type": "integer",
"description": "Due steps whose own mailbox has nothing left today; each contact keeps the address they first heard from."
},
"new_leads_started_today": {
"type": "integer"
},
@@ -23340,7 +23345,7 @@
},
"day": {
"type": "string",
"description": "The sending day, YYYY-MM-DD in the campaign's timezone."
"description": "The budget day, YYYY-MM-DD in UTC: every daily counter resets at midnight UTC whatever the campaign's timezone. The window's times are in the campaign's timezone."
},
"timezone": {
"type": "string"
@@ -25858,7 +25863,14 @@
},
"folder": {
"type": "string",
"enum": ["inbox", "sent", "drafts", "archive", "spam", "trash"],
"enum": [
"inbox",
"sent",
"drafts",
"archive",
"spam",
"trash"
],
"description": "Canonical folder the message sits in."
},
"parent_id": {
+4 -26
View File
@@ -568,36 +568,14 @@ func (s *analyticsService) coldRampInfo(ctx context.Context, email *models.Email
return nil
}
state, ok := states[email.ID]
if !ok || state.WarmupStartedAt == nil {
if !ok {
return nil
}
now := time.Now()
warmupDays := int(now.Sub(*state.WarmupStartedAt).Hours() / 24)
if warmupDays < 0 {
warmupDays = 0
}
var rampStart time.Time
if state.ColdRampStartedAt != nil {
rampStart = *state.ColdRampStartedAt
}
ceiling := warmupramp.ColdCeiling(warmupDays, rampStart, state.Placements, now, email.CampaignLimit)
if ceiling >= email.CampaignLimit {
info := warmupramp.Notice(state.WarmupStartedAt, state.ColdRampStartedAt, state.Placements, email.CampaignLimit, time.Now())
if info == nil || info.Ceiling >= email.CampaignLimit {
return nil
}
remaining := email.CampaignLimit - ceiling
days := remaining / warmupramp.ColdRampIncrement
if remaining%warmupramp.ColdRampIncrement != 0 {
days++
}
return &models.ColdRampInfo{
Ceiling: ceiling,
MailboxCap: email.CampaignLimit,
DaysToFullCap: days,
Held: warmupramp.ColdHeldUntil(rampStart, state.Placements, now, warmupramp.FreezeWindow) != nil,
}
return info
}
// sendLifecycleInfo reports the mailbox's cold-rotation state, but only when
+17
View File
@@ -1293,11 +1293,20 @@ 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()
if s.planCache != nil {
if plan, ok := s.planCache.get(key); ok {
return plan, nil
}
}
plan, err := planner.PlanCampaignDay(ctx, campaign.ID, s.orgDailyLimit(ctx, orgID))
if err != nil {
errs.CaptureException(err)
return nil, errx.InternalError()
}
if s.planCache != nil {
s.planCache.put(key, plan)
}
return plan, nil
}
@@ -1306,6 +1315,11 @@ func (s *campaignService) WorkspaceCapacity(ctx context.Context, orgID uuid.UUID
if !ok {
return nil, errx.New(errx.Internal, "send planning is not available")
}
if s.capacityCache != nil {
if out, ok := s.capacityCache.get(orgID.String()); ok {
return out, nil
}
}
accounts, xerr := s.emailRepo.GetAllActiveInScope(ctx, repository.NewAccountScope(&orgID))
if xerr != nil {
return nil, xerr
@@ -1322,5 +1336,8 @@ func (s *campaignService) WorkspaceCapacity(ctx context.Context, orgID uuid.UUID
out.Remaining = min(out.Remaining, max(0, limit-sent))
}
}
if s.capacityCache != nil {
s.capacityCache.put(orgID.String(), out)
}
return out, nil
}
+53
View File
@@ -0,0 +1,53 @@
package campaign
import (
"sync"
"time"
)
// readCache holds a derived read for a few seconds. Both the send plan and
// the workspace capacity walk every mailbox (and the plan every lead) through
// the scheduler's gates, and both are polled by every open dashboard tab, so
// two viewers of one campaign share one computation instead of doubling it.
// The clock is the only input a tick would change, and it does not change
// inside the window.
type readCache[T any] struct {
mu sync.Mutex
ttl time.Duration
m map[string]cacheEntry[T]
}
type cacheEntry[T any] struct {
v T
exp time.Time
}
func newReadCache[T any](ttl time.Duration) *readCache[T] {
return &readCache[T]{ttl: ttl, m: map[string]cacheEntry[T]{}}
}
// get returns the cached value while it is fresh.
func (c *readCache[T]) get(key string) (T, bool) {
c.mu.Lock()
defer c.mu.Unlock()
e, ok := c.m[key]
if !ok || time.Now().After(e.exp) {
var zero T
return zero, false
}
return e.v, true
}
func (c *readCache[T]) put(key string, v T) {
c.mu.Lock()
defer c.mu.Unlock()
now := time.Now()
// Sweep expired entries so a long-lived process does not keep every
// campaign it ever planned.
for k, e := range c.m {
if now.After(e.exp) {
delete(c.m, k)
}
}
c.m[key] = cacheEntry[T]{v: v, exp: now.Add(c.ttl)}
}
+6
View File
@@ -120,6 +120,10 @@ type campaignService struct {
// segments counts an audience for Estimate. Optional: without it an
// estimate reports zero recipients.
segments SegmentCounter
// planCache and capacityCache hold the two derived reads for a few
// seconds each; see readCache.
planCache *readCache[*models.CampaignSendPlan]
capacityCache *readCache[*models.WorkspaceSendCapacity]
}
// SegmentCounter is the slice of the segment service Estimate needs.
@@ -180,6 +184,8 @@ func NewService(
streamingPublisher *pubsub.StreamingPublisher,
) CampaignService {
return &campaignService{
planCache: newReadCache[*models.CampaignSendPlan](15 * time.Second),
capacityCache: newReadCache[*models.WorkspaceSendCapacity](60 * time.Second),
campaignRepository: campaignRepository,
taskRepo: taskRepo,
emailRepo: emailRepo,
+36 -1
View File
@@ -1,6 +1,10 @@
package warmupramp
import "time"
import (
"time"
"github.com/warmbly/warmbly/internal/models"
)
const (
// ColdRampIncrement is how much a graduating mailbox may add per clean day.
@@ -44,6 +48,37 @@ func ColdCeiling(warmupDays int, rampStart time.Time, placements []time.Time, no
return ceiling
}
// Notice is the graduation notice the mailbox drawer and a campaign's send
// plan both show: today's ceiling, the configured cap, the clean days left
// before the two meet, and whether a placement is pausing the climb. One
// builder, so the two surfaces cannot round differently. Nil when the
// mailbox never warmed.
func Notice(warmupStartedAt, coldRampStartedAt *time.Time, placements []time.Time, mailboxCap int, now time.Time) *models.ColdRampInfo {
if warmupStartedAt == nil {
return nil
}
warmupDays := int(now.Sub(*warmupStartedAt).Hours() / 24)
if warmupDays < 0 {
warmupDays = 0
}
var rampStart time.Time
if coldRampStartedAt != nil {
rampStart = *coldRampStartedAt
}
ceiling := ColdCeiling(warmupDays, rampStart, placements, now, mailboxCap)
left := max(0, mailboxCap-ceiling)
days := left / ColdRampIncrement
if left%ColdRampIncrement != 0 {
days++
}
return &models.ColdRampInfo{
Ceiling: ceiling,
MailboxCap: mailboxCap,
DaysToFullCap: days,
Held: ColdHeldUntil(rampStart, placements, now, FreezeWindow) != nil,
}
}
// ColdHeldUntil is when the cold ramp resumes climbing, or nil when it already
// is. It filters placements the same way ColdCeiling does, so the number the
// dashboard shows and the reason it gives cannot disagree.
+10 -2
View File
@@ -16,13 +16,17 @@ import (
type CampaignSendPlan struct {
CampaignID uuid.UUID `json:"campaign_id"`
Status string `json:"status"`
// Day is the sending day the plan is for, in the campaign's timezone.
// Day is the budget day the plan counts: every daily counter resets at
// midnight UTC, whatever the campaign's timezone, so this is the UTC
// date. The window's times are in the campaign's own timezone.
Day string `json:"day"`
Timezone string `json:"timezone"`
ComputedAt time.Time `json:"computed_at"`
// ConfiguredCeiling is the sum of the attached mailboxes' own daily caps:
// the number the settings suggest before anything else is applied.
// the number the settings suggest before anything else is applied. A
// mailbox that already sent past a cap lowered during the day counts
// what it sent, so the waterfall still adds up.
ConfiguredCeiling int `json:"configured_ceiling"`
// Projected is today's total: what has gone out plus what is still
// expected to.
@@ -116,6 +120,10 @@ type CampaignLeadSupply struct {
WaitingOnCondition int `json:"waiting_on_condition"`
// Held is the leads paused (out of office, or by hand).
Held int `json:"held"`
// WaitingOnSender is the due steps whose own mailbox has nothing left
// today. Each contact keeps the address they first heard from, so these
// wait for it rather than going out from another mailbox.
WaitingOnSender int `json:"waiting_on_sender"`
// NewLeadsStartedToday and MaxNewLeadsPerDay are the new-lead throttle;
// the cap is 0 when unlimited.
NewLeadsStartedToday int `json:"new_leads_started_today"`
+21 -3
View File
@@ -8,6 +8,8 @@ import (
"github.com/google/uuid"
)
var _ = uuid.Nil
// The day's plan reads the leads through the same routing the send path
// uses, and the per-campaign sender ledger through the same predicate the
// per-mailbox budget uses. Both are SQL, so both are checked live.
@@ -31,7 +33,7 @@ func TestLiveLeadSupplyCountsWhereEveryLeadStands(t *testing.T) {
t.Fatal(err)
}
got, err := repo.LeadSupply(ctx, f.campaign, time.Now().Add(2*time.Hour))
got, err := repo.LeadSupply(ctx, f.campaign, time.Now().Add(2*time.Hour), nil)
if err != nil {
t.Fatal(err)
}
@@ -44,20 +46,36 @@ func TestLiveLeadSupplyCountsWhereEveryLeadStands(t *testing.T) {
if _, err := pool.Exec(ctx, `UPDATE campaigns SET entry_delay_minutes = 60 WHERE id = $1`, f.campaign); err != nil {
t.Fatal(err)
}
got, err = repo.LeadSupply(ctx, f.campaign, time.Now().Add(2*time.Hour))
got, err = repo.LeadSupply(ctx, f.campaign, time.Now().Add(2*time.Hour), nil)
if err != nil {
t.Fatal(err)
}
if got.DueNow != 0 || got.DueLaterToday != 2 || got.DueLaterTodayNewLeads != 2 || got.NextDueAt == nil {
t.Fatalf("with a 60-minute entry delay and a day ending in 2 hours, got %+v, want 2 due later today", got)
}
got, err = repo.LeadSupply(ctx, f.campaign, time.Now().Add(30*time.Minute))
got, err = repo.LeadSupply(ctx, f.campaign, time.Now().Add(30*time.Minute), nil)
if err != nil {
t.Fatal(err)
}
if got.DueNow != 0 || got.DueLaterToday != 0 || got.WaitingOnStep != 2 {
t.Fatalf("with a 60-minute entry delay and a day ending in 30 minutes, got %+v, want 2 waiting on a step", got)
}
// A lead bound to a mailbox that has nothing left today waits for it,
// exactly as routing parks it, rather than counting as a send.
if _, err := pool.Exec(ctx, `UPDATE campaigns SET entry_delay_minutes = 0 WHERE id = $1`, f.campaign); err != nil {
t.Fatal(err)
}
if _, err := pool.Exec(ctx, `UPDATE campaign_leads SET email_account_id = $1 WHERE campaign_id = $2 AND contact_id = $3`, f.mailbox, f.campaign, f.leads[0]); err != nil {
t.Fatal(err)
}
got, err = repo.LeadSupply(ctx, f.campaign, time.Now().Add(2*time.Hour), map[uuid.UUID]bool{f.mailbox: true})
if err != nil {
t.Fatal(err)
}
if got.DueNow != 1 || got.WaitingOnSender != 1 {
t.Fatalf("with lead 0 bound to a spent mailbox, got %+v, want 1 due and 1 waiting on its sender", got)
}
}
func TestLiveCountCampaignSendsTodayBySenderMatchesTheMailboxLedger(t *testing.T) {
+14 -2
View File
@@ -281,7 +281,10 @@ type CampaignProgressRepository interface {
// each one stands, so a day's plan knows how many sends the leads can
// take rather than only how many the mailboxes can give. until is the
// end of the day being planned.
LeadSupply(ctx context.Context, campaignID uuid.UUID, until time.Time) (*LeadSupply, error)
// unavailable is the pool mailboxes that cannot send for the rest of the
// day: a lead bound to one waits for it, as FindRoutedPairs makes it,
// rather than counting as a send another mailbox could make.
LeadSupply(ctx context.Context, campaignID uuid.UUID, until time.Time, unavailable map[uuid.UUID]bool) (*LeadSupply, error)
// CountUndeliverableLeads counts the leads FindRoutedPairs excludes
// because address verification refused them. Reported when a campaign
@@ -2337,13 +2340,16 @@ type LeadSupply struct {
WaitingOnCondition int
// Held is the leads under a live hold with no end.
Held int
// WaitingOnSender is the due email steps whose own mailbox has nothing
// left today; each contact keeps the address they first heard from.
WaitingOnSender int
// NextDueAt is the soonest moment a waiting lead becomes due.
NextDueAt *time.Time
}
// LeadSupply walks every routable lead through the campaign's routing and
// tallies where each one stands relative to now and `until`.
func (r *campaignProgressRepository) LeadSupply(ctx context.Context, campaignID uuid.UUID, until time.Time) (*LeadSupply, error) {
func (r *campaignProgressRepository) LeadSupply(ctx context.Context, campaignID uuid.UUID, until time.Time, unavailable map[uuid.UUID]bool) (*LeadSupply, error) {
out := &LeadSupply{}
router, err := r.loadRouter(ctx, campaignID)
if err != nil {
@@ -2385,10 +2391,16 @@ func (r *campaignProgressRepository) LeadSupply(ctx context.Context, campaignID
out.WaitingOnStep++
continue
}
if in.sender != nil && unavailable[*in.sender] {
out.WaitingOnSender++
continue
}
out.DueLaterToday++
if res.IsNewLead {
out.DueLaterTodayNewLeads++
}
case in.sender != nil && unavailable[*in.sender]:
out.WaitingOnSender++
default:
out.DueNow++
if res.IsNewLead {
+36
View File
@@ -139,6 +139,9 @@ type TaskRepository interface {
// mailbox they went out from. A mailbox's daily budget is shared by every
// campaign it is on, so a plan has to know which campaign spent it.
CountCampaignSendsTodayBySender(ctx context.Context, campaignID uuid.UUID) (map[uuid.UUID]int, error)
// CountCampaignEmailsSentTodayByAccounts is CountCampaignEmailsSentToday
// for a whole pool in one read; an id with no sends is absent.
CountCampaignEmailsSentTodayByAccounts(ctx context.Context, accountIDs []uuid.UUID) (map[uuid.UUID]int, error)
CountWarmupEmailsSentToday(ctx context.Context, accountID uuid.UUID) (int, error)
// Create user-initiated email task (transactional)
@@ -450,6 +453,39 @@ func (r *taskRepository) CountCampaignEmailsSentToday(ctx context.Context, accou
return count, 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) {
out := make(map[uuid.UUID]int, len(accountIDs))
if len(accountIDs) == 0 {
return out, nil
}
query := `
SELECT t.email_account_id, COUNT(*)
FROM tasks t
WHERE t.email_account_id = ANY($1)
AND t.status = 'completed'
AND t.task_type = 'campaign'
AND DATE(t.completed_at) = CURRENT_DATE
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 n int
if err := rows.Scan(&id, &n); err != nil {
return nil, err
}
out[id] = n
}
return out, rows.Err()
}
// CountCampaignSendsTodayBySender is CountCampaignEmailsSentToday for one
// campaign, split by mailbox. Same ledger and the same day boundary, so the
// two agree on what a mailbox has spent.
+85 -58
View File
@@ -9,7 +9,6 @@ import (
"github.com/warmbly/warmbly/internal/app/behavior"
"github.com/warmbly/warmbly/internal/app/warmupramp"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
// CampaignSendPlanner is satisfied by the scheduler service.
@@ -34,24 +33,23 @@ type mailboxDay struct {
health healthRead
gate mailboxGate
configured int
sentThis, sentOther int
remaining int
byCampaignLimit int
byRamp int
byGraduation int
byRisk int
byGate int
byOther int
byHealthPace int
byHours int
byBehavior int
bySpacing int
state string
reopensAt time.Time
minGap int
graduation *models.ColdRampInfo
sendsTodayIfReopened bool
configured int
sentThis, sentOther int
remaining int
byCampaignLimit int
byRamp int
byGraduation int
byRisk int
byGate int
byOther int
byHealthPace int
byHours int
byBehavior int
bySpacing int
state string
reopensAt time.Time
minGap int
graduation *models.ColdRampInfo
}
// stagedCap is explainCap with every intermediate cap kept, so a plan can say
@@ -100,10 +98,12 @@ func room(capv, sentThis int) int {
// 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 {
d := mailboxDay{acct: acct, configured: acct.CampaignLimit, sentThis: sentThis, sentOther: max(0, sentAll-sentThis), minGap: acct.MinWaitTime}
// 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}
stages, limitedBy := stagedCap(pass, acct)
d.cap = capClamp{Cap: stages[4], LimitedBy: limitedBy}
r := room(stages[0], sentThis)
r := room(d.configured, sentThis)
step := func(capv int) int {
next := room(capv, sentThis)
delta := r - next
@@ -114,8 +114,8 @@ func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass,
d.byRamp = step(stages[2])
d.byGraduation = step(stages[3])
d.byRisk = step(stages[4])
if st, ok := pass.coldRamp[acct.ID]; ok && st.WarmupStartedAt != nil && stages[3] < stages[2] {
d.graduation = coldRampInfo(st, stages[2], now)
if st, ok := pass.coldRamp[acct.ID]; ok && stages[3] < stages[2] {
d.graduation = warmupramp.Notice(st.WarmupStartedAt, st.ColdRampStartedAt, st.Placements, stages[2], now)
}
// Standing gates: authentication, cold rotation, warmup health. Asked with
@@ -161,6 +161,17 @@ func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass,
secondsLeft := windowSecondsLeft
bhv := pass.behaviors[acct.ID]
if bhv.Enabled {
// Today's rolled budget first: placeWithinBehavior walks to tomorrow
// when it is spent, and that is the plan binding, not the hours.
today := bhv.PlanOn(behavior.PlanDateFor(now, bhv.Loc))
if today.IsWorkingDay && behavior.MinuteOfDay(now, bhv.Loc) < today.WorkEndMinute {
if s.behaviorDailyCap(ctx, bhv, r, now) == 0 {
d.state = models.MailboxPlanBudgetSpent
d.gate = mailboxGate{reason: gateBudget, paced: true}
d.byBehavior = r
return d
}
}
openAt, ok := s.placeWithinBehavior(ctx, bhv, now)
if !ok {
d.state = models.MailboxPlanNoWorkingDay
@@ -180,7 +191,7 @@ func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass,
secondsLeft = min(secondsLeft, int(windowClosesAt.Sub(openAt).Seconds()))
}
next = min(r, s.behaviorDailyCap(ctx, bhv, r, openAt))
d.byBehavior = r - next
d.byBehavior += r - next
r = next
if r == 0 {
d.state = models.MailboxPlanBudgetSpent
@@ -189,18 +200,23 @@ func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass,
}
d.minGap = s.behaviorGapFloor(bhv, openAt, acct.MinWaitTime)
plan := bhv.PlanOn(behavior.PlanDateFor(openAt, bhv.Loc))
workEnd := time.Date(openAt.In(bhv.Loc).Year(), openAt.In(bhv.Loc).Month(), openAt.In(bhv.Loc).Day(), 0, 0, 0, 0, bhv.Loc).Add(time.Duration(plan.WorkEndMinute) * time.Minute)
local := openAt.In(bhv.Loc)
workEnd := time.Date(local.Year(), local.Month(), local.Day(), 0, 0, 0, 0, bhv.Loc).Add(time.Duration(plan.WorkEndMinute) * time.Minute)
secondsLeft = min(secondsLeft, max(0, int(workEnd.Sub(openAt).Seconds())))
if plan.HourlyLimit > 0 && secondsLeft > 0 {
// The hourly ceiling is the plan's own clamp, so it is charged to it.
if plan.HourlyLimit > 0 {
hours := (secondsLeft + 3599) / 3600
if byHour := plan.HourlyLimit * hours; byHour < r {
// Counted with the spacing below: both are pace, not budget.
secondsLeft = min(secondsLeft, byHour*max(1, d.minGap))
d.byBehavior += r - byHour
r = byHour
}
}
} else if acct.Timezone != "" && acct.Timezone != pass.campaign.Timezone {
// The 8am-8pm band in the mailbox's own timezone, both ends: the
// placer moves any send past 8pm to the next morning.
loc := loadLocation(acct.Timezone)
if h := now.In(loc).Hour(); h < 8 || h >= 20 {
local := now.In(loc)
if h := local.Hour(); h < 8 || h >= 20 {
open := businessHoursReopen(now, loc)
d.reopensAt = open
if !sameLocalDay(open, now, loc) || !open.Before(windowClosesAt) {
@@ -211,6 +227,10 @@ func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass,
}
secondsLeft = min(secondsLeft, int(windowClosesAt.Sub(open).Seconds()))
}
bandEnd := time.Date(local.Year(), local.Month(), local.Day(), 20, 0, 0, 0, loc)
if bandEnd.Before(windowClosesAt) {
secondsLeft = min(secondsLeft, max(0, int(bandEnd.Sub(now).Seconds())))
}
}
if windowSecondsLeft <= 0 {
@@ -255,31 +275,6 @@ func (s *schedulerService) planMailbox(ctx context.Context, pass *campaignPass,
return d
}
// coldRampInfo is the mailbox drawer's graduation notice, computed from the
// pass's own read so the plan and the drawer cannot disagree.
func coldRampInfo(st repository.ColdRampState, mailboxCap int, now time.Time) *models.ColdRampInfo {
warmupDays := int(now.Sub(*st.WarmupStartedAt).Hours() / 24)
if warmupDays < 0 {
warmupDays = 0
}
var rampStart time.Time
if st.ColdRampStartedAt != nil {
rampStart = *st.ColdRampStartedAt
}
ceiling := warmupramp.ColdCeiling(warmupDays, rampStart, st.Placements, now, mailboxCap)
left := mailboxCap - ceiling
days := left / warmupramp.ColdRampIncrement
if left%warmupramp.ColdRampIncrement != 0 {
days++
}
return &models.ColdRampInfo{
Ceiling: ceiling,
MailboxCap: mailboxCap,
DaysToFullCap: days,
Held: warmupramp.ColdHeldUntil(rampStart, st.Placements, now, warmupramp.FreezeWindow) != nil,
}
}
// dayWindow is the campaign's calendar for today, in its own timezone.
func dayWindow(campaign *models.Campaign, now time.Time) (models.CampaignSendWindow, int, time.Time) {
tz := loadLocation(campaign.Timezone)
@@ -379,7 +374,7 @@ func (s *schedulerService) PlanCampaignDay(ctx context.Context, campaignID uuid.
plan := &models.CampaignSendPlan{
CampaignID: campaign.ID,
Status: campaign.Status,
Day: now.In(tz).Format("2006-01-02"),
Day: now.UTC().Format("2006-01-02"),
Timezone: tz.String(),
ComputedAt: now,
Window: window,
@@ -392,6 +387,9 @@ 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 {
return nil, err
}
sentBySender, err := s.taskRepo.CountCampaignSendsTodayBySender(ctx, campaignID)
if err != nil {
return nil, err
@@ -513,8 +511,15 @@ func (s *schedulerService) PlanCampaignDay(ctx context.Context, campaignID uuid.
}
}
// The leads: mailboxes can only send to a step that is due today.
supply, err := s.campaignProgressRepo.LeadSupply(ctx, campaignID, closesAt)
// 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.
unavailable := map[uuid.UUID]bool{}
for _, d := range days {
if d.remaining == 0 {
unavailable[d.acct.ID] = true
}
}
supply, err := s.campaignProgressRepo.LeadSupply(ctx, campaignID, closesAt, unavailable)
if err != nil {
return nil, err
}
@@ -528,6 +533,7 @@ func (s *schedulerService) PlanCampaignDay(ctx context.Context, campaignID uuid.
DueNow: supply.DueNow, DueLaterToday: supply.DueLaterToday,
NewLeadsDueToday: supply.DueNowNewLeads + supply.DueLaterTodayNewLeads,
WaitingOnStep: supply.WaitingOnStep, WaitingOnCondition: supply.WaitingOnCondition, Held: supply.Held,
WaitingOnSender: supply.WaitingOnSender,
NewLeadsStartedToday: newLeadsToday, MaxNewLeadsPerDay: campaign.MaxNewLeadsPerDay, NextDueAt: supply.NextDueAt,
}
followUps := supply.DueNow + supply.DueLaterToday - plan.Leads.NewLeadsDueToday
@@ -562,6 +568,10 @@ func (s *schedulerService) PoolCapacityToday(ctx context.Context, campaign *mode
campaign = &models.Campaign{}
}
pass := s.newCampaignPass(ctx, campaign, accounts)
// 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)
out := &models.WorkspaceSendCapacity{}
for _, acct := range accounts {
out.ConfiguredCeiling += acct.CampaignLimit
@@ -585,9 +595,26 @@ func (s *schedulerService) PoolCapacityToday(ctx context.Context, campaign *mode
out.Capacity += capv
sent, err := s.sentTodayFor(ctx, pass, acct.ID)
if err != nil {
return nil, err
continue
}
out.Remaining += max(0, capv-sent)
}
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 {
ids := make([]uuid.UUID, 0, len(accounts))
for _, a := range accounts {
ids = append(ids, a.ID)
}
counts, err := s.taskRepo.CountCampaignEmailsSentTodayByAccounts(ctx, ids)
if err != nil {
return err
}
for _, id := range ids {
pass.sentToday[id] = counts[id]
}
return nil
}
@@ -63,6 +63,19 @@ function route(url: string): unknown {
if (/^\/campaigns\/[^/?]+$/.test(url)) {
return { id: "camp-1", name: "Probe campaign", status: "draft", description: "" };
}
if (url.endsWith("/send-plan")) {
return {
campaign_id: "camp-1", status: "draft", day: "2026-09-19", timezone: "UTC",
computed_at: new Date().toISOString(), configured_ceiling: 0, projected_today: 0,
sent_today: 0, expected_remaining: 0, bottleneck: "", limits: [],
window: { sending_day: true, open_now: true, minutes_left: 60 },
leads: {
due_now: 0, due_later_today: 0, new_leads_due_today: 0, waiting_on_step: 0,
waiting_on_condition: 0, held: 0, waiting_on_sender: 0, new_leads_started_today: 0, max_new_leads_per_day: 0,
},
mailboxes: [],
};
}
if (url === "/auth/me" || url === "/me") {
return {
id: "u1", email: "d@w.com", first_name: "D", last_name: "W",
@@ -211,6 +211,7 @@ function LeadsLine({ plan }: { plan: SendPlan }) {
if (l.due_later_today) parts.push(`${l.due_later_today.toLocaleString()} later today`);
if (l.waiting_on_step) parts.push(`${l.waiting_on_step.toLocaleString()} waiting on a step`);
if (l.waiting_on_condition) parts.push(`${l.waiting_on_condition.toLocaleString()} in a branch window`);
if (l.waiting_on_sender) parts.push(`${l.waiting_on_sender.toLocaleString()} waiting for their own mailbox`);
if (l.held) parts.push(`${l.held.toLocaleString()} held`);
if (l.max_new_leads_per_day > 0) parts.push(`${l.new_leads_started_today}/${l.max_new_leads_per_day} new leads today`);
return (
@@ -221,12 +222,34 @@ function LeadsLine({ plan }: { plan: SendPlan }) {
);
}
// A payload with a list missing must not take the whole overview down with
// it: the strip is a hint above the page, not the page.
function withDefaults(plan: SendPlan): SendPlan {
return {
...plan,
limits: plan.limits ?? [],
mailboxes: plan.mailboxes ?? [],
window: plan.window ?? { sending_day: false, open_now: false, minutes_left: 0 },
leads: plan.leads ?? {
due_now: 0, due_later_today: 0, new_leads_due_today: 0, waiting_on_step: 0,
waiting_on_condition: 0, held: 0, waiting_on_sender: 0, new_leads_started_today: 0, max_new_leads_per_day: 0,
},
configured_ceiling: plan.configured_ceiling ?? 0,
projected_today: plan.projected_today ?? 0,
sent_today: plan.sent_today ?? 0,
expected_remaining: plan.expected_remaining ?? 0,
bottleneck: plan.bottleneck ?? "",
timezone: plan.timezone || "UTC",
status: plan.status ?? "",
};
}
export default function SendPlanCard({ campaignId }: { campaignId: string }) {
const q = useCampaignSendPlan(campaignId);
// One strip by default; the working is behind a toggle so the analytics
// below it stay above the fold.
const [open, setOpen] = useState(false);
const plan = q.data;
const plan = q.data && typeof q.data === "object" && "campaign_id" in q.data ? withDefaults(q.data) : undefined;
if (q.isPending) {
return (
@@ -56,6 +56,7 @@ export interface LeadSupply {
waiting_on_step: number;
waiting_on_condition: number;
held: number;
waiting_on_sender: number;
new_leads_started_today: number;
max_new_leads_per_day: number;
next_due_at?: string;