MCPcopy Create free account
hub / github.com/compozy/agh / ClaimNextRun

Method ClaimNextRun

internal/task/manager_test.go:1215–1295  ·  view source on GitHub ↗
(
	_ context.Context,
	criteria ClaimCriteria,
)

Source from the content-addressed store, hash-verified

1213}
1214
1215func (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

Callers

nothing calls this directly

Calls 12

taskPriorityMinFunction · 0.85
capabilitySetContainsAllFunction · 0.85
appendFunction · 0.85
cloneTaskRunFunction · 0.85
NewClaimTokenFunction · 0.85
ClaimTokenHashFunction · 0.85
cloneTaskFunction · 0.85
cloneActorIdentityFunction · 0.70
NormalizeMethod · 0.45
NowMethod · 0.45
IsZeroMethod · 0.45
AddMethod · 0.45

Tested by

no test coverage detected