返回 DeepSeek-Reasonix
multisession.go
根目录 / internal / serve / multisession.go
1 package serve
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "net/http"
9 "os"
10 "strings"
11 "sync"
12 "time"
13
14 "reasonix/internal/agent"
15 "reasonix/internal/boot"
16 "reasonix/internal/config"
17 "reasonix/internal/control"
18 "reasonix/internal/event"
19 "reasonix/internal/plugin"
20 )
21
22 // sessionTagSink stamps every event from one controller with that
23 // controller's current session path. This lets one Serve process keep several
24 // turns alive without sending a background session's frames to the foreground
25 // browser.
26 type sessionTagSink struct {
27 pendingRuntimeState *event.RuntimeStateSnapshot
28 bc *Broadcaster
29 mu sync.Mutex
30 path string
31 sessionID string
32 active bool
33 runtimeActive bool
34 pending []event.Event
35 }
36
37 func newSessionTagSink(bc *Broadcaster) *sessionTagSink {
38 return &sessionTagSink{bc: bc, runtimeActive: true}
39 }
40
41 // SessionTagSink is exported for the CLI, which builds Serve's initial
42 // controller before the Server exists.
43 type SessionTagSink = sessionTagSink
44
45 func NewSessionTagSink(bc *Broadcaster) *SessionTagSink {
46 return newSessionTagSink(bc)
47 }
48
49 func (s *sessionTagSink) SetPath(path string) {
50 s.SetIdentity(path, "")
51 }
52
53 func (s *sessionTagSink) SetIdentity(path, sessionID string) {
54 s.mu.Lock()
55 if s.path != "" && s.path != canonicalSessionPath(path) {
56 s.runtimeActive = false
57 }
58 s.path = canonicalSessionPath(path)
59 s.sessionID = strings.TrimSpace(sessionID)
60 s.activateLocked()
61 s.mu.Unlock()
62 }
63
64 // PrimePath assigns a replacement controller's route without publishing boot
65 // events. Activate is called only after the controller swap fully commits.
66 func (s *sessionTagSink) PrimePath(path string) {
67 s.mu.Lock()
68 s.path = canonicalSessionPath(path)
69 s.runtimeActive = s.bc.CurrentSession() == s.path
70 s.mu.Unlock()
71 }
72
73 // PrimeIdentity assigns a replacement controller's complete route without
74 // publishing its buffered boot events. Canonical v3 sessions have no legacy
75 // path, so PrimePath alone would drop the session ID from those events.
76 func (s *sessionTagSink) PrimeIdentity(path, sessionID string) {
77 s.mu.Lock()
78 s.path = canonicalSessionPath(path)
79 s.sessionID = strings.TrimSpace(sessionID)
80 s.runtimeActive = s.bc.CurrentSession() == s.path
81 s.mu.Unlock()
82 }
83
84 // BufferPath retags synchronous in-place Resume events but withholds them until
85 // Serve publishes the matching foreground route. Unlike PrimePath, it also
86 // pauses a sink that was already active for the previous session.
87 func (s *sessionTagSink) BufferPath(path string) {
88 s.mu.Lock()
89 s.path = canonicalSessionPath(path)
90 s.active = false
91 s.runtimeActive = false
92 s.mu.Unlock()
93 }
94
95 func canonicalSessionPath(path string) string {
96 if path != "" {
97 return agent.CanonicalSessionPath(path)
98 }
99 return ""
100 }
101
102 func (s *sessionTagSink) Activate() {
103 s.mu.Lock()
104 s.activateLocked()
105 s.mu.Unlock()
106 }
107
108 func (s *sessionTagSink) activateLocked() {
109 if s.active {
110 return
111 }
112 s.active = true
113 if s.runtimeActive && s.pendingRuntimeState != nil {
114 s.bc.publishRuntimeState(s.path, *s.pendingRuntimeState)
115 s.pendingRuntimeState = nil
116 }
117 for _, e := range s.pending {
118 if s.path != "" {
119 e.SessionPath = s.path
120 }
121 if s.sessionID != "" {
122 e.SessionID = s.sessionID
123 }
124 s.bc.Emit(e)
125 }
126 s.pending = nil
127 }
128
129 func (s *sessionTagSink) Path() string {
130 s.mu.Lock()
131 defer s.mu.Unlock()
132 return s.path
133 }
134
135 func (s *sessionTagSink) Emit(e event.Event) {
136 s.mu.Lock()
137 defer s.mu.Unlock()
138 if !s.active {
139 s.pending = append(s.pending, e)
140 return
141 }
142 if s.path != "" {
143 e.SessionPath = s.path
144 }
145 if s.sessionID != "" {
146 e.SessionID = s.sessionID
147 }
148 s.bc.Emit(e)
149 }
150
151 type detachedSession struct {
152 admissionMu sync.Mutex // new-run refresh and close-on-idle ownership
153 modelSettings *config.ModelRuntimeSettings
154 modelSettingsOfferID string
155 buildOptions boot.Options
156 path string
157 ctrl control.SessionAPI
158 keeper *control.SessionLeaseKeeper
159 tag *sessionTagSink
160 retiring bool // guarded by Server.detachedMu; blocks reattach during Close
161 force chan struct{}
162 reattach chan struct{}
163 done chan struct{}
164 }
165
166 // RegisterSessionTag associates a controller built outside Server with its
167 // tagging sink. In-place /new and /resume operations can then advance the tag.
168 func (s *Server) RegisterSessionTag(ctrl *control.Controller, tag *sessionTagSink) {
169 if ctrl == nil || tag == nil {
170 return
171 }
172 s.tagsMu.Lock()
173 if s.tags == nil {
174 s.tags = map[*control.Controller]*sessionTagSink{}
175 }
176 s.tags[ctrl] = tag
177 s.tagsMu.Unlock()
178 }
179
180 func (s *Server) tagFor(ctrl *control.Controller) *sessionTagSink {
181 if ctrl == nil {
182 return nil
183 }
184 s.tagsMu.Lock()
185 defer s.tagsMu.Unlock()
186 return s.tags[ctrl]
187 }
188
189 func (s *Server) forgetSessionTag(ctrl *control.Controller) {
190 if ctrl == nil {
191 return
192 }
193 s.tagsMu.Lock()
194 delete(s.tags, ctrl)
195 s.tagsMu.Unlock()
196 s.setControllerLeaseOwner(ctrl, nil)
197 }
198
199 func (s *Server) closeTaggedController(ctrl *control.Controller) {
200 if ctrl == nil {
201 return
202 }
203 ctrl.Close()
204 s.forgetSessionTag(ctrl)
205 }
206
207 func (s *Server) setControllerPath(ctrl *control.Controller, path string) {
208 if path != "" {
209 path = agent.CanonicalSessionPath(path)
210 }
211 if tag := s.tagFor(ctrl); tag != nil {
212 sessionID := ""
213 if ref, ok := ctrl.SessionRef(); ok {
214 sessionID = ref.SessionID
215 }
216 tag.SetIdentity(path, sessionID)
217 }
218 s.bc.SetCurrentSession(path)
219 }
220
221 // buildTagged creates a controller whose frames are session-tagged. The legacy
222 // two-argument test builder remains supported; tests that need to assert the
223 // complete boot contract can inject buildControllerWithOptions.
224 func (s *Server) buildTagged(ctx context.Context, ref string, inheritTemp bool) (*control.Controller, *sessionTagSink, error) {
225 return s.buildTaggedMode(ctx, ref, inheritTemp, false)
226 }
227
228 func (s *Server) buildTaggedMode(ctx context.Context, ref string, inheritTemp, nativeLegacy bool) (*control.Controller, *sessionTagSink, error) {
229 tag := newSessionTagSink(s.bc)
230 opts := s.buildOptions
231 opts.NativeLegacySession = nativeLegacy
232 if nativeLegacy {
233 opts.SessionRuntime = nil
234 }
235 if s.managedModels != nil {
236 opts.ModelSettings = s.managedModels
237 }
238 opts.Model = ref
239 if cur := s.ctl(); cur != nil {
240 opts.EffortModel = currentModelRef(cur)
241 }
242 opts.BeforeInboxDispatch = s.beforeInboxDispatch
243 opts.Sink = tag
244 opts.BrowserExecutor = s.sessionBrowserExecutor(tag)
245 if opts.Stderr == nil {
246 opts.Stderr = os.Stderr
247 }
248 opts.StatsSource = "serve"
249 opts.MCPHostProfile = plugin.HostProfileInteractive
250 if cur, ok := s.ctl().(*control.Controller); ok && cur != nil {
251 opts.SessionDir = cur.SessionDir()
252 opts.WorkspaceRoot = cur.WorkspaceRoot()
253 if inheritTemp {
254 opts.SessionTemp = cur.SessionTemp()
255 opts.PersistentShell = cur.PersistentShell()
256 // Model/effort switches rebuild the Agent, not the logical session:
257 // without the bound runtime the rebuilt controller's first submit
258 // allocates a new ID while the desktop fences the old one (HTTP 409).
259 if service, runtime, bound := cur.SessionBinding(); bound {
260 opts.SessionService = service
261 opts.SessionRuntime = runtime
262 opts.SessionHostID = runtime.Ref().HostID
263 }
264 }
265 }
266
267 ctrl, err := s.buildModelCandidate(ctx, ref, opts, inheritTemp)
268 if err != nil {
269 return nil, nil, err
270 }
271 ctrl.EnableServeSessionPermissionPresets(true)
272 s.RegisterSessionTag(ctrl, tag)
273 slog.Info("serve: controller built", "model", ref, "sessionDir", opts.SessionDir)
274 return ctrl, tag, nil
275 }
276
277 func (s *Server) detachedBusy(path string) bool {
278 path = sessionRouteKey(path)
279 s.detachedMu.Lock()
280 defer s.detachedMu.Unlock()
281 _, ok := s.detached[path]
282 return ok
283 }
284
285 // takeDetached transfers ownership from the close-on-idle watcher back to the
286 // request goroutine. Waiting for done is essential: without the acknowledgement
287 // the watcher can close an idle controller just after it is re-attached.
288 func (s *Server) takeDetached(path string) *detachedSession {
289 path = sessionRouteKey(path)
290 s.detachedMu.Lock()
291 d := s.detached[path]
292 if d != nil && !d.retiring {
293 delete(s.detached, path)
294 } else {
295 d = nil
296 }
297 s.detachedMu.Unlock()
298 if d == nil {
299 return nil
300 }
301 close(d.reattach)
302 <-d.done
303 return d
304 }
305
306 func (s *Server) registerDetached(ctrl control.SessionAPI, keeper *control.SessionLeaseKeeper, tag *sessionTagSink) (*detachedSession, error) {
307 if ctrl == nil {
308 return nil, fmt.Errorf("cannot detach a nil controller")
309 }
310 if tag == nil {
311 if concrete, ok := ctrl.(*control.Controller); ok {
312 tag = s.tagFor(concrete)
313 }
314 }
315 if tag == nil {
316 return nil, errSessionTagUnavailable
317 }
318 if keeper != nil {
319 if concrete, ok := ctrl.(*control.Controller); ok {
320 concrete.SetOnSessionRecovered(s.sessionRecoveryHandler(concrete, keeper))
321 }
322 }
323 if registerDetachedHookForTest != nil {
324 registerDetachedHookForTest()
325 }
326 d := &detachedSession{
327 ctrl: ctrl, keeper: keeper, tag: tag,
328 modelSettings: s.managedModels, modelSettingsOfferID: s.modelSettingsOfferID, buildOptions: s.buildOptions,
329 force: make(chan struct{}), reattach: make(chan struct{}), done: make(chan struct{}),
330 }
331 s.detachedMu.Lock()
332 path := agent.CanonicalSessionPath(ctrl.SessionPath())
333 if path == "" {
334 // Identity (exclusive v3) sessions deliberately carry no legacy path;
335 // their detached key is the immutable session-id route. Transcript and
336 // runtime reads already resolve detached controllers by identity.
337 if ref, ok := sessionAPIRef(ctrl); ok && ref.SessionID != "" {
338 path = remoteSessionIDQueryPrefix + ref.SessionID
339 }
340 }
341 if path == "" {
342 s.detachedMu.Unlock()
343 return nil, fmt.Errorf("cannot detach a session without a path")
344 }
345 d.path = path
346 if s.detached == nil {
347 s.detached = map[string]*detachedSession{}
348 }
349 if _, exists := s.detached[path]; exists {
350 s.detachedMu.Unlock()
351 return nil, fmt.Errorf("session is already running in the background")
352 }
353 s.detached[path] = d
354 s.detachedMu.Unlock()
355 slog.Info("serve: session detached", "session", path, "running", controllerHasActiveRuntimeWork(ctrl))
356 go s.watchDetached(d)
357 if concrete, ok := ctrl.(*control.Controller); ok {
358 concrete.NotifyInboxRuntimeReady()
359 }
360 return d, nil
361 }
362
363 func (s *Server) watchDetached(d *detachedSession) {
364 interval := 200 * time.Millisecond
365 forced := false
366 for s.detachedHasPendingWork(d) && !forced {
367 timer := time.NewTimer(interval)
368 select {
369 case <-d.reattach:
370 if !timer.Stop() {
371 <-timer.C
372 }
373 close(d.done)
374 return
375 case <-d.force:
376 if !timer.Stop() {
377 <-timer.C
378 }
379 forced = true
380 case <-timer.C:
381 if interval < 2*time.Second {
382 interval *= 2
383 }
384 }
385 }
386
387 // Claim close ownership only while the registry still points at d. Keep the
388 // retiring entry visible until Close and lease release finish so deletion
389 // cannot race final controller writes. takeDetached refuses retiring entries.
390 d.admissionMu.Lock()
391 s.detachedMu.Lock()
392 owns := s.detached[d.path] == d
393 if owns {
394 d.retiring = true
395 }
396 s.detachedMu.Unlock()
397 d.admissionMu.Unlock()
398 if !owns {
399 close(d.done)
400 return
401 }
402 d.ctrl.Close()
403 if d.keeper != nil {
404 d.keeper.Release()
405 }
406 if concrete, ok := d.ctrl.(*control.Controller); ok {
407 s.forgetSessionTag(concrete)
408 }
409 s.detachedMu.Lock()
410 closedPath := d.path
411 if s.detached[d.path] == d {
412 delete(s.detached, d.path)
413 }
414 s.detachedMu.Unlock()
415 slog.Info("serve: background session closed", "session", closedPath, "forced", forced)
416 close(d.done)
417 }
418
419 func (s *Server) WaitForDetachedIdle() {
420 s.detachedMu.Lock()
421 detached := make([]*detachedSession, 0, len(s.detached))
422 for _, d := range s.detached {
423 detached = append(detached, d)
424 }
425 s.detachedMu.Unlock()
426 for _, d := range detached {
427 <-d.done
428 }
429 }
430
431 func (s *Server) CloseBackground() {
432 s.detachedMu.Lock()
433 detached := make([]*detachedSession, 0, len(s.detached))
434 for _, d := range s.detached {
435 detached = append(detached, d)
436 }
437 s.detachedMu.Unlock()
438 for _, d := range detached {
439 select {
440 case <-d.force:
441 default:
442 close(d.force)
443 }
444 }
445 for _, d := range detached {
446 <-d.done
447 }
448 }
449
450 // Close stops every controller the server owns, including a foreground
451 // replacement created after the CLI's original controller was constructed.
452 func (s *Server) Close() {
453 s.CloseBackground()
454 cur := s.ctl()
455 cur.Close()
456 if concrete, ok := cur.(*control.Controller); ok {
457 s.forgetSessionTag(concrete)
458 }
459 }
460
461 // busyDetach publishes a fresh controller before demoting a busy controller.
462 // Every failure before publication restores the original lease ownership.
463 func (s *Server) busyDetach(ctx context.Context, cur *control.Controller, targetPath string, loadTarget func(*control.Controller) error) error {
464 if s.tagFor(cur) == nil {
465 return errSessionTagUnavailable
466 }
467 newCtrl, tag, err := s.buildTaggedMode(ctx, currentModelRef(cur), false, strings.TrimSpace(targetPath) != "")
468 if err != nil {
469 return err
470 }
471 if targetPath == "" {
472 newCtrl.EnsureSessionPath()
473 targetPath = newCtrl.SessionPath()
474 }
475 targetPath = agent.CanonicalSessionPath(targetPath)
476 if targetPath == "" {
477 s.closeTaggedController(newCtrl)
478 return fmt.Errorf("replacement session has no path")
479 }
480
481 demoted, err := s.leases.RebindDetaching(targetPath)
482 if err != nil {
483 s.closeTaggedController(newCtrl)
484 return err
485 }
486 if loadTarget != nil {
487 if err := loadTarget(newCtrl); err != nil {
488 s.closeTaggedController(newCtrl)
489 s.rollbackDetach(demoted, cur)
490 return err
491 }
492 }
493 // PrimeIdentity keeps the session id alive across the swap: exclusive
494 // sessions have no path, and PrimePath alone would strip the id from every
495 // frame the replacement controller emits.
496 if ref, ok := newCtrl.SessionRef(); ok {
497 tag.PrimeIdentity(targetPath, ref.SessionID)
498 } else {
499 tag.PrimePath(targetPath)
500 }
501 newCtrl.EnableInteractiveApproval()
502 newCtrl.SetOnSessionRecovered(s.sessionRecoveryHandler(newCtrl, s.leases))
503 if s.leases != nil {
504 if err := s.leases.BindControllerAuthority(newCtrl); err != nil {
505 s.closeTaggedController(newCtrl)
506 s.rollbackDetach(demoted, cur)
507 return err
508 }
509 }
510
511 if !s.publishControllerSwap(cur, newCtrl, targetPath) {
512 s.closeTaggedController(newCtrl)
513 s.rollbackDetach(demoted, cur)
514 return errReplacedDuringBind
515 }
516
517 if _, err := s.registerDetached(cur, demoted, nil); err != nil {
518 // bindMu prevents another foreground swap here. Roll publication back so
519 // a registry failure cannot strand a running controller.
520 _ = s.publishControllerSwap(newCtrl, cur, cur.SessionPath())
521 s.closeTaggedController(newCtrl)
522 s.rollbackDetach(demoted, cur)
523 return err
524 }
525 tag.Activate()
526 s.bc.ResetSessionPath(targetPath)
527 return nil
528 }
529
530 func (s *Server) announceSessionChanged(path string, reset bool) {
531 e := event.Event{Kind: event.SessionChanged, SessionPath: path, SessionReset: reset}
532 if identity, ok := s.ctl().(control.IdentityLifecycle); ok {
533 if ref, bound := identity.SessionRef(); bound {
534 e.SessionID = ref.SessionID
535 }
536 }
537 s.bc.Emit(e)
538 if ctrl, ok := s.ctl().(*control.Controller); ok {
539 if tag := s.tagFor(ctrl); tag != nil {
540 tag.ActivateRuntime()
541 }
542 }
543 }
544
545 // The routing barrier precedes the new instance's runtime projection. Content
546 // boot notices retain their historical order relative to session_changed.
547 func (s *sessionTagSink) ActivateRuntime() {
548 s.mu.Lock()
549 defer s.mu.Unlock()
550 s.runtimeActive = true
551 if s.active && s.pendingRuntimeState != nil {
552 s.bc.publishRuntimeState(s.path, *s.pendingRuntimeState)
553 s.pendingRuntimeState = nil
554 }
555 }
556
557 var errReplacedDuringBind = &replacedDuringBindError{}
558 var errSessionTagUnavailable = errors.New("multi-session switching requires a session-tagged Serve controller")
559 var errIdentityServiceUnavailable = errors.New("identity switching requires a session-service-capable Serve controller")
560
561 type replacedDuringBindError struct{}
562
563 func (*replacedDuringBindError) Error() string { return "session changed during switch" }
564
565 func (s *Server) rollbackDetach(demoted *control.SessionLeaseKeeper, ctrl *control.Controller) {
566 if demoted != nil && s.leases != nil {
567 s.leases.Adopt(demoted)
568 }
569 if ctrl != nil {
570 ctrl.SetOnSessionRecovered(s.sessionRecoveryHandler(ctrl, s.leases))
571 }
572 }
573
574 func (s *Server) resumeActiveSession(w http.ResponseWriter, r *http.Request, cur control.SessionAPI, realPath string) bool {
575 if agent.CanonicalSessionPath(cur.SessionPath()) == agent.CanonicalSessionPath(realPath) {
576 s.bc.SetCurrentSession(realPath)
577 s.announceSessionChanged(realPath, false)
578 w.WriteHeader(http.StatusNoContent)
579 return true
580 }
581 if detached := s.takeDetached(realPath); detached != nil {
582 if err := s.reattachDetached(cur, detached); err != nil {
583 s.renderBindError(w, err)
584 return true
585 }
586 s.announceSessionChanged(realPath, false)
587 w.WriteHeader(http.StatusNoContent)
588 s.replayPendingPromptsBroadcast()
589 return true
590 }
591 if s.detachedBusy(realPath) {
592 http.Error(w, "session is finishing background teardown; retry shortly", http.StatusConflict)
593 return true
594 }
595 if !controllerHasActiveRuntimeWork(cur) {
596 return false
597 }
598 curCtrl, ok := cur.(*control.Controller)
599 if !ok {
600 http.Error(w, "cannot switch session while active work or background jobs are running", http.StatusConflict)
601 return true
602 }
603 err := s.busyDetach(r.Context(), curCtrl, realPath, func(next *control.Controller) error {
604 loaded, err := agent.LoadSession(realPath)
605 if err == nil {
606 next.Resume(loaded, realPath)
607 }
608 return err
609 })
610 if err != nil {
611 s.renderBindError(w, err)
612 return true
613 }
614 s.announceSessionChanged(realPath, false)
615 w.WriteHeader(http.StatusNoContent)
616 s.replayPendingPromptsBroadcast()
617 return true
618 }
619
620 // reattachDetached promotes a controller owned by the background registry.
621 // bindMu is held, so publication and lease ownership move as one transaction.
622 func (s *Server) reattachDetached(cur control.SessionAPI, detached *detachedSession) error {
623 curCtrl, _ := cur.(*control.Controller)
624 demoted := s.leases.Split()
625 s.leases.Adopt(detached.keeper)
626 detached.keeper = nil
627 if detached.tag != nil {
628 if ref, ok := sessionAPIRef(detached.ctrl); ok && ref.SessionID != "" && detached.ctrl.SessionPath() == "" {
629 detached.tag.SetIdentity("", ref.SessionID)
630 } else {
631 detached.tag.SetPath(detached.ctrl.SessionPath())
632 }
633 }
634 if concrete, ok := detached.ctrl.(*control.Controller); ok {
635 concrete.SetOnSessionRecovered(s.sessionRecoveryHandler(concrete, s.leases))
636 }
637 if !s.publishControllerSwap(cur, detached.ctrl, detached.ctrl.SessionPath()) {
638 detached.keeper = s.leases.Split()
639 s.leases.Adopt(demoted)
640 if curCtrl != nil {
641 curCtrl.SetOnSessionRecovered(s.sessionRecoveryHandler(curCtrl, s.leases))
642 }
643 _, _ = s.registerDetached(detached.ctrl, detached.keeper, detached.tag)
644 return errReplacedDuringBind
645 }
646 if controllerHasActiveRuntimeWork(cur) {
647 if curCtrl == nil {
648 s.restoreReattach(cur, detached, demoted)
649 return fmt.Errorf("cannot switch session while active work or background jobs are running")
650 }
651 if _, err := s.registerDetached(curCtrl, demoted, nil); err != nil {
652 s.restoreReattach(cur, detached, demoted)
653 return err
654 }
655 } else {
656 if err := cur.Snapshot(); err != nil {
657 slog.Warn("serve: snapshot before background reattach", "err", err)
658 }
659 cur.Close()
660 if curCtrl != nil {
661 s.forgetSessionTag(curCtrl)
662 }
663 if demoted != nil {
664 demoted.Release()
665 }
666 }
667 s.managedModels, s.modelSettingsOfferID = detached.modelSettings, detached.modelSettingsOfferID
668 s.buildOptions = detached.buildOptions
669 slog.Info("serve: background session re-attached", "session", detached.path, "running", controllerHasActiveRuntimeWork(detached.ctrl))
670 return nil
671 }
672
673 func (s *Server) restoreReattach(cur control.SessionAPI, detached *detachedSession, demoted *control.SessionLeaseKeeper) {
674 _ = s.publishControllerSwap(detached.ctrl, cur, cur.SessionPath())
675 detached.keeper = s.leases.Split()
676 s.leases.Adopt(demoted)
677 if concrete, ok := cur.(*control.Controller); ok {
678 concrete.SetOnSessionRecovered(s.sessionRecoveryHandler(concrete, s.leases))
679 }
680 _, _ = s.registerDetached(detached.ctrl, detached.keeper, detached.tag)
681 }
682
683 // publishControllerSwap makes the command target and current-only SSE route
684 // visible as one generation. Readers cannot observe next while the broadcaster
685 // still filters against expect's session.
686 func (s *Server) publishControllerSwap(expect, next control.SessionAPI, path string) bool {
687 s.mu.Lock()
688 defer s.mu.Unlock()
689 if s.ctrl != expect {
690 return false
691 }
692 if err := control.ActivateSessionAPIReplacement(expect, next); err != nil {
693 slog.Warn("serve: activate controller replacement", "err", err)
694 return false
695 }
696 s.ctrl = next
697 s.bc.SetCurrentSession(path)
698 if ctrl, ok := next.(*control.Controller); ok {
699 ctrl.SetBeforeInboxDispatch(s.beforeInboxDispatch)
700 ctrl.NotifyInboxRuntimeReady()
701 }
702 return true
703 }
704
705 func (s *Server) replayPendingPromptsBroadcast() {
706 cur := s.ctl()
707 path := cur.SessionPath()
708 cur.ReplayPendingPromptsWith(func() event.Sink {
709 return event.FuncSink(func(e event.Event) {
710 e.SessionPath = path
711 s.bc.Emit(e)
712 })
713 })
714 }
715
716 func (s *Server) renderBindError(w http.ResponseWriter, err error) {
717 switch {
718 case errors.Is(err, agent.ErrSessionLeaseHeld):
719 http.Error(w, sessionInUseError(err), http.StatusConflict)
720 case errors.Is(err, errReplacedDuringBind), errors.Is(err, errSessionTagUnavailable):
721 http.Error(w, err.Error(), http.StatusConflict)
722 default:
723 http.Error(w, "switch session: "+err.Error(), http.StatusInternalServerError)
724 }
725 }
726
727 // retireDetachedForProviderHeal is called with bindMu held. It makes every
728 // detached controller unreattachable and waits for its provider generation to
729 // close before credential reload is acknowledged.
730 func (s *Server) retireDetachedForProviderHeal() {
731 s.detachedMu.Lock()
732 detached := make([]*detachedSession, 0, len(s.detached))
733 for _, d := range s.detached {
734 detached = append(detached, d)
735 select {
736 case <-d.force:
737 default:
738 close(d.force)
739 }
740 }
741 s.detachedMu.Unlock()
742 for _, d := range detached {
743 slog.Info("serve: provider heal retires background session", "session", d.path)
744 <-d.done
745 }
746 }
747
747 lines GO