diff --git a/apps/api/internal/runner/admission.go b/apps/api/internal/runner/admission.go index 1697e3a..1d07937 100644 --- a/apps/api/internal/runner/admission.go +++ b/apps/api/internal/runner/admission.go @@ -603,19 +603,25 @@ func (s *Service) SubmitAsyncTask(ctx context.Context, task store.GatewayTask) e return nil } -func (s *Service) dispatchWaitingAsyncTask(ctx context.Context, task store.GatewayTask) error { +func (s *Service) dispatchWaitingAsyncTask(ctx context.Context, task store.GatewayTask) (bool, error) { user := authUserFromTask(task) plan, err := s.buildTaskAdmissionPlan(ctx, task, user) if err != nil { - return err + return false, err } if !plan.Eligible { - return s.EnqueueAsyncTask(ctx, task) + if err := s.EnqueueAsyncTask(ctx, task); err != nil { + return false, err + } + return true, nil } - _, err = s.tryTaskAdmissionWithAdmittedHook(ctx, task, plan, "", func(tx pgx.Tx) error { + result, err := s.tryTaskAdmissionWithAdmittedHook(ctx, task, plan, "", func(tx pgx.Tx) error { return s.enqueueAsyncTaskTx(ctx, tx, task.ID, asyncTaskInsertOpts(task)) }) - return err + if err != nil { + return false, err + } + return result.Admitted, nil } func (s *Service) cancelAsyncSubmissionIfDisconnected(ctx context.Context, taskID string) bool { @@ -660,8 +666,17 @@ func (s *Service) dispatchWaitingAsyncAdmissions(ctx context.Context) { if task.Status != "queued" { continue } - if err := s.dispatchWaitingAsyncTask(ctx, task); err != nil { + admitted, err := s.dispatchWaitingAsyncTask(ctx, task) + if err != nil { s.logger.Warn("dispatch waiting async admission failed", "taskID", taskID, "error", err) + continue + } + if !admitted { + // Every asynchronous task shares the global worker-capacity FIFO + // scope. Once its current head cannot be admitted, later tasks + // cannot make progress until a lease is released. + s.observeTaskAdmission("dispatch_saturated") + break } } }