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