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