返回 DeepSeek-Reasonix
scheduler_test.go
根目录 / internal / agent / scheduler_test.go
1 package agent
2
3 import (
4 "context"
5 "os"
6 "path/filepath"
7 "sync"
8 "sync/atomic"
9 "testing"
10 "time"
11 )
12
13 func TestSchedulerTotalConcurrencyQueues(t *testing.T) {
14 s := NewSubagentScheduler(2, 2)
15 root := t.TempDir()
16 var started atomic.Int32
17 var max atomic.Int32
18 var wg sync.WaitGroup
19 barrier := make(chan struct{})
20
21 for range 4 {
22 wg.Go(func() {
23 release, err := s.Acquire(context.Background(), AcquireRequest{Writer: false})
24 if err != nil {
25 t.Errorf("acquire: %v", err)
26 return
27 }
28 cur := started.Add(1)
29 for {
30 old := max.Load()
31 if cur <= old || max.CompareAndSwap(old, cur) {
32 break
33 }
34 }
35 <-barrier
36 started.Add(-1)
37 release()
38 })
39 }
40
41 // Wait until at least 2 are running, then release them.
42 deadline := time.Now().Add(2 * time.Second)
43 for time.Now().Before(deadline) {
44 if max.Load() >= 2 {
45 break
46 }
47 time.Sleep(5 * time.Millisecond)
48 }
49 if got := max.Load(); got > 2 {
50 t.Fatalf("max concurrent = %d, want <= 2", got)
51 }
52 close(barrier)
53 wg.Wait()
54 _ = root
55 }
56
57 func TestSchedulerNestedFailsFast(t *testing.T) {
58 s := NewSubagentScheduler(1, 1)
59 release, err := s.Acquire(context.Background(), AcquireRequest{Writer: false})
60 if err != nil {
61 t.Fatal(err)
62 }
63 defer release()
64 _, err = s.Acquire(context.Background(), AcquireRequest{Writer: false, Nested: true})
65 if err == nil {
66 t.Fatal("nested acquire should fail fast at limit")
67 }
68 }
69
70 func TestSchedulerWriterPathConflictQueues(t *testing.T) {
71 s := NewSubagentScheduler(4, 2)
72 root := t.TempDir()
73 claim, err := NormalizeWritePaths(root, []string{"a.md"})
74 if err != nil {
75 t.Fatal(err)
76 }
77 release, err := s.Acquire(context.Background(), AcquireRequest{Writer: true, WritePaths: claim})
78 if err != nil {
79 t.Fatal(err)
80 }
81
82 ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
83 defer cancel()
84 // Same path cannot start while the first claim is held — with Nested it fails.
85 _, err = s.Acquire(ctx, AcquireRequest{Writer: true, WritePaths: claim, Nested: true})
86 if err == nil {
87 t.Fatal("expected path conflict for nested acquire")
88 }
89 release()
90
91 // After release, same path is free.
92 release2, err := s.Acquire(context.Background(), AcquireRequest{Writer: true, WritePaths: claim})
93 if err != nil {
94 t.Fatal(err)
95 }
96 release2()
97 }
98
99 func TestSchedulerDirectoryClaimsStartInParallel(t *testing.T) {
100 s := NewSubagentScheduler(4, 2)
101 root := t.TempDir()
102 if err := os.MkdirAll(filepath.Join(root, "src"), 0o755); err != nil {
103 t.Fatal(err)
104 }
105 claim, err := NormalizeWritePaths(root, []string{"src/"})
106 if err != nil {
107 t.Fatal(err)
108 }
109 release1, id1, err := s.AcquireWithID(context.Background(), AcquireRequest{Writer: true, WritePaths: claim})
110 if err != nil {
111 t.Fatal(err)
112 }
113 defer release1()
114 release2, id2, err := s.AcquireWithID(context.Background(), AcquireRequest{Writer: true, WritePaths: claim, Nested: true})
115 if err != nil {
116 t.Fatalf("second directory claim must start: %v", err)
117 }
118 defer release2()
119 if id1 == 0 || id2 == 0 || id1 == id2 {
120 t.Fatalf("claim ids = %d, %d", id1, id2)
121 }
122 }
123
124 func TestSchedulerWholeClaimCannotStartBehindUnrealizedDirectoryWriter(t *testing.T) {
125 s := NewSubagentScheduler(4, 2)
126 root := t.TempDir()
127 if err := os.MkdirAll(filepath.Join(root, "src"), 0o755); err != nil {
128 t.Fatal(err)
129 }
130 dir, err := NormalizeWritePaths(root, []string{"src/"})
131 if err != nil {
132 t.Fatal(err)
133 }
134 releaseDir, _, err := s.AcquireWithID(context.Background(), AcquireRequest{Writer: true, WritePaths: dir})
135 if err != nil {
136 t.Fatal(err)
137 }
138 defer releaseDir()
139 whole, err := WholeWorkspaceWriteClaim(root)
140 if err != nil {
141 t.Fatal(err)
142 }
143 releaseWhole, _, err := s.AcquireWithID(context.Background(), AcquireRequest{
144 Writer: true, WritePaths: whole, Nested: true,
145 })
146 if err == nil {
147 releaseWhole()
148 t.Fatal("whole-workspace claim bypassed an active unrealized directory writer")
149 }
150 releaseDir()
151 releaseWhole, _, err = s.AcquireWithID(context.Background(), AcquireRequest{
152 Writer: true, WritePaths: whole, Nested: true,
153 })
154 if err != nil {
155 t.Fatalf("whole-workspace claim after directory writer release: %v", err)
156 }
157 releaseWhole()
158 }
159
160 func TestSchedulerRealizeSameFileConflicts(t *testing.T) {
161 s := NewSubagentScheduler(4, 2)
162 root := t.TempDir()
163 if err := os.MkdirAll(filepath.Join(root, "src"), 0o755); err != nil {
164 t.Fatal(err)
165 }
166 claim, err := NormalizeWritePaths(root, []string{"src/"})
167 if err != nil {
168 t.Fatal(err)
169 }
170 _, id1, err := s.AcquireWithID(context.Background(), AcquireRequest{Writer: true, WritePaths: claim})
171 if err != nil {
172 t.Fatal(err)
173 }
174 _, id2, err := s.AcquireWithID(context.Background(), AcquireRequest{Writer: true, WritePaths: claim})
175 if err != nil {
176 t.Fatal(err)
177 }
178 file, err := NormalizeWritePaths(root, []string{"src/a.go"})
179 if err != nil {
180 t.Fatal(err)
181 }
182 if err := s.Realize(id1, file); err != nil {
183 t.Fatalf("first realize: %v", err)
184 }
185 if err := s.Realize(id2, file); err == nil {
186 t.Fatal("second realize of the same file must fail")
187 }
188 other, err := NormalizeWritePaths(root, []string{"src/b.go"})
189 if err != nil {
190 t.Fatal(err)
191 }
192 if err := s.Realize(id2, other); err != nil {
193 t.Fatalf("disjoint realize: %v", err)
194 }
195 }
196
197 func TestSchedulerMarkOpaqueBlocksRealize(t *testing.T) {
198 s := NewSubagentScheduler(4, 2)
199 root := t.TempDir()
200 if err := os.MkdirAll(filepath.Join(root, "src"), 0o755); err != nil {
201 t.Fatal(err)
202 }
203 claim, err := NormalizeWritePaths(root, []string{"src/"})
204 if err != nil {
205 t.Fatal(err)
206 }
207 _, id1, err := s.AcquireWithID(context.Background(), AcquireRequest{Writer: true, WritePaths: claim})
208 if err != nil {
209 t.Fatal(err)
210 }
211 _, id2, err := s.AcquireWithID(context.Background(), AcquireRequest{Writer: true, WritePaths: claim})
212 if err != nil {
213 t.Fatal(err)
214 }
215 if err := s.MarkOpaque(id1); err != nil {
216 t.Fatal(err)
217 }
218 file, err := NormalizeWritePaths(root, []string{"src/a.go"})
219 if err != nil {
220 t.Fatal(err)
221 }
222 if err := s.Realize(id2, file); err == nil {
223 t.Fatal("realize must fail after sibling goes opaque")
224 }
225 }
226
227 func TestSchedulerParentFileWriteAfterChildRealize(t *testing.T) {
228 s := NewSubagentScheduler(4, 2)
229 root := t.TempDir()
230 whole, err := WholeWorkspaceWriteClaim(root)
231 if err != nil {
232 t.Fatal(err)
233 }
234 _, id, err := s.AcquireWithID(context.Background(), AcquireRequest{Writer: true, WritePaths: whole})
235 if err != nil {
236 t.Fatal(err)
237 }
238 before, err := NormalizeWritePaths(root, []string{"b.go"})
239 if err != nil {
240 t.Fatal(err)
241 }
242 if _, err := s.ReserveParentWrite(before); err == nil {
243 t.Fatal("parent file write must wait while child still claims the whole workspace")
244 }
245 fileA, err := NormalizeWritePaths(root, []string{"a.go"})
246 if err != nil {
247 t.Fatal(err)
248 }
249 if err := s.Realize(id, fileA); err != nil {
250 t.Fatal(err)
251 }
252 release, err := s.ReserveParentWrite(before)
253 if err != nil {
254 t.Fatalf("parent write of disjoint file after realize: %v", err)
255 }
256 release()
257 if err := s.MarkOpaque(id); err != nil {
258 t.Fatal(err)
259 }
260 if _, err := s.ReserveParentWrite(before); err == nil {
261 t.Fatal("parent write must fail after child goes opaque")
262 }
263 }
264
265 func TestSchedulerWholeClaimNarrowsForNewSiblingsAndOpaqueRestoresExclusion(t *testing.T) {
266 s := NewSubagentScheduler(4, 3)
267 root := t.TempDir()
268 whole, err := WholeWorkspaceWriteClaim(root)
269 if err != nil {
270 t.Fatal(err)
271 }
272 releaseWhole, id, err := s.AcquireWithID(context.Background(), AcquireRequest{Writer: true, WritePaths: whole})
273 if err != nil {
274 t.Fatal(err)
275 }
276 defer releaseWhole()
277
278 fileA, err := NormalizeWritePaths(root, []string{"a.go"})
279 if err != nil {
280 t.Fatal(err)
281 }
282 fileB, err := NormalizeWritePaths(root, []string{"b.go"})
283 if err != nil {
284 t.Fatal(err)
285 }
286 if _, _, err := s.AcquireWithID(context.Background(), AcquireRequest{
287 Writer: true, WritePaths: fileB, Nested: true,
288 }); err == nil {
289 t.Fatal("new sibling must wait before the whole claim realizes a path")
290 }
291 if err := s.Realize(id, fileA); err != nil {
292 t.Fatal(err)
293 }
294 releaseB, _, err := s.AcquireWithID(context.Background(), AcquireRequest{
295 Writer: true, WritePaths: fileB, Nested: true,
296 })
297 if err != nil {
298 t.Fatalf("new disjoint sibling after realize: %v", err)
299 }
300 if _, _, err := s.AcquireWithID(context.Background(), AcquireRequest{
301 Writer: true, WritePaths: fileA, Nested: true,
302 }); err == nil {
303 t.Fatal("new same-file sibling must remain blocked")
304 }
305 releaseB()
306 if err := s.MarkOpaque(id); err != nil {
307 t.Fatal(err)
308 }
309 if _, _, err := s.AcquireWithID(context.Background(), AcquireRequest{
310 Writer: true, WritePaths: fileB, Nested: true,
311 }); err == nil {
312 t.Fatal("opaque mutation must restore whole-workspace exclusion")
313 }
314 }
315
316 func TestSchedulerTryClaimWritePaths(t *testing.T) {
317 s := NewSubagentScheduler(4, 2)
318 root := t.TempDir()
319 claim, _ := NormalizeWritePaths(root, []string{"a.md"})
320 release, err := s.Acquire(context.Background(), AcquireRequest{Writer: true, WritePaths: claim})
321 if err != nil {
322 t.Fatal(err)
323 }
324 defer release()
325 if err := s.TryClaimWritePaths(claim); err == nil {
326 t.Fatal("parent should see active claim")
327 }
328 other, _ := NormalizeWritePaths(root, []string{"b.md"})
329 if err := s.TryClaimWritePaths(other); err != nil {
330 t.Fatalf("disjoint claim should be free: %v", err)
331 }
332 }
333
333 lines GO