diff --git a/apps/api/internal/runner/admission.go b/apps/api/internal/runner/admission.go index fd91566..1697e3a 100644 --- a/apps/api/internal/runner/admission.go +++ b/apps/api/internal/runner/admission.go @@ -113,14 +113,10 @@ func (s *Service) buildTaskAdmissionPlan(ctx context.Context, task store.Gateway if !available { continue } - scopes, groupID := s.admissionScopes(ctx, user, candidate) - if task.AsyncMode { - scopes, err = s.withAsyncWorkerCapacityScope(ctx, scopes) - if err != nil { - return taskAdmissionPlan{}, err - } + scopes, groupID, scopeErr := s.taskAdmissionScopes(ctx, task, user, candidate) + if scopeErr != nil { + return taskAdmissionPlan{}, scopeErr } - scopes = acceptanceAdmissionScopes(task, scopes) hasConcurrentLimit := false for _, scope := range scopes { if scope.ConcurrentLimit > 0 { @@ -172,6 +168,23 @@ func (s *Service) withAsyncWorkerCapacityScope(ctx context.Context, scopes []sto return out, nil } +func (s *Service) taskAdmissionScopes( + ctx context.Context, + task store.GatewayTask, + user *auth.User, + candidate store.RuntimeModelCandidate, +) ([]store.AdmissionScope, string, error) { + scopes, groupID := s.admissionScopes(ctx, user, candidate) + if task.AsyncMode { + var err error + scopes, err = s.withAsyncWorkerCapacityScope(ctx, scopes) + if err != nil { + return nil, "", err + } + } + return acceptanceAdmissionScopes(task, scopes), groupID, nil +} + func acceptanceAdmissionScopes(task store.GatewayTask, scopes []store.AdmissionScope) []store.AdmissionScope { if task.RunMode != "acceptance" && task.RunMode != "acceptance_canary" { return scopes @@ -357,15 +370,10 @@ func (s *Service) ensureCandidateAdmission( body map[string]any, candidate store.RuntimeModelCandidate, ) (store.TaskAdmissionResult, bool, error) { - scopes, groupID := s.admissionScopes(ctx, user, candidate) - if task.AsyncMode { - var err error - scopes, err = s.withAsyncWorkerCapacityScope(ctx, scopes) - if err != nil { - return store.TaskAdmissionResult{}, true, err - } + scopes, groupID, err := s.taskAdmissionScopes(ctx, task, user, candidate) + if err != nil { + return store.TaskAdmissionResult{}, true, err } - scopes = acceptanceAdmissionScopes(task, scopes) hasConcurrentLimit := false for _, scope := range scopes { if scope.ConcurrentLimit > 0 { diff --git a/apps/api/internal/runner/service.go b/apps/api/internal/runner/service.go index 970698d..dc21835 100644 --- a/apps/api/internal/runner/service.go +++ b/apps/api/internal/runner/service.go @@ -544,8 +544,10 @@ func (s *Service) executeWithToken(ctx context.Context, task store.GatewayTask, var admittedLeases []store.ConcurrencyLease admittedPlatformModelID := "" if distributedAdmission && len(candidates) > 0 { - admissionScopes, groupID := s.admissionScopes(ctx, user, candidates[0]) - admissionScopes = acceptanceAdmissionScopes(task, admissionScopes) + admissionScopes, groupID, scopeErr := s.taskAdmissionScopes(ctx, task, user, candidates[0]) + if scopeErr != nil { + return Result{}, scopeErr + } hasConcurrentLimit := false for _, scope := range admissionScopes { if scope.ConcurrentLimit > 0 {