mirror of
https://github.com/warmbly/warmbly.git
synced 2026-10-03 08:02:04 +00:00
feat: take back a replayed campaign pass the task queue refused and keep its dead letter, seed the next pass on a fresh context when the one after a failed hand-off could not be written, resume a campaign whose mailboxes lack a worker at the earliest reopening, and log that wait under its own daily line
This commit is contained in:
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user