( _ context.Context, criteria ClaimCriteria, )
| 1213 | } |
| 1214 | |
| 1215 | func (s *inMemoryManagerStore) ClaimNextRun( |
| 1216 | _ context.Context, |
| 1217 | criteria ClaimCriteria, |
| 1218 | ) (ClaimResult, error) { |
| 1219 | normalized, err := criteria.Normalize(time.Now().UTC()) |
| 1220 | if err != nil { |
| 1221 | return ClaimResult{}, err |
| 1222 | } |
| 1223 | for _, run := range s.runs { |
| 1224 | if run.SessionID == normalized.ClaimerSessionID && |
| 1225 | (run.Status == TaskRunStatusClaimed || run.Status == TaskRunStatusStarting || run.Status == TaskRunStatusRunning) && |
| 1226 | (run.LeaseUntil.IsZero() || run.LeaseUntil.After(normalized.Now)) { |
| 1227 | return ClaimResult{}, ErrActiveRunLease |
| 1228 | } |
| 1229 | } |
| 1230 | |
| 1231 | candidates := make([]Run, 0) |
| 1232 | for _, run := range s.runs { |
| 1233 | taskRecord, ok := s.tasks[run.TaskID] |
| 1234 | if !ok || run.Status.Normalize() != TaskRunStatusQueued { |
| 1235 | continue |
| 1236 | } |
| 1237 | switch taskRecord.Status.Normalize() { |
| 1238 | case TaskStatusDraft, TaskStatusBlocked, TaskStatusNeedsAttention, TaskStatusCanceled: |
| 1239 | continue |
| 1240 | } |
| 1241 | if taskRecord.Scope.Normalize() != normalized.Scope { |
| 1242 | continue |
| 1243 | } |
| 1244 | if normalized.Scope == ScopeWorkspace && taskRecord.WorkspaceID != normalized.WorkspaceID { |
| 1245 | continue |
| 1246 | } |
| 1247 | if taskPriorityMin(taskRecord.Priority) < normalized.PriorityMin { |
| 1248 | continue |
| 1249 | } |
| 1250 | if !capabilitySetContainsAll(normalized.RequiredCapabilities, run.RequiredCapabilities) { |
| 1251 | continue |
| 1252 | } |
| 1253 | candidates = append(candidates, cloneTaskRun(run)) |
| 1254 | } |
| 1255 | sort.Slice(candidates, func(i int, j int) bool { |
| 1256 | leftTask := s.tasks[candidates[i].TaskID] |
| 1257 | rightTask := s.tasks[candidates[j].TaskID] |
| 1258 | if taskPriorityMin(leftTask.Priority) != taskPriorityMin(rightTask.Priority) { |
| 1259 | return taskPriorityMin(leftTask.Priority) > taskPriorityMin(rightTask.Priority) |
| 1260 | } |
| 1261 | if !candidates[i].QueuedAt.Equal(candidates[j].QueuedAt) { |
| 1262 | return candidates[i].QueuedAt.Before(candidates[j].QueuedAt) |
| 1263 | } |
| 1264 | return candidates[i].ID < candidates[j].ID |
| 1265 | }) |
| 1266 | if len(candidates) == 0 { |
| 1267 | return ClaimResult{}, ErrNoClaimableRun |
| 1268 | } |
| 1269 | |
| 1270 | token, err := NewClaimToken() |
| 1271 | if err != nil { |
| 1272 | return ClaimResult{}, err |
nothing calls this directly
no test coverage detected