返回 DeepSeek-Reasonix
session_binding.go
根目录 / internal / control / session_binding.go
1 package control
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "log/slog"
10 "strings"
11
12 "reasonix/internal/agent"
13 "reasonix/internal/event"
14 "reasonix/internal/extension"
15 "reasonix/internal/extension/dispatch"
16 "reasonix/internal/permissionpreset"
17 "reasonix/internal/provider"
18 "reasonix/internal/session"
19 )
20
21 // NativeLegacySession identifies a path-backed controller across hot rebuilds.
22 func (c *Controller) NativeLegacySession() bool {
23 if c == nil {
24 return false
25 }
26 c.v3BindingMu.RLock()
27 defer c.v3BindingMu.RUnlock()
28 return c.nativeLegacySession
29 }
30
31 // ResumeNativeSession switches to an already validated, path-backed transcript.
32 // Hosts serialize admission and acquire its path lease before calling this.
33 func (c *Controller) ResumeNativeSession(s *agent.Session, path string) error {
34 if s == nil || strings.TrimSpace(path) == "" {
35 return errors.New("native session and path are required")
36 }
37 if service := c.SessionService(); service != nil {
38 if ref, found, err := service.ExistingCanonicalForLegacy(path, s); err != nil {
39 return err
40 } else if found {
41 _, err = c.OpenSession(context.Background(), ref)
42 return err
43 }
44 }
45 if err := c.Snapshot(); err != nil {
46 return err
47 }
48 _, previousRuntime, _ := c.v3Binding()
49 if err := c.ReleaseSessionRuntimeBinding(); err != nil {
50 return err
51 }
52 c.v3BindingMu.Lock()
53 c.exclusiveSession = false
54 c.nativeLegacySession = true
55 c.v3BindingMu.Unlock()
56 // The retired canonical cache is owned by its service, never by the path
57 // adapter. Drop the borrowed pointer before rebinding legacy state.
58 c.turnEvents.mu.Lock()
59 if previousRuntime != nil {
60 c.turnEvents.v3 = nil
61 c.turnEvents.v3Runtime = nil
62 c.turnEvents.v3Path = ""
63 c.turnEvents.v3Release = nil
64 }
65 c.turnEvents.mu.Unlock()
66 c.Resume(s, path)
67 return c.turnEventLedgerError()
68 }
69
70 func bindInitialSessionRuntime(opts Options) (*session.Runtime, *session.ClientBinding) {
71 runtime := opts.SessionRuntime
72 if opts.SessionService == nil || runtime == nil {
73 return runtime, nil
74 }
75 binding, err := opts.SessionService.Bind(runtime)
76 if err != nil {
77 return nil, nil
78 }
79 return runtime, binding
80 }
81
82 // ReleaseSessionRuntimeBinding drops this controller's client reference without
83 // tearing down the controller itself. Hosts use it when a session is handed
84 // back to another runtime while keeping the current process alive.
85 func (c *Controller) ReleaseSessionRuntimeBinding() error {
86 if c == nil {
87 return nil
88 }
89 c.v3BindingMu.Lock()
90 binding := c.sessionBinding
91 runtime := c.sessionRuntime
92 c.sessionBinding = nil
93 c.sessionRuntime = nil
94 c.v3BindingMu.Unlock()
95 c.unbindExecutionControl(runtime)
96 if binding != nil {
97 return binding.Release(context.Background())
98 }
99 return nil
100 }
101
102 // ReleaseSessionForHandoff hands the bound identity to another runtime without
103 // allocating a replacement: it flushes the runtime, drops this controller's
104 // client binding and empties the in-memory transcript, leaving the controller in
105 // the never-bound exclusive state whose next turn or NewSession allocates a
106 // fresh identity lazily. Closing the runtime, which drops the writer lock, stays
107 // with the host: the host owns the rollback (OpenSession) when that close is
108 // refused, and this controller is exactly re-attachable until then.
109 func (c *Controller) ReleaseSessionForHandoff() error {
110 if c == nil {
111 return nil
112 }
113 if _, runtime, exclusive := c.v3Binding(); !exclusive || runtime == nil {
114 return session.ErrSessionNotRunning
115 }
116 if err := c.Snapshot(); err != nil {
117 return err
118 }
119 if err := c.ReleaseSessionRuntimeBinding(); err != nil {
120 return err
121 }
122 // Under snapshotMu so the swap cannot interleave with an in-flight save.
123 // Emptying the transcript keeps a later Snapshot a no-op and keeps the
124 // handed-off conversation out of the identity the next turn allocates.
125 c.snapshotMu.Lock()
126 if c.executor != nil {
127 c.executor.SetSession(agent.NewSession(c.basePrompt()))
128 }
129 c.snapshotMu.Unlock()
130 // With no runtime bound an exclusive controller has no event store, so
131 // history, transcript pages and admission answer from nothing — the released
132 // runtime's cached store must not keep serving the handed-off conversation.
133 c.rebindTurnEvents("")
134 return nil
135 }
136
137 // releaseSessionRuntimeBinding is the final controller teardown wrapper. It
138 // keeps the handoff-only release available without making normal Close paths
139 // responsible for surfacing a late binding-release error.
140 func (c *Controller) releaseSessionRuntimeBinding(service *session.Service) {
141 if err := c.ReleaseSessionRuntimeBinding(); err != nil {
142 slog.Warn("controller: release exclusive v3 binding", "err", err)
143 } else if service == nil {
144 slog.Warn("controller: exclusive v3 runtime has no service binding")
145 }
146 }
147
148 // BindFreshSession creates and publishes a fresh identity-bound session. The caller
149 // may provide an id allocated by its protocol; an empty id lets persistence
150 // allocate one. Publication happens only after the initial event batch is
151 // accepted, so failure leaves the currently-bound session usable.
152 func (c *Controller) BindFreshSession(ctx context.Context, sessionID string) (session.SessionRef, error) {
153 return c.BindFreshSessionWithOptions(ctx, session.CreateOptions{SessionID: sessionID})
154 }
155
156 // BindFreshSessionWithOptions creates a fresh identity with immutable host
157 // ownership metadata before publishing the runtime.
158 func (c *Controller) BindFreshSessionWithOptions(ctx context.Context, options session.CreateOptions) (session.SessionRef, error) {
159 return c.bindFreshSessionWithCommit(ctx, options, nil)
160 }
161
162 func (c *Controller) bindFreshSessionWithCommit(ctx context.Context, options session.CreateOptions, commit func(context.Context, session.SessionRef) error) (session.SessionRef, error) {
163 service := c.SessionCreationService()
164 if c == nil || service == nil || c.executor == nil {
165 return session.SessionRef{}, errors.New("v3 session service is unavailable")
166 }
167 prepared, err := service.PrepareCreate(ctx, options)
168 if err != nil {
169 return session.SessionRef{}, err
170 }
171 candidate := prepared.Runtime()
172 fresh := agent.NewSession(c.basePrompt())
173 if err := seedRuntimeSession(ctx, candidate, "session-create", fresh.Snapshot(), c.ModelRef(), c.ModelSelectionIdentity()); err != nil {
174 _ = service.Discard(context.Background(), prepared)
175 return session.SessionRef{}, err
176 }
177 if _, err := candidate.Session().Flush(ctx); err != nil {
178 _ = service.Discard(context.Background(), prepared)
179 return session.SessionRef{}, err
180 }
181 owner, err := service.Publish(prepared)
182 if err != nil {
183 _ = service.Discard(context.Background(), prepared)
184 return session.SessionRef{}, err
185 }
186 if _, err = c.publishSessionRuntimeWithCommit(ctx, candidate, fresh, true, commit, service); err != nil {
187 // This attempt published the identity, so an owner-scoped close is the
188 // correct cleanup. It still refuses while any client is bound.
189 _ = owner.Close(context.Background())
190 return session.SessionRef{}, err
191 }
192 return candidate.Ref(), nil
193 }
194
195 // ContinueLegacySession freezes and migrates the selected legacy head, then
196 // publishes the returned immutable v3 identity. The source remains only as a
197 // display/import locator and is never rebound as the execution store.
198 func (c *Controller) ContinueLegacySession(ctx context.Context, sourcePath, headID string) (session.SessionRef, error) {
199 return c.continueLegacySession(ctx, sourcePath, headID, true, session.CreateOptions{})
200 }
201
202 // ContinueLegacySessionWithOptions installs immutable Desktop ownership in
203 // the same publication that materializes the imported session.
204 func (c *Controller) ContinueLegacySessionWithOptions(ctx context.Context, sourcePath, headID string, options session.CreateOptions) (session.SessionRef, error) {
205 return c.continueLegacySession(ctx, sourcePath, headID, true, options)
206 }
207
208 // ContinueLegacySessionForRebuildWithOptions performs the same fail-atomic
209 // import while an Agent generation is being replaced for the same logical
210 // session, and publishes host-owned immutable metadata in that same
211 // transaction. The SessionTemp generation belongs to the logical session, so
212 // this path must not rotate it merely because persistence crossed the
213 // legacy/v3 boundary.
214 func (c *Controller) ContinueLegacySessionForRebuildWithOptions(ctx context.Context, sourcePath, headID string, options session.CreateOptions) (session.SessionRef, error) {
215 return c.continueLegacySession(ctx, sourcePath, headID, false, options)
216 }
217
218 func (c *Controller) continueLegacySession(ctx context.Context, sourcePath, headID string, rotateSessionTemp bool, options session.CreateOptions) (session.SessionRef, error) {
219 service, _, _ := c.v3Binding()
220 if c == nil || service == nil || c.executor == nil {
221 return session.SessionRef{}, errors.New("v3 session service is unavailable")
222 }
223 restoreLegacyEvents, err := c.releaseLegacyEventStoreForImport(ctx)
224 if err != nil {
225 return session.SessionRef{}, fmt.Errorf("freeze legacy event source: %w", err)
226 }
227 published := false
228 defer func() {
229 if !published {
230 restoreLegacyEvents()
231 }
232 }()
233 candidate, _, err := service.ContinueImportedWithHeader(ctx, sourcePath, headID, options)
234 if err != nil {
235 return session.SessionRef{}, err
236 }
237 owner, err := service.Owner(candidate)
238 if err != nil {
239 return session.SessionRef{}, err
240 }
241 if err := seedRuntimeConfig(ctx, candidate, "legacy-import-config", c.ModelRef(), c.ModelSelectionIdentity()); err != nil {
242 _ = owner.Close(context.Background())
243 return session.SessionRef{}, err
244 }
245 messages := candidate.Session().ExecutionSnapshot().Projection.ModelMessages
246 prepared := agent.NewSession("").CloneWithMessages(messages)
247 if _, err = c.publishSessionRuntime(candidate, prepared, rotateSessionTemp); err != nil {
248 _ = owner.Close(context.Background())
249 return session.SessionRef{}, err
250 }
251 published = true
252 return candidate.Ref(), nil
253 }
254
255 // ContinuePrototypeSession imports the retired sidecar codec through the restricted
256 // fail-closed bridge, then publishes the final linear session identity.
257 func (c *Controller) ContinuePrototypeSession(ctx context.Context, sourceDir string) (session.SessionRef, error) {
258 service, _, _ := c.v3Binding()
259 if c == nil || service == nil || c.executor == nil {
260 return session.SessionRef{}, errors.New("v3 session service is unavailable")
261 }
262 candidate, _, err := service.ContinuePrototype(ctx, sourceDir)
263 if err != nil {
264 return session.SessionRef{}, err
265 }
266 owner, err := service.Owner(candidate)
267 if err != nil {
268 return session.SessionRef{}, err
269 }
270 if err := seedRuntimeConfig(ctx, candidate, "prototype-import-config", c.ModelRef(), c.ModelSelectionIdentity()); err != nil {
271 _ = owner.Close(context.Background())
272 return session.SessionRef{}, err
273 }
274 prepared := agent.NewSession("").CloneWithMessages(candidate.Session().ExecutionSnapshot().Projection.ModelMessages)
275 if _, err = c.publishSessionRuntime(candidate, prepared, true); err != nil {
276 _ = owner.Close(context.Background())
277 return session.SessionRef{}, err
278 }
279 return candidate.Ref(), nil
280 }
281
282 // OpenSession attaches this Controller to an existing immutable session identity.
283 // Opening never creates a missing session and publication retains the current
284 // binding until the target projection and writer are ready.
285 //
286 // Attaching grants only a ClientBinding, so a failed publication withdraws this
287 // client's own grant instead of disposing a runtime another client may already
288 // be using. A retired stored codec is the one exception: importing it publishes
289 // a brand-new identity that this attempt owns outright.
290 func (c *Controller) OpenSession(ctx context.Context, ref session.SessionRef) (session.SessionRef, error) {
291 service, current, _ := c.v3Binding()
292 if c == nil || service == nil || c.executor == nil {
293 return session.SessionRef{}, errors.New("v3 session service is unavailable")
294 }
295 if current != nil && current.Ref() == ref {
296 // After a reclaim the controller still renders a runtime whose store
297 // the service closed, so the next turn append would hit a closed
298 // recovery database. Only a still-active instance may skip the re-open.
299 if active, ok := service.Runtime(ref); ok && active == current {
300 return ref, nil
301 }
302 }
303 binding, err := service.Open(ctx, ref)
304 if err != nil {
305 return session.SessionRef{}, err
306 }
307 target := binding.Runtime()
308 published, err := c.publishAttachedSession(target, "attach-existing", nil)
309 if err == nil {
310 // publishSessionRuntime installs the controller's own client grant; this
311 // temporary attach grant is no longer needed.
312 _ = binding.Release(context.Background())
313 return published, nil
314 }
315 // Only this client's grant is withdrawn. A runtime another client still
316 // holds keeps its binding count above zero and is left untouched.
317 if releaseErr := binding.Release(context.Background()); releaseErr != nil {
318 slog.Warn("controller: release failed v3 attach binding", "err", releaseErr)
319 }
320 return session.SessionRef{}, err
321 }
322
323 // publishAttachedSession publishes the prepared projection for an already-resolved
324 // runtime. retire is used only when this attempt owns a newly published
325 // identity; pass nil to withdraw a client grant instead.
326 func (c *Controller) publishAttachedSession(candidate *session.Runtime, reason string, retire func(context.Context) error) (session.SessionRef, error) {
327 if candidate == nil {
328 return session.SessionRef{}, errors.New("v3 session runtime is unavailable")
329 }
330 prepared := agent.NewSession("").CloneWithMessages(candidate.Session().ExecutionSnapshot().Projection.ModelMessages)
331 if _, err := c.publishSessionRuntime(candidate, prepared, true); err != nil {
332 if retire != nil {
333 _ = retire(context.Background())
334 }
335 return session.SessionRef{}, err
336 }
337 return candidate.Ref(), nil
338 }
339
340 // SetSessionTitle records mutable title state in the canonical event stream.
341 func (c *Controller) SetSessionTitle(ctx context.Context, title string) error {
342 _, runtime, exclusive := c.v3Binding()
343 if !exclusive || runtime == nil {
344 return session.ErrSessionNotRunning
345 }
346 payload, err := json.Marshal(map[string]string{"title": title})
347 if err != nil {
348 return err
349 }
350 snapshot := runtime.Session().ExecutionSnapshot()
351 _, err = c.appendSessionBatch(ctx, runtime.Session(), session.Batch{
352 OperationID: "session-title:" + agent.NewMessageID(),
353 TurnID: snapshot.Projection.TurnID,
354 Events: []session.Event{{Kind: "session/title", Payload: payload}},
355 })
356 return err
357 }
358
359 func seedRuntimeSession(ctx context.Context, runtime *session.Runtime, operationID string, messages []provider.Message, modelRef, modelIdentity string) error {
360 if runtime == nil {
361 return nil
362 }
363 events := make([]session.Event, 0, len(messages)+1)
364 for _, message := range messages {
365 if message.ID == "" {
366 return errors.New("initial v3 message has no stable id")
367 }
368 payload, err := json.Marshal(map[string]any{"message": message})
369 if err != nil {
370 return err
371 }
372 events = append(events, session.Event{Kind: "message/complete", Payload: payload})
373 }
374 if strings.TrimSpace(modelRef) != "" {
375 config, err := sessionConfigEvent(modelRef, modelIdentity)
376 if err != nil {
377 return err
378 }
379 events = append(events, config)
380 }
381 if len(events) == 0 {
382 return nil
383 }
384 _, err := runtime.Session().AppendBatch(ctx, operationID, events)
385 return err
386 }
387
388 func seedRuntimeConfig(ctx context.Context, runtime *session.Runtime, operationID, modelRef, modelIdentity string) error {
389 if runtime == nil || strings.TrimSpace(modelRef) == "" {
390 return nil
391 }
392 event, err := sessionConfigEvent(modelRef, modelIdentity)
393 if err != nil {
394 return err
395 }
396 digest := sha256.Sum256(event.Payload)
397 _, err = runtime.Session().AppendBatch(ctx, fmt.Sprintf("%s:%x", operationID, digest[:16]), []session.Event{event})
398 return err
399 }
400
401 func sessionConfigEvent(modelRef, modelIdentity string) (session.Event, error) {
402 payload, err := json.Marshal(map[string]string{"modelRef": modelRef, "modelIdentity": modelIdentity})
403 if err != nil {
404 return session.Event{}, err
405 }
406 return session.Event{Kind: "session/config", Payload: payload}, nil
407 }
408
409 func (c *Controller) publishSessionRuntime(candidate *session.Runtime, prepared *agent.Session, rotateSessionTemp bool) (*session.Runtime, error) {
410 return c.publishSessionRuntimeWithCommit(context.Background(), candidate, prepared, rotateSessionTemp, nil)
411 }
412
413 func (c *Controller) publishSessionRuntimeWithCommit(ctx context.Context, candidate *session.Runtime, prepared *agent.Session, rotateSessionTemp bool, commit func(context.Context, session.SessionRef) error, creationService ...*session.Service) (*session.Runtime, error) {
414 if candidate == nil || prepared == nil {
415 return nil, errors.New("v3 runtime publication candidate is unavailable")
416 }
417 service, _, _ := c.v3Binding()
418 if len(creationService) != 0 {
419 service = creationService[0]
420 }
421 if service == nil {
422 return nil, errors.New("v3 session service is unavailable")
423 }
424 if current, ok := service.Runtime(candidate.Ref()); !ok || current != candidate {
425 return nil, errors.New("v3 runtime candidate is not the exact published service instance")
426 }
427 c.permissionMu.Lock()
428 defer c.permissionMu.Unlock()
429 snapshot := candidate.Session().ExecutionSnapshot()
430 projection := snapshot.Projection
431 if err := validateSessionDomainProjection(projection); err != nil {
432 return nil, err
433 }
434 preset, err := c.presetForSessionPublication(ctx, candidate, snapshot)
435 if err != nil {
436 return nil, err
437 }
438 binding, err := service.Bind(candidate)
439 if err != nil {
440 return nil, fmt.Errorf("bind v3 runtime: %w", err)
441 }
442 published := false
443 defer func() {
444 if !published {
445 _ = binding.Release(context.Background())
446 }
447 }()
448 // Commit desktop membership before publication; failures preserve the old binding.
449 if commit != nil {
450 if err := commit(ctx, candidate.Ref()); err != nil {
451 return nil, err
452 }
453 }
454 c.snapshotMu.Lock()
455 defer c.snapshotMu.Unlock()
456 // Domain parsing was validated above. Restore it before swapping the
457 // client binding so a future validation failure cannot expose a partially
458 // published controller or require closing a shared runtime.
459 if err := c.restoreSessionDomainProjection(projection); err != nil {
460 return nil, err
461 }
462 c.mu.Lock()
463 oldGen := c.turns.generation
464 c.mu.Unlock()
465 // Serve restores the selected remote session's choice. Other frontends keep
466 // their current permission posture when changing canonical identities.
467 c.promptResolveMu.Lock()
468 c.permissionStateMu.Lock()
469 c.v3BindingMu.Lock()
470 old := c.sessionRuntime
471 oldBinding := c.sessionBinding
472 c.sessionRuntime = candidate
473 c.sessionService = service
474 c.sessionBinding = binding
475 c.exclusiveSession = true
476 c.nativeLegacySession = false
477 c.v3BindingMu.Unlock()
478 if c.servePresetRestore {
479 c.approval.setMode(preset)
480 if c.subagentGate != nil {
481 c.subagentGate.Update(preset)
482 }
483 c.refreshInteractiveGate()
484 }
485 c.permissionRevision.Add(1)
486 c.permissionStateMu.Unlock()
487 c.promptResolveMu.Unlock()
488 c.bindAttachmentService()
489 c.bindExecutionControl()
490 if old != nil && old != candidate {
491 old.UnbindExecution(oldGen)
492 }
493 c.mu.Lock()
494 // Legacy paths are import inputs only. Retaining one as the live path lets
495 // unrelated compatibility helpers recreate sidecars beside a read-only
496 // source. The immutable SessionRef is the sole execution identity.
497 c.sessionPath = ""
498 c.mu.Unlock()
499 c.executor.SetSession(prepared)
500 // The immutable session projection is the only Goal restore source. A true
501 // session switch installs a fresh, disarmed lifecycle; OpenSession's exact-runtime
502 // fast path returns before this point and therefore preserves live activation.
503 c.installGoalLifecycle(candidate)
504 // Transcript pages are a derived cache. A session switch invalidates the
505 // prior identity immediately; the next query rebuilds from the exact v3
506 // projection without reading or writing a legacy sidecar.
507 c.turnEvents.mu.Lock()
508 if c.turnEvents.projection != nil {
509 c.turnEvents.projection.CloseFollowers()
510 }
511 c.turnEvents.projection = nil
512 c.turnEvents.projectionErr = nil
513 c.turnEvents.mu.Unlock()
514 c.rebindCheckpoints("")
515 c.ResetPlannerSession()
516 // The inbox belongs to the live runtime generation, not to the imported
517 // legacy path. Close the pre-bind queue before rotating the session temp so
518 // later Agent rebuilds attach to the same current generation.
519 c.pauseInboxOnRotate()
520 if rotateSessionTemp {
521 c.rotateSessionTemp()
522 }
523 c.rebindInbox()
524 c.refreshRuntimeState(event.Event{})
525 published = true
526 if oldBinding != nil && oldBinding != binding {
527 if err := oldBinding.Release(context.Background()); err != nil {
528 slog.Warn("controller: retire previous v3 binding after publication", "err", err)
529 }
530 }
531 return old, nil
532 }
533
534 func explicitSessionPermissionPreset(source *session.Session, projection session.Projection) (string, uint64) {
535 sequence := projection.PermissionPresetSequence
536 manifest := source.Manifest()
537 if sequence == 0 || (manifest.Source != nil && sequence <= manifest.InheritedEvents) {
538 // A fork does not inherit a broad permission grant. It may start at
539 // the narrower inherited read-only posture until the child makes an
540 // explicit choice of its own.
541 if manifest.Source != nil && projection.PermissionPreset == ToolApprovalReadOnly {
542 return ToolApprovalReadOnly, 0
543 }
544 return "", 0
545 }
546 return projection.PermissionPreset, sequence
547 }
548
549 func (c *Controller) presetForSessionPublication(ctx context.Context, candidate *session.Runtime, snapshot session.Snapshot) (string, error) {
550 if !c.servePresetRestore {
551 return "", nil
552 }
553 preset, sequence := explicitSessionPermissionPreset(candidate.Session(), snapshot.Projection)
554 if sequence > snapshot.DurableSequence {
555 // Failed writes leave accepted events in memory; never publish an unflushed preset.
556 if _, err := candidate.Session().FlushThrough(ctx, sequence); err != nil {
557 return "", fmt.Errorf("persist session permission preset: %w", err)
558 }
559 }
560 return string(permissionpreset.NormalizeDefault(preset)), nil
561 }
562
563 // EnableServeSessionPermissionPresets opts a Serve-owned controller into
564 // restoring per-session choices. A same-session rebuild keeps its already
565 // migrated live mode by passing restoreCurrent=false.
566 func (c *Controller) EnableServeSessionPermissionPresets(restoreCurrent bool) {
567 c.permissionMu.Lock()
568 defer c.permissionMu.Unlock()
569 c.servePresetRestore = true
570 if !restoreCurrent {
571 return
572 }
573 _, runtime, bound := c.v3Binding()
574 if !bound || runtime == nil {
575 return
576 }
577 snapshot := runtime.Session().StateSnapshot()
578 preset, sequence := explicitSessionPermissionPreset(runtime.Session(), snapshot.Projection)
579 if sequence > snapshot.DurableSequence {
580 preset = ToolApprovalReadOnly
581 }
582 c.applyToolApprovalModeLocked(string(permissionpreset.NormalizeDefault(preset)))
583 }
584
585 func validateSessionDomainProjection(projection session.Projection) error {
586 if len(projection.PlanState) > 0 {
587 var plan struct {
588 Enabled bool `json:"enabled"`
589 }
590 if err := json.Unmarshal(projection.PlanState, &plan); err != nil {
591 return fmt.Errorf("restore v3 plan state: %w", err)
592 }
593 }
594 if len(projection.GoalState) > 0 {
595 var goal goalState
596 if err := json.Unmarshal(projection.GoalState, &goal); err != nil {
597 return fmt.Errorf("restore v3 goal state: %w", err)
598 }
599 }
600 return nil
601 }
602
603 func (c *Controller) restoreSessionDomainProjection(projection session.Projection) error {
604 var plan struct {
605 Enabled bool `json:"enabled"`
606 }
607 if len(projection.PlanState) > 0 {
608 if err := json.Unmarshal(projection.PlanState, &plan); err != nil {
609 return fmt.Errorf("restore v3 plan state: %w", err)
610 }
611 }
612 c.mu.Lock()
613 c.sessionSettings.planMode = plan.Enabled
614 c.mu.Unlock()
615 if setter, ok := c.runner.(interface{ SetPlanMode(bool) }); ok {
616 setter.SetPlanMode(plan.Enabled)
617 } else if c.executor != nil {
618 c.executor.SetPlanMode(plan.Enabled)
619 }
620 if err := c.goals.restoreGoalEvent(projection.GoalState); err != nil {
621 return fmt.Errorf("restore v3 goal state: %w", err)
622 }
623 if c.executor != nil {
624 c.executor.RestoreDeliveryCheckpoint(c.goals.deliveryState())
625 }
626 return nil
627 }
628
629 func (c *Controller) v3Binding() (*session.Service, *session.Runtime, bool) {
630 if c == nil {
631 return nil, nil, false
632 }
633 c.v3BindingMu.RLock()
634 service, runtime, exclusive := c.sessionService, c.sessionRuntime, c.exclusiveSession
635 c.v3BindingMu.RUnlock()
636 return service, runtime, exclusive
637 }
638
639 // SessionBinding exposes the host-owned service/runtime pair for an Agent
640 // rebuild. Callers must attach the pair to the replacement Controller; they
641 // must not close or republish the writer themselves.
642 func (c *Controller) SessionBinding() (*session.Service, *session.Runtime, bool) {
643 service, runtime, exclusive := c.v3Binding()
644 return service, runtime, exclusive && service != nil && runtime != nil
645 }
646
647 // SessionService exposes the host query/management owner without requiring
648 // an active runtime. Cold history listing must not create an Agent or writer.
649 func (c *Controller) SessionService() *session.Service {
650 service, _, _ := c.v3Binding()
651 return service
652 }
653
654 // SessionCreationService keeps newly created sessions in the current store
655 // even while this controller is attached to a historical directory.
656 func (c *Controller) SessionCreationService() *session.Service {
657 if c == nil {
658 return nil
659 }
660 c.v3BindingMu.RLock()
661 defer c.v3BindingMu.RUnlock()
662 if c.sessionCreateService != nil {
663 return c.sessionCreateService
664 }
665 return c.sessionService
666 }
667
668 // UsesExclusiveSession reports the configured execution contract even when
669 // a lazy fresh session has not yet been allocated. Hosts use it to avoid
670 // manufacturing a legacy path during rebuild preparation.
671 func (c *Controller) UsesExclusiveSession() bool {
672 service, _, exclusive := c.v3Binding()
673 return exclusive && service != nil
674 }
675
676 func (c *Controller) sessionEngineEnabled() bool {
677 _, _, exclusive := c.v3Binding()
678 return exclusive
679 }
680
681 type SessionRotationRequest struct {
682 Source session.SessionRef
683 SourcePath string
684 Reason string
685 }
686
687 type SessionRotationPlan struct {
688 CreateOptions session.CreateOptions
689 Commit func(context.Context, session.SessionRef) error
690 }
691
692 // rotateExclusiveSession implements /new and /clear without allocating a
693 // legacy transcript path. clear additionally deletes the closed source v3
694 // directory; new leaves it available in history.
695 func (c *Controller) rotateExclusiveSession(clear bool) error {
696 service, runtime, _ := c.v3Binding()
697 if service == nil {
698 return errors.New("exclusive v3 session runtime is unavailable")
699 }
700 reason := "new"
701 if clear {
702 reason = "clear"
703 }
704 if runtime == nil {
705 // A handoff released the identity without a replacement: with no source
706 // to flush, end or plan from, allocation is the whole rotation — the
707 // step the next turn would otherwise take lazily.
708 ref, err := c.bindFreshSessionWithCommit(context.Background(), session.CreateOptions{}, nil)
709 if err != nil {
710 return err
711 }
712 c.startExclusiveSession(ref, reason)
713 return nil
714 }
715 oldRef := runtime.Ref()
716 if err := c.Snapshot(); err != nil {
717 return err
718 }
719 if err := c.extensionSessionPhase(context.Background(), extension.PointSessionRotate, dispatch.PhaseRotate, oldRef.SessionID); err != nil {
720 return err
721 }
722 c.hooks.SessionEnd(context.Background(), reason)
723 c.extensionSessionEvent(extension.PointSessionEnd, dispatch.PhaseEnd, oldRef.SessionID)
724 createOptions := session.CreateOptions{}
725 var commitRotation func(context.Context, session.SessionRef) error
726 if c.onSessionRotation != nil {
727 plan, planErr := c.onSessionRotation(context.Background(), SessionRotationRequest{Source: oldRef, Reason: reason})
728 if planErr != nil {
729 return planErr
730 }
731 createOptions, commitRotation = plan.CreateOptions, plan.Commit
732 }
733 ref, err := c.bindFreshSessionWithCommit(context.Background(), createOptions, commitRotation)
734 if err != nil {
735 return err
736 }
737 if commitRotation == nil && clear {
738 if err := service.Delete(context.Background(), oldRef); err != nil {
739 return fmt.Errorf("new session %s is active; delete cleared session: %w", ref.SessionID, err)
740 }
741 }
742 c.startExclusiveSession(ref, reason)
743 return nil
744 }
745
746 // startExclusiveSession runs the session-start side of a rotation once the
747 // fresh identity is published.
748 func (c *Controller) startExclusiveSession(ref session.SessionRef, reason string) {
749 c.ClearGoal()
750 c.mu.Lock()
751 c.startedOnce = true
752 c.mu.Unlock()
753 c.hooks.SetSessionID(ref.SessionID)
754 c.enqueueHookContexts(c.hooks.SessionStart(context.Background(), reason))
755 c.extensionSessionEvent(extension.PointSessionStart, dispatch.PhaseStart, ref.SessionID)
756 c.clearSessionWriteAccess()
757 }
758
758 lines GO