diff --git a/cmd/cli/specs.go b/cmd/cli/specs.go index f70d003c7..587ed3b9f 100644 --- a/cmd/cli/specs.go +++ b/cmd/cli/specs.go @@ -238,6 +238,14 @@ ask before they do it.`, Method: http.MethodGet, Path: "/campaigns/{id}/ab-analysis", Args: []argSpec{{Name: "id", Help: "The campaign's id"}}, }, + { + Name: "plan", Aliases: []string{"send-plan"}, Short: "Today's sending plan: what goes out and every limit that decided it", + Method: http.MethodGet, Path: "/campaigns/{id}/send-plan", + Args: []argSpec{{Name: "id", Help: "The campaign's id"}}, + Long: `Work out how many emails the campaign sends today through the +scheduler's own gates, and list every limit that took the number below what the +mailbox caps add up to. Nothing is sent or written.`, + }, { Name: "preflight", Short: "Run the pre-send checks without sending", Method: http.MethodPost, Path: "/campaigns/{id}/preflight", Body: bodyOptional, diff --git a/cmd/warmblyctl/api_resources.go b/cmd/warmblyctl/api_resources.go index 162786493..07cdbacb5 100644 --- a/cmd/warmblyctl/api_resources.go +++ b/cmd/warmblyctl/api_resources.go @@ -65,6 +65,7 @@ var apiSpecs = []apiSpec{ {name: "campaign start", summary: "Start the campaign. This sends real mail", method: "POST", path: "/campaigns/{id}/start", sends: true}, {name: "campaign stop", summary: "Stop the campaign", method: "POST", path: "/campaigns/{id}/stop"}, {name: "campaign logs", summary: "The campaign's send log", method: "GET", path: "/campaigns/{id}/logs", query: []string{"limit", "cursor"}}, + {name: "campaign plan", summary: "Today's sending plan and every limit that decided it", method: "GET", path: "/campaigns/{id}/send-plan"}, // Per-lead hold: park ONE contact's flow in THIS campaign without // unsubscribing them or removing them from it. --data carries // {"until": "", "reason": "..."}; no until holds with no end. diff --git a/docs/content/docs/api/cli.mdx b/docs/content/docs/api/cli.mdx index 544291826..99d08345a 100644 --- a/docs/content/docs/api/cli.mdx +++ b/docs/content/docs/api/cli.mdx @@ -238,7 +238,7 @@ Run `warmbly --help` for the flags, and `warmbly Anything above `50`/day per cold mailbox needs positive reputation signals and a low complaint rate behind it. Adding mailboxes is safer than forcing a few to send more. diff --git a/docs/public/openapi.json b/docs/public/openapi.json index cef0f2f79..c49112b34 100644 --- a/docs/public/openapi.json +++ b/docs/public/openapi.json @@ -5459,6 +5459,54 @@ } } }, + "/campaigns/{id}/send-plan": { + "get": { + "operationId": "campaigns_send_plan", + "summary": "Today's sending plan", + "description": "What the campaign sends today and every limit that decided it, worked out through the scheduler's own gates. Nothing is stored or written. Scope READ_CAMPAIGNS, org permission view_campaigns.", + "tags": [ + "campaigns" + ], + "security": [ + { + "bearerAuth": [] + } + ], + "parameters": [ + { + "name": "id", + "in": "path", + "required": true, + "schema": { + "type": "string", + "format": "uuid" + } + } + ], + "responses": { + "200": { + "description": "Today's plan.", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/CampaignSendPlan" + } + } + } + }, + "404": { + "description": "Campaign not found.", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Error" + } + } + } + } + } + } + }, "/campaigns/{id}/senders": { "get": { "operationId": "campaigns_list_senders", @@ -23024,6 +23072,357 @@ } } }, + "CampaignSendLimit": { + "type": "object", + "description": "One clamp in the waterfall and what it took off the ceiling.", + "required": [ + "kind", + "emails" + ], + "properties": { + "kind": { + "type": "string", + "description": "The clamp.", + "enum": [ + "campaign_daily_limit", + "campaign_ramp", + "warmup_graduation", + "workspace_risk", + "domain_auth", + "resting", + "warmup_health_hold", + "other_campaigns", + "warmup_health_pace", + "mailbox_hours", + "sending_behavior", + "spacing", + "sending_window", + "not_running", + "org_daily_limit", + "new_lead_cap", + "leads" + ] + }, + "emails": { + "type": "integer", + "description": "How many of today's sends this clamp removed." + }, + "mailboxes": { + "type": "integer", + "description": "How many mailboxes it touched; absent for a campaign-level clamp." + } + } + }, + "CampaignSendWindow": { + "type": "object", + "required": [ + "sending_day", + "open_now", + "minutes_left" + ], + "properties": { + "sending_day": { + "type": "boolean", + "description": "False when the schedule has no window today." + }, + "open_now": { + "type": "boolean" + }, + "opens_at": { + "type": "string", + "description": "The next opening when closed now; may be on a later day.", + "format": "date-time" + }, + "closes_at": { + "type": "string", + "description": "The end of the last window today.", + "format": "date-time" + }, + "minutes_left": { + "type": "integer", + "description": "Sending time still ahead today." + }, + "starts_at": { + "type": "string", + "description": "The campaign start date when still ahead.", + "format": "date-time" + }, + "ends_at": { + "type": "string", + "description": "The campaign end date, when set.", + "format": "date-time" + } + } + }, + "CampaignLeadSupply": { + "type": "object", + "required": [ + "due_now", + "due_later_today", + "new_leads_due_today", + "waiting_on_step", + "waiting_on_condition", + "held", + "new_leads_started_today", + "max_new_leads_per_day" + ], + "properties": { + "due_now": { + "type": "integer", + "description": "Email steps that could go this minute." + }, + "due_later_today": { + "type": "integer", + "description": "Email steps whose wait elapses before the day ends." + }, + "new_leads_due_today": { + "type": "integer", + "description": "How many of due_now plus due_later_today are first emails." + }, + "waiting_on_step": { + "type": "integer", + "description": "Leads whose next step is due after today." + }, + "waiting_on_condition": { + "type": "integer", + "description": "Leads inside an undecided branch window." + }, + "held": { + "type": "integer", + "description": "Leads paused with no end." + }, + "new_leads_started_today": { + "type": "integer" + }, + "max_new_leads_per_day": { + "type": "integer", + "description": "0 when unlimited." + }, + "next_due_at": { + "type": "string", + "description": "The soonest a waiting lead becomes due.", + "format": "date-time" + } + } + }, + "CampaignMailboxPlan": { + "type": "object", + "required": [ + "id", + "email", + "provider", + "configured_cap", + "cap_today", + "limited_by", + "sent_today", + "sent_by_other_campaigns", + "expected_remaining", + "state", + "min_gap_seconds" + ], + "properties": { + "id": { + "type": "string", + "format": "uuid" + }, + "email": { + "type": "string" + }, + "provider": { + "type": "string" + }, + "configured_cap": { + "type": "integer", + "description": "The mailbox's own daily cold cap." + }, + "cap_today": { + "type": "integer", + "description": "The cap this campaign gives it today, after the cap clamps." + }, + "limited_by": { + "type": "string", + "description": "The clamp that set cap_today.", + "enum": [ + "mailbox_daily_cap", + "campaign_daily_limit", + "campaign_ramp", + "warmup_graduation", + "workspace_risk" + ] + }, + "sent_today": { + "type": "integer", + "description": "This campaign's sends from the mailbox today." + }, + "sent_by_other_campaigns": { + "type": "integer", + "description": "What other campaigns took from the same cap today." + }, + "expected_remaining": { + "type": "integer" + }, + "state": { + "type": "string", + "enum": [ + "sending", + "budget_spent", + "hours_closed", + "no_working_day", + "domain_auth", + "resting", + "health_hold", + "window_closed" + ] + }, + "reopens_at": { + "type": "string", + "description": "When a mailbox outside its own hours is next open.", + "format": "date-time" + }, + "health": { + "type": "string", + "description": "The warmup health band when not healthy." + }, + "min_gap_seconds": { + "type": "integer" + }, + "graduation": { + "$ref": "#/components/schemas/ColdRampInfo" + } + } + }, + "CampaignOrgAllowance": { + "type": "object", + "description": "The workspace's plan-level daily campaign limit.", + "required": [ + "daily_limit", + "sent_today", + "remaining" + ], + "properties": { + "daily_limit": { + "type": "integer" + }, + "sent_today": { + "type": "integer" + }, + "remaining": { + "type": "integer" + } + } + }, + "CampaignSendPlan": { + "type": "object", + "description": "Today's sending plan for a campaign, derived through the scheduler's gates on every read.", + "required": [ + "campaign_id", + "status", + "day", + "timezone", + "computed_at", + "configured_ceiling", + "projected_today", + "sent_today", + "expected_remaining", + "bottleneck", + "limits", + "window", + "leads", + "mailboxes" + ], + "properties": { + "campaign_id": { + "type": "string", + "format": "uuid" + }, + "status": { + "type": "string" + }, + "day": { + "type": "string", + "description": "The sending day, YYYY-MM-DD in the campaign's timezone." + }, + "timezone": { + "type": "string" + }, + "computed_at": { + "type": "string", + "format": "date-time" + }, + "configured_ceiling": { + "type": "integer", + "description": "The attached mailboxes' own daily caps added up." + }, + "projected_today": { + "type": "integer", + "description": "sent_today plus expected_remaining." + }, + "sent_today": { + "type": "integer" + }, + "expected_remaining": { + "type": "integer" + }, + "bottleneck": { + "type": "string", + "description": "The limit that decides projected_today: a limits[].kind, budget_spent, or empty when nothing binds." + }, + "limits": { + "type": "array", + "description": "Only the clamps that removed something, in the order the scheduler applies them. configured_ceiling minus their emails minus sent_today is expected_remaining.", + "items": { + "$ref": "#/components/schemas/CampaignSendLimit" + } + }, + "window": { + "$ref": "#/components/schemas/CampaignSendWindow" + }, + "leads": { + "$ref": "#/components/schemas/CampaignLeadSupply" + }, + "mailboxes": { + "type": "array", + "items": { + "$ref": "#/components/schemas/CampaignMailboxPlan" + } + }, + "organization": { + "$ref": "#/components/schemas/CampaignOrgAllowance" + }, + "next_wake_at": { + "type": "string", + "description": "When the campaign's chain next runs.", + "format": "date-time" + } + } + }, + "ColdRampInfo": { + "type": "object", + "required": [ + "ceiling", + "mailbox_cap", + "days_to_full_cap", + "held" + ], + "properties": { + "ceiling": { + "type": "integer", + "description": "Today's allowance." + }, + "mailbox_cap": { + "type": "integer", + "description": "What the owner configured." + }, + "days_to_full_cap": { + "type": "integer", + "description": "Clean days until the ceiling reaches the cap." + }, + "held": { + "type": "boolean", + "description": "A recent spam placement is pausing the climb." + } + } + }, "CampaignAdvancedSettings": { "type": "object", "required": [ diff --git a/internal/api/handler/analytics.go b/internal/api/handler/analytics.go index a9a47801d..173b8632e 100644 --- a/internal/api/handler/analytics.go +++ b/internal/api/handler/analytics.go @@ -250,6 +250,14 @@ func (h *Handler) GetDashboardAnalytics(c *gin.Context) { errx.Handle(c, xerr) return } + // The sidebar meter's denominator: the mailboxes' day under the + // scheduler's own clamps, not their caps added up. Best effort; the + // dashboard still renders without it. + if h.CampaignService != nil { + if capacity, cerr := h.CampaignService.WorkspaceCapacity(c.Request.Context(), *orgID); cerr == nil { + analytics.CapacityToday = capacity + } + } c.JSON(http.StatusOK, analytics) } diff --git a/internal/api/handler/campaign.go b/internal/api/handler/campaign.go index ec8682eac..fd93ffb33 100644 --- a/internal/api/handler/campaign.go +++ b/internal/api/handler/campaign.go @@ -459,6 +459,27 @@ func (h *Handler) GetCampaignLogs(c *gin.Context) { c.JSON(http.StatusOK, result) } +// GetCampaignSendPlan returns today's sending plan for a campaign: what will go +// out today and every limit that decided it, read through the scheduler's own +// gates. GET /campaigns/:id/send-plan +func (h *Handler) GetCampaignSendPlan(c *gin.Context) { + orgID := middleware.GetOrganizationID(c) + if orgID == nil { + errx.JSON(c, errx.ErrNoOrganization) + return + } + if _, err := uuid.Parse(c.Param("id")); err != nil { + errx.JSON(c, errx.ErrNotFound) + return + } + plan, xerr := h.CampaignService.SendPlan(c.Request.Context(), *orgID, c.Param("id")) + if xerr != nil { + errx.JSON(c, xerr) + return + } + c.JSON(http.StatusOK, plan) +} + // ListCampaignSenders returns a campaign's explicit sender pool. // GET /campaigns/:id/senders func (h *Handler) ListCampaignSenders(c *gin.Context) { diff --git a/internal/api/routes.go b/internal/api/routes.go index 8c20e7914..299a6c5d0 100644 --- a/internal/api/routes.go +++ b/internal/api/routes.go @@ -579,6 +579,9 @@ func Run( campaigns.POST("/:id/start", m.RequireOrganization(), m.RequireAccess(models.PermSendCampaigns, models.APIPermSendCampaigns), h.StartCampaign) campaigns.POST("/:id/stop", m.RequireOrganization(), m.RequireAccess(models.PermSendCampaigns, models.APIPermSendCampaigns), h.StopCampaign) campaigns.GET("/:id/logs", m.RequireAccess(models.PermViewCampaigns, models.APIPermReadCampaigns), h.GetCampaignLogs) + // Today's sending plan: derived through the scheduler's gates on + // every read, never stored. + campaigns.GET("/:id/send-plan", m.RateLimitMiddleware(models.RateLimitRead), m.RequireOrganization(), m.RequireAccess(models.PermViewCampaigns, models.APIPermReadCampaigns), h.GetCampaignSendPlan) // Form performance for this campaign's recipients. campaigns.GET("/:id/forms", m.RequireOrganization(), m.RequireAccess(models.PermViewCampaigns, models.APIPermReadCampaigns), h.GetCampaignForms) diff --git a/internal/app/campaign/handlers.go b/internal/app/campaign/handlers.go index e23048585..5981064bf 100644 --- a/internal/app/campaign/handlers.go +++ b/internal/app/campaign/handlers.go @@ -1096,20 +1096,38 @@ func (s *campaignService) Estimate(ctx context.Context, orgID uuid.UUID, in *mod return nil, xerr } out.Mailboxes = len(accounts) - for _, acct := range accounts { - lim := min(acct.CampaignLimit, dailyLimit) - if lim < 0 { - lim = 0 - } - out.DailyCapacity += lim - sent, err := s.taskRepo.CountCampaignEmailsSentToday(ctx, acct.ID) + // The pool's day under the same clamps the scheduler applies (the + // graduation ceiling, the workspace's risk band, a warmup health hold, + // domain authentication, cold rotation), so the wizard promises what the + // send path will honour rather than the caps added up. + if planner, ok := s.planner(); ok { + capacity, err := planner.PoolCapacityToday(ctx, &models.Campaign{OrganizationID: &orgID, DailyLimit: dailyLimit}, accounts) if err != nil { - // A counter blip must not blank the whole estimate, but it must - // not flatter it either: a mailbox whose sends today are unknown - // contributes nothing to today and only counts from tomorrow. - continue + errs.CaptureException(err) + return nil, errx.InternalError() + } + out.DailyCapacity, out.RemainingToday = capacity.Capacity, capacity.Remaining + } else { + for _, acct := range accounts { + lim := max(0, min(acct.CampaignLimit, dailyLimit)) + out.DailyCapacity += lim + sent, err := s.taskRepo.CountCampaignEmailsSentToday(ctx, acct.ID) + if err != nil { + // A counter blip must not blank the whole estimate, but it must + // not flatter it either: a mailbox whose sends today are unknown + // contributes nothing to today and only counts from tomorrow. + continue + } + out.RemainingToday += max(0, lim-sent) + } + } + if limit := s.orgDailyLimit(ctx, orgID); limit >= 0 { + out.DailyCapacity = min(out.DailyCapacity, limit) + if s.campaignProgressRepo != nil { + if sent, err := s.campaignProgressRepo.CountEmailsSentTodayByOrganization(ctx, orgID); err == nil { + out.RemainingToday = min(out.RemainingToday, max(0, limit-sent)) + } } - out.RemainingToday += max(0, lim-sent) } if out.Recipients == 0 || out.DailyCapacity == 0 { return out, nil @@ -1245,3 +1263,64 @@ func (s *campaignService) ownedCampaign(ctx context.Context, orgID, campaignID u _, _, xerr := s.campaignForOrg(ctx, orgID, campaignID.String()) return xerr } + +// planner is the scheduler's send-plan face, when the wired scheduler has one. +func (s *campaignService) planner() (scheduler.CampaignSendPlanner, bool) { + p, ok := s.scheduler.(scheduler.CampaignSendPlanner) + return p, ok && p != nil +} + +// orgDailyLimit is the workspace's plan-level daily campaign limit, negative +// when unlimited or unknown. A gate that cannot be read clamps nothing here; +// the send path asks again for itself. +func (s *campaignService) orgDailyLimit(ctx context.Context, orgID uuid.UUID) int { + if s.featureGate == nil { + return -1 + } + limit, xerr := s.featureGate.GetDailyEmailLimit(ctx, orgID) + if xerr != nil { + return -1 + } + return limit +} + +func (s *campaignService) SendPlan(ctx context.Context, orgID uuid.UUID, campaignID string) (*models.CampaignSendPlan, *errx.Error) { + campaign, xerr := s.Get(ctx, orgID.String(), campaignID) + if xerr != nil { + return nil, xerr + } + planner, ok := s.planner() + if !ok { + return nil, errx.New(errx.Internal, "send planning is not available") + } + plan, err := planner.PlanCampaignDay(ctx, campaign.ID, s.orgDailyLimit(ctx, orgID)) + if err != nil { + errs.CaptureException(err) + return nil, errx.InternalError() + } + return plan, nil +} + +func (s *campaignService) WorkspaceCapacity(ctx context.Context, orgID uuid.UUID) (*models.WorkspaceSendCapacity, *errx.Error) { + planner, ok := s.planner() + if !ok { + return nil, errx.New(errx.Internal, "send planning is not available") + } + accounts, xerr := s.emailRepo.GetAllActiveInScope(ctx, repository.NewAccountScope(&orgID)) + if xerr != nil { + return nil, xerr + } + out, err := planner.PoolCapacityToday(ctx, &models.Campaign{OrganizationID: &orgID}, accounts) + if err != nil { + errs.CaptureException(err) + return nil, errx.InternalError() + } + if limit := s.orgDailyLimit(ctx, orgID); limit >= 0 && s.campaignProgressRepo != nil { + // The plan's daily allowance caps the workspace as a whole. + if sent, err := s.campaignProgressRepo.CountEmailsSentTodayByOrganization(ctx, orgID); err == nil { + out.Capacity = min(out.Capacity, limit) + out.Remaining = min(out.Remaining, max(0, limit-sent)) + } + } + return out, nil +} diff --git a/internal/app/campaign/service.go b/internal/app/campaign/service.go index 35392da0e..1055d3489 100644 --- a/internal/app/campaign/service.go +++ b/internal/app/campaign/service.go @@ -28,6 +28,13 @@ type CampaignService interface { // many sending days a mailbox pool needs under the per-mailbox caps. // Read-only; the wizard shows it before a one-time email is created. Estimate(ctx context.Context, orgID uuid.UUID, in *models.CampaignEstimate) (*models.CampaignEstimateResult, *errx.Error) + // SendPlan is today's sending plan for one of orgID's campaigns: what + // will go out today and every limit that decided it, derived through the + // scheduler's own gates. Read-only. + SendPlan(ctx context.Context, orgID uuid.UUID, campaignID string) (*models.CampaignSendPlan, *errx.Error) + // WorkspaceCapacity is what the workspace's active mailboxes can send + // today between them, under the same clamps. Read-only. + WorkspaceCapacity(ctx context.Context, orgID uuid.UUID) (*models.WorkspaceSendCapacity, *errx.Error) Update(ctx context.Context, orgID, id string, data *models.UpdateCampaign) (*models.Campaign, *errx.Error) // Delete removes an organization's campaign outright. A running campaign // is stopped as part of it: its pending tasks are cancelled in the same diff --git a/internal/models/analytics.go b/internal/models/analytics.go index 444a1fa25..724b689e1 100644 --- a/internal/models/analytics.go +++ b/internal/models/analytics.go @@ -291,6 +291,10 @@ type DashboardAnalytics struct { TopCampaigns []TopCampaignStats `json:"top_campaigns"` AccountHealth AccountHealthSummary `json:"account_health"` DailyTrend []DashboardDailyStats `json:"daily_trend"` + // CapacityToday is what the workspace's mailboxes can send today under + // the scheduler's clamps; the sidebar meter's denominator. Absent when + // it could not be computed. + CapacityToday *WorkspaceSendCapacity `json:"capacity_today,omitempty"` } // DashboardOverallStats contains aggregate statistics for the dashboard diff --git a/internal/models/campaign_send_plan.go b/internal/models/campaign_send_plan.go new file mode 100644 index 000000000..6546e2a10 --- /dev/null +++ b/internal/models/campaign_send_plan.go @@ -0,0 +1,205 @@ +package models + +import ( + "time" + + "github.com/google/uuid" +) + +// CampaignSendPlan is what one campaign will send today and why that number is +// what it is. It is derived on read through the scheduler's own gates and never +// stored: the same clamps that decide a real send decide these figures, so the +// plan cannot promise a volume the send path would refuse. +// +// The arithmetic is a waterfall that adds up: ConfiguredCeiling minus every +// Limits entry minus SentToday equals ExpectedRemaining. +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 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. + ConfiguredCeiling int `json:"configured_ceiling"` + // Projected is today's total: what has gone out plus what is still + // expected to. + Projected int `json:"projected_today"` + // SentToday is this campaign's sends so far today. + SentToday int `json:"sent_today"` + // ExpectedRemaining is what the pool can still send today for this + // campaign after every limit, and after the leads that are actually due. + ExpectedRemaining int `json:"expected_remaining"` + // Bottleneck names the limit that decides Projected. One of the + // SendLimit* constants, or "" when nothing binds below the ceiling. + Bottleneck string `json:"bottleneck"` + + Limits []CampaignSendLimit `json:"limits"` + Window CampaignSendWindow `json:"window"` + Leads CampaignLeadSupply `json:"leads"` + Mailboxes []CampaignMailboxPlan `json:"mailboxes"` + // Organization is the workspace's plan-level daily allowance. Absent when + // the plan is unlimited. + Organization *CampaignOrgAllowance `json:"organization,omitempty"` + // NextWakeAt is when the campaign's chain next runs. Nil when it has no + // pending wakeup (not running, or being re-seeded). + NextWakeAt *time.Time `json:"next_wake_at,omitempty"` +} + +// CampaignSendLimit is one clamp in the waterfall and how many of today's +// emails it took off the ceiling. Only clamps that removed something are +// listed, in the order the scheduler applies them. +type CampaignSendLimit struct { + Kind string `json:"kind"` + // Emails is how many sends this clamp removed from today's total. + Emails int `json:"emails"` + // Mailboxes is how many of the pool's mailboxes it touched. Zero for a + // campaign-level clamp. + Mailboxes int `json:"mailboxes,omitempty"` +} + +// Limit kinds, in waterfall order. The first five are the cap clamps the +// activity feed already names; keep the strings stable, the dashboard and the +// docs key on them. +const ( + SendLimitCampaignDailyLimit = "campaign_daily_limit" + SendLimitCampaignRamp = "campaign_ramp" + SendLimitWarmupGraduation = "warmup_graduation" + SendLimitWorkspaceRisk = "workspace_risk" + SendLimitDomainAuth = "domain_auth" + SendLimitResting = "resting" + SendLimitHealthHold = "warmup_health_hold" + SendLimitOtherCampaigns = "other_campaigns" + SendLimitHealthPace = "warmup_health_pace" + SendLimitMailboxHours = "mailbox_hours" + SendLimitSpacing = "spacing" + SendLimitSendingWindow = "sending_window" + SendLimitNotRunning = "not_running" + SendLimitOrgDailyLimit = "org_daily_limit" + SendLimitNewLeadCap = "new_lead_cap" + SendLimitLeads = "leads" +) + +// CampaignSendWindow is the campaign's calendar for today. +type CampaignSendWindow struct { + // SendingDay is false when the schedule has no window today. + SendingDay bool `json:"sending_day"` + OpenNow bool `json:"open_now"` + // OpensAt is the next opening when the window is closed now; it may be on + // a later day. + OpensAt *time.Time `json:"opens_at,omitempty"` + // ClosesAt is the end of the last window today, when there is one. + ClosesAt *time.Time `json:"closes_at,omitempty"` + // MinutesLeft is the sending time still ahead today. + MinutesLeft int `json:"minutes_left"` + // StartsAt is the campaign start date when it is still ahead. + StartsAt *time.Time `json:"starts_at,omitempty"` + // EndsAt is the campaign end date, when one is set. + EndsAt *time.Time `json:"ends_at,omitempty"` +} + +// CampaignLeadSupply is the other half of the number: mailboxes can only send +// to leads whose step is due. Counted through the campaign's own routing. +type CampaignLeadSupply struct { + // DueNow is the email steps that could go this minute. + DueNow int `json:"due_now"` + // DueLaterToday is the email steps whose wait elapses before the day ends. + DueLaterToday int `json:"due_later_today"` + // NewLeadsDueToday is how many of DueNow plus DueLaterToday are first + // emails, which the new-lead cap governs. + NewLeadsDueToday int `json:"new_leads_due_today"` + // WaitingOnStep is the leads whose next step is due after today. + WaitingOnStep int `json:"waiting_on_step"` + // WaitingOnCondition is the leads inside an undecided branch window. + WaitingOnCondition int `json:"waiting_on_condition"` + // Held is the leads paused (out of office, or by hand). + Held int `json:"held"` + // NewLeadsStartedToday and MaxNewLeadsPerDay are the new-lead throttle; + // the cap is 0 when unlimited. + NewLeadsStartedToday int `json:"new_leads_started_today"` + MaxNewLeadsPerDay int `json:"max_new_leads_per_day"` + // NextDueAt is the soonest moment a waiting lead becomes due, when + // nothing is due right now. + NextDueAt *time.Time `json:"next_due_at,omitempty"` +} + +// CampaignMailboxPlan is one mailbox's day on this campaign. +type CampaignMailboxPlan struct { + ID uuid.UUID `json:"id"` + Email string `json:"email"` + Provider string `json:"provider"` + // ConfiguredCap is the mailbox's own daily cold cap. + ConfiguredCap int `json:"configured_cap"` + // CapToday is the cap this campaign gives it today, after the cap clamps. + CapToday int `json:"cap_today"` + // LimitedBy names the clamp that set CapToday: mailbox_daily_cap, + // campaign_daily_limit, campaign_ramp, warmup_graduation or + // workspace_risk. + LimitedBy string `json:"limited_by"` + // SentToday is this campaign's sends from the mailbox today; + // SentByOtherCampaigns is what other campaigns took from the same cap. + SentToday int `json:"sent_today"` + SentByOtherCampaigns int `json:"sent_by_other_campaigns"` + // ExpectedRemaining is what this mailbox is expected to still send today + // for this campaign. + ExpectedRemaining int `json:"expected_remaining"` + // State is why the mailbox is or is not sending right now. One of + // "sending", "budget_spent", "hours_closed", "no_working_day", + // "domain_auth", "resting", "health_hold". + State string `json:"state"` + // ReopensAt is when a mailbox outside its own hours is next open. + ReopensAt *time.Time `json:"reopens_at,omitempty"` + // Health is the warmup health band when it is not healthy. + Health string `json:"health,omitempty"` + // MinGapSeconds is the spacing between two of its sends, which warmup + // mail shares. + MinGapSeconds int `json:"min_gap_seconds"` + // Graduation is set while the warmup graduation ceiling holds the mailbox + // below its own cap. + Graduation *ColdRampInfo `json:"graduation,omitempty"` +} + +// Mailbox states in a send plan. +const ( + MailboxPlanSending = "sending" + MailboxPlanBudgetSpent = "budget_spent" + MailboxPlanHoursClosed = "hours_closed" + MailboxPlanNoWorkingDay = "no_working_day" + MailboxPlanDomainAuth = "domain_auth" + MailboxPlanResting = "resting" + MailboxPlanHealthHold = "health_hold" + MailboxPlanWindowClosed = "window_closed" +) + +// SendBottleneckBudgetSpent is the Bottleneck of a campaign that has sent +// everything its mailboxes had today: no limit binds, the day is simply used. +const SendBottleneckBudgetSpent = "budget_spent" + +// SendLimitSendingBehavior is a mailbox's rolled sending-behaviour plan +// lowering its day below the cold cap. +const SendLimitSendingBehavior = "sending_behavior" + +// CampaignOrgAllowance is the workspace's plan-level daily campaign limit. +type CampaignOrgAllowance struct { + DailyLimit int `json:"daily_limit"` + SentToday int `json:"sent_today"` + Remaining int `json:"remaining"` +} + +// WorkspaceSendCapacity is what a workspace's mailboxes can send today between +// them, under the same clamps a campaign's plan applies. The dashboard's +// "sent today" meter reads its denominator from here. +type WorkspaceSendCapacity struct { + // Capacity is today's cold sends across every mailbox that can send; + // Remaining is what is left of it after what has already gone out. + Capacity int `json:"capacity"` + Remaining int `json:"remaining_today"` + // ConfiguredCeiling is the same mailboxes' own caps added up. + ConfiguredCeiling int `json:"configured_ceiling"` + // Mailboxes is how many mailboxes contribute; Held is how many are + // attached but cannot send today. + Mailboxes int `json:"mailboxes"` + Held int `json:"held"` +} diff --git a/internal/repository/lead_supply_live_test.go b/internal/repository/lead_supply_live_test.go new file mode 100644 index 000000000..ff1639449 --- /dev/null +++ b/internal/repository/lead_supply_live_test.go @@ -0,0 +1,99 @@ +package repository + +import ( + "context" + "testing" + "time" + + "github.com/google/uuid" +) + +// 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. +// +// WARMBLY_TEST_DB=postgres://warmbly:warmbly@localhost:15432/warmbly_dev?sslmode=disable \ +// go test ./internal/repository/ -run LiveLeadSupply -v + +func TestLiveLeadSupplyCountsWhereEveryLeadStands(t *testing.T) { + _, pool := liveContactDB(t) + f := newRoutedPairsFixture(t, pool, 4) + ctx := context.Background() + repo := NewCampaignProgressRepository(pool) + + // Lead 1 has had the only step: their flow is over. Lead 2 is held with + // no end. Leads 0 and 3 are due now. + if _, err := pool.Exec(ctx, `INSERT INTO campaign_contact_progress (campaign_id, contact_id, sequence_id, sent_at) + VALUES ($1, $2, $3, NOW() - interval '1 hour')`, f.campaign, f.leads[1], f.step); err != nil { + t.Fatal(err) + } + if _, err := repo.HoldLead(ctx, f.campaign, f.leads[2], nil, "away", "manual"); err != nil { + t.Fatal(err) + } + + got, err := repo.LeadSupply(ctx, f.campaign, time.Now().Add(2*time.Hour)) + if err != nil { + t.Fatal(err) + } + if got.DueNow != 2 || got.DueNowNewLeads != 2 || got.Held != 1 || got.DueLaterToday != 0 || got.WaitingOnStep != 0 { + t.Fatalf("got %+v, want 2 due now (both new), 1 held", got) + } + + // An entry delay moves every new lead past "now"; whether that is later + // today or after today depends on where the day ends. + 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)) + 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)) + 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) + } +} + +func TestLiveCountCampaignSendsTodayBySenderMatchesTheMailboxLedger(t *testing.T) { + _, pool := liveContactDB(t) + f := newRoutedPairsFixture(t, pool, 1) + ctx := context.Background() + tasks := NewTaskRepository(pool) + + sent, wakeup, warmup := uuid.New(), uuid.New(), uuid.New() + exec := func(sql string, args ...any) { + t.Helper() + if _, err := pool.Exec(ctx, sql, args...); err != nil { + t.Fatalf("%q: %v", sql[:min(60, len(sql))], err) + } + } + t.Cleanup(func() { + c := context.Background() + _, _ = pool.Exec(c, `DELETE FROM campaign_tasks WHERE task_id = ANY($1)`, []uuid.UUID{sent, wakeup, warmup}) + _, _ = pool.Exec(c, `DELETE FROM tasks WHERE id = ANY($1)`, []uuid.UUID{sent, wakeup, warmup}) + }) + // A real send, a bare wake-up of the same chain, and a warmup send. + exec(`INSERT INTO tasks (id, task_type, email_account_id, status, message_id, completed_at) VALUES ($1, 'campaign', $2, 'completed', 'm-1', NOW())`, sent, f.mailbox) + exec(`INSERT INTO campaign_tasks (task_id, campaign_id) VALUES ($1, $2)`, sent, f.campaign) + exec(`INSERT INTO tasks (id, task_type, email_account_id, status, message_id, completed_at) VALUES ($1, 'campaign', $2, 'completed', '', NOW())`, wakeup, f.mailbox) + exec(`INSERT INTO campaign_tasks (task_id, campaign_id) VALUES ($1, $2)`, wakeup, f.campaign) + exec(`INSERT INTO tasks (id, task_type, email_account_id, status, message_id, completed_at) VALUES ($1, 'warmup', $2, 'completed', 'm-2', NOW())`, warmup, f.mailbox) + + bySender, err := tasks.CountCampaignSendsTodayBySender(ctx, f.campaign) + if err != nil { + t.Fatal(err) + } + all, err := tasks.CountCampaignEmailsSentToday(ctx, f.mailbox) + if err != nil { + t.Fatal(err) + } + if bySender[f.mailbox] != 1 || all != 1 { + t.Fatalf("campaign ledger %v, mailbox ledger %d: want exactly the one real send in both", bySender, all) + } +} diff --git a/internal/repository/pg_campaign_progress.go b/internal/repository/pg_campaign_progress.go index 4ef0ea3fd..38d94414a 100644 --- a/internal/repository/pg_campaign_progress.go +++ b/internal/repository/pg_campaign_progress.go @@ -277,6 +277,11 @@ type CampaignProgressRepository interface { // their flow goes next, plus the pre-send gate that excludes them, so a // per-contact preview reads the facts the send path reads. RouteContact(ctx context.Context, campaignID, contactID uuid.UUID) (*ContactRoute, error) + // LeadSupply runs the same routing over every lead and counts where + // 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) // CountUndeliverableLeads counts the leads FindRoutedPairs excludes // because address verification refused them. Reported when a campaign @@ -1343,56 +1348,7 @@ func (r *campaignProgressRepository) FindRoutedPairs(ctx context.Context, campai orderPrefix = "(lp.sequence_id IS NULL) DESC, " } - query := ` - SELECT cl.contact_id, cl.added_at, cl.email_account_id, - lp.sequence_id, lp.sent_at, lp.opened_at, lp.clicked_at, lp.replied_at, COALESCE(lp.reply_class, ''), COALESCE(lp.ai_label, ''), - COALESCE(ss.ids, '{}') AS sent_ids, - EXISTS ( - SELECT 1 FROM campaign_contact_progress rp - WHERE rp.campaign_id = $1 AND rp.contact_id = cl.contact_id AND rp.replied_at IS NOT NULL - ) AS has_replied, - ` + leadHoldColumns + ` - FROM campaign_leads cl - JOIN contacts c ON c.id = cl.contact_id - LEFT JOIN LATERAL ( - SELECT sequence_id, sent_at, - CASE WHEN p.opened_machine THEN NULL ELSE p.opened_at END AS opened_at, - clicked_at, replied_at, reply_class, ai_label - FROM campaign_contact_progress p - WHERE p.campaign_id = $1 AND p.contact_id = cl.contact_id AND p.sent_at IS NOT NULL - ORDER BY p.sent_at DESC LIMIT 1 - ) lp ON true - LEFT JOIN LATERAL ( - SELECT array_agg(sequence_id) AS ids - FROM campaign_contact_progress p2 - WHERE p2.campaign_id = $1 AND p2.contact_id = cl.contact_id - AND (p2.sent_at IS NOT NULL OR p2.dispatched_at IS NOT NULL) - ) ss ON true - WHERE cl.campaign_id = $1 - AND NOT EXISTS ( - SELECT 1 FROM campaign_contact_progress b - WHERE b.contact_id = cl.contact_id AND b.bounced_at IS NOT NULL - ) - AND NOT EXISTS ( - SELECT 1 FROM campaign_contact_progress f - WHERE f.campaign_id = $1 AND f.contact_id = cl.contact_id - AND f.sent_at IS NULL AND f.failed_at IS NOT NULL - AND f.send_attempts >= $2 - ) - -- The workspace suppression list (addresses and domains) and the - -- contact's own subscription flag are both send gates; the audience - -- count applies the same two, so the number shown is the number sent. - AND NOT recipient_suppressed((SELECT organization_id FROM campaigns WHERE id = $1), c.email) - AND c.subscribed IS NOT FALSE - -- Addresses the pre-send gates in the campaign task would refuse. - -- Without this the finder keeps handing back the same undeliverable - -- contact, the task skips it, and the campaign never reaches the - -- healthy leads behind it (issue #200). Read from the contact's - -- CURRENT verification state, so re-verifying an address puts it - -- straight back into routing. - AND NOT ` + undeliverableClause("$1") + ` - ORDER BY ` + orderPrefix + contactOrder + ` ` + dir + ` - ` + query := routedLeadsQuery(orderPrefix + contactOrder + ` ` + dir) rows, err := r.db.Query(ctx, query, args...) if err != nil { @@ -2304,3 +2260,141 @@ func (r *campaignProgressRepository) CountHeldLeads(ctx context.Context, campaig campaignID).Scan(&n) return n, err } + +// routedLeadsQuery is the lead scan FindRoutedPairs and LeadSupply share: the +// campaign's leads with their last-sent step and engagement, minus every lead +// a pre-send gate would refuse. Both walk it through the same router, so a +// count can never include a lead the send path would not offer. order is the +// ORDER BY expression; $1 is the campaign and $2 the send-attempt ceiling. +func routedLeadsQuery(order string) string { + return ` + SELECT cl.contact_id, cl.added_at, cl.email_account_id, + lp.sequence_id, lp.sent_at, lp.opened_at, lp.clicked_at, lp.replied_at, COALESCE(lp.reply_class, ''), COALESCE(lp.ai_label, ''), + COALESCE(ss.ids, '{}') AS sent_ids, + EXISTS ( + SELECT 1 FROM campaign_contact_progress rp + WHERE rp.campaign_id = $1 AND rp.contact_id = cl.contact_id AND rp.replied_at IS NOT NULL + ) AS has_replied, + ` + leadHoldColumns + ` + FROM campaign_leads cl + JOIN contacts c ON c.id = cl.contact_id + LEFT JOIN LATERAL ( + SELECT sequence_id, sent_at, + CASE WHEN p.opened_machine THEN NULL ELSE p.opened_at END AS opened_at, + clicked_at, replied_at, reply_class, ai_label + FROM campaign_contact_progress p + WHERE p.campaign_id = $1 AND p.contact_id = cl.contact_id AND p.sent_at IS NOT NULL + ORDER BY p.sent_at DESC LIMIT 1 + ) lp ON true + LEFT JOIN LATERAL ( + SELECT array_agg(sequence_id) AS ids + FROM campaign_contact_progress p2 + WHERE p2.campaign_id = $1 AND p2.contact_id = cl.contact_id + AND (p2.sent_at IS NOT NULL OR p2.dispatched_at IS NOT NULL) + ) ss ON true + WHERE cl.campaign_id = $1 + AND NOT EXISTS ( + SELECT 1 FROM campaign_contact_progress b + WHERE b.contact_id = cl.contact_id AND b.bounced_at IS NOT NULL + ) + AND NOT EXISTS ( + SELECT 1 FROM campaign_contact_progress f + WHERE f.campaign_id = $1 AND f.contact_id = cl.contact_id + AND f.sent_at IS NULL AND f.failed_at IS NOT NULL + AND f.send_attempts >= $2 + ) + -- The workspace suppression list (addresses and domains) and the + -- contact's own subscription flag are both send gates; the audience + -- count applies the same two, so the number shown is the number sent. + AND NOT recipient_suppressed((SELECT organization_id FROM campaigns WHERE id = $1), c.email) + AND c.subscribed IS NOT FALSE + -- Addresses the pre-send gates in the campaign task would refuse. + -- Without this the finder keeps handing back the same undeliverable + -- contact, the task skips it, and the campaign never reaches the + -- healthy leads behind it (issue #200). Read from the contact's + -- CURRENT verification state, so re-verifying an address puts it + -- straight back into routing. + AND NOT ` + undeliverableClause("$1") + ` + ORDER BY ` + order + ` + ` +} + +// LeadSupply is where a campaign's routable leads stand, counted for a day's +// plan. Only email steps count as sends; an action or wait node due now is +// executed without a mailbox and is not counted anywhere here. +type LeadSupply struct { + // DueNow is the email steps that could go this minute, and + // DueNowNewLeads how many of them are first emails. + DueNow int + DueNowNewLeads int + // DueLaterToday is the email steps whose wait elapses before `until`, + // and DueLaterTodayNewLeads how many of those are first emails. + DueLaterTodayNewLeads int + DueLaterToday int + // WaitingOnStep is the leads whose next step is due after `until`. + WaitingOnStep int + // WaitingOnCondition is the leads inside an undecided branch window. + WaitingOnCondition int + // Held is the leads under a live hold with no end. + Held 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) { + out := &LeadSupply{} + router, err := r.loadRouter(ctx, campaignID) + if err != nil { + return nil, err + } + if router == nil { + return out, nil + } + rows, err := r.db.Query(ctx, routedLeadsQuery("c.created_at ASC"), campaignID, config.CampaignSendMaxAttempts) + if err != nil { + return nil, err + } + defer rows.Close() + noteDue := func(at time.Time) { + if out.NextDueAt == nil || at.Before(*out.NextDueAt) { + t := at + out.NextDueAt = &t + } + } + for rows.Next() { + var in routeInput + var contactID uuid.UUID + if serr := rows.Scan(&contactID, &in.addedAt, &in.sender, &in.lastSeq, &in.sentAt, &in.openedAt, &in.clickedAt, &in.repliedAt, &in.replyClass, &in.aiLabel, &in.sentIDs, &in.hasReplied, + &in.pausedAt, &in.pausedUntil, &in.pauseReason, &in.pauseSource); serr != nil { + return nil, serr + } + res := router.route(campaignID, contactID, in) + switch { + case res.Hold != nil && res.Hold.Until == nil: + out.Held++ + case res.WaitUntil != nil: + out.WaitingOnCondition++ + noteDue(*res.WaitUntil) + case res.Target == nil || !router.isEmailStep(*res.Target): + // Finished, or a node that sends nothing. + case res.DueAt != nil && res.DueAt.After(router.dueBy): + noteDue(*res.DueAt) + if res.DueAt.After(until) { + out.WaitingOnStep++ + continue + } + out.DueLaterToday++ + if res.IsNewLead { + out.DueLaterTodayNewLeads++ + } + default: + out.DueNow++ + if res.IsNewLead { + out.DueNowNewLeads++ + } + } + } + return out, rows.Err() +} diff --git a/internal/repository/pg_task.go b/internal/repository/pg_task.go index 55e70d62f..373f79903 100644 --- a/internal/repository/pg_task.go +++ b/internal/repository/pg_task.go @@ -135,6 +135,10 @@ type TaskRepository interface { // Count only campaign tasks completed today (excludes warmup) CountCampaignEmailsSentToday(ctx context.Context, accountID uuid.UUID) (int, error) + // CountCampaignSendsTodayBySender is one campaign's sends today, by the + // 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) CountWarmupEmailsSentToday(ctx context.Context, accountID uuid.UUID) (int, error) // Create user-initiated email task (transactional) @@ -446,6 +450,38 @@ func (r *taskRepository) CountCampaignEmailsSentToday(ctx context.Context, accou return count, 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. +func (r *taskRepository) CountCampaignSendsTodayBySender(ctx context.Context, campaignID uuid.UUID) (map[uuid.UUID]int, error) { + query := ` + SELECT t.email_account_id, COUNT(*) + FROM tasks t + JOIN campaign_tasks ct ON ct.task_id = t.id + WHERE ct.campaign_id = $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, campaignID) + if err != nil { + return nil, err + } + defer rows.Close() + out := map[uuid.UUID]int{} + 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() +} + // CreateEmailTaskFull creates a task and email task entry in a single transaction func (r *taskRepository) CreateEmailTaskFull(ctx context.Context, task *Task, emailTask *EmailTask) error { tx, err := r.db.Begin(ctx) diff --git a/internal/scheduler/campaign_scheduler.go b/internal/scheduler/campaign_scheduler.go index 13140f0a5..30e6d3600 100644 --- a/internal/scheduler/campaign_scheduler.go +++ b/internal/scheduler/campaign_scheduler.go @@ -781,7 +781,12 @@ func (s *schedulerService) placeCampaignSend(ctx context.Context, campaign *mode // metronomed sends are a pattern even with additive jitter. varied := float64(remainingMinutes/remainingEmails) * (0.55 + rand.Float64()*0.9) idealInterval := time.Duration(varied * float64(time.Minute)) - minInterval := time.Second * time.Duration(gapSeconds) + // The floor is the POOL's spacing, not the chosen mailbox's. Each + // mailbox's own gap is enforced against its own last send at STEP + // 10; flooring the chain's next tick at one mailbox's whole gap + // held a three-mailbox campaign to one mailbox's rate however many + // of them were free. + minInterval := time.Second * time.Duration(gapSeconds) / time.Duration(max(1, len(pool))) if idealInterval < minInterval { idealInterval = minInterval } diff --git a/internal/scheduler/send_plan.go b/internal/scheduler/send_plan.go new file mode 100644 index 000000000..fd8445eac --- /dev/null +++ b/internal/scheduler/send_plan.go @@ -0,0 +1,593 @@ +package scheduler + +import ( + "context" + "time" + + "github.com/google/uuid" + + "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. +type CampaignSendPlanner interface { + // PlanCampaignDay is today's sending plan for one campaign: what will go + // out and every limit that decided it. orgDailyLimit is the workspace's + // plan-level daily campaign limit, negative when unlimited. Read-only. + PlanCampaignDay(ctx context.Context, campaignID uuid.UUID, orgDailyLimit int) (*models.CampaignSendPlan, error) + // PoolCapacityToday is what a mailbox pool can send today under a + // campaign's clamps, for a campaign that may not exist yet (the wizard's + // estimate) or for the whole workspace (the dashboard meter). A campaign + // with no daily limit clamps nothing at the campaign level. + PoolCapacityToday(ctx context.Context, campaign *models.Campaign, accounts []models.Email) (*models.WorkspaceSendCapacity, error) +} + +// mailboxDay is one mailbox's day on a campaign with the working shown: every +// clamp in the order the send path applies it, and how many sends it took. +// The deltas add up: room minus the deltas is remaining. +type mailboxDay struct { + acct models.Email + cap capClamp + 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 +} + +// stagedCap is explainCap with every intermediate cap kept, so a plan can say +// how many sends each clamp took. The final value is explainCap's; the test +// holds the two together. +func stagedCap(p *campaignPass, acct models.Email) (stages [5]int, limitedBy string) { + c := p.campaign + stages[0] = acct.CampaignLimit + limitedBy = capByMailbox + cur := stages[0] + clamp := func(i int, v int, why string) { + if v < cur { + cur, limitedBy = v, why + } + stages[i] = cur + } + if c.DailyLimit > 0 { + clamp(1, c.DailyLimit, capByCampaign) + } else { + stages[1] = cur + } + if c.RampEnabled { + clamp(2, campaignRampCeiling(true, c.RampStart, c.RampIncrement, c.RampCeiling, c.RampLevel), capByRamp) + } else { + stages[2] = cur + } + clamp(3, coldCeilingFor(p.coldRamp[acct.ID], cur), capByGraduation) + if m := p.risk.CapMultiplier(); m < 1 { + risked := int(float64(cur)*m + 0.5) + if risked < 1 { + risked = 1 + } + clamp(4, risked, capByRisk) + } else { + stages[4] = cur + } + return stages, limitedBy +} + +// room is what a cap leaves after this campaign's own sends. +func room(capv, sentThis int) int { + return max(0, capv-sentThis) +} + +// 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 { + d := mailboxDay{acct: acct, configured: acct.CampaignLimit, 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) + step := func(capv int) int { + next := room(capv, sentThis) + delta := r - next + r = next + return delta + } + d.byCampaignLimit = step(stages[1]) + 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) + } + + // Standing gates: authentication, cold rotation, warmup health. Asked with + // a nominal budget so the answer is the gate alone; the budget is applied + // below where its delta can be named. + gate, _ := s.gateFor(ctx, pass, acct, 1) + d.health = s.healthFor(ctx, pass, acct.ID) + if !gate.open() && gate.reason != gateBudget { + d.gate = gate + d.byGate = r + r = 0 + switch gate.reason { + case gateAuth: + d.state = models.MailboxPlanDomainAuth + case gateResting: + d.state = models.MailboxPlanResting + default: + d.state = models.MailboxPlanHealthHold + } + return d + } + // Other campaigns share this mailbox's cap; what they sent is gone. + next := max(0, r-d.sentOther) + d.byOther = r - next + r = next + if r == 0 { + d.state = models.MailboxPlanBudgetSpent + d.gate = mailboxGate{reason: gateBudget, paced: true} + return d + } + // A watch or throttled band spaces sends wider; the pass reads the + // dampened budget, so the day lands near this. + if d.health.known { + if m := adjustmentFor(d.health.state).volumeMultiplier; m < 1 { + next = int(float64(r) * m) + d.byHealthPace = r - next + r = next + } + } + + // The mailbox's own clock: a behaviour profile's rolled workday, or the + // 8am-8pm band in its own timezone when that differs from the campaign's. + secondsLeft := windowSecondsLeft + bhv := pass.behaviors[acct.ID] + if bhv.Enabled { + openAt, ok := s.placeWithinBehavior(ctx, bhv, now) + if !ok { + d.state = models.MailboxPlanNoWorkingDay + d.gate = mailboxGate{reason: gateNoWorkday} + d.byHours = r + return d + } + if !sameLocalDay(openAt, now, bhv.Loc) { + d.state = models.MailboxPlanHoursClosed + d.reopensAt = openAt + d.gate = mailboxGate{reason: gateHours, paced: true, reopensAt: openAt} + d.byHours = r + return d + } + if openAt.After(now) { + d.reopensAt = openAt + secondsLeft = min(secondsLeft, int(windowClosesAt.Sub(openAt).Seconds())) + } + next = min(r, s.behaviorDailyCap(ctx, bhv, r, openAt)) + d.byBehavior = r - next + r = next + if r == 0 { + d.state = models.MailboxPlanBudgetSpent + d.gate = mailboxGate{reason: gateBudget, paced: true} + return d + } + 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) + secondsLeft = min(secondsLeft, max(0, int(workEnd.Sub(openAt).Seconds()))) + if plan.HourlyLimit > 0 && secondsLeft > 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)) + } + } + } else if acct.Timezone != "" && acct.Timezone != pass.campaign.Timezone { + loc := loadLocation(acct.Timezone) + if h := now.In(loc).Hour(); h < 8 || h >= 20 { + open := businessHoursReopen(now, loc) + d.reopensAt = open + if !sameLocalDay(open, now, loc) || !open.Before(windowClosesAt) { + d.state = models.MailboxPlanHoursClosed + d.gate = mailboxGate{reason: gateHours, paced: true, reopensAt: open} + d.byHours = r + return d + } + secondsLeft = min(secondsLeft, int(windowClosesAt.Sub(open).Seconds())) + } + } + + if windowSecondsLeft <= 0 { + // The campaign's window is closed for the day; that clamp is reported + // once for the whole pool rather than on every mailbox. + d.state = models.MailboxPlanWindowClosed + d.remaining = r + return d + } + + // 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) { + earliest = at + secondsLeft -= int(at.Sub(now).Seconds()) + } + } + paceMax := 0 + if secondsLeft >= 0 { + paceMax = 1 + if d.minGap > 0 { + paceMax += secondsLeft / d.minGap + } else { + paceMax = r + } + } + next = min(r, paceMax) + d.bySpacing = r - next + r = next + d.remaining = r + if r > 0 { + d.state = models.MailboxPlanSending + return d + } + // Its next allowed send falls after the day's sending time ends. + d.state = models.MailboxPlanHoursClosed + if !earliest.IsZero() { + d.reopensAt = earliest + } + 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) + windows := effectiveWindows(campaign) + local := now.In(tz) + y, m, d := local.Date() + midnight := time.Date(y, m, d, 0, 0, 0, 0, tz) + out := models.CampaignSendWindow{} + if campaign.EndDate != nil { + t := *campaign.EndDate + out.EndsAt = &t + } + if campaign.StartDate != nil && campaign.StartDate.After(now) { + t := *campaign.StartDate + out.StartsAt = &t + } + if campaign.EndDate != nil && campaign.EndDate.Before(now) { + return out, 0, midnight + } + if windows.IsEmpty() { + out.SendingDay = true + out.OpenNow = out.StartsAt == nil || !campaign.StartDate.After(now) + closes := midnight.AddDate(0, 0, 1) + out.ClosesAt = &closes + from := now + if out.StartsAt != nil { + from = *out.StartsAt + } + left := 0 + if from.Before(closes) { + left = int(closes.Sub(from).Seconds()) + } + if !out.OpenNow { + out.OpensAt = out.StartsAt + } + out.MinutesLeft = left / 60 + return out, left, closes + } + nowMin := local.Hour()*60 + local.Minute() + closes := midnight + left := 0 + for _, iv := range windows[int(local.Weekday())] { + out.SendingDay = true + end := midnight.Add(time.Duration(iv.End) * time.Minute) + if end.After(closes) { + closes = end + } + if nowMin >= iv.Start && nowMin < iv.End { + out.OpenNow = true + } + if iv.End > nowMin { + left += (iv.End - max(iv.Start, nowMin)) * 60 + } + } + if out.SendingDay { + out.ClosesAt = &closes + } + // A start date still ahead holds the whole window, whatever the clock says. + 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 + } + } + } else { + left = 0 + } + } + if !out.OpenNow { + from := now + if out.StartsAt != nil && out.StartsAt.After(from) { + from = *out.StartsAt + } + if open := nextScheduleSlot(from, windows, tz); open.After(now) { + out.OpensAt = &open + } + } + out.MinutesLeft = left / 60 + return out, left, closes +} + +// PlanCampaignDay implements CampaignSendPlanner. +func (s *schedulerService) PlanCampaignDay(ctx context.Context, campaignID uuid.UUID, orgDailyLimit int) (*models.CampaignSendPlan, error) { + campaign, err := s.campaignRepo.GetByID(ctx, campaignID) + if err != nil { + return nil, err + } + now := time.Now() + projectRampLevel(campaign, now) + tz := loadLocation(campaign.Timezone) + window, windowSecondsLeft, closesAt := dayWindow(campaign, now) + + plan := &models.CampaignSendPlan{ + CampaignID: campaign.ID, + Status: campaign.Status, + Day: now.In(tz).Format("2006-01-02"), + Timezone: tz.String(), + ComputedAt: now, + Window: window, + Limits: []models.CampaignSendLimit{}, + Mailboxes: []models.CampaignMailboxPlan{}, + } + + accounts, _, err := s.campaignSenders(ctx, campaign) + if err != nil { + return nil, err + } + pass := s.newCampaignPass(ctx, campaign, accounts) + sentBySender, err := s.taskRepo.CountCampaignSendsTodayBySender(ctx, campaignID) + if err != nil { + return nil, err + } + + days := make([]mailboxDay, 0, len(accounts)) + for _, acct := range accounts { + sentAll, err := s.sentTodayFor(ctx, pass, acct.ID) + if err != nil { + return nil, err + } + days = append(days, s.planMailbox(ctx, pass, acct, sentBySender[acct.ID], sentAll, now, windowSecondsLeft, closesAt)) + } + + // The pool's waterfall: each clamp's deltas summed over the mailboxes. + type tally struct { + kind string + emails int + mailboxes int + } + tallies := []tally{ + {kind: models.SendLimitCampaignDailyLimit}, {kind: models.SendLimitCampaignRamp}, + {kind: models.SendLimitWarmupGraduation}, {kind: models.SendLimitWorkspaceRisk}, + {kind: models.SendLimitDomainAuth}, {kind: models.SendLimitResting}, {kind: models.SendLimitHealthHold}, + {kind: models.SendLimitOtherCampaigns}, {kind: models.SendLimitHealthPace}, + {kind: models.SendLimitMailboxHours}, {kind: models.SendLimitSendingBehavior}, {kind: models.SendLimitSpacing}, + } + add := func(kind string, n int) { + if n <= 0 { + return + } + for i := range tallies { + if tallies[i].kind == kind { + tallies[i].emails += n + tallies[i].mailboxes++ + return + } + } + } + poolRemaining := 0 + for _, d := range days { + plan.ConfiguredCeiling += d.configured + plan.SentToday += d.sentThis + poolRemaining += d.remaining + add(models.SendLimitCampaignDailyLimit, d.byCampaignLimit) + add(models.SendLimitCampaignRamp, d.byRamp) + add(models.SendLimitWarmupGraduation, d.byGraduation) + add(models.SendLimitWorkspaceRisk, d.byRisk) + switch d.gate.reason { + case gateAuth: + add(models.SendLimitDomainAuth, d.byGate) + case gateResting: + add(models.SendLimitResting, d.byGate) + case gateHealth: + add(models.SendLimitHealthHold, d.byGate) + } + add(models.SendLimitOtherCampaigns, d.byOther) + add(models.SendLimitHealthPace, d.byHealthPace) + add(models.SendLimitMailboxHours, d.byHours) + add(models.SendLimitSendingBehavior, d.byBehavior) + add(models.SendLimitSpacing, d.bySpacing) + + mp := models.CampaignMailboxPlan{ + ID: d.acct.ID, Email: d.acct.Email, Provider: d.acct.Provider, + ConfiguredCap: d.configured, CapToday: d.cap.Cap, LimitedBy: d.cap.LimitedBy, + SentToday: d.sentThis, SentByOtherCampaigns: d.sentOther, + ExpectedRemaining: d.remaining, State: d.state, MinGapSeconds: d.minGap, Graduation: d.graduation, + } + if !d.reopensAt.IsZero() { + t := d.reopensAt + mp.ReopensAt = &t + } + if d.health.known && d.health.state != "" && d.health.state != models.WarmupHealthHealthy { + mp.Health = string(d.health.state) + } + plan.Mailboxes = append(plan.Mailboxes, mp) + } + for _, t := range tallies { + if t.emails > 0 { + plan.Limits = append(plan.Limits, models.CampaignSendLimit{Kind: t.kind, Emails: t.emails, Mailboxes: t.mailboxes}) + } + } + // The biggest mailbox-level clamp is the story when nothing below binds. + bottleneck := "" + biggest := 0 + for _, t := range tallies { + if t.emails > biggest { + biggest, bottleneck = t.emails, t.kind + } + } + campaignLimit := func(kind string, n int) { + if n <= 0 { + return + } + plan.Limits = append(plan.Limits, models.CampaignSendLimit{Kind: kind, Emails: n}) + bottleneck = kind + } + + // Campaign-level clamps, in the order the send path meets them. + remaining := poolRemaining + if windowSecondsLeft <= 0 { + campaignLimit(models.SendLimitSendingWindow, remaining) + remaining = 0 + } + if campaign.Status != "active" || pass.risk.BlocksSending() { + campaignLimit(models.SendLimitNotRunning, remaining) + remaining = 0 + } + if orgDailyLimit >= 0 && campaign.OrganizationID != nil { + sent, err := s.campaignProgressRepo.CountEmailsSentTodayByOrganization(ctx, *campaign.OrganizationID) + if err != nil { + return nil, err + } + left := max(0, orgDailyLimit-sent) + plan.Organization = &models.CampaignOrgAllowance{DailyLimit: orgDailyLimit, SentToday: sent, Remaining: left} + if left < remaining { + campaignLimit(models.SendLimitOrgDailyLimit, remaining-left) + remaining = left + } + } + + // The leads: mailboxes can only send to a step that is due today. + supply, err := s.campaignProgressRepo.LeadSupply(ctx, campaignID, closesAt) + if err != nil { + return nil, err + } + newLeadsToday := 0 + if campaign.MaxNewLeadsPerDay > 0 { + if n, err := s.campaignRepo.CountNewLeadsStartedToday(ctx, campaignID); err == nil { + newLeadsToday = n + } + } + plan.Leads = models.CampaignLeadSupply{ + DueNow: supply.DueNow, DueLaterToday: supply.DueLaterToday, + NewLeadsDueToday: supply.DueNowNewLeads + supply.DueLaterTodayNewLeads, + WaitingOnStep: supply.WaitingOnStep, WaitingOnCondition: supply.WaitingOnCondition, Held: supply.Held, + NewLeadsStartedToday: newLeadsToday, MaxNewLeadsPerDay: campaign.MaxNewLeadsPerDay, NextDueAt: supply.NextDueAt, + } + followUps := supply.DueNow + supply.DueLaterToday - plan.Leads.NewLeadsDueToday + newDue := plan.Leads.NewLeadsDueToday + if avail := followUps + newDue; avail < remaining { + campaignLimit(models.SendLimitLeads, remaining-avail) + remaining = avail + } + if campaign.MaxNewLeadsPerDay > 0 { + roomForNew := max(0, campaign.MaxNewLeadsPerDay-newLeadsToday) + if capped := followUps + min(newDue, roomForNew); capped < remaining { + campaignLimit(models.SendLimitNewLeadCap, remaining-capped) + remaining = capped + } + } + + plan.ExpectedRemaining = remaining + plan.Projected = plan.SentToday + remaining + if remaining == 0 && bottleneck == "" && plan.SentToday > 0 { + bottleneck = models.SendBottleneckBudgetSpent + } + plan.Bottleneck = bottleneck + if campaign.Status == "active" { + plan.NextWakeAt = s.campaignWakeup(ctx, campaignID) + } + return plan, nil +} + +// PoolCapacityToday implements CampaignSendPlanner. +func (s *schedulerService) PoolCapacityToday(ctx context.Context, campaign *models.Campaign, accounts []models.Email) (*models.WorkspaceSendCapacity, error) { + if campaign == nil { + campaign = &models.Campaign{} + } + pass := s.newCampaignPass(ctx, campaign, accounts) + out := &models.WorkspaceSendCapacity{} + for _, acct := range accounts { + out.ConfiguredCeiling += acct.CampaignLimit + stages, _ := stagedCap(pass, acct) + if gate, _ := s.gateFor(ctx, pass, acct, 1); !gate.open() && gate.reason != gateBudget { + out.Held++ + continue + } + if bhv := pass.behaviors[acct.ID]; bhv.Enabled { + if _, ok := s.placeWithinBehavior(ctx, bhv, time.Now()); !ok { + out.Held++ + continue + } + } + capv := stages[4] + if capv <= 0 { + out.Held++ + continue + } + out.Mailboxes++ + out.Capacity += capv + sent, err := s.sentTodayFor(ctx, pass, acct.ID) + if err != nil { + return nil, err + } + out.Remaining += max(0, capv-sent) + } + return out, nil +} diff --git a/internal/scheduler/send_plan_live_test.go b/internal/scheduler/send_plan_live_test.go new file mode 100644 index 000000000..51b61ad7d --- /dev/null +++ b/internal/scheduler/send_plan_live_test.go @@ -0,0 +1,92 @@ +package scheduler + +import ( + "context" + "testing" + + "github.com/warmbly/warmbly/internal/models" +) + +// The plan is read through the real repositories: the sender ledger, the +// routing over every lead, the organization's day. What it says has to add +// up, and has to agree with what the fixture put on the books. Same harness +// and env var as live_integration_test.go. +func TestLivePlanCampaignDayAddsUp(t *testing.T) { + _, pool := liveDB(t) + f := newLiveFixture(t, pool, "UTC") + f.addSentLead(t) + planner, ok := loggedScheduler(t, f).(CampaignSendPlanner) + if !ok { + t.Fatal("the scheduler service does not plan") + } + + plan, err := planner.PlanCampaignDay(context.Background(), f.campaign, 5) + if err != nil { + t.Fatal(err) + } + if plan.ConfiguredCeiling != 50 || plan.SentToday != 1 { + t.Fatalf("ceiling %d sent %d, want 50 and the one send on the books", plan.ConfiguredCeiling, plan.SentToday) + } + removed := 0 + for _, l := range plan.Limits { + if l.Emails <= 0 { + t.Fatalf("limit %q listed with nothing removed", l.Kind) + } + removed += l.Emails + } + if plan.ConfiguredCeiling-removed-plan.SentToday != plan.ExpectedRemaining { + t.Fatalf("the waterfall does not add up: %d - %d - %d != %d (%+v)", + plan.ConfiguredCeiling, removed, plan.SentToday, plan.ExpectedRemaining, plan.Limits) + } + if plan.Projected != plan.SentToday+plan.ExpectedRemaining { + t.Fatalf("projected %d != sent %d + remaining %d", plan.Projected, plan.SentToday, plan.ExpectedRemaining) + } + // The fixture's own lead is the only step due: the sent lead's flow is + // over, so however much the mailbox could send, one is all there is. + if plan.Leads.DueNow != 1 || plan.ExpectedRemaining > 1 { + t.Fatalf("leads %+v remaining %d, want exactly one due and at most one to go", plan.Leads, plan.ExpectedRemaining) + } + if plan.ExpectedRemaining == 1 && plan.Bottleneck != models.SendLimitLeads { + t.Fatalf("bottleneck %q, want the leads to be what binds", plan.Bottleneck) + } + if plan.Organization == nil || plan.Organization.DailyLimit != 5 || plan.Organization.SentToday != 1 || plan.Organization.Remaining != 4 { + t.Fatalf("organization %+v, want the plan's allowance of 5 with one used", plan.Organization) + } + if len(plan.Mailboxes) != 1 { + t.Fatalf("got %d mailboxes, want the fixture's one", len(plan.Mailboxes)) + } + mb := plan.Mailboxes[0] + if mb.ID != f.mailbox || mb.CapToday != 50 || mb.LimitedBy != capByMailbox || mb.SentToday != 1 || mb.SentByOtherCampaigns != 0 || mb.MinGapSeconds != 600 { + t.Fatalf("mailbox row %+v", mb) + } + if !plan.Window.SendingDay { + t.Fatalf("an always-open campaign reported no sending day: %+v", plan.Window) + } +} + +// A campaign that is not running still says what it would send, and says +// that the campaign not running is what stops it. +func TestLivePlanCampaignDayForAPausedCampaign(t *testing.T) { + _, pool := liveDB(t) + f := newLiveFixture(t, pool, "UTC") + if _, err := pool.Exec(context.Background(), `UPDATE campaigns SET status = 'paused' WHERE id = $1`, f.campaign); err != nil { + t.Fatal(err) + } + planner := loggedScheduler(t, f).(CampaignSendPlanner) + plan, err := planner.PlanCampaignDay(context.Background(), f.campaign, -1) + if err != nil { + t.Fatal(err) + } + if plan.ExpectedRemaining != 0 || plan.Bottleneck != models.SendLimitNotRunning || plan.Organization != nil || plan.NextWakeAt != nil { + t.Fatalf("paused campaign planned %+v", plan) + } + held := 0 + for _, l := range plan.Limits { + if l.Kind == models.SendLimitNotRunning { + held = l.Emails + } + } + if held < 1 { + t.Fatalf("the not-running clamp removed %d, want what the pool would otherwise send (%+v)", held, plan.Limits) + } +} diff --git a/internal/scheduler/send_plan_test.go b/internal/scheduler/send_plan_test.go new file mode 100644 index 000000000..2bf5aee3a --- /dev/null +++ b/internal/scheduler/send_plan_test.go @@ -0,0 +1,121 @@ +package scheduler + +import ( + "testing" + "time" + + "github.com/google/uuid" + + "github.com/warmbly/warmbly/internal/models" + "github.com/warmbly/warmbly/internal/repository" +) + +// The plan explains the cap with the working shown; explainCap decides the +// send. They read the same clamps, and this keeps them from drifting apart. +func TestStagedCapAgreesWithExplainCap(t *testing.T) { + id := uuid.New() + warmed := time.Now().Add(-2 * 24 * time.Hour) + for _, tc := range []struct { + name string + pass *campaignPass + acct models.Email + }{ + {"plain", &campaignPass{campaign: &models.Campaign{DailyLimit: 50}, risk: models.OrgRiskTrusted}, models.Email{ID: id, CampaignLimit: 50}}, + {"campaign limit", &campaignPass{campaign: &models.Campaign{DailyLimit: 20}, risk: models.OrgRiskTrusted}, models.Email{ID: id, CampaignLimit: 50}}, + {"ramp", &campaignPass{campaign: &models.Campaign{DailyLimit: 50, RampEnabled: true, RampStart: 10, RampCeiling: 50, RampLevel: 15}, risk: models.OrgRiskTrusted}, models.Email{ID: id, CampaignLimit: 50}}, + {"graduation", &campaignPass{campaign: &models.Campaign{DailyLimit: 50}, coldRamp: map[uuid.UUID]repository.ColdRampState{id: {WarmupStartedAt: &warmed}}, risk: models.OrgRiskTrusted}, models.Email{ID: id, CampaignLimit: 50}}, + {"restricted", &campaignPass{campaign: &models.Campaign{DailyLimit: 50}, risk: models.OrgRiskRestricted}, models.Email{ID: id, CampaignLimit: 50}}, + {"everything", &campaignPass{campaign: &models.Campaign{DailyLimit: 30, RampEnabled: true, RampStart: 10, RampCeiling: 50, RampLevel: 25}, coldRamp: map[uuid.UUID]repository.ColdRampState{id: {WarmupStartedAt: &warmed}}, risk: models.OrgRiskRestricted}, models.Email{ID: id, CampaignLimit: 80}}, + } { + stages, by := stagedCap(tc.pass, tc.acct) + want := tc.pass.explainCap(tc.acct) + if stages[4] != want.Cap || by != want.LimitedBy { + t.Errorf("%s: staged %d by %q, explainCap %d by %q", tc.name, stages[4], by, want.Cap, want.LimitedBy) + } + for i := 1; i < len(stages); i++ { + if stages[i] > stages[i-1] { + t.Errorf("%s: stage %d (%d) above stage %d (%d): a clamp can only lower", tc.name, i, stages[i], i-1, stages[i-1]) + } + } + } +} + +// A campaign with no daily limit (the workspace meter, the wizard before a +// limit is chosen) must not be clamped to zero by it. +func TestStagedCapIgnoresAnUnsetCampaignLimit(t *testing.T) { + pass := &campaignPass{campaign: &models.Campaign{}, risk: models.OrgRiskTrusted} + stages, by := stagedCap(pass, models.Email{ID: uuid.New(), CampaignLimit: 50}) + if stages[4] != 50 || by != capByMailbox { + t.Fatalf("got %d by %q, want 50 by %q", stages[4], by, capByMailbox) + } +} + +func TestDayWindow(t *testing.T) { + tz := "Europe/Paris" + loc, _ := time.LoadLocation(tz) + // Wednesday 14:30 Paris. + now := time.Date(2026, 9, 16, 14, 30, 0, 0, loc) + nineToFive := models.ScheduleWindows{} + for wd := 1; wd <= 5; wd++ { + nineToFive[wd] = []models.TimeInterval{{Start: 9 * 60, End: 17 * 60}} + } + + t.Run("open, with the rest of the window ahead", func(t *testing.T) { + w, secs, closes := dayWindow(&models.Campaign{Timezone: tz, ScheduleWindows: nineToFive}, now) + if !w.SendingDay || !w.OpenNow || w.MinutesLeft != 150 || secs != 150*60 { + t.Fatalf("got %+v secs %d", w, secs) + } + if closes.In(loc).Hour() != 17 || closes.In(loc).Minute() != 0 { + t.Fatalf("closes at %v", closes.In(loc)) + } + }) + t.Run("closed for the day, opens tomorrow", func(t *testing.T) { + late := time.Date(2026, 9, 16, 18, 0, 0, 0, loc) + w, secs, _ := dayWindow(&models.Campaign{Timezone: tz, ScheduleWindows: nineToFive}, late) + if !w.SendingDay || w.OpenNow || secs != 0 || w.OpensAt == nil || w.OpensAt.In(loc).Day() != 17 || w.OpensAt.In(loc).Hour() != 9 { + t.Fatalf("got %+v secs %d", w, secs) + } + }) + t.Run("weekend is not a sending day", func(t *testing.T) { + sat := time.Date(2026, 9, 19, 10, 0, 0, 0, loc) + w, secs, _ := dayWindow(&models.Campaign{Timezone: tz, ScheduleWindows: nineToFive}, sat) + if w.SendingDay || w.OpenNow || secs != 0 || w.OpensAt == nil || w.OpensAt.In(loc).Weekday() != time.Monday { + t.Fatalf("got %+v secs %d", w, secs) + } + }) + t.Run("before the window opens today", func(t *testing.T) { + early := time.Date(2026, 9, 16, 8, 0, 0, 0, loc) + w, secs, _ := dayWindow(&models.Campaign{Timezone: tz, ScheduleWindows: nineToFive}, early) + if !w.SendingDay || w.OpenNow || w.MinutesLeft != 8*60 || secs != 8*3600 || w.OpensAt == nil || w.OpensAt.In(loc).Hour() != 9 { + t.Fatalf("got %+v secs %d", w, secs) + } + }) + t.Run("unconstrained schedule runs to midnight", func(t *testing.T) { + w, secs, _ := dayWindow(&models.Campaign{Timezone: tz}, now) + if !w.SendingDay || !w.OpenNow || w.MinutesLeft != 9*60+30 || secs != (9*60+30)*60 { + t.Fatalf("got %+v secs %d", w, secs) + } + }) + t.Run("a start date ahead holds the day", func(t *testing.T) { + start := time.Date(2026, 9, 18, 9, 0, 0, 0, loc) + w, secs, _ := dayWindow(&models.Campaign{Timezone: tz, ScheduleWindows: nineToFive, StartDate: &start}, now) + if w.OpenNow || secs != 0 || w.StartsAt == nil || w.OpensAt == nil || !w.OpensAt.Equal(start) { + t.Fatalf("got %+v secs %d", w, secs) + } + }) + 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) + if w.OpenNow || w.SendingDay || secs != 0 { + t.Fatalf("got %+v secs %d", w, secs) + } + }) +} + +// The waterfall's arithmetic on one mailbox, with no repositories behind it: +// the cap clamps and the sends already made must add up to what is left. +func TestRoomAddsUp(t *testing.T) { + if room(50, 20) != 30 || room(10, 20) != 0 || room(0, 0) != 0 { + t.Fatal("room is cap minus sent, never negative") + } +} diff --git a/skills/warmbly-api/SKILL.md b/skills/warmbly-api/SKILL.md index dec7aefb8..3ab204124 100644 --- a/skills/warmbly-api/SKILL.md +++ b/skills/warmbly-api/SKILL.md @@ -36,7 +36,7 @@ Run `warmblyctl --help` for subcommands and `warmblyctl | Family | Covers | |---|---| | `me` | Identity and granted scopes | -| `campaign` | list, get, create, update, delete, steps, senders, preflight, start, stop, test-email, logs, pause-lead / resume-lead | +| `campaign` | list, get, create, update, delete, steps, senders, preflight, start, stop, test-email, logs, plan, pause-lead / resume-lead | | `contact` | list (search), get, lookup, create, update, delete, notes, timeline, import, export | | `mailbox` | list, get, update, delete, auth-check, sync, identity, refresh-identity, behavior, verify, send, warmup-start/pause/resume/stop/status | | `inbox` | list, count, thread, seen, reply, compose, agent drafts, scheduled sends | @@ -83,6 +83,9 @@ These commands put real mail on the wire: `campaign start`, - Run `campaign preflight --id ` before `campaign start` and act on what it reports. It costs nothing and catches missing senders, empty audiences and broken tracking. +- When a campaign sends less than expected, `campaign plan --id ` is the + answer: today's projected sends and every limit that lowered them. Read it + before touching a cap. - Never raise a mailbox's daily cap casually. The platform default is 50 campaign emails per mailbox per day with 600 seconds between sends; a fresh mailbox should start around 10-20. Do not set a cap above 50 unless the diff --git a/skills/warmbly-cli/SKILL.md b/skills/warmbly-cli/SKILL.md index 872dce404..4577f1f0c 100644 --- a/skills/warmbly-cli/SKILL.md +++ b/skills/warmbly-cli/SKILL.md @@ -72,7 +72,7 @@ gives the arguments and flags. Ids are positional, not flags. | Command | Covers | |---|---| | `status` | one call for "what is happening": mailboxes needing attention, what is sending, what is unread | -| `campaign` | list, view, create, edit, delete, steps, senders, segments, preflight, test, start, stop, logs, pause-lead / resume-lead | +| `campaign` | list, view, create, edit, delete, steps, senders, segments, preflight, test, start, stop, logs, plan, pause-lead / resume-lead | | `contact` | list, view, create, edit, delete, lookup, timeline, emails, notes, import, export, verify | | `mailbox` | list, view, edit, check, sync, identity, refresh-identity, behavior, warmup, hold, release, send | | `inbox` | list, view, thread, read, reply, compose, drafts, scheduled, snooze | @@ -111,6 +111,10 @@ These put real mail on the wire and prompt before doing so: - Run `warmbly campaign preflight CAMPAIGN_ID` before `campaign start` and act on what it reports. It costs nothing and catches missing senders, empty audiences and broken tracking. +- When a campaign sends less than expected, `warmbly campaign plan CAMPAIGN_ID` + is the answer: today's projected sends and every limit that lowered them + (easing out of warmup, the campaign limit, other campaigns on the same + mailboxes, the plan's allowance, leads due). Read it before touching a cap. - Never raise a mailbox's daily cap casually. The default is 50 campaign emails per mailbox per day with 600 seconds between sends; a fresh mailbox starts around 10-20. Do not go above 50 unless the user asked and the diff --git a/web/src/app/app/campaigns/[id]/page.tsx b/web/src/app/app/campaigns/[id]/page.tsx index c36a934a3..285a8e671 100644 --- a/web/src/app/app/campaigns/[id]/page.tsx +++ b/web/src/app/app/campaigns/[id]/page.tsx @@ -16,6 +16,7 @@ import { TONE_DOT } from "@/components/ui/tones"; import type { DitherTone } from "@/components/ui/dither"; import AnalyticsShareButton from "@/components/app/analytics/AnalyticsShareButton"; import TaskPreview from "@/components/app/campaigns/TaskPreview"; +import SendPlanCard from "@/components/app/campaigns/SendPlanCard"; import CampaignFormsPanel from "@/components/app/campaigns/CampaignFormsPanel"; import AnimatedNumber from "@/components/ui/AnimatedNumber"; import AdvisorStrip from "@/components/app/advisor/AdvisorStrip"; @@ -115,6 +116,11 @@ export default function CampaignOverview() { that motivated it. Renders nothing when there is nothing wrong. */} + {/* What will actually go out today and every limit that decided + it, read through the scheduler's own gates. This is the number + the caps added up used to misstate (issue #606). */} + +
{/* Main analytics column */}
diff --git a/web/src/components/app/advisor/AdvisorCard.tsx b/web/src/components/app/advisor/AdvisorCard.tsx index aa166e222..6e59a6547 100644 --- a/web/src/components/app/advisor/AdvisorCard.tsx +++ b/web/src/components/app/advisor/AdvisorCard.tsx @@ -88,13 +88,17 @@ export default function AdvisorCard({ finding, onFix, compact = false, defaultOp return (
-
+
-
- {finding.title} +
+ {finding.title} + {compact && !open ? ( + {finding.detail} + ) : null} {/* Autopilot acting silently would be the thing people resent about it. Say up front which findings it is allowed to take, whether or not it is switched on. */} @@ -127,7 +134,7 @@ export default function AdvisorCard({ finding, onFix, compact = false, defaultOp ) : null}
- {!open ? ( + {!open && !compact ? (

{finding.detail}

) : null} @@ -224,15 +231,22 @@ export default function AdvisorCard({ finding, onFix, compact = false, defaultOp transition={{ duration: 0.16, ease: "easeOut" }} className="overflow-hidden" > -
-

{finding.detail}

+
+

{finding.detail}

-
-

- What to do + {compact ? ( +

+ What to do: + {finding.remedy}

-

{finding.remedy}

-
+ ) : ( +
+

+ What to do +

+

{finding.remedy}

+
+ )} {!compact && evidence.length > 0 ? (
diff --git a/web/src/components/app/advisor/AdvisorStrip.tsx b/web/src/components/app/advisor/AdvisorStrip.tsx index df8f818a2..e6c86de20 100644 --- a/web/src/components/app/advisor/AdvisorStrip.tsx +++ b/web/src/components/app/advisor/AdvisorStrip.tsx @@ -8,13 +8,13 @@ import { useMemo, useState } from "react"; import { AnimatePresence, motion } from "framer-motion"; -import { SparklesIcon } from "lucide-react"; +import { ChevronDownIcon, SparklesIcon } from "lucide-react"; import type { AdvisorCategory, AdvisorFinding, AdvisorSurface, } from "@/lib/api/models/app/advisor/Advisor"; -import { SEVERITY_RANK, groupFindings } from "@/lib/api/models/app/advisor/Advisor"; +import { SEVERITY_DOT, SEVERITY_RANK, groupFindings } from "@/lib/api/models/app/advisor/Advisor"; import { SURFACE_FETCH_LIMIT, useAdvisorFindings } from "@/lib/api/hooks/app/advisor/useAdvisor"; import AdvisorCard from "./AdvisorCard"; import AdvisorFixDrawer from "./AdvisorFixDrawer"; @@ -46,6 +46,10 @@ export default function AdvisorStrip({ title, }: Props) { const [fixing, setFixing] = useState(null); + // A compact strip folds to one line: the count and the worst finding. + // Above a page's own content, a list of suggestions is a hint, not the + // page, so it stays closed until asked. + const [expanded, setExpanded] = useState(false); // An entity-scoped strip must not fire before its id exists, or it renders // the whole org's findings for a beat while the page hydrates. @@ -96,7 +100,42 @@ export default function AdvisorStrip({
) : null} -
+ {/* Compact strips sit above a page's own content, so they are + one bordered list of single lines rather than a stack of + cards, folded to a summary line, and nothing opens on its + own. */} + {compact ? ( + + ) : null} + + {!compact || expanded ? ( + +
{groups.map((group) => ( )} ))}
+
+ ) : null} +
diff --git a/web/src/components/app/campaigns/SendPlanCard.tsx b/web/src/components/app/campaigns/SendPlanCard.tsx new file mode 100644 index 000000000..0e031ca28 --- /dev/null +++ b/web/src/components/app/campaigns/SendPlanCard.tsx @@ -0,0 +1,429 @@ +import { useState } from "react"; +import { Link } from "react-router-dom"; +import { AnimatePresence, motion } from "framer-motion"; +import { ChevronDownIcon, ClockIcon, UsersIcon } from "lucide-react"; +import useCampaignSendPlan from "@/lib/api/hooks/app/campaigns/useCampaignSendPlan"; +import type SendPlan from "@/lib/api/models/app/campaigns/SendPlan"; +import type { MailboxPlan, SendBottleneck, SendLimitKind } from "@/lib/api/models/app/campaigns/SendPlan"; +import { DitherMeter } from "@/components/ui/dither"; +import AnimatedNumber from "@/components/ui/AnimatedNumber"; + +// One line per limit: what it is called, what it did, and where to change it. +// `to` is relative to the campaign (settings, schedule, leads) or absolute. +interface LimitMeta { + label: string; + hint: string; + to?: string; +} + +const LIMIT_META: Record = { + campaign_daily_limit: { + label: "Campaign daily limit", + hint: "This campaign caps each mailbox below the mailbox's own daily cap.", + to: "preferences#senders", + }, + campaign_ramp: { + label: "Campaign ramp-up", + hint: "Daily ramp-up is still climbing toward its ceiling.", + to: "preferences#rotation", + }, + warmup_graduation: { + label: "Easing out of warmup", + hint: "A mailbox that warmed starts cold at 5 to 20 a day and adds 5 each clean day until it reaches its cap.", + to: "/app/emails", + }, + workspace_risk: { + label: "Workspace sending posture", + hint: "The workspace is restricted, so every mailbox sends a fraction of its cap.", + to: "/app/settings/deliverability", + }, + domain_auth: { + label: "Domain authentication failing", + hint: "A sending domain has failed SPF or DMARC for longer than the grace period; nothing is sent from it.", + to: "/app/emails", + }, + resting: { + label: "Resting or in reserve", + hint: "A mailbox is out of cold rotation; it keeps its warmup and sends nothing cold.", + to: "/app/emails", + }, + warmup_health_hold: { + label: "Held by warmup health", + hint: "A mailbox is quarantined or blocked by its warmup health and sends nothing cold until that lifts.", + to: "/app/emails", + }, + other_campaigns: { + label: "Used by other campaigns", + hint: "A mailbox's daily cap is shared by every campaign it is on; these sends went to another one today.", + to: "/app/campaigns", + }, + warmup_health_pace: { + label: "Slowed by warmup health", + hint: "A mailbox on watch or throttled has its sends spaced wider today.", + to: "/app/emails", + }, + mailbox_hours: { + label: "Outside the mailbox's hours", + hint: "A mailbox in another timezone, or with its own workday, is closed for the rest of today.", + to: "/app/emails", + }, + sending_behavior: { + label: "Sending behaviour plan", + hint: "A mailbox's rolled workday gives it fewer sends today than its cap.", + to: "/app/emails", + }, + spacing: { + label: "Minimum gap between sends", + hint: "With the sending time left today, the gap between two sends from one mailbox does not fit any more.", + to: "/app/emails", + }, + sending_window: { + label: "Outside the sending window", + hint: "The campaign's schedule has no sending time left today.", + to: "schedule", + }, + not_running: { + label: "Campaign not running", + hint: "The campaign is not active, so nothing goes out until it is started.", + }, + org_daily_limit: { + label: "Plan's daily allowance", + hint: "The workspace's plan caps campaign emails per day across every campaign.", + to: "/app/settings/billing", + }, + new_lead_cap: { + label: "New leads per day", + hint: "Only this many contacts may receive their first email today; follow-ups keep going.", + to: "preferences#leadflow", + }, + leads: { + label: "Not enough leads due", + hint: "The mailboxes could send more, but no more steps are due today.", + to: "leads", + }, +}; + +const STATE_META: Record = { + sending: { label: "Sending", tone: "bg-emerald-50 text-emerald-700 ring-emerald-200" }, + budget_spent: { label: "Budget used", tone: "bg-slate-100 text-slate-600 ring-slate-200" }, + hours_closed: { label: "Closed", tone: "bg-amber-50 text-amber-700 ring-amber-200" }, + no_working_day: { label: "Day off", tone: "bg-slate-100 text-slate-600 ring-slate-200" }, + domain_auth: { label: "Auth failing", tone: "bg-rose-50 text-rose-700 ring-rose-200" }, + resting: { label: "Resting", tone: "bg-slate-100 text-slate-600 ring-slate-200" }, + health_hold: { label: "Health hold", tone: "bg-rose-50 text-rose-700 ring-rose-200" }, + window_closed: { label: "Window closed", tone: "bg-amber-50 text-amber-700 ring-amber-200" }, +}; + +const LIMITED_BY: Record = { + mailbox_daily_cap: "mailbox cap", + campaign_daily_limit: "campaign limit", + campaign_ramp: "ramp-up", + warmup_graduation: "easing out of warmup", + workspace_risk: "workspace posture", +}; + +function fmtTime(iso: string | undefined, tz: string): string { + if (!iso) return ""; + try { + return new Date(iso).toLocaleTimeString("en-US", { hour: "2-digit", minute: "2-digit", timeZone: tz }); + } catch { + return new Date(iso).toLocaleTimeString("en-US", { hour: "2-digit", minute: "2-digit" }); + } +} + +function fmtDay(iso: string | undefined, tz: string): string { + if (!iso) return ""; + const d = new Date(iso); + const today = new Date(); + const opts: Intl.DateTimeFormatOptions = { timeZone: tz }; + const sameDay = d.toLocaleDateString("en-US", opts) === today.toLocaleDateString("en-US", opts); + if (sameDay) return fmtTime(iso, tz); + return `${d.toLocaleDateString("en-US", { weekday: "short", timeZone: tz })} ${fmtTime(iso, tz)}`; +} + +function shortTz(tz: string): string { + const city = tz.split("/").pop() ?? tz; + return city.replace(/_/g, " "); +} + +// The one sentence under the number: what decided it. +function headline(plan: SendPlan): string { + const b: SendBottleneck = plan.bottleneck; + const n = plan.projected_today; + if (plan.mailboxes.length === 0) return "No mailbox is attached to this campaign, so nothing can go out."; + if (plan.status !== "active") { + // What the pool would send: the not-running clamp put back. + const held = plan.limits.find((l) => l.kind === "not_running")?.emails ?? 0; + return `Once started, this campaign would send about ${(plan.expected_remaining + held).toLocaleString()} today.`; + } + switch (b) { + case "": + return n === plan.configured_ceiling + ? "Nothing is holding it below your settings." + : "Every mailbox is sending at its cap."; + case "budget_spent": + return "Every mailbox has used its budget for today; sending resumes tomorrow."; + case "leads": + return `Only ${(plan.leads.due_now + plan.leads.due_later_today).toLocaleString()} steps are due today; the mailboxes could send more.`; + case "sending_window": + return plan.window.opens_at + ? `Outside the sending window; it opens ${fmtDay(plan.window.opens_at, plan.timezone)}.` + : "Outside the sending window."; + case "new_lead_cap": + return `The new-leads-per-day limit (${plan.leads.max_new_leads_per_day}) is what holds it; follow-ups keep going.`; + case "org_daily_limit": + return `Your plan allows ${plan.organization?.daily_limit.toLocaleString() ?? ""} campaign emails a day across the workspace.`; + default: { + const meta = LIMIT_META[b]; + return meta ? `${meta.label} is holding it below the ${plan.configured_ceiling.toLocaleString()} your settings allow.` : ""; + } + } +} + +function WindowLine({ plan }: { plan: SendPlan }) { + const w = plan.window; + const tz = shortTz(plan.timezone); + let text: string; + if (w.starts_at && new Date(w.starts_at).getTime() > Date.now()) { + text = `Starts ${fmtDay(w.starts_at, plan.timezone)}`; + } else if (!w.sending_day) { + text = w.opens_at ? `Not a sending day. Opens ${fmtDay(w.opens_at, plan.timezone)}` : "Not a sending day"; + } else if (w.open_now) { + text = w.closes_at ? `Window open until ${fmtTime(w.closes_at, plan.timezone)}` : "Window open"; + } else if (w.opens_at) { + text = `Window opens ${fmtDay(w.opens_at, plan.timezone)}`; + } else { + text = "Window closed for today"; + } + return ( + + + {text} + · {tz} + + ); +} + +function LeadsLine({ plan }: { plan: SendPlan }) { + const l = plan.leads; + const parts: string[] = []; + parts.push(`${l.due_now.toLocaleString()} due now`); + 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.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 ( + + + Leads: {parts.join(" · ")} + + ); +} + +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; + + if (q.isPending) { + return ( +
+
+
+
+ ); + } + if (q.isError || !plan) { + return ( +
+ Today's sending plan + couldn't be worked out. It retries on its own. +
+ ); + } + + const frac = plan.projected_today > 0 ? plan.sent_today / plan.projected_today : 0; + const campaignBase = `/app/campaigns/${campaignId}`; + const linkFor = (to?: string) => (!to ? undefined : to.startsWith("/") ? to : `${campaignBase}/${to}`); + const showWaterfall = plan.limits.length > 0 || plan.sent_today > 0; + + return ( +
+ {/* The strip: number, meter, the one sentence, the window, and the toggle. */} + + + + {open ? ( + +
+
+ {/* Left: from the settings to today. */} +
+
+ From your settings to today +
+ {showWaterfall ? ( +
+
+ Mailbox caps added up + + {plan.configured_ceiling.toLocaleString()} + +
+ {plan.limits.map((l) => { + const meta = LIMIT_META[l.kind]; + const to = linkFor(meta?.to); + const body = ( + <> + {meta?.label ?? l.kind} + {l.mailboxes ? ( + + {l.mailboxes} {l.mailboxes === 1 ? "mailbox" : "mailboxes"} + + ) : null} + + −{l.emails.toLocaleString()} + + + ); + return to ? ( + + {body} + + ) : ( +
+ {body} +
+ ); + })} + {plan.sent_today > 0 && ( +
+ Sent today + + −{plan.sent_today.toLocaleString()} + +
+ )} +
+ Still to go today + + {plan.expected_remaining.toLocaleString()} + +
+
+ ) : ( +

Nothing lowers the number today.

+ )} +
+ + + + + {plan.organization && ( + + Plan allowance {plan.organization.sent_today.toLocaleString()}/{plan.organization.daily_limit.toLocaleString()} today + + )} + {plan.next_wake_at && ( + Next pass {fmtDay(plan.next_wake_at, plan.timezone)} + )} +
+
+ + {/* Right: each mailbox's day. */} +
+
+ {plan.mailboxes.length} {plan.mailboxes.length === 1 ? "mailbox" : "mailboxes"} +
+ {plan.mailboxes.length === 0 ? ( +

No mailbox is attached to this campaign.

+ ) : ( +
+ {plan.mailboxes.map((m) => { + const st = STATE_META[m.state] ?? STATE_META.sending; + const sentTitle = m.sent_by_other_campaigns + ? `${m.sent_today} from this campaign, ${m.sent_by_other_campaigns} from other campaigns` + : `${m.sent_today} from this campaign`; + return ( +
+ + {m.email} + + {m.limited_by !== "mailbox_daily_cap" ? `${LIMITED_BY[m.limited_by]} · ` : ""} + gap {Math.round(m.min_gap_seconds / 60)}m + {m.graduation + ? ` · ${m.graduation.days_to_full_cap} clean ${m.graduation.days_to_full_cap === 1 ? "day" : "days"} to ${m.graduation.mailbox_cap}${m.graduation.held ? ", climb paused" : ""}` + : ""} + {m.health ? ` · warmup ${m.health}` : ""} + {m.reopens_at ? ` · reopens ${fmtDay(m.reopens_at, plan.timezone)}` : ""} + + + + {m.sent_today} + {m.sent_by_other_campaigns ? +{m.sent_by_other_campaigns} : null} + / + {m.cap_today} + {m.cap_today !== m.configured_cap && of {m.configured_cap}} + + + {st.label} + +
+ ); + })} +
+ )} +
+
+
+
+ ) : null} +
+
+ ); +} diff --git a/web/src/components/layout/AppNav.tsx b/web/src/components/layout/AppNav.tsx index 0f09beee0..b6636a3f7 100644 --- a/web/src/components/layout/AppNav.tsx +++ b/web/src/components/layout/AppNav.tsx @@ -728,8 +728,9 @@ function Section({ * (shares the dashboard page's query cache; realtime invalidation keeps * it current) * - * The capacity denominator sums each mailbox's configured campaign_limit - * (default 50/day, from internal/config/constants.go). + * The capacity denominator is the dashboard payload's capacity_today: what + * the mailboxes can send today under the scheduler's clamps, not their caps + * added up. */ function LivePanel({ collapsed = false }: { collapsed?: boolean }) { const emails = useAppStore((s) => s.emails); @@ -737,18 +738,27 @@ function LivePanel({ collapsed = false }: { collapsed?: boolean }) { const dash = useDashboard("30d"); const [hovered, setHovered] = useState(null); + const serverCapacity = dash.data?.capacity_today?.capacity; const { active, mailboxes, capacity } = useMemo(() => { const m = emails.length; const a = emails.filter((e) => { const st = mailboxDisplayStatus(e); return st === "healthy" || st === "warming"; }).length; - // Capacity = the sum of each mailbox's configured daily campaign - // limit (default 50/day), not a flat count × 50 — a tuned-down or - // raised mailbox should move the meter's denominator. - const cap = emails.reduce((sum, e) => sum + (e.campaign_limit ?? 50), 0); + // The denominator is what the mailboxes can send today under the + // scheduler's own clamps (the easing-out-of-warmup ceiling, the + // workspace's risk band, a health hold, the plan's daily allowance), + // read from the dashboard payload. Adding up configured caps promised + // a volume the scheduler never intended to send. Until the payload + // arrives, the caps of the mailboxes that can send stand in. + const cap = + serverCapacity ?? + emails.reduce((sum, e) => { + const st = mailboxDisplayStatus(e); + return st === "healthy" || st === "warming" ? sum + (e.campaign_limit ?? 50) : sum; + }, 0); return { active: a, mailboxes: m, capacity: cap }; - }, [emails]); + }, [emails, serverCapacity]); const { sentToday, trend } = useMemo(() => { // daily_trend only contains days that had sends; rebuild a continuous diff --git a/web/src/lib/api/client/app/campaigns/getCampaignSendPlan.ts b/web/src/lib/api/client/app/campaigns/getCampaignSendPlan.ts new file mode 100644 index 000000000..6e6c0fe21 --- /dev/null +++ b/web/src/lib/api/client/app/campaigns/getCampaignSendPlan.ts @@ -0,0 +1,10 @@ +import Request from "../../Request"; +import type SendPlan from "@/lib/api/models/app/campaigns/SendPlan"; + +export default async function getCampaignSendPlan(id: string): Promise { + return Request({ + method: "GET", + url: `/campaigns/${id}/send-plan`, + authorization: true, + }); +} diff --git a/web/src/lib/api/hooks/app/campaigns/useCampaignSendPlan.ts b/web/src/lib/api/hooks/app/campaigns/useCampaignSendPlan.ts new file mode 100644 index 000000000..292d9972a --- /dev/null +++ b/web/src/lib/api/hooks/app/campaigns/useCampaignSendPlan.ts @@ -0,0 +1,16 @@ +import { useQuery } from "@tanstack/react-query"; +import getCampaignSendPlan from "@/lib/api/client/app/campaigns/getCampaignSendPlan"; + +// Today's sending plan. Realtime invalidation covers the sends and the +// settings; the clock is the one input no event announces (a window opening, +// a mailbox's hours, spacing running down), so the plan also re-reads itself +// once a minute while the page is open. +export default function useCampaignSendPlan(id: string) { + return useQuery({ + queryKey: ["campaigns", id, "send-plan"], + queryFn: () => getCampaignSendPlan(id), + enabled: !!id, + refetchInterval: 60_000, + refetchOnWindowFocus: true, + }); +} diff --git a/web/src/lib/api/models/app/analytics/DashboardOverview.ts b/web/src/lib/api/models/app/analytics/DashboardOverview.ts index 4366c169c..abe0cfa2e 100644 --- a/web/src/lib/api/models/app/analytics/DashboardOverview.ts +++ b/web/src/lib/api/models/app/analytics/DashboardOverview.ts @@ -1,3 +1,5 @@ +import type { WorkspaceSendCapacity } from "@/lib/api/models/app/campaigns/SendPlan" + // GET /analytics/dashboard?period=7d|30d|90d — a single (un-enveloped) object // mirroring the backend models.DashboardAnalytics. The previous flat shape // (total_campaigns/total_contacts…) did not match the wire body. @@ -62,4 +64,7 @@ export default interface DashboardOverview { top_campaigns: TopCampaignStats[] account_health: AccountHealthSummary daily_trend: DashboardDailyStats[] + // What the workspace's mailboxes can send today under the scheduler's + // clamps; the sidebar meter's denominator. Absent when not computed. + capacity_today?: WorkspaceSendCapacity } diff --git a/web/src/lib/api/models/app/campaigns/SendPlan.ts b/web/src/lib/api/models/app/campaigns/SendPlan.ts new file mode 100644 index 000000000..cedaf2bc9 --- /dev/null +++ b/web/src/lib/api/models/app/campaigns/SendPlan.ts @@ -0,0 +1,121 @@ +// GET /campaigns/:id/send-plan. Derived through the scheduler's own gates on +// every read; nothing here is stored. The arithmetic adds up: +// configured_ceiling - sum(limits.emails) - sent_today = expected_remaining. + +export type SendLimitKind = + | "campaign_daily_limit" + | "campaign_ramp" + | "warmup_graduation" + | "workspace_risk" + | "domain_auth" + | "resting" + | "warmup_health_hold" + | "other_campaigns" + | "warmup_health_pace" + | "mailbox_hours" + | "sending_behavior" + | "spacing" + | "sending_window" + | "not_running" + | "org_daily_limit" + | "new_lead_cap" + | "leads"; + +export type SendBottleneck = SendLimitKind | "budget_spent" | ""; + +export type MailboxPlanState = + | "sending" + | "budget_spent" + | "hours_closed" + | "no_working_day" + | "domain_auth" + | "resting" + | "health_hold" + | "window_closed"; + +export interface SendLimit { + kind: SendLimitKind; + emails: number; + mailboxes?: number; +} + +export interface SendWindow { + sending_day: boolean; + open_now: boolean; + opens_at?: string; + closes_at?: string; + minutes_left: number; + starts_at?: string; + ends_at?: string; +} + +export interface LeadSupply { + due_now: number; + due_later_today: number; + new_leads_due_today: number; + waiting_on_step: number; + waiting_on_condition: number; + held: number; + new_leads_started_today: number; + max_new_leads_per_day: number; + next_due_at?: string; +} + +export interface ColdRampInfo { + ceiling: number; + mailbox_cap: number; + days_to_full_cap: number; + held: boolean; +} + +export interface MailboxPlan { + id: string; + email: string; + provider: string; + configured_cap: number; + cap_today: number; + limited_by: "mailbox_daily_cap" | "campaign_daily_limit" | "campaign_ramp" | "warmup_graduation" | "workspace_risk"; + sent_today: number; + sent_by_other_campaigns: number; + expected_remaining: number; + state: MailboxPlanState; + reopens_at?: string; + health?: string; + min_gap_seconds: number; + graduation?: ColdRampInfo; +} + +export interface OrgAllowance { + daily_limit: number; + sent_today: number; + remaining: number; +} + +export default interface SendPlan { + campaign_id: string; + status: string; + day: string; + timezone: string; + computed_at: string; + configured_ceiling: number; + projected_today: number; + sent_today: number; + expected_remaining: number; + bottleneck: SendBottleneck; + limits: SendLimit[]; + window: SendWindow; + leads: LeadSupply; + mailboxes: MailboxPlan[]; + organization?: OrgAllowance; + next_wake_at?: string; +} + +// What the workspace's mailboxes can send today between them, on the +// dashboard payload as capacity_today. +export interface WorkspaceSendCapacity { + capacity: number; + remaining_today: number; + configured_ceiling: number; + mailboxes: number; + held: number; +}