返回 DeepSeek-Reasonix
scheduler.go
根目录 / internal / agent / scheduler.go
1 package agent
2
3 import (
4 "context"
5 "fmt"
6 "sync"
7 )
8
9 // SubagentSlotStatus is the queue lifecycle shown for background task/fleet
10 // items that share the session scheduler.
11 type SubagentSlotStatus string
12
13 const (
14 SubagentSlotQueued SubagentSlotStatus = "queued"
15 SubagentSlotRunning SubagentSlotStatus = "running"
16 SubagentSlotDone SubagentSlotStatus = "done"
17 SubagentSlotFailed SubagentSlotStatus = "failed"
18 )
19
20 // AcquireRequest describes a sub-agent slot request against the session pool.
21 type AcquireRequest struct {
22 // Writer is true for writer-capable runs (task without read_only, profile
23 // that is not read-only, fleet items that can write).
24 Writer bool
25 // WritePaths is the claim held while the slot is active. Empty for
26 // read-only work. Whole-workspace claims count as writers and serialize
27 // against every other writer.
28 WritePaths WritePathSet
29 // Nested fails immediately when no capacity is free instead of queueing.
30 // Nested sub-agents must not block waiting for a parent-held slot.
31 Nested bool
32 // Label is optional diagnostics text.
33 Label string
34 // callerParentClaim is the parent write claim held by the tool call that
35 // issued this request, read from the Acquire context.
36 callerParentClaim int64
37 }
38
39 // SubagentScheduler is a session-scoped concurrency controller shared by task,
40 // fleet, parallel_tasks, profile skills, and nested sub-agents.
41 type SubagentScheduler struct {
42 mu sync.Mutex
43
44 maxTotal int
45 maxWriters int
46
47 activeTotal int
48 activeWriters int
49 activeLive []liveClaim
50 nextClaimID int64
51 // parentClaims are write paths held by the parent agent during a write-tool
52 // Execute. They block overlapping subagent claims without consuming a
53 // subagent concurrency slot (parent is not a subagent).
54 parentClaims []parentWriteClaim
55
56 // waiters are FIFO waiters for non-nested acquires.
57 waiters []*schedulerWaiter
58 }
59
60 type schedulerWaiter struct {
61 req AcquireRequest
62 ready chan struct{}
63 failed error
64 id int64
65 }
66
67 // NewSubagentScheduler builds a scheduler with the given limits (normalized).
68 func NewSubagentScheduler(maxTotal, maxWriters int) *SubagentScheduler {
69 maxTotal, maxWriters = NormalizeConcurrencyLimits(maxTotal, maxWriters)
70 return &SubagentScheduler{maxTotal: maxTotal, maxWriters: maxWriters}
71 }
72
73 // Limits returns the effective total/writer caps.
74 func (s *SubagentScheduler) Limits() (total, writers int) {
75 if s == nil {
76 return DefaultMaxSubagentConcurrency, DefaultMaxParallelWriters
77 }
78 return s.maxTotal, s.maxWriters
79 }
80
81 // Acquire reserves a concurrency slot (and optional write claim). Nested
82 // requests fail immediately when capacity is exhausted. Non-nested requests
83 // queue until capacity is free or ctx is cancelled.
84 //
85 // The returned release function must be called exactly once when the sub-agent
86 // finishes. release is safe to call even if Acquire returns an error (no-op).
87 func (s *SubagentScheduler) Acquire(ctx context.Context, req AcquireRequest) (release func(), err error) {
88 release, _, err = s.AcquireWithID(ctx, req)
89 return release, err
90 }
91
92 // AcquireWithID is Acquire plus the live claim id used by Realize/MarkOpaque.
93 func (s *SubagentScheduler) AcquireWithID(ctx context.Context, req AcquireRequest) (release func(), claimID int64, err error) {
94 noop := func() {}
95 if s == nil {
96 return noop, 0, nil
97 }
98 if ctx == nil {
99 ctx = context.Background()
100 }
101 req.callerParentClaim = ParentWriteClaimID(ctx)
102
103 s.mu.Lock()
104 if ok, reason := s.canStartIncomingLocked(req); ok {
105 id := s.activateLocked(req)
106 s.mu.Unlock()
107 return s.makeReleaseID(id), id, nil
108 } else if req.Nested {
109 s.mu.Unlock()
110 return noop, 0, fmt.Errorf("subagent concurrency limit reached (%s); nested subagents fail fast to avoid parent/child slot deadlock", reason)
111 }
112
113 w := &schedulerWaiter{req: req, ready: make(chan struct{})}
114 s.waiters = append(s.waiters, w)
115 s.mu.Unlock()
116
117 select {
118 case <-w.ready:
119 if w.failed != nil {
120 return noop, 0, w.failed
121 }
122 return s.makeReleaseID(w.id), w.id, nil
123 case <-ctx.Done():
124 s.mu.Lock()
125 s.removeWaiterLocked(w)
126 s.pumpWaitersLocked()
127 s.mu.Unlock()
128 select {
129 case <-w.ready:
130 if w.failed == nil {
131 s.makeReleaseID(w.id)()
132 }
133 default:
134 }
135 return noop, 0, ctx.Err()
136 }
137 }
138
139 // TryClaimWritePaths checks whether paths conflict with active claims without
140 // taking a concurrency slot. Used for diagnostics; prefer ReserveParentWrite
141 // for parent agent writes so the check is not TOCTOU with subagent Acquire.
142 func (s *SubagentScheduler) TryClaimWritePaths(paths WritePathSet) error {
143 if s == nil || paths.Empty() {
144 return nil
145 }
146 s.mu.Lock()
147 defer s.mu.Unlock()
148 return s.conflictLocked(paths)
149 }
150
151 // Realize records path-bound writes against an active claim. Directory and
152 // whole-workspace declarations shrink to the realized files when no opaque
153 // mutation has occurred. Same-file realizes from two live writers fail.
154 func (s *SubagentScheduler) Realize(id int64, paths WritePathSet) error {
155 if s == nil || id == 0 || paths.Empty() {
156 return nil
157 }
158 s.mu.Lock()
159 defer s.mu.Unlock()
160 idx := s.liveIndexLocked(id)
161 if idx < 0 {
162 return fmt.Errorf("write path is claimed by a running background subagent; wait for it to finish before writing the same path")
163 }
164 claim := s.activeLive[idx]
165 if claim.opaque {
166 return nil
167 }
168 nextPaths := mergeRealized(claim.realized, paths)
169 next := fileReservation(claim.declared.WorkspaceRoot, nextPaths)
170 if err := s.conflictAgainstOthersLocked(id, next); err != nil {
171 return err
172 }
173 claim.realized = nextPaths
174 s.activeLive[idx] = claim
175 s.pumpWaitersLocked()
176 return nil
177 }
178
179 // MarkOpaque upgrades a live claim to a whole-workspace reservation (bash/MCP).
180 func (s *SubagentScheduler) MarkOpaque(id int64) error {
181 if s == nil || id == 0 {
182 return nil
183 }
184 s.mu.Lock()
185 defer s.mu.Unlock()
186 idx := s.liveIndexLocked(id)
187 if idx < 0 {
188 return fmt.Errorf("write path is claimed by a running background subagent; wait for it to finish before writing the same path")
189 }
190 claim := s.activeLive[idx]
191 if claim.opaque {
192 return nil
193 }
194 next := wholeReservation(claim.declared.WorkspaceRoot)
195 if err := s.conflictAgainstOthersLocked(id, next); err != nil {
196 return err
197 }
198 claim.opaque = true
199 s.activeLive[idx] = claim
200 return nil
201 }
202
203 // ReserveParentWrite holds paths against overlapping subagent claims for the
204 // duration of a parent write-tool Execute. It does not consume subagent
205 // concurrency slots. On conflict it fails immediately (parent cannot queue
206 // behind background jobs mid-tool-call). release must be called once when the
207 // write finishes so queued subagents can proceed.
208 func (s *SubagentScheduler) ReserveParentWrite(paths WritePathSet) (release func(), err error) {
209 release, _, err = s.ReserveParentWriteWithID(paths)
210 return release, err
211 }
212
213 // ReserveParentWriteWithID is ReserveParentWrite plus the claim id the tool
214 // call carries into its Execute context (see WithParentWriteClaimID).
215 func (s *SubagentScheduler) ReserveParentWriteWithID(paths WritePathSet) (release func(), claimID int64, err error) {
216 noop := func() {}
217 if s == nil || paths.Empty() {
218 return noop, 0, nil
219 }
220 s.mu.Lock()
221 if err := s.conflictLocked(paths); err != nil {
222 s.mu.Unlock()
223 return noop, 0, err
224 }
225 s.nextClaimID++
226 id := s.nextClaimID
227 s.parentClaims = append(s.parentClaims, parentWriteClaim{id: id, paths: paths})
228 s.mu.Unlock()
229
230 var once sync.Once
231 return func() {
232 once.Do(func() {
233 s.mu.Lock()
234 s.parentClaims = removeParentClaim(s.parentClaims, id)
235 s.pumpWaitersLocked()
236 s.mu.Unlock()
237 })
238 }, id, nil
239 }
240
241 // ActiveWriterClaims returns a snapshot of subagent + parent write claims.
242 func (s *SubagentScheduler) ActiveWriterClaims() []WritePathSet {
243 if s == nil {
244 return nil
245 }
246 s.mu.Lock()
247 defer s.mu.Unlock()
248 out := make([]WritePathSet, 0, len(s.activeLive)+len(s.parentClaims))
249 for _, live := range s.activeLive {
250 if !live.writer {
251 continue
252 }
253 res := live.reservation()
254 if res.Empty() {
255 if live.declared.Empty() {
256 continue
257 }
258 res = live.declared
259 }
260 out = append(out, res)
261 }
262 for _, parent := range s.parentClaims {
263 out = append(out, parent.paths)
264 }
265 return out
266 }
267
268 func (s *SubagentScheduler) conflictLocked(paths WritePathSet) error {
269 if paths.WholeWorkspace {
270 for _, live := range s.activeLive {
271 if live.writer {
272 return fmt.Errorf("write path is claimed by a running background subagent; wait for it to finish before writing the same path")
273 }
274 }
275 }
276 return s.conflictAgainstOthersLocked(0, paths)
277 }
278
279 func (s *SubagentScheduler) conflictAgainstOthersLocked(skipID int64, paths WritePathSet) error {
280 if paths.Empty() {
281 return nil
282 }
283 var callerParent int64
284 if idx := s.liveIndexLocked(skipID); skipID != 0 && idx >= 0 {
285 callerParent = s.activeLive[idx].callerParentClaim
286 }
287 for _, live := range s.activeLive {
288 if live.id == skipID {
289 continue
290 }
291 if ScheduleOverlaps(live.reservation(), paths) {
292 return fmt.Errorf("write path is claimed by a running background subagent; wait for it to finish before writing the same path")
293 }
294 }
295 for _, active := range s.parentClaims {
296 if active.id != callerParent && ScheduleOverlaps(active.paths, paths) {
297 return fmt.Errorf("write path is claimed by another parent write in progress")
298 }
299 }
300 return nil
301 }
302
303 func (s *SubagentScheduler) makeReleaseID(id int64) func() {
304 var once sync.Once
305 return func() {
306 once.Do(func() {
307 s.mu.Lock()
308 s.deactivateIDLocked(id)
309 s.pumpWaitersLocked()
310 s.mu.Unlock()
311 })
312 }
313 }
314
315 func (s *SubagentScheduler) liveIndexLocked(id int64) int {
316 for i, live := range s.activeLive {
317 if live.id == id {
318 return i
319 }
320 }
321 return -1
322 }
323
324 func (s *SubagentScheduler) canStartLocked(req AcquireRequest) (bool, string) {
325 if s.activeTotal >= s.maxTotal {
326 return false, fmt.Sprintf("total concurrency %d/%d", s.activeTotal, s.maxTotal)
327 }
328 if !req.Writer {
329 return true, ""
330 }
331 if s.activeWriters >= s.maxWriters {
332 return false, fmt.Sprintf("writer concurrency %d/%d", s.activeWriters, s.maxWriters)
333 }
334 if req.WritePaths.WholeWorkspace {
335 for _, live := range s.activeLive {
336 if live.writer {
337 return false, "whole-workspace claim conflicts with a running writer"
338 }
339 }
340 }
341 for _, live := range s.activeLive {
342 if ScheduleOverlaps(req.WritePaths, live.reservation()) {
343 return false, "write path conflict with a running subagent"
344 }
345 }
346 for _, active := range s.parentClaims {
347 if active.id != req.callerParentClaim && ScheduleOverlaps(req.WritePaths, active.paths) {
348 return false, "write path conflict with a parent write in progress"
349 }
350 }
351 return true, ""
352 }
353
354 func (s *SubagentScheduler) activateLocked(req AcquireRequest) int64 {
355 s.activeTotal++
356 s.nextClaimID++
357 id := s.nextClaimID
358 if req.Writer {
359 s.activeWriters++
360 }
361 s.activeLive = append(s.activeLive, liveClaim{id: id, writer: req.Writer, declared: req.WritePaths, callerParentClaim: req.callerParentClaim})
362 return id
363 }
364
365 func (s *SubagentScheduler) deactivateIDLocked(id int64) {
366 idx := s.liveIndexLocked(id)
367 if idx < 0 {
368 return
369 }
370 if s.activeTotal > 0 {
371 s.activeTotal--
372 }
373 if s.activeLive[idx].writer && s.activeWriters > 0 {
374 s.activeWriters--
375 }
376 s.activeLive = append(s.activeLive[:idx], s.activeLive[idx+1:]...)
377 }
378
379 func (s *SubagentScheduler) pumpWaitersLocked() {
380 if len(s.waiters) == 0 {
381 return
382 }
383 remaining := s.waiters[:0]
384 // A blocked whole-workspace writer is a FIFO barrier for later writers,
385 // while read-only work may still use otherwise available capacity.
386 wholeWriterPending := false
387 for _, w := range s.waiters {
388 if wholeWriterPending && w.req.Writer && !s.holdsParentClaimLocked(w.req.callerParentClaim) {
389 remaining = append(remaining, w)
390 continue
391 }
392 if ok, _ := s.canStartLocked(w.req); ok {
393 w.id = s.activateLocked(w.req)
394 close(w.ready)
395 continue
396 }
397 remaining = append(remaining, w)
398 if w.req.Writer && w.req.WritePaths.WholeWorkspace {
399 wholeWriterPending = true
400 }
401 }
402 s.waiters = remaining
403 }
404
405 func (s *SubagentScheduler) removeWaiterLocked(target *schedulerWaiter) {
406 if len(s.waiters) == 0 {
407 return
408 }
409 out := s.waiters[:0]
410 for _, w := range s.waiters {
411 if w == target {
412 continue
413 }
414 out = append(out, w)
415 }
416 s.waiters = out
417 }
418
419 func removeParentClaim(claims []parentWriteClaim, id int64) []parentWriteClaim {
420 for i, c := range claims {
421 if c.id == id {
422 return append(claims[:i], claims[i+1:]...)
423 }
424 }
425 return claims
426 }
427
427 lines GO