mirror of
https://github.com/warmbly/warmbly.git
synced 2026-10-03 08:02:04 +00:00
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
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -35,7 +35,7 @@ const LIMIT_META: Record<SendLimitKind, LimitMeta> = {
|
||||
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",
|
||||
|
||||
Reference in New Issue
Block a user