package models import ( "context" "log" "time" "gorm.io/gorm" ) // ReapExpiredTasks is the periodic cleanup function described in // docs/plans/worker-notifier-mvp.md section 8.5: // // - tasks in state='leased' whose lease_expires_at is past are returned to // state='queued' and have their lease_owner cleared, so the next selector // poll can pick them up. // - tasks in state='failed_retry' whose not_before is past AND // attempts >= max_attempts are moved to state='dead' so they show up on // the admin dead-letter page and stop consuming selector bandwidth. // // It returns the number of rows it touched so the caller can log a metric. // Cheap enough to run from the web process every 30s. func ReapExpiredTasks() (reaped int, deaded int, err error) { now := time.Now() err = DB().Transaction(func(tx *gorm.DB) error { if err := expireQueuedNotificationTasksTx(tx, now); err != nil { return err } var expired []Task if err := tx.Where("state = ? AND lease_expires_at IS NOT NULL AND lease_expires_at < ?", TaskStateLeased, now).Find(&expired).Error; err != nil { return err } for i := range expired { if expired[i].Attempts >= expired[i].MaxAttempts { result := tx.Model(&Task{}).Where("id = ? AND state = ?", expired[i].ID, TaskStateLeased).Updates(map[string]interface{}{"state": TaskStateDead, "lease_owner": "", "lease_token": "", "lease_expires_at": nil, "last_error": "lease expired after max attempts", "updated_at": now}) if result.Error != nil { return result.Error } if result.RowsAffected == 1 { deaded++ if err := FinalizeNotificationTaskTx(tx, &expired[i], "dead", "lease expired after max attempts"); err != nil { return err } } continue } result := tx.Model(&Task{}).Where("id = ? AND state = ?", expired[i].ID, TaskStateLeased).Updates(map[string]interface{}{"state": TaskStateQueued, "lease_owner": "", "lease_token": "", "lease_expires_at": nil, "updated_at": now}) if result.Error != nil { return result.Error } reaped += int(result.RowsAffected) } var exhausted []Task if err := tx.Where("state = ? AND not_before <= ? AND attempts >= max_attempts", TaskStateFailedRetry, now).Find(&exhausted).Error; err != nil { return err } for i := range exhausted { result := tx.Model(&Task{}).Where("id = ? AND state = ?", exhausted[i].ID, TaskStateFailedRetry).Updates(map[string]interface{}{"state": TaskStateDead, "updated_at": now}) if result.Error != nil { return result.Error } if result.RowsAffected == 1 { deaded++ if err := FinalizeNotificationTaskTx(tx, &exhausted[i], "dead", exhausted[i].LastError); err != nil { return err } } } return nil }) return reaped, deaded, err } // FinalizeNotificationTaskTx makes a terminal notification task customer-visible // and auditable. The caller owns the task state transition in this transaction. func FinalizeNotificationTaskTx(tx *gorm.DB, task *Task, status, reason string) error { if task == nil || task.Kind != TaskKindNotification || task.MessageID == nil { return nil } if err := tx.Model(&Message{}).Where("id = ? AND state NOT IN ?", *task.MessageID, []string{"sent", "error"}).Updates(map[string]interface{}{"state": "error", "error": reason}).Error; err != nil { return err } return tx.Create(&NotificationDelivery{MessageID: *task.MessageID, TaskID: task.ID, Status: status, Error: reason}).Error } func expireQueuedNotificationTasksTx(tx *gorm.DB, now time.Time) error { var tasks []Task if err := tx.Where("state = ? AND kind = ? AND deadline IS NOT NULL AND deadline <= ?", TaskStateQueued, TaskKindNotification, now).Find(&tasks).Error; err != nil { return err } for i := range tasks { result := tx.Model(&Task{}).Where("id = ? AND state = ?", tasks[i].ID, TaskStateQueued).Updates(map[string]interface{}{ "state": TaskStateDead, "last_error": "notification deadline expired", "updated_at": now, }) if result.Error != nil || result.RowsAffected == 0 { if result.Error != nil { return result.Error } continue } if err := FinalizeNotificationTaskTx(tx, &tasks[i], "expired", "notification deadline expired"); err != nil { return err } } return nil } // StartTaskReaper launches a goroutine that runs ReapExpiredTasks on the given // interval. It honors ctx.Done() so the caller can wind it down without // leaking. The function is safe to call once per process; the control plane // runs the reaper from main.init() so only one ticker ever exists in a single // web process. func StartTaskReaper(ctx context.Context, interval time.Duration) { if interval <= 0 { interval = 30 * time.Second } go func() { ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: reaped, deaded, err := ReapExpiredTasks() if err != nil { log.Printf("task_reaper: error: %v", err) continue } if reaped > 0 || deaded > 0 { log.Printf("task_reaper: reaped=%d dead=%d", reaped, deaded) } } } }() }