From 150dc7df6e43ad54201b15ae4e4f0cf7898ccf46 Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Sat, 19 Sep 2026 09:31:46 -0700 Subject: [PATCH] feat: add a Today's sending plan to the campaign overview (GET /campaigns/:id/send-plan, derived through the scheduler's own gates: per-mailbox cap clamps, warmup graduation, health bands, other campaigns on the same mailbox, hours, spacing, window, plan allowance, new-lead cap and leads due, as a waterfall that adds up), feed the sidebar meter and the wizard estimate from the same clamps instead of summing configured caps, floor the campaign chain's next tick at the pool's spacing rather than one mailbox's whole gap, fold the compact Advisor strip to one line, add warmbly campaign plan and warmblyctl campaign plan, and document it (issue #606) --- cmd/cli/specs.go | 8 + cmd/warmblyctl/api_resources.go | 1 + docs/content/docs/api/cli.mdx | 2 +- docs/content/docs/api/endpoints.mdx | 1 + docs/content/docs/api/reference/campaigns.mdx | 71 ++- docs/content/docs/guides/campaigns.mdx | 26 + docs/public/openapi.json | 399 ++++++++++++ internal/api/handler/analytics.go | 8 + internal/api/handler/campaign.go | 21 + internal/api/routes.go | 3 + internal/app/campaign/handlers.go | 103 ++- internal/app/campaign/service.go | 7 + internal/models/analytics.go | 4 + internal/models/campaign_send_plan.go | 205 ++++++ internal/repository/lead_supply_live_test.go | 99 +++ internal/repository/pg_campaign_progress.go | 194 ++++-- internal/repository/pg_task.go | 36 ++ internal/scheduler/campaign_scheduler.go | 7 +- internal/scheduler/send_plan.go | 593 ++++++++++++++++++ internal/scheduler/send_plan_live_test.go | 92 +++ internal/scheduler/send_plan_test.go | 121 ++++ skills/warmbly-api/SKILL.md | 5 +- skills/warmbly-cli/SKILL.md | 6 +- web/src/app/app/campaigns/[id]/page.tsx | 6 + .../components/app/advisor/AdvisorCard.tsx | 44 +- .../components/app/advisor/AdvisorStrip.tsx | 50 +- .../components/app/campaigns/SendPlanCard.tsx | 429 +++++++++++++ web/src/components/layout/AppNav.tsx | 24 +- .../app/campaigns/getCampaignSendPlan.ts | 10 + .../app/campaigns/useCampaignSendPlan.ts | 16 + .../models/app/analytics/DashboardOverview.ts | 5 + .../lib/api/models/app/campaigns/SendPlan.ts | 121 ++++ 32 files changed, 2623 insertions(+), 94 deletions(-) create mode 100644 internal/models/campaign_send_plan.go create mode 100644 internal/repository/lead_supply_live_test.go create mode 100644 internal/scheduler/send_plan.go create mode 100644 internal/scheduler/send_plan_live_test.go create mode 100644 internal/scheduler/send_plan_test.go create mode 100644 web/src/components/app/campaigns/SendPlanCard.tsx create mode 100644 web/src/lib/api/client/app/campaigns/getCampaignSendPlan.ts create mode 100644 web/src/lib/api/hooks/app/campaigns/useCampaignSendPlan.ts create mode 100644 web/src/lib/api/models/app/campaigns/SendPlan.ts 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; +}