diff --git a/apps/api/internal/runner/admission.go b/apps/api/internal/runner/admission.go index 1d07937..23b68d3 100644 --- a/apps/api/internal/runner/admission.go +++ b/apps/api/internal/runner/admission.go @@ -3,6 +3,7 @@ package runner import ( "context" "errors" + "fmt" "hash/fnv" "strings" "time" @@ -279,25 +280,42 @@ func (s *Service) tryTaskAdmissionWithAdmittedHook( waiterID string, onAdmitted func(pgx.Tx) error, ) (store.TaskAdmissionResult, error) { - input := taskAdmissionInput(task, plan, waiterID) - current, currentErr := s.store.GetTaskAdmission(ctx, task.ID) - if currentErr != nil && !errors.Is(currentErr, pgx.ErrNoRows) { - return store.TaskAdmissionResult{}, currentErr - } - if currentErr == nil && current.Status == "waiting" && - (current.PlatformID != input.PlatformID || current.PlatformModelID != input.PlatformModelID || current.UserGroupID != input.UserGroupID) { - if _, rebindErr := s.store.RebindWaitingTaskAdmission(ctx, input); rebindErr != nil { - return store.TaskAdmissionResult{}, rebindErr - } - s.observeTaskAdmission("candidate_migrated") + input, err := s.prepareTaskAdmissionInput(ctx, task, plan, waiterID) + if err != nil { + return store.TaskAdmissionResult{}, err } var result store.TaskAdmissionResult - var err error if onAdmitted == nil { result, err = s.store.TryTaskAdmission(ctx, input) } else { result, err = s.store.TryTaskAdmissionWithAdmittedHook(ctx, input, onAdmitted) } + s.observeTaskAdmissionAttempt(result, err) + return result, err +} + +func (s *Service) prepareTaskAdmissionInput( + ctx context.Context, + task store.GatewayTask, + plan taskAdmissionPlan, + waiterID string, +) (store.TaskAdmissionInput, error) { + input := taskAdmissionInput(task, plan, waiterID) + current, currentErr := s.store.GetTaskAdmission(ctx, task.ID) + if currentErr != nil && !errors.Is(currentErr, pgx.ErrNoRows) { + return store.TaskAdmissionInput{}, currentErr + } + if currentErr == nil && current.Status == "waiting" && + (current.PlatformID != input.PlatformID || current.PlatformModelID != input.PlatformModelID || current.UserGroupID != input.UserGroupID) { + if _, rebindErr := s.store.RebindWaitingTaskAdmission(ctx, input); rebindErr != nil { + return store.TaskAdmissionInput{}, rebindErr + } + s.observeTaskAdmission("candidate_migrated") + } + return input, nil +} + +func (s *Service) observeTaskAdmissionAttempt(result store.TaskAdmissionResult, err error) { if result.NewlyAdmitted { s.observeTaskAdmission("admitted") s.observeTaskAdmissionWait(time.Since(result.Admission.EnqueuedAt)) @@ -309,7 +327,6 @@ func (s *Service) tryTaskAdmissionWithAdmittedHook( case errors.Is(err, store.ErrQueueTimeout): s.observeTaskAdmission("timeout") } - return result, err } func taskAdmissionInput(task store.GatewayTask, plan taskAdmissionPlan, waiterID string) store.TaskAdmissionInput { @@ -603,25 +620,61 @@ func (s *Service) SubmitAsyncTask(ctx context.Context, task store.GatewayTask) e return nil } -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 false, err +func (s *Service) dispatchWaitingAsyncTasks(ctx context.Context, tasks []store.GatewayTask) (bool, error) { + if len(tasks) == 0 { + return false, nil } - if !plan.Eligible { - if err := s.EnqueueAsyncTask(ctx, task); err != nil { + inputs := make([]store.TaskAdmissionInput, 0, len(tasks)) + tasksByID := make(map[string]store.GatewayTask, len(tasks)) + for _, task := range tasks { + user := authUserFromTask(task) + plan, err := s.buildTaskAdmissionPlan(ctx, task, user) + if err != nil { return false, err } - return true, nil + if !plan.Eligible { + if err := s.EnqueueAsyncTask(ctx, task); err != nil { + return false, err + } + continue + } + input, err := s.prepareTaskAdmissionInput(ctx, task, plan, "") + if err != nil { + return false, err + } + inputs = append(inputs, input) + tasksByID[task.ID] = task } - result, err := s.tryTaskAdmissionWithAdmittedHook(ctx, task, plan, "", func(tx pgx.Tx) error { - return s.enqueueAsyncTaskTx(ctx, tx, task.ID, asyncTaskInsertOpts(task)) - }) + outcomes, err := s.store.TryTaskAdmissionBatchWithAdmittedHook( + ctx, + inputs, + func(tx pgx.Tx, input store.TaskAdmissionInput) error { + task, ok := tasksByID[input.TaskID] + if !ok { + return fmt.Errorf("async admission batch task %s was not prepared", input.TaskID) + } + return s.enqueueAsyncTaskTx(ctx, tx, task.ID, asyncTaskInsertOpts(task)) + }, + ) if err != nil { return false, err } - return result.Admitted, nil + saturated := false + for _, outcome := range outcomes { + s.observeTaskAdmissionAttempt(outcome.Result, outcome.Err) + if outcome.Err != nil { + if errors.Is(outcome.Err, store.ErrTaskExecutionFinished) || + errors.Is(outcome.Err, store.ErrQueueTimeout) { + continue + } + return false, outcome.Err + } + if !outcome.Result.Admitted { + saturated = true + break + } + } + return saturated, nil } func (s *Service) cancelAsyncSubmissionIfDisconnected(ctx context.Context, taskID string) bool { @@ -650,11 +703,19 @@ func (s *Service) dispatchWaitingAsyncAdmissions(ctx context.Context) { case <-s.asyncAdmissionWake: case <-ticker.C: } - taskIDs, err := s.store.ListWaitingAsyncAdmissionTaskIDs(ctx, 1000) + batchLimit := s.cfg.AsyncWorkerHardLimit + if batchLimit <= 0 { + batchLimit = 64 + } + if batchLimit > 128 { + batchLimit = 128 + } + taskIDs, err := s.store.ListWaitingAsyncAdmissionTaskIDs(ctx, batchLimit) if err != nil { s.logger.Warn("list waiting async admissions failed", "error", err) continue } + tasks := make([]store.GatewayTask, 0, len(taskIDs)) for _, taskID := range taskIDs { task, err := s.store.GetTask(ctx, taskID) if err != nil { @@ -666,18 +727,18 @@ func (s *Service) dispatchWaitingAsyncAdmissions(ctx context.Context) { if task.Status != "queued" { continue } - 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 - } + tasks = append(tasks, task) + } + saturated, err := s.dispatchWaitingAsyncTasks(ctx, tasks) + if err != nil { + s.logger.Warn("dispatch waiting async admission batch failed", "tasks", len(tasks), "error", err) + continue + } + if saturated { + // 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") } } } diff --git a/apps/api/internal/store/admission_queue.go b/apps/api/internal/store/admission_queue.go index 0a80d7b..21426ed 100644 --- a/apps/api/internal/store/admission_queue.go +++ b/apps/api/internal/store/admission_queue.go @@ -85,6 +85,12 @@ type TaskAdmissionResult struct { Leases []ConcurrencyLease } +type TaskAdmissionBatchOutcome struct { + TaskID string + Result TaskAdmissionResult + Err error +} + type TaskAdmissionMetricsSnapshot struct { QueueDepth int WaitingSync int @@ -254,23 +260,44 @@ func (s *Store) tryTaskAdmissionOnce( } defer rollbackTransaction(tx) - if err := tryAdmissionTransactionLock(ctx, tx, "task-admission:"+input.TaskID); err != nil { + outcome, err := s.tryTaskAdmissionTx(ctx, tx, input, onAdmitted) + if err != nil { return TaskAdmissionResult{}, err } + if err := tx.Commit(ctx); err != nil { + return TaskAdmissionResult{}, err + } + return outcome.Result, outcome.PostCommitErr +} + +type taskAdmissionTxOutcome struct { + Result TaskAdmissionResult + PostCommitErr error +} + +func (s *Store) tryTaskAdmissionTx( + ctx context.Context, + tx pgx.Tx, + input TaskAdmissionInput, + onAdmitted func(pgx.Tx) error, +) (taskAdmissionTxOutcome, error) { + if err := tryAdmissionTransactionLock(ctx, tx, "task-admission:"+input.TaskID); err != nil { + return taskAdmissionTxOutcome{}, err + } var taskActive bool if err := tx.QueryRow(ctx, ` SELECT status IN ('queued', 'running') FROM gateway_tasks WHERE id = $1::uuid`, input.TaskID).Scan(&taskActive); err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } if !taskActive { - return TaskAdmissionResult{}, ErrTaskExecutionFinished + return taskAdmissionTxOutcome{}, ErrTaskExecutionFinished } admission, found, err := loadTaskAdmissionTx(ctx, tx, input.TaskID) if err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } scopes := normalizedAdmissionScopes(input.Scopes) lockScopes := append([]AdmissionScope{}, scopes...) @@ -282,13 +309,13 @@ WHERE id = $1::uuid`, input.TaskID).Scan(&taskActive); err != nil { } for _, scope := range normalizedAdmissionScopes(lockScopes) { if err := tryAdmissionTransactionLock(ctx, tx, admissionLockKey(scope)); err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } } if found && admission.Status == "admitted" { activeLeases, leaseErr := activeTaskAdmissionLeaseCountTx(ctx, tx, input.TaskID) if leaseErr != nil { - return TaskAdmissionResult{}, leaseErr + return taskAdmissionTxOutcome{}, leaseErr } bindingChanged := admission.PlatformID != input.PlatformID || admission.PlatformModelID != input.PlatformModelID || @@ -297,78 +324,72 @@ WHERE id = $1::uuid`, input.TaskID).Scan(&taskActive); err != nil { result := TaskAdmissionResult{Admission: admission, Admitted: true} if onAdmitted != nil { if err := onAdmitted(tx); err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } } - if err := tx.Commit(ctx); err != nil { - return TaskAdmissionResult{}, err - } - return result, nil + return taskAdmissionTxOutcome{Result: result}, nil } if bindingChanged { targetStates, stateErr := admissionScopeStatesTx(ctx, tx, scopes) if stateErr != nil { - return TaskAdmissionResult{}, stateErr + return taskAdmissionTxOutcome{}, stateErr } if admissionErr := validateNewAdmissionCapacity(targetStates); admissionErr != nil { - return TaskAdmissionResult{}, admissionErr + return taskAdmissionTxOutcome{}, admissionErr } } admission, err = resetTaskAdmissionToWaitingTx(ctx, tx, input) if err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } } if found && !admission.WaitDeadlineAt.After(time.Now()) { if err := expireTaskAdmissionTx(ctx, tx, input.TaskID); err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } if err := notifyTaskAdmissionTx(ctx, tx, "*"); err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } - if err := tx.Commit(ctx); err != nil { - return TaskAdmissionResult{}, err - } - return TaskAdmissionResult{}, &QueueTimeoutError{TaskID: input.TaskID} + return taskAdmissionTxOutcome{ + PostCommitErr: &QueueTimeoutError{TaskID: input.TaskID}, + }, nil } scopeStates, err := admissionScopeStatesTx(ctx, tx, scopes) if err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } if !found { if admissionErr := validateNewAdmissionCapacity(scopeStates); admissionErr != nil { - return TaskAdmissionResult{}, admissionErr + return taskAdmissionTxOutcome{}, admissionErr } deadline := admissionDeadline(scopes) admission, err = insertTaskAdmissionTx(ctx, tx, input, deadline) if err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } found = true } else if input.Mode == "sync" && strings.TrimSpace(input.WaiterID) != "" { admission, err = renewTaskAdmissionWaiterTx(ctx, tx, input.TaskID, input.WaiterID) if err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } } head, err := isTaskAdmissionHeadTx(ctx, tx, admission, scopes) if err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } if !head { - if err := tx.Commit(ctx); err != nil { - return TaskAdmissionResult{}, err - } - return TaskAdmissionResult{Admission: admission}, nil + return taskAdmissionTxOutcome{ + Result: TaskAdmissionResult{Admission: admission}, + }, nil } for _, state := range scopeStates { if state.Saturated { - if err := tx.Commit(ctx); err != nil { - return TaskAdmissionResult{}, err - } - return TaskAdmissionResult{Admission: admission}, nil + return taskAdmissionTxOutcome{ + Result: TaskAdmissionResult{Admission: admission}, + }, nil } } @@ -390,34 +411,98 @@ WHERE id = $1::uuid`, input.TaskID).Scan(&taskActive); err != nil { }) if err != nil { if errors.Is(err, ErrRateLimited) { - if err := tx.Commit(ctx); err != nil { - return TaskAdmissionResult{}, err - } - return TaskAdmissionResult{Admission: admission}, nil + return taskAdmissionTxOutcome{ + Result: TaskAdmissionResult{Admission: admission}, + }, nil } - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } leases = append(leases, lease) } admission, err = markTaskAdmissionAdmittedTx(ctx, tx, input.TaskID) if err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } // Continue the FIFO chain after one task is admitted. Every API process // only probes the global head it owns, so free capacity is filled without // broadcasting an advisory-lock attempt to every waiter. if err := notifyTaskAdmissionTx(ctx, tx, "*"); err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } if onAdmitted != nil { if err := onAdmitted(tx); err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{}, err } } - if err := tx.Commit(ctx); err != nil { - return TaskAdmissionResult{}, err + return taskAdmissionTxOutcome{ + Result: TaskAdmissionResult{ + Admission: admission, + Admitted: true, + NewlyAdmitted: true, + Leases: leases, + }, + }, nil +} + +// TryTaskAdmissionBatchWithAdmittedHook fills a FIFO capacity window in one +// transaction. The hook inserts every unique River job in that same +// transaction, preserving the admission/job crash boundary while avoiding one +// synchronous replication round trip per task. +func (s *Store) TryTaskAdmissionBatchWithAdmittedHook( + ctx context.Context, + inputs []TaskAdmissionInput, + onAdmitted func(pgx.Tx, TaskAdmissionInput) error, +) ([]TaskAdmissionBatchOutcome, error) { + if len(inputs) == 0 { + return nil, nil } - return TaskAdmissionResult{Admission: admission, Admitted: true, NewlyAdmitted: true, Leases: leases}, nil + lockKeys := make([]string, 0, len(inputs)*3) + for _, input := range inputs { + if err := validateTaskAdmissionInput(input); err != nil { + return nil, err + } + lockKeys = append(lockKeys, admissionOperationLockKeys(input)...) + } + return retryAdmissionOperation(ctx, lockKeys, func() ([]TaskAdmissionBatchOutcome, error) { + tx, err := s.pool.Begin(ctx) + if err != nil { + return nil, err + } + defer rollbackTransaction(tx) + + outcomes := make([]TaskAdmissionBatchOutcome, 0, len(inputs)) + for _, input := range inputs { + hook := func(tx pgx.Tx) error { + if onAdmitted == nil { + return nil + } + return onAdmitted(tx, input) + } + outcome, admissionErr := s.tryTaskAdmissionTx(ctx, tx, input, hook) + if errors.Is(admissionErr, ErrTaskExecutionFinished) { + outcomes = append(outcomes, TaskAdmissionBatchOutcome{ + TaskID: input.TaskID, + Err: admissionErr, + }) + continue + } + if admissionErr != nil { + return nil, admissionErr + } + outcomes = append(outcomes, TaskAdmissionBatchOutcome{ + TaskID: input.TaskID, + Result: outcome.Result, + Err: outcome.PostCommitErr, + }) + if outcome.PostCommitErr == nil && !outcome.Result.Admitted { + break + } + } + if err := tx.Commit(ctx); err != nil { + return nil, err + } + return outcomes, nil + }) } func validateNewAdmissionCapacity(states []admissionScopeState) error { diff --git a/apps/api/internal/store/admission_queue_integration_test.go b/apps/api/internal/store/admission_queue_integration_test.go index b1c90ac..f2522e3 100644 --- a/apps/api/internal/store/admission_queue_integration_test.go +++ b/apps/api/internal/store/admission_queue_integration_test.go @@ -501,6 +501,140 @@ WHERE id = $1::uuid`, recoveryGraceTask.ID); err != nil { if !recoveryGraceFound { t.Fatal("generic River recovery did not claim a task beyond the recovery grace period") } + + batchScope := AdmissionScope{ + ScopeType: "worker_capacity", + ScopeKey: "batch-" + suffix, + ScopeName: "batch capacity", + ConcurrentLimit: 3, + Amount: 1, + LeaseTTLSeconds: 120, + QueueLimit: 100, + MaxWaitSeconds: 600, + } + batchTasks := []GatewayTask{createTask(true), createTask(true), createTask(true)} + batchInputs := make([]TaskAdmissionInput, 0, len(batchTasks)) + batchJobIDs := make(map[string]int64, len(batchTasks)) + for index, task := range batchTasks { + input := inputFor(task, 100+index, "") + input.Scopes = []AdmissionScope{batchScope} + if _, err := first.QueueTaskAdmissionWithHook(ctx, input, nil); err != nil { + t.Fatalf("queue batch task %d: %v", index, err) + } + batchInputs = append(batchInputs, input) + batchJobIDs[task.ID] = int64(987654400 + index) + } + batchOutcomes, err := first.TryTaskAdmissionBatchWithAdmittedHook( + ctx, + batchInputs, + func(tx pgx.Tx, input TaskAdmissionInput) error { + _, hookErr := tx.Exec(ctx, ` +UPDATE gateway_tasks +SET river_job_id = $2 +WHERE id = $1::uuid`, input.TaskID, batchJobIDs[input.TaskID]) + return hookErr + }, + ) + if err != nil || len(batchOutcomes) != len(batchTasks) { + t.Fatalf("batch admission outcomes=%+v err=%v", batchOutcomes, err) + } + for index, outcome := range batchOutcomes { + if outcome.Err != nil || !outcome.Result.Admitted || !outcome.Result.NewlyAdmitted || + len(outcome.Result.Leases) != 1 { + t.Fatalf("batch outcome %d=%+v", index, outcome) + } + } + var batchAdmissions, batchLeases, batchRiverJobs int + if err := first.pool.QueryRow(ctx, ` +SELECT + (SELECT count(*) + FROM gateway_task_admissions + WHERE task_id = ANY($1::uuid[]) AND status = 'admitted'), + (SELECT count(*) + FROM gateway_concurrency_leases + WHERE task_id = ANY($1::uuid[]) AND released_at IS NULL AND expires_at > now()), + (SELECT count(*) + FROM gateway_tasks + WHERE id = ANY($1::uuid[]) AND river_job_id IS NOT NULL)`, + []string{batchTasks[0].ID, batchTasks[1].ID, batchTasks[2].ID}, + ).Scan(&batchAdmissions, &batchLeases, &batchRiverJobs); err != nil { + t.Fatalf("read committed batch admission: %v", err) + } + if batchAdmissions != 3 || batchLeases != 3 || batchRiverJobs != 3 { + t.Fatalf( + "batch admissions=%d leases=%d River jobs=%d, want 3/3/3", + batchAdmissions, + batchLeases, + batchRiverJobs, + ) + } + for _, task := range batchTasks { + if err := first.DeleteTaskAdmission(ctx, task.ID); err != nil { + t.Fatalf("release batch task %s: %v", task.ID, err) + } + } + + rollbackScope := batchScope + rollbackScope.ScopeKey = "batch-rollback-" + suffix + rollbackScope.ConcurrentLimit = 2 + rollbackTasks := []GatewayTask{createTask(true), createTask(true)} + rollbackInputs := make([]TaskAdmissionInput, 0, len(rollbackTasks)) + for index, task := range rollbackTasks { + input := inputFor(task, 200+index, "") + input.Scopes = []AdmissionScope{rollbackScope} + if _, err := first.QueueTaskAdmissionWithHook(ctx, input, nil); err != nil { + t.Fatalf("queue rollback batch task %d: %v", index, err) + } + rollbackInputs = append(rollbackInputs, input) + } + batchHookFailure := errors.New("synthetic batch hook failure") + _, err = first.TryTaskAdmissionBatchWithAdmittedHook( + ctx, + rollbackInputs, + func(tx pgx.Tx, input TaskAdmissionInput) error { + if input.TaskID == rollbackTasks[1].ID { + return batchHookFailure + } + _, hookErr := tx.Exec(ctx, ` +UPDATE gateway_tasks +SET river_job_id = 987654499 +WHERE id = $1::uuid`, input.TaskID) + return hookErr + }, + ) + if !errors.Is(err, batchHookFailure) { + t.Fatalf("batch hook failure=%v, want synthetic failure", err) + } + var rollbackWaiting, rollbackLeases, rollbackRiverJobs int + if err := first.pool.QueryRow(ctx, ` +SELECT + (SELECT count(*) + FROM gateway_task_admissions + WHERE task_id = ANY($1::uuid[]) AND status = 'waiting'), + (SELECT count(*) + FROM gateway_concurrency_leases + WHERE task_id = ANY($1::uuid[]) AND released_at IS NULL), + (SELECT count(*) + FROM gateway_tasks + WHERE id = ANY($1::uuid[]) AND river_job_id IS NOT NULL)`, + []string{rollbackTasks[0].ID, rollbackTasks[1].ID}, + ).Scan(&rollbackWaiting, &rollbackLeases, &rollbackRiverJobs); err != nil { + t.Fatalf("read rolled back batch admission: %v", err) + } + if rollbackWaiting != 2 || rollbackLeases != 0 || rollbackRiverJobs != 0 { + t.Fatalf( + "rolled back batch waiting=%d leases=%d River jobs=%d, want 2/0/0", + rollbackWaiting, + rollbackLeases, + rollbackRiverJobs, + ) + } + for _, task := range rollbackTasks { + if err := first.DeleteTaskAdmission(ctx, task.ID); err != nil { + t.Fatalf("release rollback batch task %s: %v", task.ID, err) + } + } + 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/scripts/cluster/run-production-acceptance.sh b/scripts/cluster/run-production-acceptance.sh index 2dd5a20..7eca470 100755 --- a/scripts/cluster/run-production-acceptance.sh +++ b/scripts/cluster/run-production-acceptance.sh @@ -911,6 +911,7 @@ run_load_profile() { load_pod=$(remote_kubectl get pods -n "$namespace" \ -l 'app.kubernetes.io/name=easyai-acceptance-emulator' -o json | jq -r '[.items[] + | select(.metadata.deletionTimestamp == null) | select(any(.status.conditions[]?; .type=="Ready" and .status=="True")) | .metadata.name][0] // empty') [[ $load_pod =~ ^easyai-acceptance-emulator-[a-z0-9-]+$ ]] || { @@ -1339,6 +1340,7 @@ WHERE acceptance_run_id='$run_id'::uuid;"); then pod=$(remote_kubectl get pods -n "$namespace" \ -l "app.kubernetes.io/name=easyai-$workload,easyai.io/site=$site" -o json | jq -r '[.items[] + | select(.metadata.deletionTimestamp == null) | select(any(.status.conditions[]?; .type=="Ready" and .status=="True")) | .metadata.name][0] // empty') || exit [[ -n $pod ]] || continue