Files
worker/app/models/task_reaper.go
Gleb Tv 2c7a0236da feat: publish standalone worker
Separate worker packaging and service lifecycle from the control plane.
2026-07-13 17:55:14 +03:00

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)
}
}
}
}()
}