diff --git a/apps/api/internal/store/admission_queue_integration_test.go b/apps/api/internal/store/admission_queue_integration_test.go index a852b47..b1c90ac 100644 --- a/apps/api/internal/store/admission_queue_integration_test.go +++ b/apps/api/internal/store/admission_queue_integration_test.go @@ -471,6 +471,36 @@ WHERE id = $1::uuid`, queuedAtomicTask.ID).Scan(&riverJobID, &queuedLeases); err t.Fatal("generic River recovery claimed a task owned by the async admission dispatcher") } } + recoveryGraceTask := createTask(true) + recoverableTasks, err = first.ListRecoverableAsyncTasks(ctx, 1000) + if err != nil { + t.Fatalf("list fresh recoverable tasks: %v", err) + } + for _, item := range recoverableTasks { + if item.ID == recoveryGraceTask.ID { + t.Fatal("generic River recovery claimed a task still preparing admission") + } + } + if _, err := first.pool.Exec(ctx, ` +UPDATE gateway_tasks +SET created_at = now() - interval '6 minutes' +WHERE id = $1::uuid`, recoveryGraceTask.ID); err != nil { + t.Fatalf("age task beyond River recovery grace period: %v", err) + } + recoverableTasks, err = first.ListRecoverableAsyncTasks(ctx, 1000) + if err != nil { + t.Fatalf("list aged recoverable tasks: %v", err) + } + recoveryGraceFound := false + for _, item := range recoverableTasks { + if item.ID == recoveryGraceTask.ID { + recoveryGraceFound = true + break + } + } + if !recoveryGraceFound { + t.Fatal("generic River recovery did not claim a task beyond the recovery grace period") + } result, err = second.TryTaskAdmission(ctx, inputFor(queuedAtomicTask, 100, "")) if err != nil || !result.Admitted || len(result.Leases) != 1 { t.Fatalf("worker-time queued admission result=%+v err=%v", result, err) diff --git a/apps/api/internal/store/tasks_runtime.go b/apps/api/internal/store/tasks_runtime.go index ae2940c..fd594c0 100644 --- a/apps/api/internal/store/tasks_runtime.go +++ b/apps/api/internal/store/tasks_runtime.go @@ -1043,11 +1043,17 @@ func (s *Store) ListRecoverableAsyncTasks(ctx context.Context, limit int) ([]Asy if limit <= 0 { limit = 500 } + // API handlers create the durable task before they finish media restoration + // and register distributed admission. Give that normal submission path + // ownership of fresh tasks so the generic crash recovery scan cannot create + // a River job inside the create-to-admission window. + const minimumRecoveryAge = 5 * time.Minute rows, err := s.pool.Query(ctx, ` SELECT task.id::text, task.priority, task.next_run_at FROM gateway_tasks task LEFT JOIN river_job job ON job.id = task.river_job_id WHERE task.async_mode = true + AND task.created_at <= now() - $2::interval AND ( task.status = 'queued' OR ( @@ -1067,7 +1073,7 @@ WHERE task.async_mode = true AND admission.status = 'waiting' ) ORDER BY task.priority ASC, task.created_at ASC -LIMIT $1`, limit) +LIMIT $1`, limit, minimumRecoveryAge.String()) if err != nil { return nil, err }