返回 DeepSeek-Reasonix
session_ownership_identity_test.go
根目录 / internal / serve / session_ownership_identity_test.go
1 package serve
2
3 import (
4 "context"
5 "encoding/json"
6 "net/http"
7 "net/http/httptest"
8 "path/filepath"
9 "strings"
10 "testing"
11 "time"
12
13 "reasonix/internal/agent"
14 "reasonix/internal/config"
15 "reasonix/internal/control"
16 "reasonix/internal/event"
17 "reasonix/internal/provider"
18 "reasonix/internal/session"
19 )
20
21 // retireExclusiveForeground releases every writer deterministically (the
22 // foreground plus any runtime a handoff rotated away from) so Windows temp-dir
23 // cleanup does not race the idle-retirement TTL.
24 func retireExclusiveForeground(t *testing.T, ctrl *control.Controller, service *session.Service) {
25 t.Helper()
26 ctrl.Close()
27 if service != nil {
28 _ = service.CloseAll(context.Background())
29 }
30 }
31
32 // newIdentityLifecycleServe wraps the lifecycle test server with the frame
33 // tag an exclusive controller's production host would register, so identity
34 // transitions can assert how live frames get stamped.
35 func newIdentityLifecycleServe(t *testing.T, ctrl *control.Controller, ref session.SessionRef) *Server {
36 t.Helper()
37 bc := NewBroadcaster()
38 lifecycle := newLifecycleTestServer(t, ctrl, bc, config.ServeConfig{})
39 tag := newSessionTagSink(bc)
40 tag.SetIdentity("", ref.SessionID)
41 lifecycle.RegisterSessionTag(ctrl, tag)
42 return lifecycle
43 }
44
45 // openIdentityWriter opens the final-format session's writer through a bare
46 // persistence handle, standing in for the local runtime that takes over. The
47 // returned session keeps the writer lock until Close.
48 func openIdentityWriter(t *testing.T, root string, ref session.SessionRef) *session.Session {
49 t.Helper()
50 handle, err := session.NewFilesystemPersistence(root).Open(ref.SessionID, session.ReadWrite)
51 if err != nil {
52 t.Fatalf("open identity writer: %v", err)
53 }
54 return handle
55 }
56
57 func identityRoot(t *testing.T, service *session.Service, ref session.SessionRef) string {
58 t.Helper()
59 dir, err := service.SessionDir(t.Context(), ref)
60 if err != nil {
61 t.Fatalf("resolve identity dir: %v", err)
62 }
63 return filepath.Dir(dir)
64 }
65
66 func serveBody(t *testing.T, method, url, body string) (*http.Response, string) {
67 t.Helper()
68 req, err := http.NewRequest(method, url, strings.NewReader(body))
69 if err != nil {
70 t.Fatal(err)
71 }
72 req.Header.Set("Content-Type", "application/json")
73 resp, err := http.DefaultClient.Do(req)
74 if err != nil {
75 t.Fatal(err)
76 }
77 defer resp.Body.Close()
78 raw := make([]byte, 0, 4<<10)
79 buf := make([]byte, 4<<10)
80 for {
81 n, readErr := resp.Body.Read(buf)
82 raw = append(raw, buf[:n]...)
83 if readErr != nil {
84 break
85 }
86 }
87 return resp, string(raw)
88 }
89
90 // TestIdentityHandoffReleasesWriterAndOwnershipTracks proves the release half:
91 // /handoff on a final-format identity rotates the foreground, drops the writer
92 // lock before the grant is answered, and /ownership reports the external
93 // holder while the taker keeps the lock.
94 func TestIdentityHandoffReleasesWriterAndOwnershipTracks(t *testing.T) {
95 _, ctrl, service, current := newExclusiveSessionServe(t)
96 root := identityRoot(t, service, current)
97 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
98 ts := httptest.NewServer(operatorHandler(lifecycle))
99 defer ts.Close()
100 route := "session-id:" + current.SessionID
101
102 resp, raw := serveBody(t, http.MethodGet, ts.URL+"/ownership?session="+route, "")
103 if resp.StatusCode != http.StatusOK {
104 t.Fatalf("ownership status = %d body %s", resp.StatusCode, raw)
105 }
106 var view ownershipView
107 if err := json.Unmarshal([]byte(raw), &view); err != nil {
108 t.Fatal(err)
109 }
110 if view.Holder != "serve" || view.Running {
111 t.Fatalf("before handoff view = %+v, want serve holder idle", view)
112 }
113 defer retireExclusiveForeground(t, ctrl, service)
114
115 resp, raw = serveBody(t, http.MethodPost, ts.URL+"/handoff", `{"sessionPath":"`+route+`","targetWriterId":"taker-writer","force":true,"mode":"wait","timeoutMs":2000}`)
116 if resp.StatusCode != http.StatusOK {
117 t.Fatalf("handoff status = %d body %s", resp.StatusCode, raw)
118 }
119 var grant mirrorGrant
120 if err := json.Unmarshal([]byte(raw), &grant); err != nil {
121 t.Fatal(err)
122 }
123 if grant.MirrorID == "" || grant.ReturnHandoffID == "" || grant.SourceWriterID == "" ||
124 grant.TargetWriterID != "taker-writer" || grant.SessionPath != route {
125 t.Fatalf("handoff grant = %+v", grant)
126 }
127 if ref, bound := ctrl.SessionRef(); bound {
128 t.Fatalf("foreground still bound to %q after handoff; release must not allocate a replacement identity", ref.SessionID)
129 }
130 // The frame tag must follow the released foreground: a tag still pointing
131 // at the handed-off identity misroutes every subsequent live frame.
132 if tag := lifecycle.tagFor(ctrl); tag == nil {
133 t.Fatal("frame tag missing after handoff release")
134 } else if tag.path != "" || tag.sessionID != "" {
135 t.Fatalf("frame tag after handoff = %+v, want no route until the next identity is allocated", tag)
136 }
137 if session.ProbeWriterHeld(filepath.Join(root, current.SessionID)) {
138 t.Fatal("writer lock still held after handoff grant")
139 }
140
141 // The taker acquires the released writer; ownership flips to external.
142 writer := openIdentityWriter(t, root, current)
143 defer writer.Close(t.Context())
144 resp, raw = serveBody(t, http.MethodGet, ts.URL+"/ownership?session="+route, "")
145 if err := json.Unmarshal([]byte(raw), &view); err != nil {
146 t.Fatal(err)
147 }
148 if resp.StatusCode != http.StatusOK || view.Holder != "external" || !view.TakenOver {
149 t.Fatalf("after takeover view = %+v (status %d)", view, resp.StatusCode)
150 }
151 }
152
153 // TestIdentityResumeMountsSpectatorWhenWriterHeld proves the attach contract:
154 // /resume for an identity another runtime writes answers 204 with the
155 // taken-over header instead of a hard failure, and /history serves the cold
156 // event log so the spectator can render.
157 func TestIdentityResumeMountsSpectatorWhenWriterHeld(t *testing.T) {
158 _, ctrl, service, current := newExclusiveSessionServe(t)
159 root := identityRoot(t, service, current)
160 ts := httptest.NewServer(operatorHandler(newLifecycleTestServer(t, ctrl, NewBroadcaster(), config.ServeConfig{})))
161 defer ts.Close()
162 route := "session-id:" + current.SessionID
163
164 // Detach the foreground from the identity first so the resume path cannot
165 // short-circuit onto the already-bound current session.
166 if _, err := ctrl.BindFreshSession(t.Context(), "spectator-fresh"); err != nil {
167 t.Fatal(err)
168 }
169 if err := service.Close(t.Context(), current); err != nil {
170 t.Fatalf("close rotated-out runtime: %v", err)
171 }
172 writer := openIdentityWriter(t, root, current)
173 defer writer.Close(t.Context())
174
175 resp, raw := serveBody(t, http.MethodPost, ts.URL+"/resume", `{"sessionId":"`+current.SessionID+`"}`)
176 if resp.StatusCode != http.StatusNoContent {
177 t.Fatalf("spectator resume status = %d body %s", resp.StatusCode, raw)
178 }
179 if resp.Header.Get(sessionTakenOverHeader) == "" {
180 t.Fatal("spectator resume omitted the taken-over header")
181 }
182 if resp.Header.Get(sessionIDHeader) != current.SessionID {
183 t.Fatalf("spectator resume session header = %q", resp.Header.Get(sessionIDHeader))
184 }
185
186 resp, raw = serveBody(t, http.MethodGet, ts.URL+"/history?session="+route, "")
187 if resp.StatusCode != http.StatusOK {
188 t.Fatalf("spectator history status = %d body %s", resp.StatusCode, raw)
189 }
190
191 resp, raw = serveBody(t, http.MethodGet, ts.URL+"/status?session="+route, "")
192 if resp.StatusCode != http.StatusOK {
193 t.Fatalf("spectator status code = %d body %s", resp.StatusCode, raw)
194 }
195 var status map[string]any
196 if err := json.Unmarshal([]byte(raw), &status); err != nil {
197 t.Fatal(err)
198 }
199 if taken, _ := status["takenOver"].(bool); !taken {
200 t.Fatalf("spectator status missing takenOver: %v", status)
201 }
202 retireExclusiveForeground(t, ctrl, service)
203 }
204
205 // TestIdentityAdoptRegistersWriterAfterServeRestart pins the canonical
206 // re-registration path used when a CLI survives the resident serve restarting.
207 func TestIdentityAdoptRegistersWriterAfterServeRestart(t *testing.T) {
208 _, ctrl, service, current := newExclusiveSessionServe(t)
209 root := identityRoot(t, service, current)
210 lifecycle := newLifecycleTestServer(t, ctrl, NewBroadcaster(), config.ServeConfig{})
211 ts := httptest.NewServer(operatorHandler(lifecycle))
212 defer ts.Close()
213 defer retireExclusiveForeground(t, ctrl, service)
214 route := "session-id:" + current.SessionID
215
216 if _, err := ctrl.BindFreshSession(t.Context(), "adopt-fresh"); err != nil {
217 t.Fatal(err)
218 }
219 if err := service.Close(t.Context(), current); err != nil {
220 t.Fatalf("close rotated-out runtime: %v", err)
221 }
222 writer := openIdentityWriter(t, root, current)
223 defer writer.Close(t.Context())
224
225 resp, raw := serveBody(t, http.MethodPost, ts.URL+"/adopt", `{"sessionPath":"`+route+`","writerId":"surviving-cli"}`)
226 if resp.StatusCode != http.StatusOK {
227 t.Fatalf("identity adopt status = %d body %s", resp.StatusCode, raw)
228 }
229 var grant mirrorGrant
230 if err := json.Unmarshal([]byte(raw), &grant); err != nil {
231 t.Fatal(err)
232 }
233 if grant.SessionPath != route || grant.MirrorID == "" || grant.TargetWriterID != "surviving-cli" {
234 t.Fatalf("identity adopt grant = %+v", grant)
235 }
236 if mirrored, ok := lifecycle.mirroredEntry(route); !ok || mirrored.targetWriterID != "surviving-cli" {
237 t.Fatalf("identity mirror after adopt = %+v, present=%v", mirrored, ok)
238 }
239 }
240
241 // TestIdentityReclaimReattachesForeground proves the return half: once the
242 // local writer drops the lock, /reclaim re-owns the identity for the serve and
243 // clears the mirror bookkeeping.
244 func TestIdentityReclaimReattachesForeground(t *testing.T) {
245 _, ctrl, service, current := newExclusiveSessionServe(t)
246 root := identityRoot(t, service, current)
247 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
248 ts := httptest.NewServer(operatorHandler(lifecycle))
249 defer ts.Close()
250 route := "session-id:" + current.SessionID
251
252 resp, raw := serveBody(t, http.MethodPost, ts.URL+"/handoff", `{"sessionPath":"`+route+`","targetWriterId":"taker-writer","force":true,"mode":"wait","timeoutMs":2000}`)
253 if resp.StatusCode != http.StatusOK {
254 t.Fatalf("handoff status = %d body %s", resp.StatusCode, raw)
255 }
256 writer := openIdentityWriter(t, root, current)
257 if err := writer.Close(t.Context()); err != nil {
258 t.Fatalf("taker release: %v", err)
259 }
260
261 resp, raw = serveBody(t, http.MethodPost, ts.URL+"/reclaim", `{"sessionPath":"`+route+`","mode":"wait","timeoutMs":5000}`)
262 if resp.StatusCode != http.StatusNoContent {
263 t.Fatalf("reclaim status = %d body %s", resp.StatusCode, raw)
264 }
265 // After the reclaim the identity-selected status must answer ownership
266 // explicitly false: clients apply present fields only, so an omitted
267 // takenOver would pin the spectator banner forever.
268 resp, raw = serveBody(t, http.MethodGet, ts.URL+"/status?session="+route, "")
269 if resp.StatusCode != http.StatusOK {
270 t.Fatalf("post-reclaim status code = %d body %s", resp.StatusCode, raw)
271 }
272 var reclaimed map[string]any
273 if err := json.Unmarshal([]byte(raw), &reclaimed); err != nil {
274 t.Fatal(err)
275 }
276 if taken, _ := reclaimed["takenOver"].(bool); taken {
277 t.Fatalf("post-reclaim status still reports takenOver: %v", reclaimed)
278 }
279 if ref, bound := ctrl.SessionRef(); !bound || ref != current {
280 t.Fatalf("foreground ref after reclaim = %+v (bound %v), want %+v", ref, bound, current)
281 }
282 // The frame tag must follow the re-owned identity: a stale tag stamps live
283 // frames with another session and identity-routed subscribers drop them.
284 if tag := lifecycle.tagFor(ctrl); tag == nil || tag.path != "" || tag.sessionID != current.SessionID {
285 t.Fatalf("frame tag after reclaim = %+v, want identity %q", tag, current.SessionID)
286 }
287 if _, still := lifecycle.mirroredEntry(route); still {
288 t.Fatal("mirror entry survived reclaim")
289 }
290 defer retireExclusiveForeground(t, ctrl, service)
291 if !session.ProbeWriterHeld(filepath.Join(root, current.SessionID)) {
292 t.Fatal("serve did not re-acquire the writer after reclaim")
293 }
294 }
295
296 // TestIdentityHandoffRefusesForeignHolder pins the guard: a handoff for an
297 // identity this serve does not run is refused instead of granting a session
298 // the serve cannot release.
299 func TestIdentityHandoffRefusesForeignHolder(t *testing.T) {
300 _, ctrl, service, _ := newExclusiveSessionServe(t)
301 ts := httptest.NewServer(operatorHandler(newLifecycleTestServer(t, ctrl, NewBroadcaster(), config.ServeConfig{})))
302 defer ts.Close()
303 resp, raw := serveBody(t, http.MethodPost, ts.URL+"/handoff", `{"sessionPath":"session-id:does-not-exist","targetWriterId":"taker","force":true}`)
304 if resp.StatusCode != http.StatusBadRequest {
305 t.Fatalf("unknown identity handoff status = %d body %s", resp.StatusCode, raw)
306 }
307 retireExclusiveForeground(t, ctrl, service)
308 }
309
310 // shortMirrorEndWait shrinks the farewell's release wait so the tests pin the
311 // protocol (which probe re-owns, how many probes the bound allows) rather than
312 // wall-clock time.
313 func shortMirrorEndWait(t *testing.T, polls int) {
314 t.Helper()
315 wait, poll := mirrorEndReleaseWait, mirrorEndReleasePoll
316 mirrorEndReleasePoll = 2 * time.Millisecond
317 mirrorEndReleaseWait = time.Duration(polls) * mirrorEndReleasePoll
318 t.Cleanup(func() {
319 mirrorEndReleaseWait, mirrorEndReleasePoll = wait, poll
320 mirrorEndProbeHookForTest = nil
321 })
322 }
323
324 // handoffIdentityForTest hands the fixture's identity to "taker-writer" and
325 // returns the grant.
326 func handoffIdentityForTest(t *testing.T, url, route string) mirrorGrant {
327 t.Helper()
328 resp, raw := serveBody(t, http.MethodPost, url+"/handoff", `{"sessionPath":"`+route+`","targetWriterId":"taker-writer","force":true,"mode":"wait","timeoutMs":2000}`)
329 if resp.StatusCode != http.StatusOK {
330 t.Fatalf("handoff status = %d body %s", resp.StatusCode, raw)
331 }
332 var grant mirrorGrant
333 if err := json.Unmarshal([]byte(raw), &grant); err != nil {
334 t.Fatal(err)
335 }
336 return grant
337 }
338
339 // TestIdentityMirrorEndAcceptsLiveWriter pins the farewell contract: the
340 // writer's return transaction sends mirror-end before process exit, so the
341 // writer lock may still be held past the release wait and the serve must accept
342 // (204) instead of 409ing a call its own protocol ordering requires. The
343 // outstanding return then belongs to the stale auto-reclaim, and the wait is
344 // bounded by the configured number of probes.
345 func TestIdentityMirrorEndAcceptsLiveWriter(t *testing.T) {
346 const polls = 3
347 shortMirrorEndWait(t, polls)
348 _, ctrl, service, current := newExclusiveSessionServe(t)
349 root := identityRoot(t, service, current)
350 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
351 ts := httptest.NewServer(operatorHandler(lifecycle))
352 defer ts.Close()
353 route := "session-id:" + current.SessionID
354 grant := handoffIdentityForTest(t, ts.URL, route)
355 writer := openIdentityWriter(t, root, current)
356 defer writer.Close(t.Context())
357 probes := 0
358 mirrorEndProbeHookForTest = func(int) { probes++ }
359
360 resp, raw := serveBody(t, http.MethodPost, ts.URL+"/mirror-end", `{"sessionPath":"`+route+`","mirrorId":"`+grant.MirrorID+`"}`)
361 if resp.StatusCode != http.StatusNoContent {
362 t.Fatalf("mirror-end with live writer status = %d body %s", resp.StatusCode, raw)
363 }
364 if ref, bound := ctrl.SessionRef(); bound && ref == current {
365 t.Fatal("mirror-end re-owned the identity under a live writer")
366 }
367 if _, mirrored := lifecycle.mirroredEntry(route); !mirrored {
368 t.Fatal("mirror entry dropped although the writer never released; the auto-reclaim has nothing to finish")
369 }
370 if probes < 1 || probes > polls+2 {
371 t.Fatalf("farewell probed %d times, want between 1 and %d (bounded wait)", probes, polls+2)
372 }
373 retireExclusiveForeground(t, ctrl, service)
374 }
375
376 // The desktop releases its runtime right after sending mirror-end; the
377 // farewell must re-own the identity on the first probe that sees the lock
378 // free instead of answering 204 and leaving the remote side read-only until
379 // the 30 s stale auto-reclaim.
380 func TestIdentityMirrorEndReclaimsOnceWriterReleases(t *testing.T) {
381 shortMirrorEndWait(t, 50)
382 _, ctrl, service, current := newExclusiveSessionServe(t)
383 root := identityRoot(t, service, current)
384 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
385 ts := httptest.NewServer(operatorHandler(lifecycle))
386 defer ts.Close()
387 route := "session-id:" + current.SessionID
388 grant := handoffIdentityForTest(t, ts.URL, route)
389 writer := openIdentityWriter(t, root, current)
390 heldProbes := 0
391 mirrorEndProbeHookForTest = func(attempt int) {
392 heldProbes++
393 if attempt == 0 {
394 // The writer's teardown lands after the farewell was sent.
395 if err := writer.Close(t.Context()); err != nil {
396 t.Errorf("release taker writer: %v", err)
397 }
398 }
399 }
400
401 resp, raw := serveBody(t, http.MethodPost, ts.URL+"/mirror-end", `{"sessionPath":"`+route+`","mirrorId":"`+grant.MirrorID+`"}`)
402 if resp.StatusCode != http.StatusNoContent {
403 t.Fatalf("mirror-end status = %d body %s", resp.StatusCode, raw)
404 }
405 if heldProbes != 1 {
406 t.Fatalf("farewell saw the lock held on %d probes, want exactly 1: the release must be picked up by the very next probe", heldProbes)
407 }
408 if ref, bound := ctrl.SessionRef(); !bound || ref != current {
409 t.Fatalf("foreground after farewell = %+v (bound %v), want the reclaimed identity %+v", ref, bound, current)
410 }
411 if _, mirrored := lifecycle.mirroredEntry(route); mirrored {
412 t.Fatal("mirror entry survived the farewell reclaim")
413 }
414 defer retireExclusiveForeground(t, ctrl, service)
415 if !session.ProbeWriterHeld(filepath.Join(root, current.SessionID)) {
416 t.Fatal("serve did not re-acquire the writer after the farewell")
417 }
418 }
419
420 // An identity history read runs outside bindMu. A rotation or handoff landing
421 // between the read and the response must be reported as a changed runtime, as
422 // transcriptBoundRead already does, instead of answering the new route with
423 // the outgoing controller's transcript.
424 func TestHistoryIdentityRouteDetectsRuntimeChangeDuringRead(t *testing.T) {
425 _, ctrl, service, current := newExclusiveSessionServe(t)
426 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
427 ts := httptest.NewServer(operatorHandler(lifecycle))
428 defer ts.Close()
429 defer retireExclusiveForeground(t, ctrl, service)
430 route := "session-id:" + current.SessionID
431
432 historyIdentityReadHookForTest = func() {
433 if _, err := ctrl.BindFreshSession(context.Background(), "rotated-mid-read"); err != nil {
434 t.Errorf("rotate during read: %v", err)
435 }
436 }
437 t.Cleanup(func() { historyIdentityReadHookForTest = nil })
438 resp, raw := serveBody(t, http.MethodGet, ts.URL+"/history?session="+route, "")
439 if resp.StatusCode != http.StatusConflict || !strings.Contains(raw, "transcript runtime changed during read") {
440 t.Fatalf("history across a mid-read rotation = %d %q, want 409 runtime changed", resp.StatusCode, raw)
441 }
442 historyIdentityReadHookForTest = nil
443 if err := service.Close(t.Context(), current); err != nil {
444 t.Fatalf("close rotated-out runtime: %v", err)
445 }
446 // With the rotation settled the route is served cold from the event log.
447 resp, raw = serveBody(t, http.MethodGet, ts.URL+"/history?session="+route, "")
448 if resp.StatusCode != http.StatusOK {
449 t.Fatalf("history after rotation = %d body %s", resp.StatusCode, raw)
450 }
451 }
452
453 // TestIdentityStatusAnswersFreeWriterWithRouteMatch pins the serve-restart
454 // recovery: a spectator identity whose writer exited and whose mirror entry
455 // was lost to the restart must get an explicit route-matching status with
456 // takenOver=false — the foreground snapshot names a different session and a
457 // pinned tab would discard it, leaving the banner stuck until re-attach.
458 func TestIdentityStatusAnswersFreeWriterWithRouteMatch(t *testing.T) {
459 _, ctrl, service, current := newExclusiveSessionServe(t)
460 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
461 ts := httptest.NewServer(operatorHandler(lifecycle))
462 defer ts.Close()
463 route := "session-id:" + current.SessionID
464
465 // Move the foreground off the identity and free its writer, mimicking a
466 // post-restart world where nothing holds the session.
467 if _, err := ctrl.BindFreshSession(t.Context(), "elsewhere"); err != nil {
468 t.Fatal(err)
469 }
470 if err := service.Close(t.Context(), current); err != nil {
471 t.Fatalf("release identity writer: %v", err)
472 }
473
474 resp, raw := serveBody(t, http.MethodGet, ts.URL+"/status?session="+route, "")
475 if resp.StatusCode != http.StatusOK {
476 t.Fatalf("status code = %d body %s", resp.StatusCode, raw)
477 }
478 var status map[string]any
479 if err := json.Unmarshal([]byte(raw), &status); err != nil {
480 t.Fatal(err)
481 }
482 if sid, _ := status["sessionId"].(string); sid != current.SessionID {
483 t.Fatalf("status sessionId = %v, want the queried identity", status["sessionId"])
484 }
485 if taken, _ := status["takenOver"].(bool); taken {
486 t.Fatalf("free-writer identity status still reports takenOver: %v", status)
487 }
488 retireExclusiveForeground(t, ctrl, service)
489 }
490
491 // countSessions returns the /sessions row count so identity lifecycle tests
492 // can pin what a transition persists.
493 func countSessions(t *testing.T, url string) int {
494 t.Helper()
495 resp, raw := serveBody(t, http.MethodGet, url+"/sessions", "")
496 if resp.StatusCode != http.StatusOK {
497 t.Fatalf("sessions status = %d body %s", resp.StatusCode, raw)
498 }
499 var rows []sessionListEntry
500 if err := json.Unmarshal([]byte(raw), &rows); err != nil {
501 t.Fatal(err)
502 }
503 return len(rows)
504 }
505
506 // A handoff releases authority; it is not a conversation. The legacy keeper
507 // only unbinds, and the identity path must match: no replacement identity is
508 // created until the user actually starts one, so handoff/reclaim cycles do not
509 // litter /sessions with empty rows. /new on the released foreground still
510 // allocates on demand.
511 func TestIdentityHandoffDoesNotPersistReplacementSession(t *testing.T) {
512 _, ctrl, service, current := newExclusiveSessionServe(t)
513 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
514 ts := httptest.NewServer(operatorHandler(lifecycle))
515 defer ts.Close()
516 defer retireExclusiveForeground(t, ctrl, service)
517 route := "session-id:" + current.SessionID
518
519 before := countSessions(t, ts.URL)
520 resp, raw := serveBody(t, http.MethodPost, ts.URL+"/handoff", `{"sessionPath":"`+route+`","targetWriterId":"taker-writer","force":true,"mode":"wait","timeoutMs":2000}`)
521 if resp.StatusCode != http.StatusOK {
522 t.Fatalf("handoff status = %d body %s", resp.StatusCode, raw)
523 }
524 if after := countSessions(t, ts.URL); after != before {
525 t.Fatalf("/sessions rows after handoff = %d, want %d (handoff persisted a replacement session)", after, before)
526 }
527 if _, bound := ctrl.SessionRef(); bound {
528 t.Fatal("foreground is bound after handoff; nothing should be allocated until the next turn")
529 }
530 for _, msg := range ctrl.History() {
531 if msg.Role != provider.RoleSystem {
532 t.Fatalf("released foreground still carries the handed-off conversation: %+v", msg)
533 }
534 }
535 // The released foreground stays usable: /new allocates exactly one fresh
536 // identity on demand.
537 resp, raw = serveBody(t, http.MethodPost, ts.URL+"/new", "")
538 if resp.StatusCode != http.StatusNoContent {
539 t.Fatalf("/new after handoff status = %d body %s", resp.StatusCode, raw)
540 }
541 fresh, bound := ctrl.SessionRef()
542 if !bound || fresh == current {
543 t.Fatalf("/new after handoff bound %+v (bound %v), want a fresh identity", fresh, bound)
544 }
545 if got := countSessions(t, ts.URL); got != before+1 {
546 t.Fatalf("/sessions rows after /new = %d, want %d", got, before+1)
547 }
548 }
549
550 // A turn admitted between the unlocked quiet probe and the locked release must
551 // be refused with the legacy path's busy-again error, not raced: otherwise the
552 // foreground is unbound mid-turn, the close fails as busy, and the caller gets
553 // a 500 with an orphaned running runtime.
554 func TestIdentityHandoffRefusesTurnAdmittedAfterQuietProbe(t *testing.T) {
555 _, ctrl, service, current := newExclusiveSessionServeWithOptions(t, func(opts *control.Options) {
556 opts.Runner = blockingRunner{}
557 })
558 root := identityRoot(t, service, current)
559 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
560 ts := httptest.NewServer(operatorHandler(lifecycle))
561 defer ts.Close()
562 defer retireExclusiveForeground(t, ctrl, service)
563 route := "session-id:" + current.SessionID
564
565 handoffIdentityBeforeLockHookForTest = func() {
566 // POST /chat was admitted right after the probe saw an idle foreground.
567 ctrl.Submit("keep running")
568 waitRunning(t, ctrl)
569 }
570 t.Cleanup(func() { handoffIdentityBeforeLockHookForTest = nil })
571 defer func() {
572 ctrl.Cancel()
573 waitNotRunning(t, ctrl)
574 }()
575
576 resp, raw := serveBody(t, http.MethodPost, ts.URL+"/handoff", `{"sessionPath":"`+route+`","targetWriterId":"taker-writer","force":true,"mode":"wait","timeoutMs":2000}`)
577 if resp.StatusCode != http.StatusConflict || !strings.Contains(raw, errHandoffBusyAgain.Error()) {
578 t.Fatalf("handoff against a freshly admitted turn = %d %q, want 409 busy-again", resp.StatusCode, raw)
579 }
580 if ref, bound := ctrl.SessionRef(); !bound || ref != current {
581 t.Fatalf("foreground binding after refused handoff = %+v (bound %v), want %+v", ref, bound, current)
582 }
583 if _, mirrored := lifecycle.mirroredEntry(route); mirrored {
584 t.Fatal("refused handoff registered a mirror entry")
585 }
586 if !session.ProbeWriterHeld(filepath.Join(root, current.SessionID)) {
587 t.Fatal("refused handoff dropped the writer lock")
588 }
589 if !ctrl.Running() {
590 t.Fatal("refused handoff interrupted the admitted turn")
591 }
592 }
593
594 // The controller can report idle while the runtime is still finalizing the
595 // turn's terminal commit; closing such a runtime is refused as busy. The
596 // handoff must report busy-again from the locked re-check and leave the
597 // binding intact, then succeed once the runtime settles.
598 func TestIdentityHandoffRefusesFinalizingRuntimeAndRecovers(t *testing.T) {
599 _, ctrl, service, current := newExclusiveSessionServe(t)
600 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
601 ts := httptest.NewServer(operatorHandler(lifecycle))
602 defer ts.Close()
603 defer retireExclusiveForeground(t, ctrl, service)
604 route := "session-id:" + current.SessionID
605 runtime, ok := service.Runtime(current)
606 if !ok {
607 t.Fatal("current runtime is not published")
608 }
609 generation := ctrl.ExecutionGeneration()
610 handoffIdentityBeforeLockHookForTest = func() {
611 runtime.NoteExecution(generation, session.RuntimeFinalizing, "terminal_commit")
612 }
613 t.Cleanup(func() { handoffIdentityBeforeLockHookForTest = nil })
614
615 body := `{"sessionPath":"` + route + `","targetWriterId":"taker-writer","force":true,"mode":"wait","timeoutMs":2000}`
616 resp, raw := serveBody(t, http.MethodPost, ts.URL+"/handoff", body)
617 if resp.StatusCode != http.StatusConflict || !strings.Contains(raw, errHandoffBusyAgain.Error()) {
618 t.Fatalf("handoff against a finalizing runtime = %d %q, want 409 busy-again", resp.StatusCode, raw)
619 }
620 if ref, bound := ctrl.SessionRef(); !bound || ref != current {
621 t.Fatalf("foreground binding after refused handoff = %+v (bound %v), want %+v", ref, bound, current)
622 }
623 if _, mirrored := lifecycle.mirroredEntry(route); mirrored {
624 t.Fatal("refused handoff registered a mirror entry")
625 }
626
627 handoffIdentityBeforeLockHookForTest = nil
628 runtime.NoteExecution(generation, session.RuntimeIdle, "")
629 resp, raw = serveBody(t, http.MethodPost, ts.URL+"/handoff", body)
630 if resp.StatusCode != http.StatusOK {
631 t.Fatalf("handoff after the runtime settled = %d %q, want 200", resp.StatusCode, raw)
632 }
633 if _, bound := ctrl.SessionRef(); bound {
634 t.Fatal("foreground still bound after the successful handoff")
635 }
636 }
637
638 // registerDetachedIdentityHolderForTest parks ctrl in the background registry
639 // the way a detached legacy controller ends up there after upgrading to an
640 // identity mid-turn: keyed by its former transcript path while SessionRef
641 // reports the identity. The close-on-idle watcher is omitted because these
642 // tests assert ownership predicates, not idle retirement; done is pre-closed so
643 // takeDetached hands the entry over immediately.
644 func (s *Server) registerDetachedIdentityHolderForTest(key string, ctrl *control.Controller, tag *sessionTagSink) *detachedSession {
645 done := make(chan struct{})
646 close(done)
647 d := &detachedSession{
648 path: agent.CanonicalSessionPath(key), ctrl: ctrl, tag: tag,
649 force: make(chan struct{}), reattach: make(chan struct{}), done: done,
650 }
651 s.detachedMu.Lock()
652 s.detached[d.path] = d
653 s.detachedMu.Unlock()
654 return d
655 }
656
657 // A detached background session that runs an identity is still this serve's
658 // writer. /ownership must say so instead of "other", /adopt must refuse the
659 // claim instead of registering a mirror over our own writer, /status must not
660 // pin the tab read-only, and /handoff must release the detached holder the way
661 // the legacy detached handoff does.
662 func TestIdentityOwnershipCoversDetachedHolder(t *testing.T) {
663 _, ctrl, service, current := newExclusiveSessionServe(t)
664 root := identityRoot(t, service, current)
665 lifecycle := newIdentityLifecycleServe(t, ctrl, current)
666 ts := httptest.NewServer(operatorHandler(lifecycle))
667 defer ts.Close()
668 defer retireExclusiveForeground(t, ctrl, service)
669
670 background, err := service.Create(t.Context(), session.CreateOptions{SessionID: "background"})
671 if err != nil {
672 t.Fatal(err)
673 }
674 if err := service.Close(t.Context(), background.Ref()); err != nil {
675 t.Fatal(err)
676 }
677 exec := agent.New(nil, nil, agent.NewSession("system"), agent.Options{}, event.Discard)
678 bg := control.New(control.Options{Runner: blockingRunner{}, Executor: exec, SessionDir: ctrl.SessionDir(), SessionService: service, ExclusiveSession: true})
679 if _, err := bg.OpenSession(t.Context(), background.Ref()); err != nil {
680 t.Fatal(err)
681 }
682 tag := newSessionTagSink(lifecycle.bc)
683 tag.SetIdentity("", background.Ref().SessionID)
684 lifecycle.RegisterSessionTag(bg, tag)
685 bg.Submit("keep running")
686 waitRunning(t, bg)
687 lifecycle.registerDetachedIdentityHolderForTest(filepath.Join(ctrl.SessionDir(), "upgraded.jsonl"), bg, tag)
688 closed := false
689 defer func() {
690 if !closed {
691 bg.Cancel()
692 waitNotRunning(t, bg)
693 bg.Close()
694 }
695 }()
696 route := "session-id:" + background.Ref().SessionID
697
698 resp, raw := serveBody(t, http.MethodGet, ts.URL+"/ownership?session="+route, "")
699 var view ownershipView
700 if err := json.Unmarshal([]byte(raw), &view); err != nil {
701 t.Fatal(err)
702 }
703 if resp.StatusCode != http.StatusOK || view.Holder != "serve" || !view.Running {
704 t.Fatalf("ownership of a detached identity = %+v (status %d), want serve holder running", view, resp.StatusCode)
705 }
706 resp, raw = serveBody(t, http.MethodPost, ts.URL+"/adopt", `{"sessionPath":"`+route+`","writerId":"impostor"}`)
707 if resp.StatusCode != http.StatusConflict {
708 t.Fatalf("adopt of our own detached identity = %d body %s, want 409", resp.StatusCode, raw)
709 }
710 if _, mirrored := lifecycle.mirroredEntry(route); mirrored {
711 t.Fatal("adopt registered a mirror over this serve's own detached writer")
712 }
713 resp, raw = serveBody(t, http.MethodGet, ts.URL+"/status?session="+route, "")
714 var status map[string]any
715 if err := json.Unmarshal([]byte(raw), &status); err != nil {
716 t.Fatal(err)
717 }
718 if taken, _ := status["takenOver"].(bool); resp.StatusCode != http.StatusOK || taken {
719 t.Fatalf("status of a detached identity reports takenOver: %v (status %d)", status, resp.StatusCode)
720 }
721
722 // The detached holder is handed off like a legacy detached session: the
723 // turn is interrupted, the writer lock drops, the controller is retired.
724 resp, raw = serveBody(t, http.MethodPost, ts.URL+"/handoff", `{"sessionPath":"`+route+`","targetWriterId":"taker-writer","force":true,"mode":"interrupt","timeoutMs":5000}`)
725 if resp.StatusCode != http.StatusOK {
726 t.Fatalf("handoff of a detached identity = %d body %s", resp.StatusCode, raw)
727 }
728 closed = true
729 if lifecycle.detachedBusy(filepath.Join(ctrl.SessionDir(), "upgraded.jsonl")) {
730 t.Fatal("handed-off detached holder is still registered")
731 }
732 if session.ProbeWriterHeld(filepath.Join(root, background.Ref().SessionID)) {
733 t.Fatal("writer lock still held after handing off the detached identity")
734 }
735 if _, mirrored := lifecycle.mirroredEntry(route); !mirrored {
736 t.Fatal("handoff of the detached identity did not register a mirror entry")
737 }
738 if ref, bound := ctrl.SessionRef(); !bound || ref != current {
739 t.Fatalf("foreground changed while handing off a detached identity: %+v (bound %v)", ref, bound)
740 }
741 }
742
742 lines GO