返回 DeepSeek-Reasonix
scheduler_test.go
根目录 / internal / agent / scheduler_test.go
1 package agent
2
3 import (
4 "context"
5 "sync"
6 "sync/atomic"
7 "testing"
8 "time"
9 )
10
11 func TestSchedulerTotalConcurrencyQueues(t *testing.T) {
12 s := NewSubagentScheduler(2, 2)
13 root := t.TempDir()
14 var started atomic.Int32
15 var max atomic.Int32
16 var wg sync.WaitGroup
17 barrier := make(chan struct{})
18
19 for i := 0; i < 4; i++ {
20 wg.Add(1)
21 go func() {
22 defer wg.Done()
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 TestSchedulerTryClaimWritePaths(t *testing.T) {
100 s := NewSubagentScheduler(4, 2)
101 root := t.TempDir()
102 claim, _ := NormalizeWritePaths(root, []string{"a.md"})
103 release, err := s.Acquire(context.Background(), AcquireRequest{Writer: true, WritePaths: claim})
104 if err != nil {
105 t.Fatal(err)
106 }
107 defer release()
108 if err := s.TryClaimWritePaths(claim); err == nil {
109 t.Fatal("parent should see active claim")
110 }
111 other, _ := NormalizeWritePaths(root, []string{"b.md"})
112 if err := s.TryClaimWritePaths(other); err != nil {
113 t.Fatalf("disjoint claim should be free: %v", err)
114 }
115 }
116
116 lines GO