返回 DeepSeek-Reasonix
runtime_fixture_lifecycle_test.go
根目录 / internal / serve / runtime_fixture_lifecycle_test.go
1 package serve
2
3 import (
4 "context"
5 "encoding/json"
6 "net/http"
7 "net/http/httptest"
8 "os"
9 "path/filepath"
10 "runtime"
11 "strings"
12 "testing"
13 "time"
14
15 "reasonix/internal/agent"
16 "reasonix/internal/boot"
17 "reasonix/internal/config"
18 "reasonix/internal/control"
19 "reasonix/internal/provider"
20 )
21
22 func TestResumeActiveSessionTreatsSymlinkAliasAsCurrent(t *testing.T) {
23 if runtime.GOOS == "windows" {
24 t.Skip("symlink alias identity is exercised on POSIX CI")
25 }
26 dir := t.TempDir()
27 realPath := filepath.Join(dir, "real.jsonl")
28 aliasPath := filepath.Join(dir, "alias.jsonl")
29 saveServeTestSession(t, realPath)
30 if err := os.Symlink(realPath, aliasPath); err != nil {
31 t.Fatal(err)
32 }
33 ctrl := control.New(control.Options{Runner: blockingRunner{}, SessionDir: dir, SessionPath: aliasPath})
34 server := New(ctrl, NewBroadcaster(), config.ServeConfig{})
35 defer server.Close()
36 ctrl.Submit("keep running")
37 waitRunning(t, ctrl)
38 defer func() {
39 ctrl.Cancel()
40 waitNotRunning(t, ctrl)
41 }()
42 rec := httptest.NewRecorder()
43 req := httptest.NewRequest(http.MethodPost, "/resume", nil)
44 if !server.resumeActiveSession(rec, req, ctrl, realPath) {
45 t.Fatal("active current-session alias was not handled")
46 }
47 if rec.Code != http.StatusNoContent {
48 t.Fatalf("resume through current-session symlink alias = %d, want 204", rec.Code)
49 }
50 if server.ctl() != control.SessionAPI(ctrl) || !ctrl.Running() {
51 t.Fatal("current-session alias detached or replaced the active controller")
52 }
53 }
54
55 func TestSessionsReportsForegroundBackgroundJobsAsRunning(t *testing.T) {
56 dir := t.TempDir()
57 path := filepath.Join(dir, "jobs.jsonl")
58 saveServeTestSession(t, path)
59 ctrl := &backgroundJobOnlyController{Controller: control.New(control.Options{SessionDir: dir, SessionPath: path})}
60 server := New(ctrl, NewBroadcaster(), config.ServeConfig{})
61 defer server.Close()
62 rec := httptest.NewRecorder()
63 server.sessions(rec, httptest.NewRequest(http.MethodGet, "/sessions", nil))
64 var rows []sessionListEntry
65 if err := json.Unmarshal(rec.Body.Bytes(), &rows); err != nil {
66 t.Fatal(err)
67 }
68 if len(rows) != 1 || !rows[0].Current || !rows[0].Running {
69 t.Fatalf("foreground background-job session = %+v, want current and running", rows)
70 }
71 }
72
73 func TestDetachedRecoveryMovesRegistryKey(t *testing.T) {
74 dir := t.TempDir()
75 oldPath := filepath.Join(dir, "old.jsonl")
76 recoveryPath := filepath.Join(dir, "old-recovery.jsonl")
77 ctrl := control.New(control.Options{SessionDir: dir, SessionPath: oldPath})
78 server := New(ctrl, NewBroadcaster(), config.ServeConfig{})
79 defer ctrl.Close()
80 detached := &detachedSession{path: oldPath, ctrl: ctrl}
81 server.detached[oldPath] = detached
82 if err := server.moveDetachedRecovery(ctrl, recoveryPath); err != nil {
83 t.Fatal(err)
84 }
85 canonical := agent.CanonicalSessionPath(recoveryPath)
86 if got := server.detached[canonical]; got != detached || detached.path != canonical {
87 t.Fatalf("recovery registry = %+v path=%q", got, detached.path)
88 }
89 if server.detached[oldPath] != nil {
90 t.Fatal("old detached registry key was retained")
91 }
92 }
93
94 func TestRegisterDetachedRevalidatesPathAtPublication(t *testing.T) {
95 dir := t.TempDir()
96 oldPath := filepath.Join(dir, "old.jsonl")
97 newPath := filepath.Join(dir, "recovery.jsonl")
98 saveServeTestSession(t, oldPath)
99 saveServeTestSession(t, newPath)
100 ctrl := control.New(control.Options{Runner: blockingRunner{}, SessionDir: dir, SessionPath: oldPath})
101 server := New(ctrl, NewBroadcaster(), config.ServeConfig{})
102 defer server.Close()
103 tag := NewSessionTagSink(server.bc)
104 server.RegisterSessionTag(ctrl, tag)
105 started, release := make(chan struct{}), make(chan struct{})
106 registerDetachedHookForTest = func() { close(started); <-release }
107 t.Cleanup(func() { registerDetachedHookForTest = nil })
108 result := make(chan *detachedSession, 1)
109 go func() { detached, _ := server.registerDetached(ctrl, nil, tag); result <- detached }()
110 <-started
111 loaded, err := agent.LoadSession(newPath)
112 if err != nil {
113 t.Fatal(err)
114 }
115 ctrl.Resume(loaded, newPath)
116 ctrl.Submit("keep running")
117 waitRunning(t, ctrl)
118 close(release)
119 detached := <-result
120 canonical := agent.CanonicalSessionPath(newPath)
121 server.detachedMu.Lock()
122 registered := server.detached[canonical]
123 server.detachedMu.Unlock()
124 if detached == nil || detached.path != canonical || registered != detached {
125 t.Fatalf("detached publication path = %q entry=%v, want %q", detached.path, registered == detached, canonical)
126 }
127 ctrl.Cancel()
128 waitNotRunning(t, ctrl)
129 server.CloseBackground()
130 }
131
132 func TestDetachedRecoveryKeepsServeRoutingWrapper(t *testing.T) {
133 t.Setenv(agent.SessionLogSchemaEnv, "v1")
134 dir := t.TempDir()
135 aPath := filepath.Join(dir, "a.jsonl")
136 bPath := filepath.Join(dir, "b.jsonl")
137 saveServeTestSession(t, aPath)
138 saveServeTestSession(t, bPath)
139 loaded, err := agent.LoadSession(aPath)
140 if err != nil {
141 t.Fatal(err)
142 }
143 bc := NewBroadcaster()
144 tag := NewSessionTagSink(bc)
145 tag.SetPath(aPath)
146 exec := agent.New(nil, nil, loaded, agent.Options{}, tag)
147 ctrlA := control.New(control.Options{Runner: blockingRunner{}, Executor: exec, Sink: tag, SessionDir: dir, SessionPath: aPath, Label: "test"})
148 server := New(ctrlA, bc, config.ServeConfig{})
149 defer server.Close()
150 server.RegisterSessionTag(ctrlA, tag)
151 leases := control.NewSessionLeaseKeeper()
152 defer leases.Release()
153 if err := leases.Rebind(aPath); err != nil {
154 t.Fatal(err)
155 }
156 if err := server.SetSessionLeases(leases); err != nil {
157 t.Fatal(err)
158 }
159 server.buildControllerWithOptions = func(_ context.Context, _ string, opts boot.Options) (*control.Controller, error) {
160 return control.New(control.Options{Sink: opts.Sink, SessionDir: opts.SessionDir, Label: "test"}), nil
161 }
162 ctrlA.Submit("keep running")
163 waitRunning(t, ctrlA)
164 if err := server.busyDetach(context.Background(), ctrlA, bPath, func(next *control.Controller) error {
165 session, loadErr := agent.LoadSession(bPath)
166 if loadErr == nil {
167 next.Resume(session, bPath)
168 }
169 return loadErr
170 }); err != nil {
171 t.Fatal(err)
172 }
173 disk, err := agent.LoadSession(aPath)
174 if err != nil {
175 t.Fatal(err)
176 }
177 disk.Add(provider.Message{Role: provider.RoleUser, Content: "disk diverged"})
178 if err := disk.Save(aPath); err != nil {
179 t.Fatal(err)
180 }
181 ctrlA.Executor().Session().Add(provider.Message{Role: provider.RoleUser, Content: "local diverged"})
182 if err := ctrlA.Snapshot(); err != nil {
183 t.Fatal(err)
184 }
185 recoveryPath := agent.CanonicalSessionPath(ctrlA.SessionPath())
186 if recoveryPath == agent.CanonicalSessionPath(aPath) {
187 t.Fatal("detached controller did not move to a recovery transcript")
188 }
189 server.detachedMu.Lock()
190 detached := server.detached[recoveryPath]
191 oldEntry := server.detached[agent.CanonicalSessionPath(aPath)]
192 server.detachedMu.Unlock()
193 if detached == nil || detached.ctrl != control.SessionAPI(ctrlA) || oldEntry != nil || tag.Path() != recoveryPath {
194 t.Fatalf("detached recovery routing = entry %v old %v tag %q want %q", detached != nil, oldEntry != nil, tag.Path(), recoveryPath)
195 }
196 ctrlA.Cancel()
197 waitNotRunning(t, ctrlA)
198 server.CloseBackground()
199 }
200
201 func TestServeSwitchEffortUsesModelRefForDuplicateModelNames(t *testing.T) {
202 writeServeModelConfig(t)
203
204 bc := NewBroadcaster()
205 ctrl := control.New(control.Options{
206 Sink: bc,
207 Label: "shared-chat",
208 ModelRef: "alternate/shared-chat",
209 SessionDir: t.TempDir(),
210 })
211 server := New(ctrl, bc, config.ServeConfig{})
212 defer server.Close()
213 var builtRef string
214 server.buildController = func(_ context.Context, ref string) (*control.Controller, error) {
215 builtRef = ref
216 return control.New(control.Options{
217 Sink: bc,
218 Label: "shared-chat",
219 ModelRef: ref,
220 SessionDir: t.TempDir(),
221 }), nil
222 }
223
224 if err := server.switchEffort(context.Background(), "high"); err != nil {
225 t.Fatalf("switchEffort: %v", err)
226 }
227 if builtRef != "alternate/shared-chat" {
228 t.Fatalf("rebuilt model ref = %q, want alternate/shared-chat", builtRef)
229 }
230 edit := config.LoadForEdit(config.UserConfigPath())
231 def, _ := edit.Provider("default")
232 if def.Effort != "" {
233 t.Fatalf("default effort = %q, want unchanged", def.Effort)
234 }
235 alt, _ := edit.Provider("alternate")
236 if alt.Effort != "high" {
237 t.Fatalf("alternate effort = %q, want high", alt.Effort)
238 }
239 }
240
241 func TestSubmitNewHoldsBindingLockUntilRotationCompletes(t *testing.T) {
242 bc := NewBroadcaster()
243 ctrl := &blockingNewSessionController{
244 Controller: control.New(control.Options{Sink: bc, SessionDir: t.TempDir()}),
245 entered: make(chan struct{}),
246 release: make(chan struct{}),
247 }
248 t.Cleanup(func() {
249 select {
250 case <-ctrl.release:
251 default:
252 close(ctrl.release)
253 }
254 })
255 s := New(ctrl, bc, config.ServeConfig{})
256 defer func() {
257 select {
258 case <-ctrl.release:
259 default:
260 close(ctrl.release)
261 }
262 s.Close()
263 }()
264 submitDone := make(chan *httptest.ResponseRecorder, 1)
265 go func() {
266 req := httptest.NewRequest(http.MethodPost, "/submit", strings.NewReader(`{"input":"/new"}`))
267 rec := httptest.NewRecorder()
268 s.submit(rec, req)
269 submitDone <- rec
270 }()
271 select {
272 case <-ctrl.entered:
273 case <-time.After(2 * time.Second):
274 t.Fatal("/submit /new did not enter synchronous rotation")
275 }
276 lockAcquired := make(chan struct{})
277 go func() {
278 s.bindMu.Lock()
279 close(lockAcquired)
280 s.bindMu.Unlock()
281 }()
282 select {
283 case <-lockAcquired:
284 t.Fatal("bindMu was released before /new finished")
285 case <-time.After(100 * time.Millisecond):
286 }
287 close(ctrl.release)
288 var rec *httptest.ResponseRecorder
289 select {
290 case rec = <-submitDone:
291 case <-time.After(2 * time.Second):
292 t.Fatal("/submit /new did not return after rotation finished")
293 }
294 if rec.Code != http.StatusNoContent {
295 t.Fatalf("/submit /new status = %d, want 204", rec.Code)
296 }
297 select {
298 case <-lockAcquired:
299 case <-time.After(2 * time.Second):
300 t.Fatal("bindMu stayed locked after /new completed")
301 }
302 }
303
304 func TestSessionSnapshotEndpointsWaitForBindingEpoch(t *testing.T) {
305 bc := NewBroadcaster()
306 ctrl := control.New(control.Options{Sink: bc, SessionDir: t.TempDir()})
307 s := New(ctrl, bc, config.ServeConfig{})
308 defer s.Close()
309
310 for _, endpoint := range []struct {
311 name string
312 handler func(http.ResponseWriter, *http.Request)
313 }{
314 {name: "history", handler: s.history},
315 {name: "status", handler: s.status},
316 } {
317 t.Run(endpoint.name, func(t *testing.T) {
318 s.bindMu.Lock()
319 done := make(chan struct{})
320 go func() {
321 rec := httptest.NewRecorder()
322 req := httptest.NewRequest(http.MethodGet, "/"+endpoint.name+"?runtime=1", nil)
323 endpoint.handler(rec, req)
324 close(done)
325 }()
326 select {
327 case <-done:
328 s.bindMu.Unlock()
329 t.Fatalf("/%s observed a controller snapshot during an active binding epoch", endpoint.name)
330 case <-time.After(100 * time.Millisecond):
331 }
332 s.bindMu.Unlock()
333 select {
334 case <-done:
335 case <-time.After(2 * time.Second):
336 t.Fatalf("/%s stayed blocked after the binding epoch completed", endpoint.name)
337 }
338 })
339 }
340 }
341
342 func TestDeleteSessionSerializesWithForegroundPromotion(t *testing.T) {
343 dir := t.TempDir()
344 active, target := filepath.Join(dir, "active.jsonl"), filepath.Join(dir, "target.jsonl")
345 for _, path := range []string{active, target} {
346 if err := os.WriteFile(path, []byte("{}\n"), 0o644); err != nil {
347 t.Fatal(err)
348 }
349 }
350 bc := NewBroadcaster()
351 first := control.New(control.Options{Sink: bc, SessionDir: dir, SessionPath: active})
352 defer first.Close()
353 promoted := control.New(control.Options{Sink: bc, SessionDir: dir, SessionPath: target})
354 server := New(first, bc, config.ServeConfig{})
355 defer server.Close()
356 reachedLock := make(chan struct{})
357 deleteSessionBeforeOwnershipLockHookForTest = func() { close(reachedLock) }
358 t.Cleanup(func() { deleteSessionBeforeOwnershipLockHookForTest = nil })
359 server.bindMu.Lock()
360 rec := httptest.NewRecorder()
361 done := make(chan struct{})
362 go func() {
363 server.deleteSession(rec, httptest.NewRequest(http.MethodPost, "/delete-session", strings.NewReader(`{"name":"target"}`)))
364 close(done)
365 }()
366 select {
367 case <-reachedLock:
368 case <-time.After(2 * time.Second):
369 server.bindMu.Unlock()
370 t.Fatal("delete did not reach ownership boundary")
371 }
372 if !server.publishControllerSwap(first, promoted, target) {
373 server.bindMu.Unlock()
374 t.Fatal("foreground promotion failed")
375 }
376 server.bindMu.Unlock()
377 select {
378 case <-done:
379 case <-time.After(2 * time.Second):
380 t.Fatal("delete remained blocked after promotion")
381 }
382 if rec.Code != http.StatusConflict {
383 t.Fatalf("delete promoted session status = %d, want 409", rec.Code)
384 }
385 if _, err := os.Stat(target); err != nil {
386 t.Fatalf("promoted session was deleted: %v", err)
387 }
388 }
389
389 lines GO