返回 DeepSeek-Reasonix
background_scope_test.go
根目录 / internal / jobs / background_scope_test.go
1 package jobs
2
3 import (
4 "context"
5 "io"
6 "sync"
7 "testing"
8 "time"
9
10 "reasonix/internal/event"
11 )
12
13 type scopeEventSink struct{ emit func(event.Event) }
14
15 func (s scopeEventSink) Emit(e event.Event) { s.emit(e) }
16
17 func TestBackgroundScopeEventsSwitchInOrderAndAllowReentrantReads(t *testing.T) {
18 var mu sync.Mutex
19 var received []string
20 m := NewManager(event.Discard)
21 scope := NewSessionBackgroundScope(m, nil)
22 defer scope.Release(false)
23 delivered := make(chan struct{})
24 release, err := m.BeginReplacement("")
25 if err != nil {
26 t.Fatal(err)
27 }
28 m.Emit(event.Event{Kind: event.Notice, Text: "one"})
29 m.Emit(event.Event{Kind: event.Notice, Text: "two"})
30 scope.Bind(scopeEventSink{func(e event.Event) {
31 _ = m.Running() // The registry lock must not cross the callback.
32 mu.Lock()
33 received = append(received, e.Text)
34 count := len(received)
35 mu.Unlock()
36 if count == 2 {
37 close(delivered)
38 }
39 }}, nil)
40 release()
41 select {
42 case <-delivered:
43 case <-time.After(time.Second):
44 t.Fatal("buffered completion events were lost")
45 }
46 mu.Lock()
47 defer mu.Unlock()
48 if len(received) != 2 || received[0] != "one" || received[1] != "two" {
49 t.Fatalf("events: %v", received)
50 }
51 }
52
53 func TestBackgroundScopeReplacementPreservesProcessAndCompletion(t *testing.T) {
54 m := NewManager(event.Discard)
55 scope := NewSessionBackgroundScope(m, nil)
56 defer scope.Release(false)
57 entered, finish := make(chan struct{}), make(chan struct{})
58 j := m.StartSessionProcess("session", "bash", "gateway", func(ctx context.Context, out io.Writer) (string, error) {
59 close(entered)
60 select {
61 case <-finish:
62 return "completed", nil
63 case <-ctx.Done():
64 return "", ctx.Err()
65 }
66 })
67 <-entered
68 release, err := m.BeginReplacement("session")
69 if err != nil {
70 t.Fatal(err)
71 }
72 if err := scope.Acquire(); err != nil {
73 t.Fatal(err)
74 }
75 // Candidate disposal must not cancel the outgoing runtime's process.
76 scope.Release(false)
77 select {
78 case <-j.done:
79 t.Fatal("candidate disposal killed gateway")
80 default:
81 }
82 blocked := m.StartForSession("session", "task", "late old-generation task", func(context.Context, io.Writer) (string, error) { t.Error("sealed task ran"); return "", nil })
83 <-blocked.done
84 release()
85 if len(m.RunningForSession("session")) != 1 || len(m.BlockingJobs("session")) != 0 {
86 t.Fatal("incorrect replacement classification")
87 }
88 close(finish)
89 select {
90 case <-j.done:
91 case <-time.After(time.Second):
92 t.Fatal("process did not finish")
93 }
94 if len(m.RunningForSession("session")) != 0 {
95 t.Fatal("finished process remained active")
96 }
97 }
98
99 func TestBackgroundScopeRuntimeTaskBlocksUntilActuallyExited(t *testing.T) {
100 m := NewManager(event.Discard)
101 defer m.Close()
102 entered, unwind := make(chan struct{}), make(chan struct{})
103 j := m.StartForSession("session", "task", "dependent", func(ctx context.Context, _ io.Writer) (string, error) {
104 close(entered)
105 <-ctx.Done()
106 <-unwind
107 return "", ctx.Err()
108 })
109 <-entered
110 m.Kill(j.ID)
111 if release, err := m.BeginReplacement("session"); err == nil {
112 release()
113 t.Error("cancelled-but-running task admitted replacement")
114 }
115 close(unwind)
116 <-j.done
117 release, err := m.BeginReplacement("session")
118 if err != nil {
119 t.Fatal(err)
120 }
121 release()
122 }
123
124 func TestBackgroundScopeLastRuntimeReleaseCancelsProcesses(t *testing.T) {
125 m := NewManager(event.Discard)
126 scope := NewSessionBackgroundScope(m, nil)
127 j := m.StartSessionProcess("session", "bash", "gateway", func(ctx context.Context, _ io.Writer) (string, error) { <-ctx.Done(); return "", ctx.Err() })
128 if err := scope.Acquire(); err != nil {
129 t.Fatal(err)
130 }
131 scope.Release(false)
132 select {
133 case <-j.done:
134 t.Fatal("old generation killed process")
135 default:
136 }
137 scope.Release(false)
138 select {
139 case <-j.done:
140 case <-time.After(time.Second):
141 t.Fatal("final release did not cancel")
142 }
143 if err := scope.Acquire(); err == nil {
144 t.Fatal("closed scope was resurrected")
145 }
146 }
147
148 func TestBackgroundScopeResamplesExitSuppressedDuringReplacement(t *testing.T) {
149 m := NewManager(event.Discard)
150 defer m.Close()
151 finish := make(chan struct{})
152 j := m.StartSessionProcess("session", "bash", "gateway", func(ctx context.Context, _ io.Writer) (string, error) {
153 select {
154 case <-finish:
155 return "", nil
156 case <-ctx.Done():
157 return "", ctx.Err()
158 }
159 })
160 release, err := m.BeginReplacement("session")
161 if err != nil {
162 t.Fatal(err)
163 }
164 defer release()
165 suppressed, published := make(chan struct{}, 1), make(chan struct{}, 1)
166 _, unsubscribe := m.SubscribeRuntime("session", func(state RuntimeState) {
167 if state.Running != 0 {
168 return
169 }
170 target := published
171 if m.ReplacementInProgress() {
172 target = suppressed
173 }
174 select {
175 case target <- struct{}{}:
176 default:
177 }
178 })
179 defer unsubscribe()
180 close(finish)
181 <-j.done
182 select {
183 case <-suppressed:
184 case <-time.After(time.Second):
185 t.Fatal("completion did not reach sealed observer")
186 }
187 release()
188 select {
189 case <-published:
190 case <-time.After(time.Second):
191 t.Fatal("publication lost the process exit suppressed during build")
192 }
193 }
194
194 lines GO