diff --git a/internal/app/advanced/campaign_replay_test.go b/internal/app/advanced/campaign_replay_test.go index 5a3916a68..da8827668 100644 --- a/internal/app/advanced/campaign_replay_test.go +++ b/internal/app/advanced/campaign_replay_test.go @@ -36,6 +36,12 @@ type replayTasks struct { created bool inserted int repended int + deleted int +} + +func (t *replayTasks) DeleteTask(context.Context, uuid.UUID) error { + t.deleted++ + return nil } func (t *replayTasks) GetTask(context.Context, uuid.UUID) (*repository.Task, error) { @@ -65,9 +71,12 @@ func (t *replayTasks) UpdateTaskStatus(context.Context, uuid.UUID, string) error return nil } -type replayClient struct{} +type replayClient struct{ fail bool } -func (replayClient) CreateTask(context.Context, *proto.ProcessTask, time.Time) (string, error) { +func (c replayClient) CreateTask(context.Context, *proto.ProcessTask, time.Time) (string, error) { + if c.fail { + return "", errors.New("queue unavailable") + } return "local", nil } func (replayClient) DeleteTask(context.Context, string) error { return nil } @@ -81,19 +90,22 @@ func TestReplayDeadLetterForACampaignPass(t *testing.T) { name string ctErr error created bool + queueFails bool wantErr bool wantInserted int wantReplayed int + wantDeleted int }{ {name: "chain has no pass", created: true, wantInserted: 1, wantReplayed: 1}, {name: "chain already has its next pass", created: false, wantReplayed: 1}, {name: "campaign task unreadable", ctErr: errors.New("connection reset"), wantErr: true}, + {name: "queue refuses the new pass", created: true, queueFails: true, wantErr: true, wantInserted: 1, wantDeleted: 1}, } { t.Run(tc.name, func(t *testing.T) { task := &repository.Task{ID: uuid.New(), TaskType: "campaign", EmailAccountID: uuid.New()} repo := &replayRepo{dlq: &models.TaskDeadLetter{ID: uuid.New(), TaskID: task.ID}} tasks := &replayTasks{task: task, campaignID: &campaignID, ctErr: tc.ctErr, created: tc.created} - svc := &service{repo: repo, taskRepo: tasks, tasksClient: replayClient{}} + svc := &service{repo: repo, taskRepo: tasks, tasksClient: replayClient{fail: tc.queueFails}} xerr := svc.ReplayDeadLetter(context.Background(), uuid.New(), repo.dlq.ID) if (xerr != nil) != tc.wantErr { @@ -102,8 +114,9 @@ func TestReplayDeadLetterForACampaignPass(t *testing.T) { if tasks.repended != 0 { t.Fatal("the old pass was put back to pending") } - if tasks.inserted != tc.wantInserted || repo.replayed != tc.wantReplayed { - t.Fatalf("inserted %d replayed %d, want %d and %d", tasks.inserted, repo.replayed, tc.wantInserted, tc.wantReplayed) + if tasks.inserted != tc.wantInserted || repo.replayed != tc.wantReplayed || tasks.deleted != tc.wantDeleted { + t.Fatalf("inserted %d replayed %d deleted %d, want %d, %d and %d", + tasks.inserted, repo.replayed, tasks.deleted, tc.wantInserted, tc.wantReplayed, tc.wantDeleted) } }) } diff --git a/internal/app/advanced/service.go b/internal/app/advanced/service.go index d58619a59..af395527d 100644 --- a/internal/app/advanced/service.go +++ b/internal/app/advanced/service.go @@ -2153,9 +2153,15 @@ func (s *service) replayCampaignPass(ctx context.Context, task *repository.Task) return false, true, nil } name, err := s.tasksClient.CreateTask(ctx, &proto.ProcessTask{TaskId: id.String()}, at) - if err == nil { - _ = s.taskRepo.UpdateTaskScheduledAt(ctx, id, at, name) + if err != nil { + // Nothing will fire the row, and while it is pending the chain can + // seed no other pass: take it back and keep the dead letter. + if derr := s.taskRepo.DeleteTask(ctx, id); derr != nil { + log.Warn().Err(derr).Str("task_id", id.String()).Msg("dead-letter replay: could not remove a pass that was never queued; overdue reconciliation will") + } + return false, true, err } + _ = s.taskRepo.UpdateTaskScheduledAt(ctx, id, at, name) return true, true, nil } diff --git a/internal/scheduler/campaign_scheduler.go b/internal/scheduler/campaign_scheduler.go index 13f74cd2d..3c12899a1 100644 --- a/internal/scheduler/campaign_scheduler.go +++ b/internal/scheduler/campaign_scheduler.go @@ -556,14 +556,21 @@ func (s *schedulerService) placeCampaignSend(ctx context.Context, campaign *mode switch { case noWorker > 0: // Mailboxes that could send once a worker holds them again, which - // the worker reconciler sees to within minutes: the soonest any of - // the pool can come back, so the pass looks again then. A deferral, - // never a pause. - logDecisionOnce("mailboxes_unavailable", + // the worker reconciler sees to within minutes, unless a mailbox + // whose hours reopen sooner comes back first. A deferral, never a + // pause, logged under its own name so an earlier line about + // budgets or hours does not hide it for the rest of the day. + resume := time.Now().Add(workerRecheck) + if hoursClosed > 0 { + if open := nextScheduleSlot(reopensAt, windows, campaignTZ); open.Before(resume) { + resume = open + } + } + logDecisionOnce("mailboxes_no_worker", fmt.Sprintf("No mailbox can send right now: %d not connected to a running sending worker; sending resumes as soon as one is", noWorker), map[string]interface{}{"no_worker": noWorker, "capped_mailboxes": budgetSpent, "hours_closed": hoursClosed, "health_held": healthHeld, "resting_mailboxes": lifecycleGated, "auth_gated": authGated, "pool_size": len(accounts)}) - return time.Now().Add(workerRecheck), nil, accounts[0].ID, ErrCampaignDeferred + return resume, nil, accounts[0].ID, ErrCampaignDeferred case budgetSpent > 0 || hoursClosed > 0: // Resume when the first of them can send again: tomorrow for a // spent budget, the reopening of the mailbox's own 8am-8pm band diff --git a/internal/scheduler/no_worker_live_test.go b/internal/scheduler/no_worker_live_test.go index 5e626a35a..bdf6fad82 100644 --- a/internal/scheduler/no_worker_live_test.go +++ b/internal/scheduler/no_worker_live_test.go @@ -24,12 +24,16 @@ func TestLivePlanShowsAMailboxWithoutAWorkerAsReconnecting(t *testing.T) { s := loggedScheduler(t, f) planner := s.(CampaignSendPlanner) - before, err := planner.PlanCampaignDay(context.Background(), f.campaign, 0) + // -1: no workspace allowance, which would otherwise clamp both to zero. + before, err := planner.PlanCampaignDay(context.Background(), f.campaign, -1) if err != nil { t.Fatal(err) } + if before.ExpectedRemaining <= 0 { + t.Fatalf("baseline remaining %d, want positive capacity", before.ExpectedRemaining) + } s.(WorkerLivenessAware).WireWorkerLiveness(noWorkers{}) - plan, err := planner.PlanCampaignDay(context.Background(), f.campaign, 0) + plan, err := planner.PlanCampaignDay(context.Background(), f.campaign, -1) if err != nil { t.Fatal(err) } @@ -47,7 +51,7 @@ func TestLivePlanShowsAMailboxWithoutAWorkerAsReconnecting(t *testing.T) { t.Cleanup(func() { _, _ = pool.Exec(context.Background(), `DELETE FROM warmup_pool_participants WHERE email_account_id = $1`, f.mailbox) }) - held, err := planner.PlanCampaignDay(context.Background(), f.campaign, 0) + held, err := planner.PlanCampaignDay(context.Background(), f.campaign, -1) if err != nil { t.Fatal(err) } diff --git a/internal/tasks/campaign_task.go b/internal/tasks/campaign_task.go index 29fa41912..fa3dfb075 100644 --- a/internal/tasks/campaign_task.go +++ b/internal/tasks/campaign_task.go @@ -909,6 +909,9 @@ func (s *tasksService) HandleCampaignTask(task *proto.ProcessTask) (result *errx // would run beside the chain that already moved on. if cerr := s.createCampaignTask(ctx, campaign.ID, account.ID, time.Now().Add(config.CampaignTickRetrySeconds*time.Second)); cerr != nil { log.Warn().Err(cerr).Str("campaign_id", campaign.ID.String()).Str("task_id", taskID.String()).Msg("Failed to schedule the campaign's next pass after a send that never left") + // The pass's deadline may have run out during the send; the + // safety net closes it and seeds the next on a fresh context. + s.keepChainAfterFailure(task.TaskId, false) } executionStatus = "completed" return nil