137 строки
5.0 KiB
Go
137 строки
5.0 KiB
Go
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)
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
}
|