| 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 |