返回 DeepSeek-Reasonix
serve_lease_test.go
根目录 / internal / serve / serve_lease_test.go
1 package serve
2
3 import (
4 "encoding/json"
5 "fmt"
6 "io"
7 "net/http"
8 "net/http/httptest"
9 "path/filepath"
10 "strings"
11 "sync"
12 "sync/atomic"
13 "testing"
14 "time"
15
16 "reasonix/internal/agent"
17 "reasonix/internal/config"
18 "reasonix/internal/control"
19 "reasonix/internal/provider"
20 )
21
22 func saveServeTestSession(t *testing.T, path string) {
23 t.Helper()
24 s := agent.NewSession("sys")
25 s.Add(provider.Message{Role: provider.RoleUser, Content: "hi"})
26 if err := s.Save(path); err != nil {
27 t.Fatal(err)
28 }
29 }
30
31 // TestResumeRefusedWhenSessionLeaseHeld proves POST /resume refuses to bind a
32 // session another runtime holds, keeps the server on its current session, and
33 // reports the shared holder wording.
34 func TestResumeRefusedWhenSessionLeaseHeld(t *testing.T) {
35 dir := t.TempDir()
36 active := filepath.Join(dir, "active.jsonl")
37 held := filepath.Join(dir, "held.jsonl")
38 saveServeTestSession(t, active)
39 saveServeTestSession(t, held)
40
41 holder, err := agent.TryAcquireSessionLease(held)
42 if err != nil {
43 t.Fatalf("test holder acquire: %v", err)
44 }
45 defer holder.Release()
46
47 bc := NewBroadcaster()
48 ctrl := control.New(control.Options{Sink: bc, SessionDir: dir, SessionPath: active})
49 server := New(ctrl, bc, config.ServeConfig{})
50 leases := control.NewSessionLeaseKeeper()
51 defer leases.Release()
52 if err := leases.Rebind(active); err != nil {
53 t.Fatalf("seed lease on active: %v", err)
54 }
55 server.SetSessionLeases(leases)
56 srv := httptest.NewServer(server.Handler())
57 defer srv.Close()
58
59 body, err := json.Marshal(map[string]string{"path": held})
60 if err != nil {
61 t.Fatal(err)
62 }
63 resp, err := http.Post(srv.URL+"/resume", "application/json", strings.NewReader(string(body)))
64 if err != nil {
65 t.Fatal(err)
66 }
67 respBody, _ := readAll(resp)
68 if resp.StatusCode != http.StatusConflict {
69 t.Fatalf("held resume status = %d, want 409 (body %q)", resp.StatusCode, respBody)
70 }
71 if !strings.Contains(respBody, "in use by another Reasonix") {
72 t.Fatalf("held resume body = %q, want holder wording", respBody)
73 }
74 if strings.Contains(respBody, held) {
75 t.Fatalf("held resume body leaks the session path: %q", respBody)
76 }
77 if got := filepath.Clean(ctrl.SessionPath()); got != filepath.Clean(active) {
78 t.Fatalf("session path after refused resume = %q, want active %q", got, active)
79 }
80 if got, want := leases.HeldPath(), agent.CanonicalSessionPath(active); got != want {
81 t.Fatalf("lease after refused resume = %q, want %q", got, want)
82 }
83 }
84
85 // TestResumeMovesSessionLease proves a successful POST /resume releases the old
86 // session's lease and holds the new one.
87 func TestResumeMovesSessionLease(t *testing.T) {
88 dir := t.TempDir()
89 active := filepath.Join(dir, "active.jsonl")
90 next := filepath.Join(dir, "next.jsonl")
91 saveServeTestSession(t, active)
92 saveServeTestSession(t, next)
93
94 bc := NewBroadcaster()
95 ctrl := control.New(control.Options{Sink: bc, SessionDir: dir, SessionPath: active})
96 server := New(ctrl, bc, config.ServeConfig{})
97 leases := control.NewSessionLeaseKeeper()
98 defer leases.Release()
99 if err := leases.Rebind(active); err != nil {
100 t.Fatalf("seed lease on active: %v", err)
101 }
102 server.SetSessionLeases(leases)
103 srv := httptest.NewServer(server.Handler())
104 defer srv.Close()
105
106 body, err := json.Marshal(map[string]string{"path": next})
107 if err != nil {
108 t.Fatal(err)
109 }
110 resp, err := http.Post(srv.URL+"/resume", "application/json", strings.NewReader(string(body)))
111 if err != nil {
112 t.Fatal(err)
113 }
114 respBody, _ := readAll(resp)
115 if resp.StatusCode != http.StatusNoContent {
116 t.Fatalf("resume status = %d, want 204 (body %q)", resp.StatusCode, respBody)
117 }
118 want, err := filepath.EvalSymlinks(next)
119 if err != nil {
120 t.Fatal(err)
121 }
122 if got, wantHeld := leases.HeldPath(), agent.CanonicalSessionPath(want); got != wantHeld {
123 t.Fatalf("lease after resume = %q, want %q", got, wantHeld)
124 }
125 lease, err := agent.TryAcquireSessionLease(active)
126 if err != nil {
127 t.Fatalf("old session lease not released by resume: %v", err)
128 }
129 lease.Release()
130 }
131
132 func readAll(resp *http.Response) (string, error) {
133 defer resp.Body.Close()
134 b, err := io.ReadAll(resp.Body)
135 return strings.TrimSpace(string(b)), err
136 }
137
138 func postServeLeaseJSON(client *http.Client, url string, body any, want int) error {
139 payload, err := json.Marshal(body)
140 if err != nil {
141 return err
142 }
143 resp, err := client.Post(url, "application/json", strings.NewReader(string(payload)))
144 if err != nil {
145 return err
146 }
147 respBody, readErr := readAll(resp)
148 if readErr != nil {
149 return readErr
150 }
151 if resp.StatusCode != want {
152 return fmt.Errorf("POST %s status = %d, want %d (body %q)", url, resp.StatusCode, want, respBody)
153 }
154 return nil
155 }
156
157 func waitServeLeaseDone(t *testing.T, done <-chan struct{}, what string, timeout time.Duration) {
158 t.Helper()
159 select {
160 case <-done:
161 case <-time.After(timeout):
162 t.Fatalf("timed out waiting for %s", what)
163 }
164 }
165
166 func waitServeLeaseResult(t *testing.T, ch <-chan error, what string, timeout time.Duration) error {
167 t.Helper()
168 select {
169 case err := <-ch:
170 return err
171 case <-time.After(timeout):
172 return fmt.Errorf("timed out waiting for %s", what)
173 }
174 }
175
176 // TestConcurrentResumesKeepControllerAndLeaseAligned hammers POST /resume from
177 // two goroutines bouncing between different targets and asserts the invariant
178 // this PR exists for: whatever session the controller ends up writing is the
179 // session the lease keeper guards. Without bindMu serializing the
180 // snapshot→rebind→resume sequence, interleaved handlers split the two (the
181 // controller on one path, the lease on another), leaving the written session
182 // unprotected and a foreign one wrongly occupied.
183 func TestConcurrentResumesKeepControllerAndLeaseAligned(t *testing.T) {
184 dir := t.TempDir()
185 active := filepath.Join(dir, "active.jsonl")
186 targetB := filepath.Join(dir, "target-b.jsonl")
187 targetC := filepath.Join(dir, "target-c.jsonl")
188 for _, p := range []string{active, targetB, targetC} {
189 saveServeTestSession(t, p)
190 }
191
192 bc := NewBroadcaster()
193 ctrl := control.New(control.Options{Sink: bc, SessionDir: dir, SessionPath: active})
194 defer ctrl.Close()
195 server := New(ctrl, bc, config.ServeConfig{})
196 leases := control.NewSessionLeaseKeeper()
197 defer leases.Release()
198 if err := leases.Rebind(active); err != nil {
199 t.Fatalf("seed lease on active: %v", err)
200 }
201 server.SetSessionLeases(leases)
202 srv := httptest.NewServer(server.Handler())
203 defer srv.Close()
204 client := &http.Client{Timeout: 10 * time.Second}
205
206 post := func(target string) error {
207 return postServeLeaseJSON(client, srv.URL+"/resume", map[string]string{"path": target}, http.StatusNoContent)
208 }
209
210 const rounds = 25
211 var wg sync.WaitGroup
212 errs := make(chan error, rounds*2)
213 wg.Add(2)
214 go func() {
215 defer wg.Done()
216 for range rounds {
217 if err := post(targetB); err != nil {
218 errs <- err
219 return
220 }
221 }
222 }()
223 go func() {
224 defer wg.Done()
225 for range rounds {
226 if err := post(targetC); err != nil {
227 errs <- err
228 return
229 }
230 }
231 }()
232 done := make(chan struct{})
233 go func() {
234 wg.Wait()
235 close(done)
236 close(errs)
237 }()
238 waitServeLeaseDone(t, done, "concurrent resume posts", 20*time.Second)
239 for err := range errs {
240 if err != nil {
241 t.Fatal(err)
242 }
243 }
244
245 got := agent.CanonicalSessionPath(server.ctl().SessionPath())
246 held := leases.HeldPath()
247 if got != held {
248 t.Fatalf("controller/lease split after concurrent resumes:\n controller writes %q\n lease guards %q", got, held)
249 }
250 }
251
252 // TestConcurrentResumeAndForkKeepAlignment interleaves /resume with /fork —
253 // the two handlers rotate the active path through different code paths — and
254 // asserts the same controller/lease alignment invariant.
255 func TestConcurrentResumeAndForkKeepAlignment(t *testing.T) {
256 dir := t.TempDir()
257 active := filepath.Join(dir, "active.jsonl")
258 target := filepath.Join(dir, "target.jsonl")
259 saveServeTestSession(t, active)
260 saveServeTestSession(t, target)
261
262 bc := NewBroadcaster()
263 ctrl := control.New(control.Options{Sink: bc, SessionDir: dir, SessionPath: active})
264 defer ctrl.Close()
265 server := New(ctrl, bc, config.ServeConfig{})
266 leases := control.NewSessionLeaseKeeper()
267 defer leases.Release()
268 if err := leases.Rebind(active); err != nil {
269 t.Fatalf("seed lease on active: %v", err)
270 }
271 server.SetSessionLeases(leases)
272 srv := httptest.NewServer(server.Handler())
273 defer srv.Close()
274 client := &http.Client{Timeout: 10 * time.Second}
275
276 var wg sync.WaitGroup
277 errs := make(chan error, 30)
278 wg.Add(2)
279 go func() {
280 defer wg.Done()
281 for range 15 {
282 if err := postServeLeaseJSON(client, srv.URL+"/resume", map[string]string{"path": target}, http.StatusNoContent); err != nil {
283 errs <- err
284 return
285 }
286 }
287 }()
288 go func() {
289 defer wg.Done()
290 for range 15 {
291 payload, _ := json.Marshal(map[string]any{"turn": 0, "name": ""})
292 resp, err := client.Post(srv.URL+"/fork", "application/json", strings.NewReader(string(payload)))
293 if err != nil {
294 errs <- err
295 return
296 }
297 _, readErr := readAll(resp)
298 if readErr != nil {
299 errs <- readErr
300 return
301 }
302 }
303 }()
304 done := make(chan struct{})
305 go func() {
306 wg.Wait()
307 close(done)
308 close(errs)
309 }()
310 waitServeLeaseDone(t, done, "concurrent resume/fork posts", 20*time.Second)
311 for err := range errs {
312 if err != nil {
313 t.Fatal(err)
314 }
315 }
316
317 got := agent.CanonicalSessionPath(server.ctl().SessionPath())
318 held := leases.HeldPath()
319 if got != held {
320 t.Fatalf("controller/lease split after resume×fork interleave:\n controller writes %q\n lease guards %q", got, held)
321 }
322 }
323
324 // TestInterleavedResumesForcedThroughBindWindow deterministically forces the
325 // interleaving bindMu prevents: request 1 is parked between its lease rebind
326 // and its controller Resume while request 2 tries to resume a different
327 // session. With bindMu, request 2 waits and both requests land aligned;
328 // without it, request 2 completes inside request 1's window and request 1
329 // then rebinds the controller to a session the lease no longer guards.
330 func TestInterleavedResumesForcedThroughBindWindow(t *testing.T) {
331 dir := t.TempDir()
332 active := filepath.Join(dir, "active.jsonl")
333 targetB := filepath.Join(dir, "target-b.jsonl")
334 targetC := filepath.Join(dir, "target-c.jsonl")
335 for _, p := range []string{active, targetB, targetC} {
336 saveServeTestSession(t, p)
337 }
338
339 bc := NewBroadcaster()
340 ctrl := control.New(control.Options{Sink: bc, SessionDir: dir, SessionPath: active})
341 defer ctrl.Close()
342 server := New(ctrl, bc, config.ServeConfig{})
343 leases := control.NewSessionLeaseKeeper()
344 defer leases.Release()
345 if err := leases.Rebind(active); err != nil {
346 t.Fatalf("seed lease on active: %v", err)
347 }
348 server.SetSessionLeases(leases)
349 srv := httptest.NewServer(server.Handler())
350 defer srv.Close()
351 client := &http.Client{Timeout: 10 * time.Second}
352
353 entered := make(chan struct{})
354 release := make(chan struct{})
355 // Park only the FIRST request through the hook. A sync.Once would not do:
356 // Once.Do holds an internal mutex while f runs, so the second request
357 // would block on the Once itself and accidentally reproduce bindMu's
358 // serialization even with the guard removed.
359 var first atomic.Bool
360 first.Store(true)
361 resumeBindHookForTest = func() {
362 if first.CompareAndSwap(true, false) {
363 close(entered)
364 <-release
365 }
366 }
367 defer func() { resumeBindHookForTest = nil }()
368 var releaseOnce sync.Once
369 unpark := func() {
370 releaseOnce.Do(func() {
371 close(release)
372 })
373 }
374 t.Cleanup(unpark)
375
376 post := func(target string) error {
377 return postServeLeaseJSON(client, srv.URL+"/resume", map[string]string{"path": target}, http.StatusNoContent)
378 }
379
380 done1 := make(chan error, 1)
381 go func() { done1 <- post(targetB) }()
382 select {
383 case <-entered:
384 // Request 1 parked inside its bind window, lease moved to B.
385 case err := <-done1:
386 if err != nil {
387 t.Fatalf("first resume finished before entering bind hook: %v", err)
388 }
389 t.Fatal("first resume finished before entering bind hook")
390 case <-time.After(5 * time.Second):
391 t.Fatal("timed out waiting for first resume to enter bind hook")
392 }
393
394 done2 := make(chan error, 1)
395 go func() { done2 <- post(targetC) }()
396 // Request 2 must not complete while request 1 owns the bind window.
397 select {
398 case err := <-done2:
399 // Unpark request 1 before failing, or httptest's Close would wait the
400 // full package timeout on the parked handler.
401 unpark()
402 if waitErr := waitServeLeaseResult(t, done1, "first resume after failed serialization check", 5*time.Second); waitErr != nil {
403 t.Fatalf("second resume completed inside the first resume's bind window; first resume cleanup failed: %v", waitErr)
404 }
405 if err != nil {
406 t.Fatalf("second resume completed inside the first resume's bind window with error: %v", err)
407 }
408 t.Fatal("second resume completed inside the first resume's bind window")
409 case <-time.After(150 * time.Millisecond):
410 }
411
412 unpark()
413 if err := waitServeLeaseResult(t, done1, "first resume", 5*time.Second); err != nil {
414 t.Fatal(err)
415 }
416 if err := waitServeLeaseResult(t, done2, "second resume", 5*time.Second); err != nil {
417 t.Fatal(err)
418 }
419
420 got := agent.CanonicalSessionPath(server.ctl().SessionPath())
421 held := leases.HeldPath()
422 if got != held {
423 t.Fatalf("controller/lease split after forced interleave:\n controller writes %q\n lease guards %q", got, held)
424 }
425 wantC := targetC
426 if resolved, err := filepath.EvalSymlinks(targetC); err == nil {
427 wantC = resolved // serve resolves the request path; match its form
428 }
429 if got != agent.CanonicalSessionPath(wantC) {
430 t.Fatalf("last resume should win: controller on %q, want %q", got, wantC)
431 }
432 }
433
433 lines GO