返回 DeepSeek-Reasonix
controller.go
根目录 / internal / control / controller.go
1 // Package control is the transport-agnostic session driver. A Controller owns
2 // the agent run loop and session lifecycle, takes commands (Send/Cancel/Approve/
3 // SetPlanMode/Compact/NewSession/…), and emits everything that happens —
4 // reasoning, tool calls, approvals, turn completion — as a typed event stream to
5 // a single event.Sink.
6 //
7 // The point is one orchestration layer behind every frontend: a terminal TUI, a
8 // desktop webview, or an HTTP/SSE server each drive the Controller identically
9 // (issue commands, render events) and none of them re-implement turn lifecycle,
10 // cancellation, or approval. The Controller depends on no frontend.
11 package control
12
13 import (
14 "context"
15 "encoding/json"
16 "errors"
17 "fmt"
18 "log/slog"
19 "net/http"
20 "os"
21 "path/filepath"
22 "reflect"
23 "slices"
24 "sort"
25 "strconv"
26 "strings"
27 "sync"
28 "sync/atomic"
29 "time"
30
31 "reasonix/internal/ablation"
32 "reasonix/internal/agent"
33 "reasonix/internal/agentpreset"
34 "reasonix/internal/autoresearch"
35 "reasonix/internal/billing"
36 "reasonix/internal/capability"
37 "reasonix/internal/checkpoint"
38 "reasonix/internal/command"
39 "reasonix/internal/config"
40 "reasonix/internal/event"
41 "reasonix/internal/evidence"
42 "reasonix/internal/extension"
43 "reasonix/internal/extension/dispatch"
44 "reasonix/internal/extension/uihub"
45 "reasonix/internal/gitcmd"
46 goaldomain "reasonix/internal/goal"
47 "reasonix/internal/guardian"
48 "reasonix/internal/hook"
49 "reasonix/internal/i18n"
50 "reasonix/internal/jobs"
51 "reasonix/internal/mcpinteraction"
52 "reasonix/internal/memory"
53 "reasonix/internal/nilutil"
54 "reasonix/internal/permission"
55 "reasonix/internal/permissionpreset"
56 "reasonix/internal/persistentshell"
57 "reasonix/internal/plugin"
58 "reasonix/internal/provider"
59 "reasonix/internal/sandbox"
60 "reasonix/internal/session"
61 "reasonix/internal/sessioncontext"
62 "reasonix/internal/sessioninbox"
63 "reasonix/internal/sessiontemp"
64 "reasonix/internal/shellrun"
65 "reasonix/internal/skill"
66 "reasonix/internal/store"
67 "reasonix/internal/taskmonitor"
68 "reasonix/internal/tool"
69 "reasonix/internal/workspacelease"
70 )
71
72 // ErrTurnRunning reports that a caller tried to start a second foreground turn
73 // while one is already active in the same Controller.
74 var ErrTurnRunning = errors.New("turn already running")
75
76 // ErrNoFinalReadinessRecovery is retained for old recovery-action clients.
77 var ErrNoFinalReadinessRecovery = errors.New("final_readiness_recovery_retired: readiness recovery actions can no longer restore evidence or replay checks")
78
79 // ErrRuntimeDraining reports that a caller targeted a controller generation
80 // superseded by a successful rebuild.
81 var ErrRuntimeDraining = errors.New("runtime is draining after rebuild")
82
83 // ErrRecoveryRequired reports that the previous foreground activity did not
84 // stop inside the cancellation grace period. The controller keeps queued input
85 // and write ownership sealed until the process is restarted.
86 var ErrRecoveryRequired = errors.New("session recovery is required before another turn can start")
87
88 // errTurnRunningRotation and errRotationInProgress are returned by the
89 // session-rotation gate (beginRotation) when a rotation cannot proceed: a turn
90 // is in flight, or another rotation already holds the gate.
91 var (
92 errTurnRunningRotation = errors.New("cannot start a new session while a turn is running")
93 errRotationInProgress = errors.New("cannot start a new session while another session change is in progress")
94 )
95
96 // errNoSessionPath is returned by snapshot when a session has content to persist
97 // but no resolved session path — a misconfiguration (e.g. an unresolvable data
98 // dir in a bot deployment) that previously dropped conversations silently
99 // (#4414). Callers log it and continue; it must never be swallowed quietly.
100 var errNoSessionPath = errors.New("session has content but no session path; conversation cannot be persisted")
101
102 // Controller drives one chat session. Construct with New; drive with the command
103 // methods; observe through the Sink passed in Options.
104 type Controller struct {
105 lifecycleDiagnostics lifecycleDiagnosticBuffer
106 providerDiagnostics providerDiagnosticBuffer
107 runtimeState controllerRuntimeState
108 controllerPromptRouting
109 authentication authenticationGate
110 runner agent.Runner
111 executor *agent.Agent
112 guardianSess *guardian.Session // nil when guardian is disabled
113 guardianPath string // persisted guardian session file ("" when disabled)
114 // taskBudget is the configured spend gate, as passed at construction.
115 taskBudget agent.TaskBudget
116 // goalTokenBudget bounds an unattended Goal loop; 0 leaves it unbounded.
117 goalTokenBudget int
118 goalResourceMu sync.Mutex
119 goalTokensUsed int
120 goalRequestsUsed int
121 goalTokenLimit int
122 goalBudgetExtensions int
123
124 // goalUsageTee accounts billable usage events into the active goal turn's
125 // observational token total. It wraps the public sink when the caller didn't provide one.
126 goalUsageTee *goalUsageTee
127 sink event.Sink
128 policy permission.Policy
129 // subagentGate is the shared gate every headless-only sub-agent surface
130 // reads from (see Options.SubagentGate). Nil when the caller didn't build
131 // one — sub-agents then keep whatever gate they were constructed with.
132 subagentGate *SharedHeadlessGate
133
134 label string
135 selection modelSelection
136 resolveSessionModel func(string, string) (string, error)
137 visionModel string
138 visionProviderResolver func(string) (provider.Provider, error)
139 visionModelSelector func(string, string) (string, bool)
140 modelCapabilityResolver func(*config.ProviderEntry) config.ResolvedModelCapability
141 frozenImageInput *bool
142 imageCapabilityChanged func() bool
143 controllerAttachmentState
144 modelSettings controllerModelSettings
145 prompt controllerPromptState
146 pinnedContextLoader PinnedContextLoader
147 sessionContextStatic sessioncontext.Sections
148 sessionDir string
149 controllerSessionBinding
150 // managedSessionEvents is set for hosts that publish controllers only after
151 // a session-lease handoff. An unpublished replacement may read the shared
152 // v3 projection, but it must not mutate that projection before its final
153 // write authority is bound.
154 managedSessionEvents atomic.Bool
155 commands atomic.Pointer[[]command.Command]
156 // skills owns the session's discovered skills (enabled subset, full set, and
157 // the reloadable stores) — the skills slice of the Capabilities concern. See
158 // skill.go.
159 skills skillSet
160 skillRunner skill.SubagentRunner
161 readOnlySkillRunner skill.SubagentRunner
162 skillProfile skill.ProfileResolver
163 disableImplicitSkillInvocation bool
164 slashSkillSeq atomic.Uint64
165 hooks *hook.Runner // session hook runner; nil-safe (no hooks configured)
166 // hookContexts carries one-shot lifecycle hook context into the next real
167 // user turn without changing the cache-stable system prompt.
168 hookContexts []string
169 // memory owns the loaded memory snapshot, the pending turn-tail notes queue,
170 // and write serialization behind its own locks, off c.mu — so a memory-panel
171 // save never stalls an approval or status poll. See memory.go.
172 memory memoryManager
173 cleanup func()
174 responseLanguage string
175 reasoningLanguage string
176 disableColdResumePrune bool // legacy; rewrite elision removed, still gates cold notice
177 headPolicy sessionHeadPolicy
178 // testCacheColdAfter overrides cacheColdAfter() in tests. Zero uses the
179 // vendor-aware resolution from config.
180 testCacheColdAfter time.Duration
181 // testCancelGrace overrides the production cancellation grace in tests.
182 // Zero uses the documented 15 second boundary.
183 testCancelGrace time.Duration
184
185 shell sandbox.Shell // interpreter for user-invoked "!" commands; zero = auto
186 startedOnce bool // guards the one-shot SessionStart hook on first turn
187 closeOnce sync.Once // makes close idempotent under racing teardown paths
188 closeFinalizeOnce sync.Once // releases persistence/resources only after the terminal boundary
189 closeFinalized chan struct{} // closes after every controller-owned resource has been released
190 closeFireSessionEnd bool
191 closeJobsMode closeJobsMode
192 onRemember func(rule string) RememberResult // set via Options; invoked when user picks "always allow"
193 onRememberPlanModeReadOnlyCommand func(prefix string) PlanModeReadOnlyCommandTrustResult
194 writeAccess controllerWriteAccess
195 sessionRecoveryMeta func(SessionRecoveryRequest) agent.BranchMeta
196 onSessionRecovered func(SessionRecoveryInfo) error
197 onSessionTransition func(SessionTransitionInfo) error
198 onSessionRotation func(context.Context, SessionRotationRequest) (SessionRotationPlan, error)
199
200 // balanceURL/balanceKey target the active provider's optional wallet-balance
201 // endpoint (empty when the provider declares none). Captured at build so a
202 // model/key switch — which rebuilds the controller — refreshes them.
203 balanceURL string
204 balanceKey string
205 balanceClient *http.Client
206
207 // jobs is the session-scoped background-job manager. The agent's background
208 // tools spawn into it; Compose drains its completion notes into the next turn;
209 // Close cancels its still-running jobs.
210 jobs *jobs.Manager
211 background controllerBackground
212 // workspaceLease is the Delivery writer owner shared with the executor.
213 // It is exposed only through a sanitized state snapshot for Desktop recovery.
214 workspaceLease *workspacelease.Owner
215
216 // mcp owns the session's live tool/plugin surface behind its own lock, off
217 // c.mu; the Controller keeps config-facing orchestration. See mcp.go.
218 mcp mcpManager
219 mcpDefaultCallTimeout time.Duration
220 mcpConfigureSpec func(*plugin.Spec)
221 capabilityRuntime *agent.MCPCapabilityRuntime
222
223 runtimeGeneration uint64 // PublishGate gen; 0 disables
224 permissionMu sync.Mutex
225 permissionStateMu sync.RWMutex
226 permissionRevision atomic.Uint64
227 runtimeOwner *extension.RuntimeOwner
228 lastResumeDecision extension.ResumeDecision
229 // extensions is the frozen extension dispatcher for this controller
230 // generation, or nil when no v2 runtime packages are installed (the
231 // universal pre-dispatch fast path). It is installed before the controller
232 // starts serving (Options.Extensions or SetExtensions) and never swapped
233 // afterwards, so wiring points read it without locking.
234 extensions *dispatch.Dispatcher
235 // extensionUI is the host extension UI hub for this controller generation
236 // (stage 8a), or nil when no v2 runtime packages started. Installed via
237 // SetExtensionUI before serving and never swapped; readers take c.mu.
238 extensionUI *uihub.Hub
239 // providerResolver is the build's merged provider catalog (extension
240 // sidecar providers over the config/broker base), or nil when no sidecar
241 // declared providers. Immutable after New; ProviderCatalog reads it.
242 providerResolver provider.Resolver
243
244 // Capability routing (Delivery hybrid route + dual-model Planner proxy).
245 // Not part of the provider-visible prefix; only seeds the turn-scoped ledger
246 // and optional semantic router.
247 pluginCfg []config.PluginEntry
248 capCachedTools map[string][]plugin.CachedTool
249 capCacheKeyOK map[string]bool
250 semanticRouter *capability.SemanticRouter
251 capabilityAudit *capability.Audit
252 // capabilityProxy directs unready MCP candidates to use_capability in the
253 // transient route block (Delivery and dual-model Planner).
254 capabilityProxy bool
255 // proxyToolsFn returns live tools observed through use_capability without
256 // entering the provider-visible registry (dual-model Planner).
257 proxyToolsFn func() map[string][]plugin.CachedTool
258 ablation ablation.Set
259
260 // goals owns the active goal's FSM (status, intercepts, idle/turn counters)
261 // and its persistence, behind its own mutex so a per-turn goal save never
262 // stalls an approval or status poll on c.mu. See goal.go.
263 goals goalMachine
264 // goalLifecycle is the versioned session goal authority. The legacy
265 // goalMachine remains only while old sidecars are imported and must not be
266 // used as the execution source once the v3 lifecycle cutover is complete.
267 goalLifecycleMu sync.RWMutex
268 goalLifecycleMutationMu sync.Mutex
269 goalLifecycle *goaldomain.Machine
270 goalLifecycleLoadErr error
271 // goalDriver is a level-triggered, process-local scheduler. It never owns a
272 // cross-turn Activity: each accepted continuation enters through the normal
273 // guarded top-level turn path.
274 goalDriverMu sync.Mutex
275 goalDriverWG sync.WaitGroup
276 goalDriverPending bool
277 goalDriverActive *goalRoundReservation
278 goalDriverControl goalDriverControl
279 // legacyResearchArchive reads explicit pre-unification task paths. It never
280 // creates or mutates archive state. See
281 // autoresearch_manager.go.
282 legacyResearchArchive legacyResearchArchive
283 legacyRestoreMu sync.Mutex
284 legacyRestore legacyGoalRestore
285
286 // workspaceRoot is the workspace root: the base for resolving @-refs and slash
287 // path refs, the working directory for user "!" shell commands and custom
288 // command discovery, and the guard root for checkpoint restore writes. It is
289 // surfaced to frontends via WorkspaceRoot().
290 workspaceRoot string
291 // workspaceRepo is workspaceRoot's git identity, resolved before any turn ran.
292 workspaceRepo gitcmd.Repo
293
294 // externalFolderRefs maps session-generated @ tokens to user-dropped
295 // directories outside workspaceRoot. It is intentionally per-controller:
296 // dragging a folder authorizes that folder for this chat session only, without
297 // widening scoped @ resolution to arbitrary absolute paths.
298 externalFolderRefsMu sync.RWMutex
299 externalFolderRefs map[string]string
300 externalFolderToolRefs externalFolderToolRefs
301
302 // checkpoints owns the snapshot-based rewind bookkeeping (the per-session
303 // store, the monotonic turn counter, and the conversation-rewind boundary map)
304 // behind its own lock, off c.mu — so a boundary read for a rewind/fork never
305 // contends on the run-state lock. The Controller keeps the rewind/fork/summarize
306 // orchestration (truncating the session, restoring code, emitting events). See
307 // checkpoint.go.
308 checkpoints checkpointManager
309 // mutationObserver is the host-side file mutation observer for v2 checkpoints.
310 mutationObserver *checkpoint.MutationObserver
311 // sessionRevision increments on successful rewind/undo and is used as a
312 // prepare/commit freshness token.
313 sessionRevision int64
314
315 // approval owns the approval/ask prompt bookkeeping and the runtime approval
316 // posture (ask/auto/yolo, session grants, the just-approved-plan window)
317 // behind its own locks, off c.mu. The Controller keeps the I/O orchestration
318 // (requestApproval/Ask emit events + fire hooks + rebuild the executor gate).
319 // See approval.go.
320 approval approvalManager
321
322 // mu guards the run state; every critical section under it is short and
323 // non-blocking.
324 mu sync.Mutex
325 // turns is the sole execution authority: phase, current cancel/done/token,
326 // and the FIFO pending queue.
327 turns turnLoop
328 // executionGeneration is the Runtime BindExecution generation for this
329 // controller. Unbind uses the exact value so a rebuilt controller cannot
330 // clear the replacement's control.
331 executionGeneration atomic.Uint64
332 // closed marks the controller as terminally torn down (close() ran). It
333 // seals turn admission: without it, a submit arriving AFTER close cleared
334 // the parked queue — but while a still-running turn's TurnDone delivery
335 // was in flight — would park again and then start against freed resources
336 // when the window closed.
337 closed bool
338 // rotating is set under mu while NewSession/ClearSession swap the executor
339 // session out. Checking running once and then swapping later leaves a
340 // TOCTOU window: a turn can start (running=false at check time) during the
341 // intervening Snapshot() and then have its live session replaced. running
342 // and rotating are mutually exclusive gates — a turn refuses to start while
343 // a rotation is in progress, and a rotation refuses to start while a turn
344 // runs — so the run loop's session reference cannot change under it.
345 rotating bool
346 // maintenance is the controller-owned foreground maintenance operation.
347 // It is deliberately separate from rotating: Stop can cancel it and normal
348 // user input can remain in the durable inbox until it reaches a terminal
349 // boundary.
350 maintenance *controllerMaintenance
351 autosaveWG sync.WaitGroup
352 // sessionSettings groups the per-session posture knobs that share one
353 // lifetime: swapped together on session rotation.
354 sessionSettings sessionSettings
355 sessionPath string
356 // sessionTemp owns the logical-session private temporary directory shared
357 // by Bash calls. Retained for this Controller's lifetime; rotated on
358 // /new, /clear, resume of another session, and branch switches.
359 sessionTemp *sessiontemp.Manager
360 persistentShell *persistentshell.Manager
361 // snapshotMu serializes the whole save/recovery handoff for this controller.
362 // Agent-level path locks protect individual files, but recovery also moves
363 // controller-owned state (sessionPath, guardianPath, checkpoints, rewrite
364 // baseline). Letting a second snapshot observe that migration halfway through
365 // can turn one conflict into a recovery cascade. Session/path swaps
366 // (new/clear/fork/branch/switch/resume/SetSessionPath) hold it for the same
367 // reason: a save that reads the old path but the new session would write one
368 // transcript's messages into another's file, or manufacture a bogus conflict.
369 // Not reentrant — never call snapshot (or anything that snapshots, such as
370 // recoverInterruptedTurn or maybeColdResumePrune) while holding it.
371 snapshotMu sync.Mutex
372 // turn counts model turns this session, passed to hooks in their payload.
373 turn int
374 turnEvents turnEventState
375 submissions submissionIdentityState
376 liveness turnLiveness
377
378 displayRecorder func(content, display string)
379
380 // inbox is the durable session-level instruction queue. Disk I/O never
381 // runs under c.mu; the store owns its own lock.
382 inbox inboxState
383 }
384
385 type approvalReply struct {
386 allow bool
387 session bool
388 persist bool // true = write "always allow" rule to config
389 onceDirs []string
390 persistErr error
391 }
392
393 type pendingApproval struct {
394 id string
395 tool string
396 subject string
397 reason string
398 rawInput json.RawMessage
399 fresh bool
400 requireHuman bool
401 autoDrain bool
402 kind string // tool | plan | recovery | write_access; empty = tool
403 recovery *event.RecoveryApproval
404 writeAccess *event.WriteAccessApproval
405 reply chan approvalReply
406 }
407
408 // pendingAsk is an in-flight ask question batch. questions is retained so the
409 // AskRequest can be re-emitted to a frontend that reconnected after the original
410 // event (see ReplayPendingPrompts).
411 type pendingAsk struct {
412 questions []event.AskQuestion
413 reply chan []event.AskAnswer
414 queued bool // registered but not yet shown; replay must skip it
415 }
416
417 type plannerSessionResetter interface {
418 ResetPlannerSession()
419 }
420
421 type controllerSessionBinding struct {
422 // sessionRuntime is the final identity-bound v3 owner. When exclusiveSession is
423 // set, SessionPath is a legacy import/display locator only and no production
424 // transcript or business sidecar may be written through it.
425 sessionService *session.Service
426 sessionCreateService *session.Service
427 sessionRuntime *session.Runtime
428 sessionBinding *session.ClientBinding
429 exclusiveSession bool
430 nativeLegacySession bool
431 // Only Serve opts into session-scoped preset restoration. TUI and Desktop
432 // keep their own live permission policy across canonical session binds.
433 servePresetRestore bool
434 v3BindingMu sync.RWMutex
435 }
436
437 type controllerPromptRouting struct {
438 // promptEpochMu protects only the routing epoch. One-shot resolution is
439 // owned by PendingPromptOwner and the typed registries; cancellation must
440 // never wait behind an answer callback.
441 promptEpochMu sync.RWMutex
442 promptRuntimeEpoch string
443 promptOwner PendingPromptOwner
444 // promptResolveMu serializes permission-generation changes with legacy
445 // resolver entry points while they transition onto PendingPromptOwner.
446 // It is never held while waiting for a user answer.
447 promptResolveMu sync.Mutex
448 }
449
450 // RuntimeStatus is the frontend-facing snapshot of foreground turn state. It is
451 // intentionally more explicit than the legacy Running bool so UI code can
452 // distinguish a cancellable foreground turn from pending prompts and background
453 // jobs.
454 type RuntimeStatus struct {
455 Running bool
456 PendingPrompt bool
457 BackgroundJobs int
458 CancelRequested bool
459 Cancellable bool
460 TurnID string
461 Status event.TurnStatus
462 TurnEventSeq uint64
463 ReplayAfterSeq uint64
464 }
465
466 const (
467 ToolApprovalReadOnly = string(permissionpreset.ReadOnly)
468 ToolApprovalWorkspaceWrite = string(permissionpreset.WorkspaceWrite)
469 ToolApprovalDangerFullAccess = string(permissionpreset.DangerFullAccess)
470 ToolApprovalDontAsk = "dontAsk"
471
472 // Deprecated source aliases. Persisted legacy strings are migrated by
473 // normalizeToolApprovalMode; these names keep older integrations building.
474 ToolApprovalAsk = ToolApprovalReadOnly
475 ToolApprovalAuto = ToolApprovalWorkspaceWrite
476 ToolApprovalYolo = ToolApprovalDangerFullAccess
477 )
478
479 const (
480 memoryRememberTool = "remember"
481 memoryForgetTool = "forget"
482 )
483
484 // RememberResult describes what happened when an approval rule was persisted.
485 type RememberResult struct {
486 Rule string
487 Path string
488 Saved bool
489 CoveredBy string
490 Err error
491 }
492
493 // PlanModeReadOnlyCommandTrustResult describes what happened when a trusted bash
494 // command prefix was persisted for plan-mode research.
495 type PlanModeReadOnlyCommandTrustResult struct {
496 Prefix string
497 Path string
498 Saved bool
499 CoveredBy string
500 Err error
501 }
502
503 type SessionRecoveryRequest struct {
504 OriginalPath string
505 Reason string
506 Mode string
507 }
508
509 type SessionRecoveryInfo struct {
510 OriginalPath string
511 RecoveryPath string
512 Existing bool
513 Reason string
514 BaseRevision int64
515 DiskRevision int64
516 Meta agent.BranchMeta
517 commit *sessionRecoveryCommit
518 }
519
520 func snapshotConflictRevisions(err error) (base, disk int64) {
521 var conflict *agent.SessionSnapshotConflictError
522 if errors.As(err, &conflict) && conflict != nil {
523 return conflict.BaseRevision, conflict.DiskRevision
524 }
525 return 0, 0
526 }
527
528 // OnCommit defers publication work until the controller has installed the
529 // recovery path and rebound all path-scoped state.
530 func (i SessionRecoveryInfo) OnCommit(fn func()) {
531 if i.commit != nil && fn != nil {
532 i.commit.hooks = append(i.commit.hooks, fn)
533 }
534 }
535
536 type sessionRecoveryCommit struct {
537 hooks []func()
538 }
539
540 func (c *sessionRecoveryCommit) publish() {
541 for _, hook := range c.hooks {
542 hook()
543 }
544 c.hooks = nil
545 }
546
547 type externalFolderToolRefs interface {
548 RegisterReadRoot(token, root string)
549 }
550
551 // Options carries the already-built pieces setup assembles. Lifecycle metadata
552 // lets the controller mint and rotate session files; Host/Commands are surfaced
553 // to frontends that resolve MCP prompts and slash commands.
554 type Options struct {
555 ImageRouteConfig *config.Config
556 Runner agent.Runner
557 Executor *agent.Agent
558 // Authentication is the frozen runtime credential snapshot's initial
559 // admission state. An empty value remains Ready for source compatibility.
560 Authentication AuthenticationState
561 AuthenticationForModel func(string) AuthenticationState
562 Guardian *guardian.Session
563 // RecoveryHeadless is decoded for source compatibility and ignored. Auto
564 // Guard cannot be re-enabled through Controller options.
565 RecoveryHeadless bool
566 // TaskBudget is the configured spend gate; unset leaves a turn unbounded.
567 TaskBudget agent.TaskBudget
568 // GoalTokenBudget bounds an unattended Goal loop by cumulative tokens.
569 GoalTokenBudget int
570 // GoalEvaluator is accepted only for source compatibility.
571 // when the working model submits no update_goal report. nil fails closed:
572 // the goal pauses instead of defaulting to continue.
573 GoalEvaluator any // deprecated: accepted but never invoked
574 Sink event.Sink
575 Policy permission.Policy
576 // SubagentGate is the shared, mutable gate every headless-only sub-agent
577 // surface (task, writer-capable skill sub-agents, planner) reads from. Nil
578 // disables gating for those surfaces same as before this field existed.
579 // SetToolApprovalMode and ApplyHeadlessApprovalMode call Update on it so a
580 // runtime approval-mode switch reaches sub-agents, not just the parent
581 // executor's own gate.
582 SubagentGate *SharedHeadlessGate
583 Label string
584 ModelRef string
585 ModelIdentity string
586 ResolveSessionModel func(string, string) (string, error)
587 // VisionModel is empty (off), "auto", or a canonical provider/model ref.
588 // The resolver and selector are assembled by boot so the controller remains
589 // transport-agnostic and tests can inject deterministic fake providers.
590 VisionModel string
591 VisionProviderResolver func(string) (provider.Provider, error)
592 VisionModelSelector func(string, string) (string, bool)
593 // ModelCapabilityResolver returns the adapter/config-resolved metadata for
594 // the exact active model. Nil keeps the legacy config-only behavior.
595 ModelCapabilityResolver func(*config.ProviderEntry) config.ResolvedModelCapability
596 // FrozenImageInput belongs to the provider instance built for this runtime.
597 FrozenImageInput *bool
598 ImageCapabilityChanged func() bool
599 ModelSettingsRevision string
600 ModelSettingsSourceRevision string
601 ModelSettingsCurrent func() (string, error)
602 ModelSettingsContinuation func() error
603 ModelConnectionTarget string
604 // BeforeInboxDispatch lets the owner reserve runtime admission before a
605 // queued message becomes a new turn. The returned release runs after claim
606 // and synchronous turn admission, outside every controller lock.
607 BeforeInboxDispatch func(*Controller) (func(), error)
608 SystemPrompt string
609 // PinnedContextLoader snapshots the current session sidecar at turn
610 // admission. The Agent persists changes as append-only user-role revisions.
611 PinnedContextLoader PinnedContextLoader
612 SessionDir string
613 SessionPath string
614 // SessionRuntime binds this Controller/Agent view to an already-published
615 // immutable v3 session identity. SessionService owns exact-instance close,
616 // fork, query and cancellation. ExclusiveSession disables legacy
617 // transcript and business-sidecar writes.
618 SessionService *session.Service
619 SessionCreateService *session.Service
620 SessionRuntime *session.Runtime
621 ExclusiveSession bool
622 NativeLegacySession bool
623 Host *plugin.Host
624 // MCPHostProfile is the surface lazily created hosts declare; injected
625 // hosts keep their own profile.
626 MCPHostProfile plugin.HostProfile
627 Commands []command.Command
628 Skills []skill.Skill
629 AllSkills []skill.Skill
630 SkillStore *skill.Store
631 AllSkillStore *skill.Store
632 // DisableImplicitSkillInvocation controls model-facing discovery only;
633 // explicit /skill commands and management remain host-side capabilities.
634 DisableImplicitSkillInvocation bool
635 // SkillRunner executes a runAs=subagent skill in an isolated child loop.
636 // ReadOnlySkillRunner is reserved for explicitly read-only entry points;
637 // Plan itself is a workflow instruction and uses SkillRunner with the shared
638 // Permissions/Sandbox gate. SkillProfile supplies model/effort display
639 // metadata for the synthetic top-level run_skill event.
640 SkillRunner skill.SubagentRunner
641 ReadOnlySkillRunner skill.SubagentRunner
642 SkillProfile skill.ProfileResolver
643 Hooks *hook.Runner
644 Memory *memory.Set
645 Cleanup func()
646 // BalanceURL/BalanceKey wire the active provider's optional wallet-balance
647 // endpoint and bearer key; empty when the provider declares no balance_url.
648 BalanceURL string
649 BalanceKey string
650 BalanceClient *http.Client
651 // Jobs is the session-scoped background-job manager (nil disables background jobs).
652 Jobs *jobs.Manager
653 BackgroundScope *jobs.SessionBackgroundScope
654 BackgroundSink event.Sink
655 // TaskStore remains a FileStore-compatible authority. Desktop injects one
656 // observed instance so recorder and task-control APIs share post-commit
657 // projection hints; nil preserves the ordinary FileStore.
658 TaskStore taskmonitor.WriteStore
659 // WorkspaceLease is the Delivery writer owner shared with the executor.
660 WorkspaceLease *workspacelease.Owner
661 // Registry is the executor's live tool set, and PluginCtx the session-scoped
662 // context; both are needed for hot-adding MCP servers via AddMCPServer.
663 Registry *tool.Registry
664 PluginCtx context.Context
665 // MCPDefaultCallTimeout is the global MCP call cap used by hot-connected
666 // servers when they do not declare a server- or tool-specific override.
667 MCPDefaultCallTimeout time.Duration
668 // MCPConfigureSpec injects host-local launch and isolation policy into every
669 // hot-connected server without persisting that state in project config.
670 MCPConfigureSpec func(*plugin.Spec)
671 // CapabilityRuntime is the controller-local authoritative MCP inventory used
672 // by stable use_capability frontends. It shares Host processes with sibling
673 // tabs but never shares their enabled/disabled state.
674 CapabilityRuntime *agent.MCPCapabilityRuntime
675 RuntimeGeneration uint64 // PublishGate generation for admission
676 // RuntimeOwner isolates publish/drain gates and receipts to one
677 // controller/session rebuild lineage. Nil preserves compatibility behavior.
678 RuntimeOwner *extension.RuntimeOwner
679 // WorkspaceRoot is the project root checkpoint restores are confined to ("" =
680 // no confinement). Frontends pass the cwd they launched the session in.
681 WorkspaceRoot string
682 WorkspaceRepo gitcmd.Repo // WorkspaceRoot's identity, resolved when the session opened
683 ExternalFolderToolRefs externalFolderToolRefs
684 // ResponseLanguage controls final-answer language preference. Empty/auto
685 // means no transient injection because the stable language policy follows the
686 // current user turn.
687 ResponseLanguage string
688 // ReasoningLanguage controls visible reasoning language preference. Empty/auto
689 // means no transient injection because the stable language policy already
690 // follows the conversation language.
691 ReasoningLanguage string
692 // SessionContextStatic carries boot-observed runtime facts that belong in a
693 // host user-turn snapshot rather than the cache-stable system prompt. Only
694 // Environment and Workspace are consumed; memory and skills stay live.
695 SessionContextStatic sessioncontext.Sections
696 // FileBranchesOnly keeps fork, branch, switch, and conversation rewind on
697 // separate session files even for schema-2 logs. Hosts that expose forks as
698 // independent conversations (desktop tabs and multi-session Serve) set it.
699 FileBranchesOnly bool
700 // DisableColdResumePrune suppresses the cold-resume cache-state notice.
701 // Resume never rewrites history regardless of this flag.
702 DisableColdResumePrune bool
703 // Shell is the interpreter user-invoked "!" commands run under, so /shell
704 // matches the agent's configured [tools.shell] choice. Zero value = auto.
705 Shell sandbox.Shell
706 // OnRemember, when set, is invoked with a new allow rule the user chose to
707 // persist to disk (e.g. "Bash(go test:*)"). The callback is wired into the
708 // permission Gate on EnableInteractiveApproval.
709 OnRemember func(rule string) RememberResult
710 // OnRememberPlanModeReadOnlyCommand persists a bash command prefix as trusted
711 // read-only when the user chooses "always allow" from the plan-mode trust
712 // prompt.
713 OnRememberPlanModeReadOnlyCommand func(prefix string) PlanModeReadOnlyCommandTrustResult
714 // OnPersistWriteAccess writes sandbox.allow_write and an optional permission
715 // rule to the workspace reasonix.toml as one transaction.
716 OnPersistWriteAccess PersistWriteAccessFunc
717 // WriteRoots is the session-scoped writable directory manager shared with
718 // built-in file tools and bash.
719 WriteRoots *sandbox.WritableRootSet
720 // BashSandboxEnforced is true when this session's bash tool actually wraps
721 // commands in an OS sandbox. Windows and bash=off leave this false so
722 // directory prompts are not implied for unisolated shell writes.
723 BashSandboxEnforced bool
724 // SessionRecoveryMeta lets a frontend attach scope/topic/profile metadata to
725 // an automatic recovery branch before it is written.
726 SessionRecoveryMeta func(SessionRecoveryRequest) agent.BranchMeta
727 // OnSessionRecovered is called after a stale runtime's transcript has been
728 // saved as a recovery branch, before the controller commits to that branch.
729 OnSessionRecovered func(SessionRecoveryInfo) error
730 // OnSessionTransition transfers write ownership before an intentional
731 // fork, branch, or switch publishes a different Session.
732 OnSessionTransition func(SessionTransitionInfo) error
733 // OnSessionRotation lets an identity-owning host durably reserve a fresh
734 // SessionID before /new or /clear publishes it. When installed, the host
735 // also owns clear archival; the controller never permanently deletes the
736 // source session.
737 OnSessionRotation func(context.Context, SessionRotationRequest) (SessionRotationPlan, error)
738 // ApprovalTimeout bounds how long a tool-approval or ask prompt blocks waiting
739 // for a user decision. Zero (default) waits forever — right for an interactive
740 // terminal. Bot/headless frontends set a positive value so an unanswered
741 // prompt can't wedge the session indefinitely (#4626, #4402).
742 ApprovalTimeout time.Duration
743 // Extensions is the frozen extension dispatcher for this controller
744 // generation (Extension Protocol v2, stage 6b1). Nil means no v2 runtime
745 // packages are installed: every extension wiring point takes an untouched
746 // fast path. Boot installs it through SetExtensions because sidecars (and
747 // therefore the dispatcher) only exist after snapshot assembly, which runs
748 // after New.
749 Extensions *dispatch.Dispatcher
750 // ProviderResolver is the build's merged provider catalog — extension
751 // sidecar providers folded over the config/broker base (stage 7). Nil when
752 // no v2 runtime sidecar declared providers; ProviderCatalog then returns
753 // nil and frontends enumerate providers from config alone, as before.
754 ProviderResolver provider.Resolver
755 // Ablation switches subsystems off for a benchmark arm. The zero value runs
756 // everything.
757 Ablation ablation.Set
758 // SessionTemp is the logical-session private temporary directory manager
759 // shared by sandboxed Bash calls. Nil creates a fresh Manager owned by this
760 // Controller. Hot rebuilds pass the previous Controller's Manager so the
761 // temporary directory survives model/settings swaps.
762 SessionTemp *sessiontemp.Manager
763 // PersistentShell is the session-scoped PTY used by ordinary foreground
764 // bash. Nil creates a fresh Manager owned by this Controller. Hot rebuilds
765 // pass the previous Controller's Manager so cwd and exported environment
766 // survive model/settings swaps.
767 PersistentShell *persistentshell.Manager
768 }
769
770 // New builds a Controller. A nil Sink becomes event.Discard; unless the caller
771 // already provided a goalUsageTee (NewGoalUsageTee), the sink is wrapped in one
772 // so billable usage can be accounted to Goal budgets.
773 func controllerSessionTemp(existing *sessiontemp.Manager) *sessiontemp.Manager {
774 if existing != nil {
775 return existing
776 }
777 return sessiontemp.New()
778 }
779
780 func controllerPersistentShell(existing *persistentshell.Manager) *persistentshell.Manager {
781 if existing != nil {
782 return existing
783 }
784 return persistentshell.New()
785 }
786
787 func New(opts Options) *Controller {
788 sink := opts.Sink
789 if nilutil.IsNil(sink) {
790 sink = event.Discard
791 }
792 usageTee, ok := sink.(*goalUsageTee)
793 if !ok {
794 usageTee = NewGoalUsageTee(sink).(*goalUsageTee)
795 sink = usageTee
796 }
797 pluginCtx := opts.PluginCtx
798 if pluginCtx == nil {
799 pluginCtx = context.Background()
800 }
801 runtimeOwner := runtimeOwnerOrDefault(opts.RuntimeOwner)
802 pluginCtx = extension.ContextWithRuntimeOwner(pluginCtx, runtimeOwner)
803 goalDriverCtx, goalDriverCancel := context.WithCancel(context.Background())
804 if opts.Hooks != nil {
805 opts.Hooks.SetSessionID(agent.BranchID(opts.SessionPath))
806 }
807 sessionRuntime, sessionBinding := bindInitialSessionRuntime(opts)
808 c := &Controller{
809 authentication: newAuthenticationGate(opts.Authentication, opts.ModelRef),
810 taskBudget: opts.TaskBudget,
811 goalTokenBudget: opts.GoalTokenBudget,
812 goalTokenLimit: opts.GoalTokenBudget,
813 goals: goalMachine{tokenBudget: opts.GoalTokenBudget},
814 runner: opts.Runner,
815 executor: opts.Executor,
816 guardianSess: opts.Guardian,
817 guardianPath: guardian.PathFor(opts.SessionPath),
818 goalUsageTee: usageTee,
819 sink: sink,
820 policy: opts.Policy,
821 subagentGate: opts.SubagentGate,
822 label: opts.Label,
823 selection: modelSelection{ref: opts.ModelRef, identity: opts.ModelIdentity},
824 resolveSessionModel: opts.ResolveSessionModel,
825 visionModel: strings.TrimSpace(opts.VisionModel),
826 visionProviderResolver: opts.VisionProviderResolver,
827 visionModelSelector: opts.VisionModelSelector,
828 modelCapabilityResolver: opts.ModelCapabilityResolver,
829 frozenImageInput: opts.FrozenImageInput,
830 imageCapabilityChanged: opts.ImageCapabilityChanged,
831 modelSettings: newControllerModelSettings(opts),
832 prompt: newControllerPromptState(opts.SystemPrompt, opts.Executor),
833 pinnedContextLoader: opts.PinnedContextLoader,
834 sessionContextStatic: opts.SessionContextStatic,
835 sessionDir: opts.SessionDir,
836 sessionPath: opts.SessionPath,
837 controllerSessionBinding: controllerSessionBinding{sessionService: opts.SessionService, sessionCreateService: opts.SessionCreateService, sessionRuntime: sessionRuntime, sessionBinding: sessionBinding, exclusiveSession: opts.ExclusiveSession, nativeLegacySession: opts.NativeLegacySession},
838 commands: atomic.Pointer[[]command.Command]{},
839 skills: newSkillSet(opts.Skills, opts.AllSkills, opts.SkillStore, opts.AllSkillStore),
840 disableImplicitSkillInvocation: opts.DisableImplicitSkillInvocation,
841 skillRunner: opts.SkillRunner,
842 readOnlySkillRunner: opts.ReadOnlySkillRunner,
843 skillProfile: opts.SkillProfile,
844 hooks: opts.Hooks,
845 memory: newMemoryManager(opts.Memory),
846 cleanup: opts.Cleanup,
847 responseLanguage: config.NormalizeLanguage(opts.ResponseLanguage),
848 reasoningLanguage: config.NormalizeReasoningLanguage(opts.ReasoningLanguage),
849 disableColdResumePrune: opts.DisableColdResumePrune,
850 headPolicy: sessionHeadPolicy{fileBranchesOnly: opts.FileBranchesOnly},
851 shell: opts.Shell,
852 onRemember: opts.OnRemember,
853 onRememberPlanModeReadOnlyCommand: opts.OnRememberPlanModeReadOnlyCommand,
854 writeAccess: newControllerWriteAccess(opts),
855 sessionRecoveryMeta: opts.SessionRecoveryMeta,
856 onSessionRecovered: opts.OnSessionRecovered,
857 onSessionTransition: opts.OnSessionTransition,
858 onSessionRotation: opts.OnSessionRotation,
859 balanceURL: opts.BalanceURL,
860 balanceKey: opts.BalanceKey,
861 balanceClient: opts.BalanceClient,
862 jobs: opts.Jobs,
863 background: controllerBackground{scope: opts.BackgroundScope, sink: opts.BackgroundSink,
864 candidate: opts.BackgroundScope != nil && opts.BackgroundScope.Manager.ReplacementInProgress()},
865 workspaceLease: opts.WorkspaceLease,
866 mcp: newMcpManager(opts.Host, opts.Registry, pluginCtx, opts.MCPHostProfile),
867 mcpDefaultCallTimeout: opts.MCPDefaultCallTimeout,
868 mcpConfigureSpec: opts.MCPConfigureSpec,
869 capabilityRuntime: opts.CapabilityRuntime,
870 ablation: opts.Ablation,
871 workspaceRoot: opts.WorkspaceRoot,
872 workspaceRepo: opts.WorkspaceRepo,
873 externalFolderToolRefs: opts.ExternalFolderToolRefs,
874 providerResolver: opts.ProviderResolver,
875 runtimeGeneration: opts.RuntimeGeneration,
876 runtimeOwner: runtimeOwner,
877 goalDriverControl: goalDriverControl{ctx: goalDriverCtx, cancel: goalDriverCancel},
878 approval: newApprovalManager(opts.Policy, ToolApprovalAsk, opts.ApprovalTimeout),
879 turns: turnLoop{phase: session.RuntimeIdle},
880 closeFinalized: make(chan struct{}),
881 }
882 c.authentication.initialForModel = opts.AuthenticationForModel
883 c.initializeOwnedResources(opts)
884 c.bindAttachmentService()
885 if opts.ImageRouteConfig != nil {
886 c.captureImageRoutes(opts.ImageRouteConfig)
887 }
888 return c
889 }
890
891 func (c *Controller) initializeOwnedResources(opts Options) {
892 if c.executor != nil {
893 c.executor.SetImageRequestResolver(c)
894 }
895 c.goalUsageTee.setLifecycleUsageRecorder(c.recordGoalLifecycleUsage)
896 c.installGoalLifecycle(opts.SessionRuntime)
897 c.managedSessionEvents.Store(opts.OnSessionTransition != nil)
898 c.permissionRevision.Store(1)
899 // Session-private temporary directory: reuse a shared Manager on hot
900 // rebuild, otherwise create one. Retain so ReleaseResources/Close drop the
901 // owner reference without racing a replacement Controller.
902 c.sessionTemp = controllerSessionTemp(opts.SessionTemp)
903 c.sessionTemp.Retain()
904 c.persistentShell = controllerPersistentShell(opts.PersistentShell)
905 c.persistentShell.Retain()
906 if strings.TrimSpace(opts.WorkspaceRoot) != "" {
907 c.legacyResearchArchive = legacyResearchArchive{store: autoresearch.NewStore(opts.WorkspaceRoot)}
908 }
909 if opts.Extensions != nil {
910 c.extensions = opts.Extensions
911 c.sink = newFrontendEventSink(c.sink, opts.Extensions)
912 if c.executor != nil {
913 c.executor.SetExtensions(opts.Extensions)
914 }
915 }
916 // Checkpoints: bind a store to the session and route writer pre-edits into it.
917 c.rebindCheckpoints(opts.SessionPath)
918 c.setActiveJobSession(opts.SessionPath)
919 c.rebindInbox()
920 // Observe Steer / unapplied-steer for durable inbox state transitions.
921 // Must wrap both the controller sink and the executor sink: agent.Steer
922 // emits on the executor path, TurnDone on the controller path.
923 c.sink = &inboxEventSink{inner: newTurnEventSink(c.sink, c), c: c}
924 if runner, ok := c.runner.(interface{ SetSink(event.Sink) }); ok {
925 runner.SetSink(c.sink)
926 }
927 // Establish mutation authority before any constructor-time session seed.
928 // A hot-rebuild candidate sharing an already-bound Runtime remains at
929 // generation zero and can restore from the projection without writing it.
930 c.bindExecutionControl()
931 if c.executor != nil {
932 c.executor.SetSink(c.sink)
933 c.executor.SetSessionCheckpointer(c)
934 if _, runtime, _ := c.v3Binding(); runtime != nil {
935 if runtime.StateSnapshot().Session.EventSequence == 0 {
936 if err := c.seedSessionEventsFromExecutor("session-open"); err != nil {
937 c.failTurnEventLedger(err)
938 }
939 }
940 c.restoreExecutorFromSessionEvents()
941 } else if err := c.seedSessionEventsFromExecutor("session-open"); err != nil {
942 c.failTurnEventLedger(err)
943 }
944 }
945 cmdsInit := opts.Commands
946 c.commands.Store(&cmdsInit)
947 if c.executor != nil {
948 c.wireMutationObserver()
949 c.executor.SetMemoryQueue(c)
950 }
951 // Task monitoring: record background-job lifecycle into the project-local
952 // task store so CLI, Desktop, scripts, and future clients observe the same
953 // state/event evidence. The recorder swallows its own failures — monitoring
954 // must never affect the agent pipeline. The session id is resolved lazily
955 // because the session path is only fixed once the first turn begins.
956 c.initializeTaskRecorder(opts.TaskStore)
957 c.initializeRuntimeState()
958 }
959
960 func (c *Controller) initializeTaskRecorder(store taskmonitor.WriteStore) {
961 if c.jobs == nil || c.workspaceRoot == "" {
962 return
963 }
964 if store == nil {
965 store = taskmonitor.NewFileStore(filepath.Join(".reasonix", "tasks"))
966 }
967 sessionID := func() string { return c.parentSessionID() }
968 if c.background.scope != nil {
969 sessionID = c.jobs.ActiveSessionID
970 }
971 recorder := taskmonitor.NewTaskRecorder(store, c.workspaceRoot, sessionID)
972 c.background.recorder = recorder
973 if c.background.scope == nil {
974 c.jobs.SetTaskRecorder(recorder)
975 }
976 }
977
978 // SetDisplayRecorder installs an optional hook used by frontends that persist a
979 // shorter user-facing transcript than the fully composed model prompt.
980 func (c *Controller) SetDisplayRecorder(fn func(content, display string)) {
981 c.mu.Lock()
982 defer c.mu.Unlock()
983 c.displayRecorder = fn
984 }
985
986 // SetExtensions installs the extension dispatcher after construction. Boot
987 // uses it because sidecars — and therefore the dispatcher — only exist after
988 // snapshot assembly, which runs after New. First non-nil install wins for the
989 // cold-start path; use ReplaceExtensions for generation-safe rebuild swaps.
990 // Nil is a no-op. The executor agent receives the same dispatcher (stage 6b2).
991 func (c *Controller) SetExtensions(d *dispatch.Dispatcher) {
992 if d == nil {
993 return
994 }
995 c.mu.Lock()
996 defer c.mu.Unlock()
997 if c.extensions != nil {
998 return
999 }
1000 c.installExtensionsLocked(d)
1001 }
1002
1003 // ReplaceExtensions atomically swaps the dispatcher for a reused controller
1004 // after a narrow rebuild. Updates sink strategy owner and executor together.
1005 func (c *Controller) ReplaceExtensions(d *dispatch.Dispatcher) {
1006 if c == nil || d == nil {
1007 return
1008 }
1009 c.mu.Lock()
1010 defer c.mu.Unlock()
1011 c.installExtensionsLocked(d)
1012 }
1013
1014 func (c *Controller) installExtensionsLocked(d *dispatch.Dispatcher) {
1015 c.extensions = d
1016 // Keep the inbox observer as the outermost sink so Steer/unapplied events
1017 // always update durable state, while still installing/updating the
1018 // frontendEventSink wrapper underneath for extension rulings.
1019 switch sink := c.sink.(type) {
1020 case *inboxEventSink:
1021 if lifecycle, ok := sink.inner.(*turnEventSink); ok {
1022 current := lifecycle.innerSnapshot()
1023 if existing, ok := current.(*frontendEventSink); ok {
1024 existing.setDispatcher(d)
1025 } else {
1026 lifecycle.setInner(newFrontendEventSink(current, d))
1027 }
1028 } else if existing, ok := sink.inner.(*frontendEventSink); ok {
1029 existing.setDispatcher(d)
1030 } else {
1031 sink.inner = newTurnEventSink(newFrontendEventSink(sink.inner, d), c)
1032 }
1033 case *frontendEventSink:
1034 sink.setDispatcher(d)
1035 // Ensure inbox observer stays outer.
1036 c.sink = &inboxEventSink{inner: newTurnEventSink(sink, c), c: c}
1037 default:
1038 c.sink = &inboxEventSink{inner: newTurnEventSink(newFrontendEventSink(c.sink, d), c), c: c}
1039 }
1040 if c.executor != nil {
1041 c.executor.SetExtensions(d)
1042 c.executor.SetSink(c.sink)
1043 }
1044 if runner, ok := c.runner.(interface{ SetSink(event.Sink) }); ok {
1045 runner.SetSink(c.sink)
1046 }
1047 }
1048
1049 // SetProviderResolver replaces the session's merged provider catalog (narrow
1050 // rebuild after sidecar Manager roll). Nil clears extension-hosted providers.
1051 func (c *Controller) SetProviderResolver(r provider.Resolver) {
1052 if c == nil {
1053 return
1054 }
1055 c.mu.Lock()
1056 c.providerResolver = r
1057 c.mu.Unlock()
1058 }
1059
1060 // SetOnSessionRecovered installs the ownership handoff invoked before the
1061 // controller commits to an automatically created recovery branch. Frontends
1062 // that acquire their session owner after controller construction (for example
1063 // reasonix serve) use this before publishing the controller.
1064 func (c *Controller) SetOnSessionRecovered(fn func(SessionRecoveryInfo) error) {
1065 if c == nil {
1066 return
1067 }
1068 c.mu.Lock()
1069 defer c.mu.Unlock()
1070 c.onSessionRecovered = fn
1071 }
1072
1073 func (c *Controller) sessionRecoveredHandler() func(SessionRecoveryInfo) error {
1074 c.mu.Lock()
1075 defer c.mu.Unlock()
1076 return c.onSessionRecovered
1077 }
1078
1079 func (c *Controller) recordDisplay(content, display string) {
1080 if strings.TrimSpace(display) == "" || content == display {
1081 return
1082 }
1083 c.mu.Lock()
1084 record := c.displayRecorder
1085 c.mu.Unlock()
1086 if record != nil {
1087 record(content, display)
1088 }
1089 }
1090
1091 // ToolContractEntries returns a stable snapshot of the executor's live tool
1092 // contract: provider-visible names, descriptions, canonical schemas, and
1093 // read-only flags. It is intended for diagnostics and regression tests.
1094 func (c *Controller) ToolContractEntries() []tool.ContractEntry {
1095 if c == nil {
1096 return nil
1097 }
1098 reg := c.mcp.registry()
1099 if reg == nil {
1100 return nil
1101 }
1102 return reg.ContractEntries()
1103 }
1104
1105 // AllToolContractEntries returns every registered tool, including those hidden
1106 // from the provider-visible schema and only reachable via use_capability.
1107 func (c *Controller) AllToolContractEntries() []tool.ContractEntry {
1108 if c == nil {
1109 return nil
1110 }
1111 reg := c.mcp.registry()
1112 if reg == nil {
1113 return nil
1114 }
1115 return reg.AllContractEntries()
1116 }
1117
1118 // ProviderCatalog returns the session's merged provider catalog: the config
1119 // (or broker) base plus every provider a live extension sidecar declared,
1120 // keyed by ref — extension refs carry their plugin/<plugin>/<provider>/<model>
1121 // namespace. Nil when no sidecar declared providers, so frontends can tell
1122 // "enumerate config only" apart from "the extension catalog is empty".
1123 func (c *Controller) ProviderCatalog() []provider.Descriptor {
1124 if c == nil {
1125 return nil
1126 }
1127 c.mu.Lock()
1128 r := c.providerResolver
1129 c.mu.Unlock()
1130 if r == nil {
1131 return nil
1132 }
1133 return r.Catalog()
1134 }
1135
1136 func (c *Controller) recordDisplayForNewUser(startMessages int, display string) {
1137 if strings.TrimSpace(display) == "" {
1138 return
1139 }
1140 msgs := c.History()
1141 if startMessages > len(msgs) {
1142 startMessages = len(msgs)
1143 }
1144 for _, m := range msgs[startMessages:] {
1145 if agent.IsUserAuthoredTurnMessage(m) {
1146 c.recordDisplay(m.Content, display)
1147 return
1148 }
1149 }
1150 }
1151
1152 func (c *Controller) markEditedForNewUser(startMessages int, original string) {
1153 if strings.TrimSpace(original) == "" || c.executor == nil {
1154 return
1155 }
1156 s := c.executor.Session()
1157 msgs := s.Snapshot()
1158 if startMessages > len(msgs) {
1159 startMessages = len(msgs)
1160 }
1161 for i := startMessages; i < len(msgs); i++ {
1162 if !agent.IsUserAuthoredTurnMessage(msgs[i]) {
1163 continue
1164 }
1165 if agent.UserMessageText(msgs[i]) == original {
1166 return
1167 }
1168 msgs[i].Edited = true
1169 msgs[i].Original = original
1170 // A periodic autosave may already contain this user message without its
1171 // local edit metadata. Classify the mutation atomically so the turn-end
1172 // save performs an owned rewrite instead of forking a bogus
1173 // same-revision recovery branch. Edited/Original are local-only display
1174 // metadata (provider requests ignore them), so this must not report a
1175 // cache-prefix change — ReplaceLocalMetadata, not Rewrite.
1176 s.ReplaceLocalMetadata(msgs)
1177 return
1178 }
1179 }
1180
1181 // ckptDir derives a session's checkpoint directory from its file path
1182 // (…/<id>.jsonl → …/<id>.ckpt). Empty path → empty (in-memory checkpoints).
1183 func ckptDir(sessionPath string) string {
1184 return store.SessionCheckpointDir(sessionPath)
1185 }
1186
1187 // rebindCheckpoints points the store at the (possibly new) session, loading any
1188 // checkpoints already on disk, and resets the turn boundaries. Called on
1189 // construction and whenever the session path changes (NewSession/Resume/SetSessionPath).
1190 // Also re-wires the mutation observer so capture targets the new store.
1191 func (c *Controller) rebindCheckpoints(sessionPath string) {
1192 if c.sessionEngineEnabled() {
1193 // Goal and runtime business state are v3 events. Legacy goal/checkpoint
1194 // sidecars must not become a second restore source in exclusive mode.
1195 c.goals.setStatePath("")
1196 c.checkpoints.rebind("", c.workspaceRoot, c.checkpointOptions()...)
1197 c.rebindTurnEvents(sessionPath)
1198 if c.executor != nil {
1199 c.wireMutationObserver()
1200 }
1201 return
1202 }
1203 c.goals.setStatePath(goalStatePath(sessionPath))
1204 c.checkpoints.rebind(ckptDir(sessionPath), c.workspaceRoot, c.checkpointOptions()...)
1205 c.rebindTurnEvents(sessionPath)
1206 if c.executor != nil {
1207 c.wireMutationObserver()
1208 }
1209 }
1210
1211 // commands (frontend → controller)
1212
1213 func (c *Controller) Send(input string) {
1214 c.SendWithRaw(input, input)
1215 }
1216
1217 // SendWithRaw starts a turn with separate model input and raw prompt text.
1218 func (c *Controller) SendWithRaw(input, raw string) {
1219 _, _ = c.submitIdentifiedWithSetup(SubmissionRequest{Input: input, Display: raw}, nil, func(admission turnAdmission) {
1220 c.runGuardedWithAdmission(func(ctx context.Context) error { return c.runGoalLoopWithRaw(ctx, input, raw) }, admission)
1221 })
1222 }
1223
1224 // planApprovalTool is the Tool name on the ApprovalRequest the controller emits
1225 // to gate a proposed plan. Frontends key their plan-approval UI on it (the
1226 // desktop renders a plan card; the chat TUI a plan banner).
1227 const planApprovalTool = "exit_plan_mode"
1228
1229 // PlanDecisionAction preserves the three user-owned meanings of the Plan card.
1230 // Revise and exit both deny execution at the approval gate, but they are not the
1231 // same product decision and must remain distinguishable in durable receipts.
1232 type PlanDecisionAction string
1233
1234 const (
1235 PlanDecisionStartExecution PlanDecisionAction = "start_execution"
1236 PlanDecisionRevisePlan PlanDecisionAction = "revise_plan"
1237 PlanDecisionExitPlan PlanDecisionAction = "exit_plan"
1238 )
1239
1240 // SandboxEscapeApprovalTool is the internal Tool name used for one-shot approval
1241 // to rerun a shell command without the OS sandbox after the sandbox failed.
1242 const SandboxEscapeApprovalTool = "sandbox_escape"
1243
1244 // ManagedConfigWriteApprovalTool is the internal Tool name used for per-write
1245 // approval when a file tool targets a Reasonix-managed config file outside the
1246 // workspace write roots. It is a fresh human decision: config files control
1247 // providers, sandbox rules, permissions, and MCP servers for future sessions,
1248 // so YOLO/auto approval must never answer it.
1249 const ManagedConfigWriteApprovalTool = "config_write"
1250
1251 // planApprovedMessage is the follow-up turn sent once the user approves a plan.
1252 // Approval grants the plan scope; the model owns any fresh todo list it chooses
1253 // to write during the execution turn.
1254 const planApprovedMessage = "Plan approved — plan mode is off. Implement the approved plan and user feedback. Explicit scope, permission and sandbox restrictions still apply. Update todos to reflect actual progress. Use the plan’s checks and acceptance notes as task instructions, and report what you changed and verified."
1255
1256 // runTurn runs one model turn, then applies the plan-approval gate. This is the
1257 // single, frontend-agnostic plan flow: in Plan the model is instructed to
1258 // research and write its plan as a normal answer, while any tool calls still use
1259 // the active Permissions/Sandbox path.
1260 // When the turn ends with a text proposal, the controller asks the user to
1261 // approve (reusing the ApprovalRequest channel both frontends already render);
1262 // on approval it exits plan mode and continues straight into execution; on
1263 // rejection it stays in plan mode so the
1264 // next turn can revise. Plan mode is only ever set interactively, so the headless
1265 // `Run` path (which doesn't call this) never blocks on a prompt.
1266 func (c *Controller) runTurn(ctx context.Context, input string) error {
1267 return c.runGoalLoopWithRaw(ctx, input, input)
1268 }
1269
1270 func (c *Controller) runTurnWithRaw(ctx context.Context, input, raw string) error {
1271 return c.runTurnWithRawDisplay(ctx, input, raw, "")
1272 }
1273
1274 func (c *Controller) runGoalLoopWithRaw(ctx context.Context, input, raw string) error {
1275 return c.runGoalLoopWithRawDisplay(ctx, input, raw, "")
1276 }
1277
1278 // withTurnFormat binds a structured-output format to the turn context
1279 // (empty is a no-op). Extracted from the runGoalLoop closure so tests can
1280 // assert the format actually reaches the agent request path.
1281 func (c *Controller) withTurnFormat(ctx context.Context, format string) context.Context {
1282 if format == "" {
1283 return ctx
1284 }
1285 return agent.WithResponseFormat(ctx, format)
1286 }
1287
1288 func (c *Controller) runGoalLoopWithRawDisplay(ctx context.Context, input, raw, display string) error {
1289 // Structured-output format is bound to the submitted turn (passed via
1290 // submitHTTPWithFormat → submitCommandOrTurn → runGoalLoop closure);
1291 // no global one-shot slot to race across concurrent requests.
1292 return newTurnOrchestrator(c).runGoalLoopWithRawDisplay(ctx, input, raw, display)
1293 }
1294
1295 func (c *Controller) runEditedGoalLoopWithRawDisplay(ctx context.Context, input, raw, display, original string) error {
1296 return newTurnOrchestrator(c).runEditedGoalLoopWithRawDisplay(ctx, input, raw, display, original)
1297 }
1298
1299 func (c *Controller) runTurnWithRawDisplay(ctx context.Context, input, raw, display string) error {
1300 return newTurnOrchestrator(c).runTurnWithRawDisplay(ctx, input, raw, display)
1301 }
1302
1303 func (c *Controller) runSubagentSkillSlash(sk skill.Skill, task, raw, display string, admission turnAdmission) {
1304 sk = c.skills.prepare(sk)
1305 c.runGuardedWithAdmission(func(ctx context.Context) error {
1306 planMode := c.PlanMode()
1307 runner := c.skillRunner
1308 if runner == nil {
1309 return fmt.Errorf("subagent skill runner is unavailable for /%s", sk.Name)
1310 }
1311 return newTurnOrchestrator(c).runSubagentSkillGoalLoop(ctx, sk, task, raw, display, runner, planMode)
1312 }, admission)
1313 }
1314
1315 func (c *Controller) stopGoal(status string) {
1316 path, data, ok := c.goals.stop(status)
1317 c.persistGoalState(path, data, ok)
1318 }
1319
1320 // lastAssistantText returns the content of the most recent assistant message with
1321 // non-empty text — the model's final answer for the turn (its plan, in plan mode).
1322 func lastAssistantText(msgs []provider.Message) string {
1323 for _, msg := range slices.Backward(msgs) {
1324 if msg.Role == provider.RoleAssistant && strings.TrimSpace(msg.Content) != "" {
1325 return msg.Content
1326 }
1327 }
1328 return ""
1329 }
1330
1331 // Submit is the one-call entry for a simple frontend: it takes raw user input
1332 // and does everything — slash-command dispatch, @-reference expansion, plan-mode
1333 // composition — emitting all output as events. The HTTP/SSE server uses this so
1334 // a browser client only POSTs the typed line.
1335 //
1336 // Slash commands route to the matching primitive: /compact, /new, and /clear
1337 // run their session op and emit a Notice; /mcp__server__prompt and custom /commands
1338 // resolve to a turn; an unknown slash emits a Notice. Anything else is a normal
1339 // turn with its @-references resolved first.
1340 func (c *Controller) Submit(input string) {
1341 c.submit(input, "", "")
1342 }
1343
1344 // SubmitHTTP accepts input from the unauthenticated localhost HTTP frontend. It
1345 // deliberately omits the trusted TUI-only "!cmd" shell shortcut and resolves file
1346 // references only through the controller's workspace root.
1347 func (c *Controller) SubmitHTTP(input string) {
1348 c.submitHTTP(input, "")
1349 }
1350
1351 // SubmitDisplay runs input as a turn while remembering the user-facing display
1352 // text for transcript replay when controller-side composition expands input.
1353 func (c *Controller) SubmitDisplay(display, input string) {
1354 c.submit(input, display, "")
1355 }
1356
1357 // SubmitInvocationDisplay executes composer-selected invocation entities
1358 // independently of slash-command parsing. Plain string submit entry points keep
1359 // their existing behavior for CLI, HTTP, and backward-compatible clients.
1360 func (c *Controller) SubmitInvocationDisplay(display, input string, invocations []InvocationRequest) {
1361 c.submitInvocations(input, display, invocations)
1362 }
1363
1364 func (c *Controller) submitInvocations(input, display string, requests []InvocationRequest) {
1365 _, _ = c.submitIdentifiedWithSetup(SubmissionRequest{Input: input, Display: display, Invocations: requests}, nil, func(admission turnAdmission) {
1366 c.submitInvocationsLocked(input, display, requests, admission)
1367 })
1368 }
1369
1370 func (c *Controller) submitInvocationsLocked(input, display string, requests []InvocationRequest, admission turnAdmission) {
1371 if len(requests) == 0 {
1372 c.submitLocked(input, display, "", admission)
1373 return
1374 }
1375 prepared, err := c.prepareInvocationTurn(input, requests)
1376 if err != nil {
1377 c.notice(err.Error())
1378 return
1379 }
1380 c.runGuardedWithAdmission(func(ctx context.Context) error {
1381 return c.runPreparedInvocationTurn(ctx, prepared, input, input, display, nil)
1382 }, admission)
1383 }
1384
1385 type preparedInvocationTurn struct {
1386 composed string
1387 subagents []skill.Skill
1388 inlineSkillNames []string
1389 }
1390
1391 func (c *Controller) prepareInvocationTurn(input string, requests []InvocationRequest) (preparedInvocationTurn, error) {
1392 ordered := append([]InvocationRequest(nil), requests...)
1393 sort.SliceStable(ordered, func(i, j int) bool { return ordered[i].Offset < ordered[j].Offset })
1394 inline := make([]skill.Skill, 0, len(ordered))
1395 subagents := make([]skill.Skill, 0, len(ordered))
1396 for _, request := range ordered {
1397 sk, _, ok := c.resolveSkillInvocation("/" + strings.TrimSpace(request.Name))
1398 if !ok {
1399 return preparedInvocationTurn{}, fmt.Errorf("unknown invocation: /%s", strings.TrimSpace(request.Name))
1400 }
1401 kind := "skill"
1402 if sk.RunAs == skill.RunSubagent {
1403 kind = "subagent"
1404 }
1405 if strings.TrimSpace(request.Kind) != "" && request.Kind != kind {
1406 return preparedInvocationTurn{}, fmt.Errorf("invocation /%s is %s, not %s", sk.SlashName(), kind, request.Kind)
1407 }
1408 if sk.RunAs == skill.RunSubagent {
1409 subagents = append(subagents, sk)
1410 } else {
1411 inline = append(inline, sk)
1412 }
1413 }
1414
1415 if strings.TrimSpace(input) == "" && len(subagents) > 0 {
1416 return preparedInvocationTurn{}, fmt.Errorf("subagent invocation requires a task")
1417 }
1418 inlineSkillNames := make([]string, 0, len(inline))
1419 for _, sk := range inline {
1420 inlineSkillNames = append(inlineSkillNames, sk.Name)
1421 }
1422 // A lone inline skill takes the typed text as its arguments, as "/name task" does.
1423 if len(inline) == 1 && len(subagents) == 0 {
1424 return preparedInvocationTurn{composed: c.skills.renderInvocation(inline[0], strings.TrimSpace(input)), inlineSkillNames: inlineSkillNames}, nil
1425 }
1426 parts := make([]string, 0, len(inline)+1)
1427 for _, sk := range inline {
1428 parts = append(parts, c.skills.renderInvocation(sk, ""))
1429 }
1430 if strings.TrimSpace(input) != "" {
1431 parts = append(parts, input)
1432 }
1433 return preparedInvocationTurn{composed: strings.Join(parts, "\n\n"), subagents: subagents, inlineSkillNames: inlineSkillNames}, nil
1434 }
1435
1436 func (c *Controller) runPreparedInvocationTurn(
1437 ctx context.Context,
1438 prepared preparedInvocationTurn,
1439 input, raw, display string,
1440 frozenImages []string,
1441 ) error {
1442 ctx = withInvokedSkills(ctx, prepared.inlineSkillNames)
1443 if len(prepared.subagents) == 0 {
1444 return c.runGoalLoopWithFrozenImagesRawDisplay(ctx, prepared.composed, raw, display, frozenImages)
1445 }
1446 runner := c.skillRunner
1447 if runner == nil {
1448 return fmt.Errorf("subagent skill runner is unavailable")
1449 }
1450 return newTurnOrchestrator(c).runSubagentSkillTurnsGoalLoop(
1451 ctx,
1452 prepared.subagents,
1453 prepared.composed,
1454 input,
1455 display,
1456 runner,
1457 c.PlanMode(),
1458 frozenImages,
1459 )
1460 }
1461
1462 // SubmitEditedDisplay is SubmitDisplay for an inline-edited prompt. The model
1463 // sees input; the saved user message also keeps the pre-edit prompt as local UI
1464 // metadata so the edit survives session rewrites.
1465 func (c *Controller) SubmitEditedDisplay(display, input, original string) {
1466 c.submit(input, display, original)
1467 }
1468
1469 // SubmitUserTurn starts a normal model turn without interpreting shell or slash
1470 // commands. It still resolves references, so callers can submit trusted
1471 // user-authored prompt text without expanding the command surface.
1472 func (c *Controller) SubmitUserTurn(input, display string) {
1473 _, _ = c.submitIdentifiedWithSetup(SubmissionRequest{Input: input, Display: display}, nil, func(admission turnAdmission) {
1474 c.runRefTurnWithAdmission(input, display, admission)
1475 })
1476 }
1477
1478 func (c *Controller) submit(input, display, editedOriginal string) {
1479 if isSessionManagementSubmission(input) {
1480 c.submissions.mu.Lock()
1481 defer c.releaseSubmissionAdmission()
1482 c.submitLocked(input, display, editedOriginal, turnAdmission{})
1483 return
1484 }
1485 _, _ = c.submitIdentifiedWithSetup(SubmissionRequest{Input: input, Display: display, Original: editedOriginal}, nil, func(admission turnAdmission) {
1486 c.submitLocked(input, display, editedOriginal, admission)
1487 })
1488 }
1489
1490 func (c *Controller) submitLocked(input, display, editedOriginal string, admission turnAdmission) {
1491 trimmed := strings.TrimSpace(input)
1492 if note, ok := MemoryQuickAddNote(trimmed); ok {
1493 c.rememberProjectNote(note)
1494 return
1495 }
1496 if note, ok := RememberCommandNote(trimmed); ok {
1497 c.rememberProjectNote(note)
1498 return
1499 }
1500 if c.applyGoalCommandWithAdmission(trimmed, display, admission) {
1501 return
1502 }
1503 if strings.HasPrefix(trimmed, "!") {
1504 c.RunShell(trimmed[1:])
1505 return
1506 }
1507 c.submitCommandOrTurn(trimmed, input, display, false, editedOriginal, "", admission)
1508 }
1509
1510 func (c *Controller) submitHTTP(input, display string) {
1511 c.submitHTTPWithFormat(input, display, "")
1512 }
1513
1514 func (c *Controller) submitHTTPWithFormat(input, display, format string) {
1515 if isSessionManagementSubmission(input) {
1516 c.submissions.mu.Lock()
1517 defer c.releaseSubmissionAdmission()
1518 c.submitHTTPWithFormatLocked(input, display, format, turnAdmission{})
1519 return
1520 }
1521 _, _ = c.submitIdentifiedWithSetup(SubmissionRequest{Input: input, Display: display, HTTP: true, Format: format}, nil, func(admission turnAdmission) {
1522 c.submitHTTPWithFormatLocked(input, display, format, admission)
1523 })
1524 }
1525
1526 func (c *Controller) submitHTTPWithFormatLocked(input, display, format string, admission turnAdmission) {
1527 trimmed := strings.TrimSpace(input)
1528 if note, ok := MemoryQuickAddNote(trimmed); ok {
1529 c.rememberProjectNote(note)
1530 return
1531 }
1532 if note, ok := RememberCommandNote(trimmed); ok {
1533 c.rememberProjectNote(note)
1534 return
1535 }
1536 if c.applyGoalCommandWithAdmission(trimmed, display, admission) {
1537 return
1538 }
1539 if strings.HasPrefix(trimmed, "!") {
1540 c.notice("shell commands are unavailable from this frontend")
1541 return
1542 }
1543 c.submitCommandOrTurn(trimmed, input, display, true, "", format, admission)
1544 }
1545
1546 func (c *Controller) submitCommandOrTurnReady(trimmed, input, display string, scopedRefsOnly bool, editedOriginal, format string, admission turnAdmission) {
1547 runRefTurn := func(input, display string) {
1548 c.runRefTurnWithFormat(input, display, format, admission)
1549 }
1550 runRefTurnWithRefs := func(input, refLine, display string) {
1551 c.runRefTurnWithRefsFormat(input, refLine, display, format, admission)
1552 }
1553 runGoalLoop := func(ctx context.Context, input, raw, display string) error {
1554 return c.runGoalLoopWithRawDisplay(c.withTurnFormat(ctx, format), input, raw, display)
1555 }
1556 if scopedRefsOnly {
1557 runRefTurn = func(input, display string) {
1558 c.runScopedRefTurnWithFormat(input, display, format, admission)
1559 }
1560 runRefTurnWithRefs = func(input, refLine, display string) {
1561 c.runScopedRefTurnWithRefsFormat(input, refLine, display, format, admission)
1562 }
1563 }
1564 if strings.TrimSpace(editedOriginal) != "" {
1565 runRefTurn = func(input, display string) {
1566 c.runEditedRefTurnWithFormat(input, display, editedOriginal, format, admission)
1567 }
1568 runRefTurnWithRefs = func(input, refLine, display string) {
1569 c.runEditedRefTurnWithRefsFormat(input, refLine, display, editedOriginal, format, admission)
1570 }
1571 runGoalLoop = func(ctx context.Context, input, raw, display string) error {
1572 return c.runEditedGoalLoopWithRawDisplay(ctx, input, raw, display, editedOriginal)
1573 }
1574 }
1575 if id, guidance, ok := ParseProtocolRecoveryCommand(trimmed); ok {
1576 c.submitProtocolRecoveryLocked(id, guidance, admission)
1577 return
1578 }
1579 if c.submitFinalReadinessCommand(trimmed, display, admission) {
1580 return
1581 }
1582 switch {
1583 case trimmed == "/compact" || strings.HasPrefix(trimmed, "/compact "):
1584 focus := strings.TrimSpace(strings.TrimPrefix(trimmed, "/compact"))
1585 // Register synchronously before returning the management-command receipt.
1586 // This removes the window in which the command looked complete while a
1587 // new turn could still enter before compaction claimed the session.
1588 if err := c.startCompactAsync(focus); err != nil {
1589 c.notice("compaction failed: " + err.Error())
1590 }
1591 case trimmed == "/context":
1592 c.noticeDetail(c.ContextReport())
1593 case trimmed == "/new":
1594 c.runSessionVerb(c.NewSession, "new session", "new session failed: ")
1595 case trimmed == "/clear":
1596 c.runSessionVerb(c.ClearSession, "context cleared", "clear context failed: ")
1597 case strings.HasPrefix(trimmed, "/mcp__"):
1598 c.runGuardedWithAdmission(func(ctx context.Context) error {
1599 sent, found, err := c.MCPPrompt(ctx, trimmed)
1600 if err != nil {
1601 return err
1602 }
1603 if !found {
1604 c.notice("unknown command: " + trimmed)
1605 return nil
1606 }
1607 return runGoalLoop(ctx, sent, sent, display)
1608 }, admission)
1609 case SlashCodeCommentLine(trimmed):
1610 // Slash-prefixed code comments are prompt text, not slash commands.
1611 runRefTurn(input, display)
1612 case strings.HasPrefix(trimmed, "/"):
1613 if ref, ok := FileRefLine(trimmed); ok {
1614 runRefTurn(ref, display)
1615 return
1616 }
1617 if ref, ok := SlashPathLineRef(trimmed, c.workspaceRoot); ok {
1618 runRefTurnWithRefs(input, ref, display)
1619 return
1620 }
1621 if SlashPathLikeLine(trimmed) {
1622 runRefTurn(input, display)
1623 return
1624 }
1625 // Management verbs (/model /memory /skills /hooks /mcp) emit a Notice, so
1626 // Submit-based frontends (desktop, HTTP) get them with no extra wiring.
1627 // The chat TUI handles these itself with richer output.
1628 fields := strings.Fields(trimmed)
1629 switch fields[0] {
1630 case "/tree":
1631 c.notice(c.BranchTreeText())
1632 return
1633 case "/branch":
1634 args := strings.TrimSpace(strings.TrimPrefix(trimmed, fields[0]))
1635 if turn, name, fromTurn, err := ParseBranchTarget(args); err != nil {
1636 c.notice(err.Error())
1637 } else if fromTurn {
1638 if _, err := c.ForkNamed(turn-1, name); err != nil {
1639 c.notice(err.Error())
1640 }
1641 } else {
1642 if _, err := c.Branch(name); err != nil {
1643 c.notice(err.Error())
1644 }
1645 }
1646 return
1647 case "/switch":
1648 ref := strings.TrimSpace(strings.TrimPrefix(trimmed, fields[0]))
1649 if _, err := c.SwitchBranch(ref); err != nil {
1650 c.notice(err.Error())
1651 }
1652 return
1653 case "/rewind":
1654 args := strings.TrimSpace(strings.TrimPrefix(trimmed, fields[0]))
1655 turn, scope, err := parseRewind(args, c.Checkpoints())
1656 if err != nil {
1657 c.notice("usage: /rewind [turn] [code|conversation|both]")
1658 return
1659 }
1660 if err := c.Rewind(turn, scope); err != nil {
1661 c.notice(err.Error())
1662 }
1663 return
1664 case "/plan-exec":
1665 c.applyPlanExec(trimmed, display)
1666 return
1667 case "/prometheus":
1668 c.applyPrometheus(trimmed, display, admission)
1669 return
1670 }
1671 if c.managementNotice(trimmed) {
1672 return
1673 }
1674 if IsBuiltinDocsSlash(fields[0], c.Commands(), c.SlashSkills()) {
1675 query := strings.TrimSpace(strings.TrimPrefix(trimmed, fields[0]))
1676 if query == "" {
1677 text, err := DocsCommandOverviewFor(fields[0])
1678 if err != nil {
1679 c.notice("docs: " + err.Error())
1680 } else {
1681 c.notice(text)
1682 }
1683 return
1684 }
1685 c.runGuardedWithAdmission(func(ctx context.Context) error {
1686 sent, err := docsCommandPrompt(ctx, query)
1687 if err != nil {
1688 return fmt.Errorf("docs: %w", err)
1689 }
1690 return runGoalLoop(ctx, sent, sent, display)
1691 }, admission)
1692 return
1693 }
1694 // A custom command wins over a skill of the same name; both resolve to a
1695 // turn. Built-ins and their explicit Reasonix namespace are handled above.
1696 if sent, ok := c.CustomCommand(trimmed); ok {
1697 c.runGuardedWithAdmission(func(ctx context.Context) error {
1698 return runGoalLoop(ctx, sent, sent, display)
1699 }, admission)
1700 return
1701 }
1702 if sk, task, ok := c.resolveSkillInvocation(trimmed); ok {
1703 if sk.RunAs == skill.RunSubagent {
1704 if strings.TrimSpace(task) == "" {
1705 c.notice("usage: /" + sk.Name + " <task>")
1706 return
1707 }
1708 c.runSubagentSkillSlash(sk, task, trimmed, display, admission)
1709 return
1710 }
1711 sent := c.skills.renderInvocation(sk, task)
1712 c.runGuardedWithAdmission(func(ctx context.Context) error {
1713 return runGoalLoop(withInvokedSkills(ctx, []string{sk.Name}), sent, input, display)
1714 }, admission)
1715 return
1716 }
1717 // Unknown slash input is prose more often than a typo ("/etc/hosts
1718 // looks wrong", pasted paths, half-remembered commands) — send it as a
1719 // regular message instead of dead-ending the submission, with a notice
1720 // so real typos are still visible (#5756).
1721 c.notice("unknown command: " + trimmed + " — sent as a regular message")
1722 runRefTurn(input, display)
1723 default:
1724 runRefTurn(input, display)
1725 }
1726 }
1727
1728 func (c *Controller) rememberProjectNote(note string) {
1729 if note == "" {
1730 c.notice("nothing to remember")
1731 return
1732 }
1733 if path, err := c.QuickAdd(memory.ScopeProject, note); err != nil {
1734 c.notice("memory: " + err.Error())
1735 } else {
1736 c.notice("remembered → " + path)
1737 }
1738 }
1739
1740 // applyPlanExec is a command tombstone. The old path coupled Plan approval,
1741 // todo state and Goal continuation and is intentionally absent from the runtime.
1742 func (c *Controller) applyPlanExec(_, _ string) {
1743 c.notice("/plan-exec is retired; approve the Plan, then let the model create a fresh todo list for the new turn")
1744 }
1745
1746 // prometheusPrompt is the strategic planner system prompt.
1747 const prometheusPrompt = "You are Prometheus, a strategic planner. Interview the user one question at a time. Cover: scope, modules, files, constraints, tests. When ready, output a numbered plan with each step tagged by module. Read the current goal with get_goal, then call update_goal with its exact ID/revision and action complete. Do not implement.\n\nFor independent research directions, use parallel_tasks before planning."
1748
1749 // applyPrometheus starts an interactive planning interview, inspired by OMO's
1750 // Prometheus agent. It enters goal mode with a structured interview prompt.
1751 func (c *Controller) applyPrometheus(input, display string, admission turnAdmission) {
1752 args := strings.TrimSpace(strings.TrimPrefix(input, "/prometheus"))
1753 if args == "" || args == "--strict" {
1754 c.notice("usage: /prometheus <your task description>")
1755 return
1756 }
1757 strict := false
1758 if strings.HasPrefix(args, "--strict ") {
1759 strict = true
1760 args = strings.TrimPrefix(args, "--strict ")
1761 }
1762 prompt := prometheusPrompt + "\n\n## User request\n\n" + args + "\n\nBegin the interview by asking your first clarifying question."
1763 c.SetPlanMode(false)
1764 c.SetGoal("plan: " + ShortGoalForNotice(args))
1765 c.GoalStrict(strict)
1766 c.notice("prometheus: starting planning interview")
1767 if c.runner != nil {
1768 c.runGuardedWithAdmission(func(ctx context.Context) error {
1769 return c.runGoalLoopWithRawDisplay(ctx, prompt, prompt, display)
1770 }, admission)
1771 }
1772 }
1773
1774 // shellTimeout is the maximum time a user-invoked "!command" may run. Matches
1775 // the bash tool's timeout so behaviour is consistent across invocation paths.
1776 const shellTimeout = 120 * time.Second
1777
1778 // shellWaitDelay bounds how long cmd.Run() waits after context cancellation for
1779 // the child's pipes to drain, matching the bash tool's WaitDelay.
1780 const shellWaitDelay = 5 * time.Second
1781
1782 func shellCommandPreview(command string) string {
1783 command = strings.TrimSpace(strings.ReplaceAll(command, "\n", " "))
1784 const max = 48
1785 r := []rune(command)
1786 if len(r) > max {
1787 return string(r[:max]) + "…"
1788 }
1789 return command
1790 }
1791
1792 // RunShell executes a shell command directly (bypassing the model) and streams
1793 // the output as ToolDispatch/ToolProgress/ToolResult events. It uses the same
1794 // bash-tool infrastructure (shell resolution, timeout) and shares the runGuarded
1795 // lock with model turns — only one can run at a time. User-invoked "!" commands
1796 // run without the OS sandbox (the user typed the command explicitly).
1797 func (c *Controller) RunShell(command string) {
1798 c.runShell(command, turnAdmission{})
1799 }
1800
1801 func (c *Controller) runShell(command string, admission turnAdmission) {
1802 command = strings.TrimSpace(command)
1803 if command == "" {
1804 c.notice(i18n.M.ShellExecEmpty)
1805 return
1806 }
1807 c.runGuardedWithAdmission(func(ctx context.Context) error {
1808 sh := c.shell
1809 if sh.Path == "" {
1810 sh = sandbox.ResolveShell("", "", nil)
1811 }
1812 argv, _ := sandbox.Command(sandbox.Spec{}, sh, command) // false = unsandboxed (user invoked)
1813
1814 preview := []rune(command)
1815 if len(preview) > 32 {
1816 preview = preview[:32]
1817 }
1818 id := "shell-" + string(preview)
1819 diagnosticPreview := shellCommandPreview(command)
1820 desc := shellrun.DescriptorFromShell(sh)
1821 toolName := "bash"
1822 if sh.Kind == sandbox.ShellPowerShell {
1823 toolName = "pwsh"
1824 }
1825
1826 if err := event.EmitChecked(c.sink, event.Event{
1827 Kind: event.ToolDispatch,
1828 Tool: event.Tool{
1829 ID: id,
1830 Name: toolName,
1831 Args: fmt.Sprintf(`{"command":%q}`, command),
1832 Execution: &event.ShellExecution{
1833 Kind: desc.Kind, Shell: desc.Shell, ShellVersion: desc.ShellVersion,
1834 Platform: desc.Platform, SupportsAndAnd: desc.SupportsAndAnd,
1835 State: tool.ShellStateRunning,
1836 },
1837 },
1838 }); err != nil {
1839 return fmt.Errorf("persist shell dispatch: %w", err)
1840 }
1841
1842 start := time.Now()
1843 res := shellrun.RunForeground(ctx, shellrun.Request{
1844 Argv: argv,
1845 Dir: c.workspaceRoot,
1846 Timeout: shellTimeout,
1847 WaitDelay: shellWaitDelay,
1848 CommandPreview: diagnosticPreview,
1849 ShellKind: sh.Kind.String(),
1850 ShellPath: sh.Path,
1851 Source: "user_shell",
1852 Track: true,
1853 Progress: func(chunk string) {
1854 c.sink.Emit(event.Event{
1855 Kind: event.ToolProgress,
1856 Tool: event.Tool{ID: id, Output: chunk},
1857 })
1858 },
1859 })
1860 durationMs := time.Since(start).Milliseconds()
1861 ex := &event.ShellExecution{
1862 Kind: desc.Kind, Shell: desc.Shell, ShellVersion: desc.ShellVersion,
1863 Platform: desc.Platform, SupportsAndAnd: desc.SupportsAndAnd,
1864 State: res.State, FailurePhase: res.FailurePhase,
1865 OutputTail: res.OutputTail, DurationMs: durationMs,
1866 MutationRisk: tool.ShellMutationNone,
1867 Verification: tool.ShellVerificationNotVerification,
1868 }
1869 if res.ExitCode != nil {
1870 code := *res.ExitCode
1871 ex.ExitCode = &code
1872 }
1873 switch res.State {
1874 case tool.ShellStateCompleted:
1875 ex.MutationRisk = tool.ShellMutationNone
1876 case tool.ShellStateNotRun:
1877 ex.MutationRisk = tool.ShellMutationNotStarted
1878 case tool.ShellStateFailed:
1879 if res.FailurePhase == tool.ShellPhaseLaunch {
1880 ex.MutationRisk = tool.ShellMutationNotStarted
1881 } else {
1882 ex.MutationRisk = tool.ShellMutationMayBePartial
1883 }
1884 case tool.ShellStateTimedOut, tool.ShellStateCancelled:
1885 ex.MutationRisk = tool.ShellMutationMayBePartial
1886 }
1887
1888 errText := ""
1889 switch res.State {
1890 case tool.ShellStateCancelled:
1891 errText = i18n.M.TurnCancelled
1892 case tool.ShellStateTimedOut:
1893 errText = fmt.Sprintf(i18n.M.ShellExecTimeoutFmt, shellTimeout)
1894 case tool.ShellStateFailed, tool.ShellStateNotRun:
1895 if res.Err != nil {
1896 errText = fmt.Sprintf(i18n.M.ShellExecFailedFmt, res.Err)
1897 }
1898 }
1899 c.sink.Emit(event.Event{
1900 Kind: event.ToolResult,
1901 Tool: event.Tool{
1902 ID: id, Name: "bash", Output: res.Combined, Err: errText,
1903 DurationMs: durationMs, Execution: ex,
1904 },
1905 })
1906 return nil
1907 }, admission)
1908 }
1909
1910 // runRefTurn resolves a line's @references into a context block and starts a
1911 // turn with it prepended (or the raw line when nothing resolved).
1912 func (c *Controller) runRefTurn(input, display string) {
1913 c.runRefTurnWithAdmission(input, display, turnAdmission{})
1914 }
1915
1916 func (c *Controller) runRefTurnWithAdmission(input, display string, admission turnAdmission) {
1917 c.runRefTurnWithRefs(input, input, display, admission)
1918 }
1919
1920 // runRefTurnWithFormat runs a reference turn with a structured-output
1921 // format bound to its context (symmetric with runGoalLoop's withTurnFormat
1922 // injection — format is a property of every accepted turn, not just the
1923 // plain-goal path; review #7234 binds format to the accepted turn).
1924 func (c *Controller) runRefTurnWithFormat(input, display, format string, admission turnAdmission) {
1925 c.runPreparedRefTurn(input, input, display, "", c.resolveUnscopedRefsForTurn, func(ctx context.Context) context.Context {
1926 return c.withTurnFormat(ctx, format)
1927 }, admission)
1928 }
1929
1930 func (c *Controller) runScopedRefTurnWithFormat(input, display, format string, admission turnAdmission) {
1931 c.runPreparedRefTurn(input, input, display, "", c.resolveScopedRefsForTurn, func(ctx context.Context) context.Context {
1932 return c.withTurnFormat(ctx, format)
1933 }, admission)
1934 }
1935
1936 func (c *Controller) runRefTurnWithRefsFormat(input, refLine, display, format string, admission turnAdmission) {
1937 c.runPreparedRefTurn(input, refLine, display, "", c.resolveUnscopedRefsForTurn, func(ctx context.Context) context.Context {
1938 return c.withTurnFormat(ctx, format)
1939 }, admission)
1940 }
1941
1942 func (c *Controller) runScopedRefTurnWithRefsFormat(input, refLine, display, format string, admission turnAdmission) {
1943 c.runPreparedRefTurn(input, refLine, display, "", c.resolveScopedRefsForTurn, func(ctx context.Context) context.Context {
1944 return c.withTurnFormat(ctx, format)
1945 }, admission)
1946 }
1947
1948 func (c *Controller) runEditedRefTurnWithFormat(input, display, original, format string, admission turnAdmission) {
1949 c.runPreparedRefTurn(input, input, display, original, c.resolveUnscopedRefsForTurn, func(ctx context.Context) context.Context {
1950 return c.withTurnFormat(ctx, format)
1951 }, admission)
1952 }
1953
1954 func (c *Controller) runEditedRefTurnWithRefsFormat(input, refLine, display, original, format string, admission turnAdmission) {
1955 c.runPreparedRefTurn(input, refLine, display, original, c.resolveUnscopedRefsForTurn, func(ctx context.Context) context.Context {
1956 return c.withTurnFormat(ctx, format)
1957 }, admission)
1958 }
1959
1960 // runRefTurnWithRefs resolves references from refLine while preserving input as
1961 // the user's actual prompt text. This lets compiler diagnostics such as
1962 // "/path/File.kt:12: error" attach @/path/File.kt without rewriting the error.
1963 func (c *Controller) runRefTurnWithRefs(input, refLine, display string, admission turnAdmission) {
1964 c.runRefTurnWithResolver(input, refLine, display, c.resolveUnscopedRefsForTurn, admission)
1965 }
1966
1967 func (c *Controller) runRefTurnWithResolver(input, refLine, display string, resolve func(context.Context, string) resolvedReferences, admission turnAdmission) {
1968 c.runPreparedRefTurn(input, refLine, display, "", resolve, func(ctx context.Context) context.Context { return ctx }, admission)
1969 }
1970
1971 func (c *Controller) runRefTurnWithResolverSync(ctx context.Context, input, refLine, display, original string, resolve func(context.Context, string) resolvedReferences) error {
1972 resolved := resolve(ctx, refLine)
1973 return c.runResolvedRefTurnSync(ctx, input, display, original, resolved)
1974 }
1975
1976 func (c *Controller) runResolvedRefTurnSync(ctx context.Context, input, display, original string, resolved resolvedReferences) error {
1977 if len(resolved.imageErrs) > 0 {
1978 return ImageReferenceFailures(resolved.imageErrs)
1979 }
1980 for _, e := range resolved.errs {
1981 c.notice(e)
1982 }
1983 sent := input
1984 if resolved.block != "" {
1985 sent = "Referenced context:\n\n" + resolved.block + "\n\n" + input
1986 }
1987 if strings.TrimSpace(original) != "" {
1988 return c.runEditedGoalLoopWithFrozenImagesRawDisplay(ctx, sent, input, display, original, resolved.images)
1989 }
1990 return c.runGoalLoopWithFrozenImagesRawDisplay(ctx, sent, input, display, resolved.images)
1991 }
1992
1993 // notice emits an informational Notice event.
1994 func (c *Controller) notice(text string) {
1995 c.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelInfo, Text: text})
1996 }
1997
1998 func (c *Controller) noticeDetail(text, detail string) {
1999 c.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelInfo, Text: text, Detail: detail})
2000 }
2001
2002 // Run executes a turn synchronously, returning the agent's error. Used by the
2003 // headless `reasonix run` path, where the Sink renders to stdout and the caller
2004 // just needs the exit status — no TurnDone event, no cancel bookkeeping.
2005 func (c *Controller) runReady(ctx context.Context, input string) (err error) {
2006 ctx = extension.ContextWithRuntimeOwner(ctx, c.RuntimeOwner())
2007 if c.RuntimePhase() == RuntimePhaseDraining {
2008 c.emitDrainingNotice()
2009 return ErrRuntimeDraining
2010 }
2011 c.maybeSessionStart(ctx)
2012 parentSession := c.parentSessionID()
2013 ctx = agent.WithParentSession(ctx, parentSession)
2014 ctx = jobs.WithSession(ctx, parentSession)
2015 rawInput := input
2016 ctx = c.withTurnImages(ctx, rawInput)
2017 ctx = agent.WithRawUserInput(ctx, rawInput)
2018 input = c.Compose(input)
2019 // input.receive: same interception seam as the orchestrated turn — the
2020 // composed headless input crosses the extension chain before it enters
2021 // the session.
2022 input, blocked, interceptErr := c.interceptInputReceive(ctx, input)
2023 if interceptErr != nil {
2024 return interceptErr
2025 }
2026 if blocked {
2027 return nil
2028 }
2029 startMessages := c.messageCount()
2030 var marker agent.InFlightTurnMeta
2031 defer func() { c.finishInFlightTurn(startMessages, marker) }()
2032 c.beginCheckpoint(ctx, rawInput)
2033 if c.guardianSess != nil {
2034 c.guardianSess.ResetTurn()
2035 }
2036 if c.hooks.Enabled() {
2037 c.mu.Lock()
2038 c.turn++
2039 turn := c.turn
2040 c.mu.Unlock()
2041 if block, _ := c.hooks.PromptSubmit(ctx, input, turn); block {
2042 return nil
2043 }
2044 defer func() { c.hooks.StopResult(context.Background(), lastAssistantText(c.History()), turn, err) }()
2045 }
2046 marker = c.markInFlightTurn(startMessages, true)
2047 ctx = c.withTurnContext(ctx, true)
2048 ctx = c.withPlannerTurnMetadata(ctx, rawInput, false, startMessages)
2049 modelInput := c.withCapabilityRoute(ctx, input, rawInput)
2050 modelInput, ctx, err = c.prepareVisionTurn(ctx, modelInput, agent.SubagentImageCandidates(ctx))
2051 if err != nil {
2052 return err
2053 }
2054 err = c.runModelTurn(ctx, modelInput)
2055 return err
2056 }
2057
2058 // beginRotation claims the session-rotation gate. It fails if a turn is running
2059 // or another rotation is already in progress, so the caller holds exclusive
2060 // rights to swap the executor session from the check here through endRotation.
2061 // This closes the TOCTOU window that a bare `if c.running` check left open:
2062 // between that check and the actual SetSession, a turn could start and then be
2063 // yanked out from under the run loop.
2064 func (c *Controller) beginRotation() error {
2065 c.mu.Lock()
2066 defer c.mu.Unlock()
2067 if c.bodyActiveLocked() || c.finalizingLocked() {
2068 return errTurnRunningRotation
2069 }
2070 if c.maintenance != nil {
2071 return ErrMaintenanceBusy
2072 }
2073 if c.rotating {
2074 return errRotationInProgress
2075 }
2076 c.rotating = true
2077 return nil
2078 }
2079
2080 // CancelRequested reports whether Cancel has been requested for the active turn.
2081 func (c *Controller) CancelRequested() bool {
2082 c.mu.Lock()
2083 defer c.mu.Unlock()
2084 return c.cancelRequestedLocked()
2085 }
2086
2087 // PendingPrompt reports whether the current turn is blocked waiting for a user
2088 // approval, plan approval, memory approval, or ask-tool answer.
2089 func (c *Controller) PendingPrompt() bool {
2090 return len(c.promptOwner.Identities()) > 0
2091 }
2092
2093 // RuntimeStatus reports the active work owned by the foreground controller.
2094 func (c *Controller) RuntimeStatus() RuntimeStatus {
2095 snapshot := c.RuntimeStateSnapshot()
2096 _, _, _, replayAfterSeq := c.turnEventRuntimeStatus()
2097 return RuntimeStatus{
2098 Running: snapshot.Running,
2099 PendingPrompt: snapshot.PendingPrompt,
2100 BackgroundJobs: snapshot.BackgroundJobs,
2101 CancelRequested: snapshot.CancelRequested,
2102 Cancellable: snapshot.Cancellable,
2103 TurnID: snapshot.TurnID,
2104 Status: snapshot.TurnStatus,
2105 TurnEventSeq: snapshot.TurnEventSeq,
2106 ReplayAfterSeq: replayAfterSeq,
2107 }
2108 }
2109
2110 // Turn returns the current turn number (0 before the first submit).
2111 func (c *Controller) Turn() int {
2112 c.mu.Lock()
2113 defer c.mu.Unlock()
2114 return c.turn
2115 }
2116
2117 func (c *Controller) recordDecisionReceipt(pending pendingApproval, outcome string) {
2118 if c == nil || c.executor == nil || pending.reply == nil {
2119 return
2120 }
2121 kind := pending.kind
2122 if kind == "" {
2123 kind = "tool"
2124 if pending.tool == planApprovalTool {
2125 kind = "plan"
2126 }
2127 }
2128 receipt := &provider.DecisionReceipt{
2129 ID: pending.id,
2130 Kind: kind,
2131 Tool: strings.TrimSpace(pending.tool),
2132 Subject: clipUTF8(strings.TrimSpace(pending.subject), 240),
2133 Outcome: strings.TrimSpace(outcome),
2134 }
2135 // Keep the receipt bounded and provider-excluded even when an older caller
2136 // omits optional approval metadata.
2137 c.executor.Session().AddDecisionReceipt(receipt)
2138 c.sink.Emit(event.Event{
2139 Kind: event.Notice,
2140 Code: event.NoticeCodeDecisionReceipt,
2141 Level: event.LevelInfo,
2142 Text: "Decision recorded: " + receipt.Outcome,
2143 DecisionReceipt: receipt,
2144 })
2145 }
2146
2147 // EnableInteractiveApproval swaps the executor's gate for one that routes
2148 // approval decisions to the frontend via ApprovalRequest events, and wires the
2149 // controller in as the executor's Asker so the `ask` tool can question the user.
2150 // Interactive frontends (chat, desktop) call this; the headless run keeps the
2151 // silent gate and a nil asker from setup.
2152 func (c *Controller) EnableInteractiveApproval() {
2153 trustGate := planModeReadOnlyTrustApprover{c}
2154 escapeApprover := sandboxEscapeApprover{c}
2155 configApprover := managedConfigWriteApprover{c}
2156 c.writeAccess.interactive = true
2157 if c.executor != nil {
2158 c.executor.SetGate(c.newInteractiveGate())
2159 c.executor.SetPlanModeReadOnlyTrustGate(trustGate)
2160 c.executor.SetSandboxEscapeApprover(escapeApprover)
2161 c.executor.SetConfigWriteApprover(configApprover)
2162 c.executor.SetWriteAccessGate(c)
2163 c.executor.SetWriteRoots(c.writeAccess.roots)
2164 c.executor.SetPermissionPresetProvider(c.ToolApprovalMode)
2165 c.executor.SetAsker(c)
2166 c.executor.SetInteractionBroker(c)
2167 }
2168 if setter, ok := c.runner.(interface {
2169 SetPlanModeReadOnlyTrustGate(agent.PlanModeReadOnlyTrustGate)
2170 }); ok {
2171 setter.SetPlanModeReadOnlyTrustGate(trustGate)
2172 }
2173 if setter, ok := c.runner.(interface {
2174 SetSandboxEscapeApprover(sandbox.EscapeApprover)
2175 }); ok {
2176 setter.SetSandboxEscapeApprover(escapeApprover)
2177 }
2178 if setter, ok := c.runner.(interface {
2179 SetConfigWriteApprover(tool.ConfigWriteApprover)
2180 }); ok {
2181 setter.SetConfigWriteApprover(configApprover)
2182 }
2183 if setter, ok := c.runner.(interface {
2184 SetWriteAccessGate(agent.WriteAccessGate)
2185 }); ok {
2186 setter.SetWriteAccessGate(c)
2187 }
2188 if setter, ok := c.runner.(interface{ SetPermissionPresetProvider(func() string) }); ok {
2189 setter.SetPermissionPresetProvider(c.ToolApprovalMode)
2190 }
2191 if setter, ok := c.runner.(interface {
2192 SetWriteRoots(*sandbox.WritableRootSet)
2193 }); ok {
2194 setter.SetWriteRoots(c.writeAccess.roots)
2195 }
2196 if setter, ok := c.runner.(interface {
2197 SetPlannerPlanApprover(agent.PlannerPlanApprover)
2198 }); ok {
2199 setter.SetPlannerPlanApprover(plannerPlanApprover{c: c})
2200 }
2201 // The planner holds the real ask tool, so it reaches the same approval
2202 // surface the executor does instead of a parallel prose-question path.
2203 if setter, ok := c.runner.(interface{ SetAsker(agent.Asker) }); ok {
2204 setter.SetAsker(c)
2205 }
2206 if setter, ok := c.runner.(interface{ SetInteractionBroker(mcpinteraction.Broker) }); ok {
2207 setter.SetInteractionBroker(c)
2208 }
2209 }
2210
2211 type plannerPlanApprover struct {
2212 c *Controller
2213 }
2214
2215 func (p plannerPlanApprover) RunWithPlannerApproval(ctx context.Context, plan string, run func(context.Context) error) error {
2216 c := p.c
2217 allow, _, err := c.requestApprovalWithReason(ctx, planApprovalTool, "", nil, "Planner requested host approval before execution.")
2218 if err != nil {
2219 return err
2220 }
2221 if !allow {
2222 return nil
2223 }
2224 c.approval.setPlanAutoApprove(true)
2225 defer c.approval.setPlanAutoApprove(false)
2226 if err := run(ctx); err != nil {
2227 return err
2228 }
2229 return nil
2230 }
2231
2232 func (c *Controller) newInteractiveGate() *permission.Gate {
2233 policy := c.policy
2234 mode := c.approval.mode()
2235 switch mode {
2236 case ToolApprovalWorkspaceWrite, ToolApprovalDangerFullAccess:
2237 policy.Mode = permission.Allow
2238 case ToolApprovalDontAsk:
2239 policy.Mode = permission.Deny
2240 default:
2241 policy.Mode = permission.Ask
2242 }
2243 // SessionAllow must not cover fresh-human tools: it is checked before Ask,
2244 // so `--allowed-tools remember` would skip the prompt. Interactive Auto and
2245 // YOLO treat remember/forget as ordinary policy decisions; Auto still
2246 // preserves an explicit configured Ask rule, while YOLO bypasses it.
2247 policy.SessionAllow = rulesWithoutFreshHumanApproval(policy.SessionAllow)
2248 if mode != ToolApprovalWorkspaceWrite && mode != ToolApprovalDangerFullAccess {
2249 policy.Ask = append(policy.Ask,
2250 permission.Rule{Tool: memoryRememberTool},
2251 permission.Rule{Tool: memoryForgetTool},
2252 )
2253 }
2254 // The OS sandbox, rather than shell syntax heuristics, owns the write
2255 // boundary for all three presets. Explicit ask and deny rules still win.
2256 var approver permission.Approver = gateApprover{c}
2257 if mode == ToolApprovalDontAsk {
2258 approver = denyPermissionApprover{}
2259 }
2260 gate := permission.NewGate(policy, approver)
2261 gate.OnRemember = func(rule string) {
2262 if c.onRemember != nil {
2263 _ = c.onRemember(rule)
2264 }
2265 }
2266 return gate
2267 }
2268
2269 func (c *Controller) allowLowRiskRemember(args json.RawMessage) bool {
2270 mem := c.Memory()
2271 if mem != nil {
2272 if assessment := memory.AssessRememberWrite(mem.Store, args); assessment.AutoAllow {
2273 c.memory.authorizeAutoRemember(args)
2274 return true
2275 }
2276 }
2277 c.memory.revokeAutoRemember(args)
2278 return false
2279 }
2280
2281 func (c *Controller) newHeadlessGate(mode string) *freshHumanHeadlessGate {
2282 gate := BuildHeadlessApprovalGate(c.policy, mode)
2283 gate.allowLowRiskFreshAction = func(toolName string, args json.RawMessage) bool {
2284 return toolName == memoryRememberTool && c.allowLowRiskRemember(args)
2285 }
2286 return gate
2287 }
2288
2289 type denyPermissionApprover struct{}
2290
2291 func (denyPermissionApprover) Approve(context.Context, string, string, json.RawMessage) (bool, bool, error) {
2292 return false, false, nil
2293 }
2294
2295 // rulesWithoutFreshHumanApproval drops any session-allow rule that targets a
2296 // tool requiring fresh human approval, so an explicit allowlist cannot bypass
2297 // the always-prompt contract for those tools.
2298 func rulesWithoutFreshHumanApproval(rules []permission.Rule) []permission.Rule {
2299 if len(rules) == 0 {
2300 return rules
2301 }
2302 filtered := make([]permission.Rule, 0, len(rules))
2303 for _, r := range rules {
2304 if RequiresFreshHumanApprovalTool(r.Tool) {
2305 continue
2306 }
2307 filtered = append(filtered, r)
2308 }
2309 return filtered
2310 }
2311
2312 // ApplyHeadlessApprovalMode configures the executor gate for a non-interactive
2313 // (`reasonix run`) session from an explicit --permission-mode. Unlike
2314 // EnableInteractiveApproval it installs no blocking approver, asker, or
2315 // fresh-approval prompt: there is no key loop to answer them, and the default
2316 // infinite approval timeout would wedge the run forever on an Ask rule, the
2317 // `ask` tool, or a sandbox/config approval. Modes map straight onto a headless
2318 // gate, and each preserves the interactive contract as closely as a run with no
2319 // one to prompt allows:
2320 //
2321 // - auto: auto-approve the writer fallback (Mode=Allow) but PRESERVE explicit
2322 // ask rules. Interactive auto prompts on those (it never auto-approves them);
2323 // headless can't prompt, so a would-ask decision fails closed (deny) rather
2324 // than running silently. Only bypass may run such a command unattended.
2325 // - yolo/bypassPermissions: skip ordinary approval-gated decisions (nil
2326 // approver); deny rules and fresh decisions still fail closed.
2327 // - dontAsk: deny anything that would ask, and deny the writer fallback too.
2328 //
2329 // Deny rules and fresh-human tools (memory, plan, sandbox, config) stay enforced
2330 // by the gate for every mode. The only exception is a controller-assessed,
2331 // create-only project/reference memory; every other memory write remains denied.
2332 func (c *Controller) ApplyHeadlessApprovalMode(mode string) {
2333 mode = normalizeToolApprovalMode(mode)
2334 c.permissionStateMu.Lock()
2335 defer c.permissionStateMu.Unlock()
2336 c.approval.setMode(mode)
2337 if c.subagentGate != nil {
2338 c.subagentGate.Update(mode)
2339 }
2340 c.writeAccess.interactive = false
2341 if c.executor != nil {
2342 c.executor.SetGate(c.newHeadlessGate(mode))
2343 c.executor.SetWriteAccessGate(c)
2344 c.executor.SetWriteRoots(c.writeAccess.roots)
2345 c.executor.SetPermissionPresetProvider(c.ToolApprovalMode)
2346 }
2347 }
2348
2349 func (c *Controller) refreshInteractiveGate() {
2350 if c.executor != nil {
2351 c.executor.SetGate(c.newInteractiveGate())
2352 }
2353 }
2354
2355 // TrySteer queues mid-turn guidance only when the active agent turn accepts it.
2356 func (c *Controller) TrySteer(text string) bool {
2357 c.mu.Lock()
2358 exec := c.executor
2359 running := c.bodyActiveLocked()
2360 c.mu.Unlock()
2361 return running && exec != nil && exec.Steer(text)
2362 }
2363
2364 // Steer is the compatibility path for callers that cannot observe admission.
2365 // Interactive hosts should call TrySteer so a rejected steer remains in their
2366 // draft/queue and can be retried as a regular follow-up.
2367 func (c *Controller) Steer(text string) {
2368 if c.TrySteer(text) {
2369 return
2370 }
2371 // No active turn accepted the steer: the frontend's runningRef was stale,
2372 // the turn exited between our running check and the enqueue, or no
2373 // executor is bound yet. Deliver it as a regular turn instead.
2374 c.submitSteerFallback(text)
2375 }
2376
2377 // submitSteerFallback records steer text that no active turn accepted as
2378 // unapplied guidance, not as a new task. This compatibility path deliberately
2379 // never opens a provider turn: replaying stale historical guidance as the
2380 // user's current request caused unintended code changes (#7045).
2381 func (c *Controller) submitSteerFallback(text string) admissionResult {
2382 return c.runGuardedOrPark(func(context.Context) error {
2383 if c.executor != nil {
2384 c.executor.RecordUnappliedSteer(text)
2385 }
2386 return nil
2387 })
2388 }
2389
2390 // SteerConsumed returns true when the steer queue is empty after the last consume.
2391 func (c *Controller) SteerConsumed() bool {
2392 c.mu.Lock()
2393 exec := c.executor
2394 c.mu.Unlock()
2395 if exec != nil {
2396 return exec.SteerConsumed()
2397 }
2398 return true
2399 }
2400
2401 // Ask implements agent.Asker: it emits an AskRequest and blocks until
2402 // AnswerQuestion(ID, …) answers or ctx is cancelled. Multiple requests may be
2403 // outstanding; the frontend presents the shared pending list one at a time.
2404 // Unlike tool-approval gates, Ask is NOT bypassed in YOLO mode — the `ask`
2405 // tool exists to get a genuine user decision, and YOLO only auto-approves
2406 // tool calls; it must not answer the user's questions for them.
2407 func (c *Controller) Ask(ctx context.Context, questions []event.AskQuestion) ([]event.AskAnswer, error) {
2408 c.approval.promptEmitMu.Lock()
2409 id, reply := c.approval.registerAsk(questions)
2410 c.registerOwnedPrompt(id, PromptAsk)
2411 turnID, _, _, _ := c.turnEventRuntimeStatus()
2412 _, runtimeEpoch := c.promptIdentitySnapshot()
2413 if identity := c.bindOwnedPromptRouting(id, turnID, runtimeEpoch); identity.TurnID != "" {
2414 turnID = identity.TurnID
2415 }
2416 if err := event.EmitChecked(c.sink, event.Event{Kind: event.AskRequest, TurnID: turnID, ItemID: id, Ask: event.Ask{ID: id, Questions: questions, TurnID: turnID}}); err != nil {
2417 c.approval.promptEmitMu.Unlock()
2418 c.cancelOwnedPrompt(id)
2419 return nil, fmt.Errorf("persist ask request: %w", err)
2420 }
2421 c.approval.markAskEmitted(id)
2422 c.approval.promptEmitMu.Unlock()
2423
2424 waitCtx, cancelWait := c.approval.waitContext(ctx)
2425 defer cancelWait()
2426
2427 select {
2428 case ans := <-reply:
2429 return ans, nil
2430 case <-waitCtx.Done():
2431 c.cancelOwnedPrompt(id)
2432 return nil, waitCtx.Err()
2433 }
2434 }
2435
2436 // AnswerQuestion resolves a pending AskRequest by ID with the user's selections.
2437 // Unknown/expired IDs are ignored.
2438 func (c *Controller) AnswerQuestion(id string, answers []event.AskAnswer) {
2439 _ = c.AnswerQuestionChecked(id, answers)
2440 }
2441
2442 // AnswerQuestionChecked persists the prompt transition before releasing the
2443 // agent loop. A failed ledger write leaves the prompt pending and retryable.
2444 func (c *Controller) AnswerQuestionChecked(id string, answers []event.AskAnswer) error {
2445 defer c.refreshRuntimeState(event.Event{})
2446 return c.answerQuestionCheckedLocked(id, answers)
2447 }
2448
2449 func (c *Controller) answerQuestionCheckedLocked(id string, answers []event.AskAnswer) error {
2450 pending, ok, err := c.approval.resolveAskAfter(id, func(p pendingAsk) error {
2451 return c.emitTurnEventChecked(event.Event{Kind: event.PromptAnswered, ItemID: id, InteractionState: string(PromptAnswered), Status: event.TurnInProgress})
2452 })
2453 if err != nil {
2454 return err
2455 }
2456 if ok {
2457 c.promptOwner.MarkIDTerminal(id, PromptAnswered)
2458 // An answer batch with no selections is the explicit "skip and continue
2459 // chat" path. End the current turn instead of feeding a prose dismissal
2460 // back to the model and trusting it not to ask again (#6869).
2461 if !askAnswersHaveSelection(answers) {
2462 c.mu.Lock()
2463 activeTurn := c.turns.cancel != nil
2464 c.mu.Unlock()
2465 if activeTurn {
2466 c.cancelLocked()
2467 return nil
2468 }
2469 }
2470 c.recordAskDecisionReceipt(id, pending, answers)
2471 pending.reply <- answers // buffered, never blocks
2472 }
2473 return nil
2474 }
2475
2476 func (c *Controller) recordAskDecisionReceipt(id string, pending pendingAsk, answers []event.AskAnswer) {
2477 if c == nil || c.executor == nil {
2478 return
2479 }
2480 selected := make(map[string][]string, len(answers))
2481 for _, answer := range answers {
2482 selected[answer.QuestionID] = append([]string(nil), answer.Selected...)
2483 }
2484 parts := make([]string, 0, len(pending.questions))
2485 for _, question := range pending.questions {
2486 answer := strings.TrimSpace(strings.Join(selected[question.ID], ", "))
2487 if answer == "" {
2488 answer = "—"
2489 }
2490 prompt := strings.TrimSpace(question.Prompt)
2491 if prompt == "" {
2492 prompt = strings.TrimSpace(question.Header)
2493 }
2494 if prompt == "" {
2495 prompt = question.ID
2496 }
2497 parts = append(parts, prompt+": "+answer)
2498 }
2499 receipt := &provider.DecisionReceipt{
2500 ID: id,
2501 Kind: "ask",
2502 Subject: clipUTF8(strings.Join(parts, " · "), 240),
2503 Outcome: "answered",
2504 }
2505 c.executor.Session().AddDecisionReceipt(receipt)
2506 c.sink.Emit(event.Event{
2507 Kind: event.Notice,
2508 Code: event.NoticeCodeDecisionReceipt,
2509 Level: event.LevelInfo,
2510 Text: "Decision recorded: answered",
2511 DecisionReceipt: receipt,
2512 })
2513 }
2514
2515 func askAnswersHaveSelection(answers []event.AskAnswer) bool {
2516 for _, answer := range answers {
2517 if len(answer.Selected) > 0 {
2518 return true
2519 }
2520 }
2521 return false
2522 }
2523
2524 // ReplayPendingPrompts re-emits the ApprovalRequest / AskRequest event for every
2525 // prompt currently blocking the run loop. A frontend that reconnected or reloaded
2526 // after the original event has no way to rebuild its approval/ask modal otherwise,
2527 // so the blocked gate goroutine stays stuck forever while the session shows a
2528 // "waiting" status with no actionable prompt. All outstanding interactions
2529 // are replayed from the same registry; the frontend presents them in order.
2530 func (c *Controller) ReplayPendingPrompts() {
2531 c.approval.promptEmitMu.Lock()
2532 noApprovals := c.replayPendingPromptsTo(c.sink)
2533 c.approval.promptEmitMu.Unlock()
2534 if noApprovals {
2535 // Retained compatibility hook; live Auto Guard cards are ordinary approvals.
2536 c.ReplayUnresolvedRecoveries()
2537 }
2538 }
2539
2540 // ReplayPendingPromptsTo re-emits pending prompts to one frontend sink. Serve
2541 // uses this for a newly attached SSE client so existing browsers do not receive
2542 // duplicate approval/ask cards when another client reconnects.
2543 func (c *Controller) ReplayPendingPromptsTo(sink event.Sink) {
2544 c.approval.promptEmitMu.Lock()
2545 defer c.approval.promptEmitMu.Unlock()
2546 c.replayPendingPromptsTo(sink)
2547 }
2548
2549 // ReplayPendingPromptsWith performs an SSE connection handoff while prompt
2550 // registration and emission are paused. The factory must subscribe the new
2551 // client and return a sink that targets it; this closes the attach race where
2552 // the original prompt could otherwise land between Subscribe and replay.
2553 func (c *Controller) ReplayPendingPromptsWith(sinkFactory func() event.Sink) {
2554 if sinkFactory == nil {
2555 return
2556 }
2557 c.approval.promptEmitMu.Lock()
2558 defer c.approval.promptEmitMu.Unlock()
2559 c.replayPendingPromptsTo(sinkFactory())
2560 }
2561
2562 func (c *Controller) replayPendingPromptsTo(sink event.Sink) bool {
2563 approvals, asks := c.approval.snapshotPrompts()
2564 interactions := c.approval.snapshotMCPInteractions()
2565 c.emitPendingPrompts(sink, approvals, asks, interactions)
2566 return len(approvals) == 0
2567 }
2568
2569 func (c *Controller) emitPendingPrompts(sink event.Sink, approvals []event.Approval, asks []event.Ask, interactions []event.MCPInteraction) {
2570 if sink == nil {
2571 return
2572 }
2573 for _, a := range approvals {
2574 if identity, ok := c.promptOwner.Identity(a.ID); ok {
2575 a.TurnID = identity.TurnID
2576 }
2577 e := c.approvalRequestEvent(a)
2578 e.Replayed = true
2579 sink.Emit(e)
2580 }
2581 for _, a := range asks {
2582 if identity, ok := c.promptOwner.Identity(a.ID); ok {
2583 a.TurnID = identity.TurnID
2584 }
2585 sink.Emit(event.Event{Kind: event.AskRequest, TurnID: a.TurnID, ItemID: a.ID, Ask: a, Replayed: true})
2586 }
2587 for _, i := range interactions {
2588 if identity, ok := c.promptOwner.Identity(i.ID); ok {
2589 i.TurnID = identity.TurnID
2590 }
2591 sink.Emit(event.Event{Kind: event.MCPInteractionRequest, TurnID: i.TurnID, ItemID: i.ID, MCPInteraction: i, Replayed: true})
2592 }
2593 }
2594
2595 // SetPlanMode flips the executor's plan-first workflow flag without touching the
2596 // cache-stable system/tool prefix, and remembers the state so Compose can prepend
2597 // the plan-mode marker to outgoing user turns.
2598 func (c *Controller) SetPlanMode(v bool) {
2599 c.applyPlanMode(v)
2600 }
2601
2602 // SetAgentPreset accepts retired role inputs for compatibility. Recognized
2603 // values no longer change runtime behavior.
2604 func (c *Controller) SetAgentPreset(preset string) {
2605 if c == nil {
2606 return
2607 }
2608 if p, err := agentpreset.Normalize(preset); err == nil {
2609 _ = c.SetQualityFloor(string(p))
2610 }
2611 }
2612
2613 // AgentPreset returns the fixed compatibility label.
2614 func (c *Controller) AgentPreset() string {
2615 return string(agentpreset.Standard)
2616 }
2617
2618 // SetResponseLanguage updates the final-answer language preference for
2619 // subsequent turns.
2620 func (c *Controller) SetResponseLanguage(lang string) {
2621 mode := config.NormalizeLanguage(lang)
2622 c.mu.Lock()
2623 c.responseLanguage = mode
2624 c.mu.Unlock()
2625 if setter, ok := c.runner.(interface{ SetResponseLanguage(string) }); ok {
2626 setter.SetResponseLanguage(mode)
2627 } else if c.executor != nil {
2628 c.executor.SetResponseLanguage(mode)
2629 }
2630 }
2631
2632 // SetReasoningLanguage updates the visible reasoning language preference for
2633 // subsequent turns.
2634 func (c *Controller) SetReasoningLanguage(lang string) {
2635 mode := config.NormalizeReasoningLanguage(lang)
2636 c.mu.Lock()
2637 c.reasoningLanguage = mode
2638 c.mu.Unlock()
2639 if setter, ok := c.runner.(interface{ SetReasoningLanguage(string) }); ok {
2640 setter.SetReasoningLanguage(mode)
2641 } else if c.executor != nil {
2642 c.executor.SetReasoningLanguage(mode)
2643 }
2644 }
2645
2646 // PlanMode reports whether outgoing turns currently receive the plan-mode
2647 // marker.
2648 func (c *Controller) PlanMode() bool {
2649 c.mu.Lock()
2650 defer c.mu.Unlock()
2651 return c.sessionSettings.planMode
2652 }
2653
2654 // GoalStrict enables or disables strict goal mode. Since the structured
2655 // protocol, every complete claim is validated against host readiness and an
2656 // incomplete-todo intercept can never be overridden, so the flag is persisted
2657 // for compatibility with older frontends but no longer changes FSM behavior.
2658 func (c *Controller) GoalStrict(strict bool) {
2659 if c.sessionEngineEnabled() {
2660 return
2661 }
2662 path, data, ok := c.goals.setStrict(strict)
2663 c.persistGoalState(path, data, ok)
2664 }
2665
2666 // SetGoal stores a session-scoped active goal. Compose injects it into outgoing
2667 // user turns, not the system prompt or tool schema, so it does not disturb the
2668 // cache-stable prefix.
2669 func (c *Controller) SetGoal(goal string) {
2670 c.SetGoalWithResearchMode(goal, GoalResearchAuto)
2671 }
2672
2673 // LoadInactiveGoal restores a legacy metadata-only objective without starting
2674 // execution or writing a sidecar during read-only history access.
2675 func (c *Controller) LoadInactiveGoal(goal string) {
2676 c.goals.mu.Lock()
2677 defer c.goals.mu.Unlock()
2678 c.goals.installGoalLocked(strings.TrimSpace(goal), ClassifyGoalBudget(goal))
2679 c.goals.disarmed = true
2680 }
2681
2682 // SetGoalDurable updates the Goal only when its sidecar can be replaced
2683 // atomically.
2684 func (c *Controller) SetGoalDurable(goal string) error {
2685 if c.sessionEngineEnabled() {
2686 goal = strings.TrimSpace(goal)
2687 current, err := c.goalLifecycleView()
2688 if err != nil {
2689 return err
2690 }
2691 if goal == "" {
2692 if current == nil {
2693 return nil
2694 }
2695 _, err = c.applyHostGoalMutation(context.Background(), "clear", func(machine *goaldomain.Machine) (*goaldomain.View, error) {
2696 if clearErr := machine.Clear(current.Ref()); clearErr != nil {
2697 return nil, clearErr
2698 }
2699 return nil, nil
2700 })
2701 if err == nil {
2702 c.resetGoalResourceBudget()
2703 }
2704 return err
2705 }
2706 if current != nil && current.Objective == goal && current.Phase == goaldomain.PhaseActive && current.Activation == goaldomain.ActivationArmed {
2707 return nil
2708 }
2709 _, err = c.applyHostGoalMutation(context.Background(), "set", func(machine *goaldomain.Machine) (*goaldomain.View, error) {
2710 created, createErr := machine.Replace(goaldomain.CreateRequest{Objective: goal})
2711 return &created, createErr
2712 })
2713 if err == nil {
2714 c.resetGoalResourceBudget()
2715 }
2716 return err
2717 }
2718 snapshot := c.goals.capture()
2719 legacySnapshot, hadLegacySnapshot := c.legacyRestoreSnapshot()
2720 resolved, setup := c.resolveGoalText(goal, GoalResearchAuto)
2721 var path string
2722 var data []byte
2723 var persist bool
2724 if setup.blockReason != "" {
2725 path, data, persist = c.goals.setLegacyArchiveBlockedWithTaskID(resolved, setup.budgetClass, setup.blockReason, setup.legacyTaskID)
2726 c.replaceLegacyRestore(legacyGoalRestore{taskID: setup.legacyTaskID, epoch: c.goals.continuationToken(), explicit: setup.explicit})
2727 } else {
2728 path, data, persist = c.goals.set(resolved, setup.budgetClass)
2729 c.replaceLegacyRestore(legacyGoalRestore{})
2730 }
2731 if persist {
2732 if err := c.goals.writeStateErr(path, data); err != nil {
2733 c.goals.restore(snapshot)
2734 if hadLegacySnapshot {
2735 legacySnapshot.epoch = c.goals.continuationToken()
2736 c.replaceLegacyRestore(legacySnapshot)
2737 } else {
2738 c.replaceLegacyRestore(legacyGoalRestore{})
2739 }
2740 return err
2741 }
2742 }
2743 if setup.notice != "" {
2744 c.notice(setup.notice)
2745 }
2746 if setup.blockReason != "" {
2747 c.notice("legacy research archive resume failed: " + setup.blockReason)
2748 }
2749 return nil
2750 }
2751
2752 func (c *Controller) SetGoalWithResearchMode(goal string, researchMode GoalResearchMode) {
2753 if c.sessionEngineEnabled() {
2754 if err := c.SetGoalDurable(goal); err != nil {
2755 c.notice("goal: " + err.Error())
2756 }
2757 return
2758 }
2759 resolved, setup := c.resolveGoalText(goal, researchMode)
2760 if setup.notice != "" {
2761 c.notice(setup.notice)
2762 }
2763 var path string
2764 var data []byte
2765 var ok bool
2766 if setup.blockReason != "" {
2767 path, data, ok = c.goals.setLegacyArchiveBlockedWithTaskID(resolved, setup.budgetClass, setup.blockReason, setup.legacyTaskID)
2768 c.replaceLegacyRestore(legacyGoalRestore{taskID: setup.legacyTaskID, epoch: c.goals.continuationToken(), explicit: setup.explicit})
2769 c.notice("legacy research archive resume failed: " + setup.blockReason)
2770 } else {
2771 path, data, ok = c.goals.set(resolved, setup.budgetClass)
2772 c.replaceLegacyRestore(legacyGoalRestore{})
2773 }
2774 c.persistGoalState(path, data, ok)
2775 }
2776
2777 // goalSetSetup is the resolved objective and budget class after archive lookup.
2778 type goalSetSetup struct {
2779 budgetClass string
2780 notice string
2781 blockReason string
2782 legacyTaskID string
2783 explicit bool
2784 }
2785
2786 func (c *Controller) resolveGoalText(goal string, researchMode GoalResearchMode) (string, goalSetSetup) {
2787 setup := goalSetSetup{budgetClass: budgetClassForLegacyMode(goal, researchMode)}
2788 legacy := c.prepareLegacyResearchTask(goal)
2789 if !legacy.explicit {
2790 return goal, setup
2791 }
2792 setup.notice, setup.blockReason, setup.legacyTaskID, setup.explicit = legacy.notice, legacy.blockReason, legacy.taskID, legacy.explicit
2793 if legacy.blockReason != "" {
2794 return goal, setup
2795 }
2796 setup.budgetClass = budgetClassResearch
2797 return legacy.goal, setup
2798 }
2799
2800 // ResumeGoal re-enters a recoverable blocked/stopped Goal without resetting its
2801 // delivery evidence scope or accumulated usage statistics.
2802 func (c *Controller) ResumeGoal() bool {
2803 if c.sessionEngineEnabled() {
2804 current, err := c.goalLifecycleView()
2805 if err != nil || current == nil {
2806 return false
2807 }
2808 _, err = c.applyHostGoalMutation(context.Background(), "resume", func(machine *goaldomain.Machine) (*goaldomain.View, error) {
2809 resumed, resumeErr := machine.Resume(current.Ref(), true)
2810 return &resumed, resumeErr
2811 })
2812 if err != nil {
2813 return false
2814 }
2815 if current.BlockedReason != nil && current.BlockedReason.Code == "resource-budget" && c.goalTokenBudget > 0 {
2816 c.goalResourceMu.Lock()
2817 c.goalTokenLimit += c.goalTokenBudget
2818 c.goalBudgetExtensions++
2819 c.goalResourceMu.Unlock()
2820 }
2821 c.kickGoalDriver()
2822 return true
2823 }
2824 if handled, resumed := c.retryBlockedLegacyGoal(); handled {
2825 return resumed
2826 }
2827 spentBudget := c.goals.runtimeView().StopCause == stopCauseBudgetSpend
2828 path, data, persist, resumed := c.goals.resume()
2829 if !resumed {
2830 return false
2831 }
2832 c.persistGoalState(path, data, persist)
2833 if c.executor != nil {
2834 if spentBudget {
2835 c.executor.ResetTaskBudget()
2836 }
2837 c.executor.RestoreDeliveryCheckpoint(c.goals.deliveryState())
2838 }
2839 return true
2840 }
2841
2842 // PauseGoal suspends a running Goal without losing its Delivery checkpoint or
2843 // runtime history; ResumeGoal restores it. Returns false when no
2844 // running Goal exists.
2845 func (c *Controller) PauseGoal() bool {
2846 if c.sessionEngineEnabled() {
2847 current, err := c.goalLifecycleView()
2848 if err != nil || current == nil || current.Phase != goaldomain.PhaseActive {
2849 return false
2850 }
2851 // Revoke automatic execution before persistence or cancellation can block.
2852 c.disarmGoalLifecycle("user-paused")
2853 _, err = c.applyHostGoalMutation(context.Background(), "pause", func(machine *goaldomain.Machine) (*goaldomain.View, error) {
2854 paused, pauseErr := machine.Pause(current.Ref())
2855 return &paused, pauseErr
2856 })
2857 if err != nil {
2858 return false
2859 }
2860 c.goalDriverMu.Lock()
2861 activeGoalRound := c.goalDriverActive != nil
2862 c.goalDriverMu.Unlock()
2863 if activeGoalRound {
2864 c.Cancel()
2865 }
2866 c.notice(i18n.M.GoalPaused)
2867 return true
2868 }
2869 if !c.goals.active() {
2870 return false
2871 }
2872 path, data, ok := c.goals.pauseFor(stopCauseManual, i18n.M.GoalPausedReason)
2873 c.persistGoalState(path, data, ok)
2874 c.notice(i18n.M.GoalPaused)
2875 return true
2876 }
2877
2878 // GoalRuntime returns the active Goal's usage/runtime summary for frontends.
2879 func (c *Controller) GoalRuntime() GoalRuntimeView {
2880 if c.sessionEngineEnabled() {
2881 view, _ := c.goalLifecycleView()
2882 if view == nil {
2883 return GoalRuntimeView{}
2884 }
2885 limit := 0
2886 if view.MaxGoalRounds != nil {
2887 limit = int(*view.MaxGoalRounds)
2888 }
2889 c.goalResourceMu.Lock()
2890 used, requests, tokenLimit, extensions := c.goalTokensUsed, c.goalRequestsUsed, c.goalTokenLimit, c.goalBudgetExtensions
2891 c.goalResourceMu.Unlock()
2892 return GoalRuntimeView{TurnsUsed: int(view.RoundsStarted), TurnsLimit: limit, TokensUsed: used,
2893 RequestsUsed: requests, TokensLimit: tokenLimit, StopCause: view.StopReason, BudgetExtensions: extensions}
2894 }
2895 return c.goals.runtimeView()
2896 }
2897
2898 func (c *Controller) ClearGoal() {
2899 if c.sessionEngineEnabled() {
2900 c.disarmGoalLifecycle("cleared")
2901 _ = c.SetGoalDurable("")
2902 c.goalDriverMu.Lock()
2903 activeGoalRound := c.goalDriverActive != nil
2904 c.goalDriverMu.Unlock()
2905 if activeGoalRound {
2906 c.Cancel()
2907 }
2908 return
2909 }
2910 c.SetGoal("")
2911 }
2912
2913 func (c *Controller) Goal() string {
2914 if c.sessionEngineEnabled() {
2915 view, _ := c.goalLifecycleView()
2916 if view == nil {
2917 return ""
2918 }
2919 return view.Objective
2920 }
2921 return c.goals.goalText()
2922 }
2923
2924 func (c *Controller) GoalStatus() string {
2925 if c.sessionEngineEnabled() {
2926 view, err := c.goalLifecycleView()
2927 if err != nil || view == nil {
2928 return GoalStatusStopped
2929 }
2930 switch view.Phase {
2931 case goaldomain.PhaseComplete:
2932 return GoalStatusComplete
2933 case goaldomain.PhaseBlocked:
2934 return GoalStatusBlocked
2935 case goaldomain.PhaseActive:
2936 if view.Activation == goaldomain.ActivationArmed {
2937 return GoalStatusRunning
2938 }
2939 }
2940 return GoalStatusStopped
2941 }
2942 return c.goals.statusForDisplay()
2943 }
2944
2945 // Compact runs one compaction pass on the executor's session on demand.
2946 // instructions is optional `/compact <focus>` guidance steering what to keep.
2947 func (c *Controller) Compact(ctx context.Context, instructions string) error {
2948 ctx = c.withAuthentication(ctx)
2949 if err := c.authentication.admissionError(); err != nil {
2950 return err
2951 }
2952 if c.executor == nil {
2953 return nil
2954 }
2955 op, runCtx, err := c.beginMaintenance(ctx, "compact")
2956 if err != nil {
2957 return err
2958 }
2959 return c.executeMaintenance(op, runCtx, func(work context.Context) error {
2960 return c.executor.CompactNow(work, instructions)
2961 })
2962 }
2963
2964 // maybeSessionStart fires the SessionStart hook exactly once per session, lazily
2965 // on the first turn — by then the sink/notify is wired, and a resumed session
2966 // fires it too (its first post-resume turn).
2967 func (c *Controller) maybeSessionStart(ctx context.Context) {
2968 c.hooks.SetSessionID(c.parentSessionID())
2969 c.mu.Lock()
2970 if c.startedOnce {
2971 c.mu.Unlock()
2972 return
2973 }
2974 c.startedOnce = true
2975 c.mu.Unlock()
2976 c.enqueueHookContexts(c.hooks.SessionStart(ctx))
2977 c.extensionSessionEvent(extension.PointSessionStart, dispatch.PhaseStart, c.SessionPath())
2978 }
2979
2980 // NewSession snapshots the current conversation, rotates to a fresh file, and
2981 // resets the executor to a clean session carrying the same base system prompt.
2982 // Session-owned pinned context intentionally starts empty. It ends the old
2983 // session and starts the new one for lifecycle hooks.
2984 func (c *Controller) NewSession() error {
2985 if c.executor == nil {
2986 return nil
2987 }
2988 // Claim the rotation gate for the whole snapshot-then-swap sequence. A bare
2989 // `if c.running` check released before Snapshot() left a window where a turn
2990 // could start during the snapshot and then have its live session replaced by
2991 // the SetSession below. Submit ("/new") and the bot gateway call this
2992 // asynchronously, so the gate is load-bearing, not defensive.
2993 if err := c.beginRotation(); err != nil {
2994 return err
2995 }
2996 defer c.endRotation()
2997 if c.NativeLegacySession() && c.SessionService() != nil {
2998 oldPath := c.SessionPath()
2999 c.flushRecoveryPersistence(oldPath)
3000 if err := c.Snapshot(); err != nil {
3001 return err
3002 }
3003 if err := c.extensionSessionPhase(context.Background(), extension.PointSessionRotate, dispatch.PhaseRotate, oldPath); err != nil {
3004 return err
3005 }
3006 c.hooks.SessionEnd(context.Background(), "new")
3007 c.extensionSessionEvent(extension.PointSessionEnd, dispatch.PhaseEnd, oldPath)
3008 plan := SessionRotationPlan{}
3009 if c.onSessionRotation != nil {
3010 var err error
3011 plan, err = c.onSessionRotation(context.Background(), SessionRotationRequest{SourcePath: c.SessionPath(), Reason: "new"})
3012 if err != nil {
3013 return err
3014 }
3015 }
3016 ref, err := c.bindFreshSessionWithCommit(context.Background(), plan.CreateOptions, plan.Commit)
3017 if err == nil {
3018 c.startExclusiveSession(ref, "new")
3019 }
3020 return err
3021 }
3022 if c.sessionEngineEnabled() {
3023 return c.rotateExclusiveSession(false)
3024 }
3025 // Retire asynchronous recovery writes before Snapshot publishes the final
3026 // old-session checkpoint. Otherwise an earlier write can outlive the path
3027 // rotation (or process teardown) and race cleanup of the old session.
3028 oldPath := c.SessionPath()
3029 c.flushRecoveryPersistence(oldPath)
3030 if err := c.Snapshot(); err != nil {
3031 return err
3032 }
3033 // session.rotate: the session_policy owner rules on the rotation before
3034 // anything is torn down, so its failure (required-class) aborts the
3035 // rotation cleanly. SessionPath is the file being rotated away from; the
3036 // fresh path arrives with the session.start event below.
3037 if err := c.extensionSessionPhase(context.Background(), extension.PointSessionRotate, dispatch.PhaseRotate, oldPath); err != nil {
3038 return err
3039 }
3040 c.hooks.SessionEnd(context.Background(), "clear")
3041 c.extensionSessionEvent(extension.PointSessionEnd, dispatch.PhaseEnd, oldPath)
3042 freshPath := oldPath
3043 if c.sessionDir != "" {
3044 freshPath = agent.NewSessionPath(c.sessionDir, c.label)
3045 }
3046 freshSession := agent.NewSession(c.basePrompt())
3047 commitTransition, err := c.prepareSessionTransition(freshPath, "new", freshSession)
3048 if err != nil {
3049 return fmt.Errorf("bind new session: %w", err)
3050 }
3051 // Hold snapshotMu across the swap so an in-flight save cannot pair the old
3052 // path with the fresh session (or the fresh path with the old session).
3053 c.snapshotMu.Lock()
3054 commitTransition.publish()
3055 c.bindExecutorProjection(c.SessionPath(), false)
3056 if c.guardianSess != nil {
3057 c.guardianSess.Reset()
3058 }
3059 c.ResetPlannerSession()
3060 c.rebindCheckpoints(freshPath)
3061 seedErr := c.seedSessionEventsFromExecutor("session-new")
3062 c.resetRecoveryForNewSession(freshPath)
3063 c.rotateSessionTemp()
3064 c.snapshotMu.Unlock()
3065 // Old session keeps its inbox (paused); the fresh session starts empty.
3066 c.pauseInboxOnRotate()
3067 c.rebindInbox()
3068 // A new session starts with no active goal: without this, a running goal's
3069 // text kept injecting into the fresh session's first turns. The old
3070 // session's goal-state sidecar was persisted before the rotation and stays
3071 // intact, so resuming it restores its goal; the cleared state below lands
3072 // on the NEW path (rebindCheckpoints just moved it).
3073 c.ClearGoal()
3074 c.mu.Lock()
3075 c.startedOnce = true // NewSession fires SessionStart itself; don't re-fire on the next turn
3076 c.mu.Unlock()
3077 c.hooks.SetSessionID(c.parentSessionID())
3078 c.enqueueHookContexts(c.hooks.SessionStart(context.Background(), "clear"))
3079 c.extensionSessionEvent(extension.PointSessionStart, dispatch.PhaseStart, c.SessionPath())
3080 c.clearSessionWriteAccess()
3081 if seedErr != nil {
3082 return fmt.Errorf("seed new session events: %w", seedErr)
3083 }
3084 return nil
3085 }
3086
3087 func (c *Controller) hasUnfinishedSessionJobs(sessionPath string) bool {
3088 if c.jobs == nil {
3089 return false
3090 }
3091 return c.jobs.HasUnfinishedForSession(agent.BranchID(sessionPath))
3092 }
3093
3094 func removeSessionArtifacts(path string) error {
3095 if path == "" {
3096 return nil
3097 }
3098 if err := jobs.RemoveArtifacts(path); err != nil {
3099 return err
3100 }
3101 remove := []string{path}
3102 // Sidecars include the event log — the authoritative transcript. Leaving
3103 // it behind would both leak the cleared conversation and let LoadSession
3104 // resurrect it on the recycled path. The guardian transcript saves through
3105 // the same session layer, so its sidecars are swept too.
3106 remove = append(remove, store.SessionSidecarFiles(path)...)
3107 remove = append(remove, guardian.PathFor(path), guardian.CursorPathFor(path))
3108 remove = append(remove, store.SessionSidecarFiles(guardian.PathFor(path))...)
3109 for _, p := range remove {
3110 if p == "" {
3111 continue
3112 }
3113 if err := os.Remove(p); err != nil && !os.IsNotExist(err) {
3114 return err
3115 }
3116 }
3117 if err := sessioninbox.RemoveDir(path); err != nil && !os.IsNotExist(err) {
3118 return err
3119 }
3120 if dir := ckptDir(path); dir != "" {
3121 if err := os.RemoveAll(dir); err != nil && !os.IsNotExist(err) {
3122 return err
3123 }
3124 }
3125 if err := agent.DeleteSubagentsByParent(filepath.Dir(path), agent.BranchID(path)); err != nil {
3126 return err
3127 }
3128 if err := agent.ClearCleanupPending(path); err != nil {
3129 return err
3130 }
3131 return nil
3132 }
3133
3134 // RemoveSessionArtifacts removes a transcript and every durable artifact owned
3135 // by it. Remote runtimes use this when a newly-created fork fails before it can
3136 // be registered as a live session.
3137 func RemoveSessionArtifacts(path string) error {
3138 return removeSessionArtifacts(path)
3139 }
3140
3141 // ReconcileCleanupPending retries physical cleanup for logically removed
3142 // sessions that were left behind by a previous process.
3143 func ReconcileCleanupPending(dir string) error {
3144 return agent.ReconcileCleanupPending(dir, func(item agent.CleanupPendingInfo) error {
3145 return removeSessionArtifacts(item.SessionPath)
3146 })
3147 }
3148
3149 // RewindScope selects what a Rewind restores.
3150 type RewindScope int
3151
3152 const (
3153 RewindCode RewindScope = iota // files only
3154 RewindConversation // message log only
3155 RewindBoth // both
3156 )
3157
3158 // Checkpoints lists the session's rewind points (one per user turn), oldest first.
3159 //
3160 // Each Meta.Prompt is reduced to what the user typed. A checkpoint opens with
3161 // the composed turn, so the stored prompt can carry the plan-mode marker and
3162 // transient blocks; every consumer of this list is a label (the rewind picker,
3163 // the desktop change list, the workbench projection) and the picker also
3164 // restores the prompt into the composer, so composed text must not reach them.
3165 // Stripping on read rather than only on write keeps checkpoints already on disk
3166 // readable — they were recorded composed.
3167 func (c *Controller) Checkpoints() []checkpoint.Meta {
3168 metas := c.checkpoints.list()
3169 for i := range metas {
3170 metas[i].Prompt = StripComposePrefixes(metas[i].Prompt)
3171 }
3172 return metas
3173 }
3174
3175 func (c *Controller) CheckpointFileState(path string) (checkpoint.FileState, bool) {
3176 return c.checkpoints.fileState(path)
3177 }
3178
3179 func (c *Controller) CheckpointTurnsByMessageIndex() map[int]int {
3180 return c.checkpoints.turnsByMessageIndex()
3181 }
3182
3183 // rewindFail emits the error as a Warn notice (so a frontend that swallows the
3184 // returned error — e.g. the desktop bridge's .catch — still shows the user why
3185 // the rewind did nothing) and returns it.
3186 func (c *Controller) rewindFail(err error) error {
3187 c.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Text: err.Error()})
3188 return err
3189 }
3190
3191 // Rewind is implemented in rewind.go (transactional conversation+file restore).
3192
3193 // SummarizeFrom and SummarizeUpTo preserve the historical turn-index API while
3194 // changing only the model-visible context projection. The canonical transcript
3195 // and checkpoint boundaries remain available for rewind, undo, and fork.
3196 func (c *Controller) SummarizeFrom(ctx context.Context, turn int) error {
3197 return c.summarizeAt(ctx, turn, true)
3198 }
3199
3200 func (c *Controller) SummarizeUpTo(ctx context.Context, turn int) error {
3201 return c.summarizeAt(ctx, turn, false)
3202 }
3203
3204 func (c *Controller) summarizeAt(ctx context.Context, turn int, from bool) error {
3205 ctx = c.withAuthentication(ctx)
3206 if err := c.authentication.admissionError(); err != nil {
3207 return err
3208 }
3209 if c.executor == nil {
3210 return c.rewindFail(fmt.Errorf("checkpoints unavailable"))
3211 }
3212 kind := "summarize_up_to"
3213 if from {
3214 kind = "summarize_from"
3215 }
3216 op, runCtx, err := c.beginMaintenance(ctx, kind)
3217 if err != nil {
3218 return c.rewindFail(err)
3219 }
3220 boundary, hasBound := c.checkpoints.boundary(turn)
3221 if !hasBound {
3222 err = fmt.Errorf("summarize unavailable for turn %d (resumed session)", turn)
3223 _ = c.executeMaintenance(op, runCtx, func(context.Context) error { return err })
3224 return c.rewindFail(err)
3225 }
3226 err = c.executeMaintenance(op, runCtx, func(work context.Context) error {
3227 if from {
3228 return c.executor.SummarizeFrom(work, boundary)
3229 }
3230 return c.executor.SummarizeUpTo(work, boundary)
3231 })
3232 if err != nil {
3233 return c.rewindFail(err)
3234 }
3235 return nil
3236 }
3237
3238 // Resume seeds the session from a loaded transcript and pins the active file to
3239 // its path so auto-save keeps appending there.
3240 //
3241 // When the controller already has a different non-empty session path, Resume
3242 // rotates the private temporary generation so the loaded conversation cannot
3243 // see the previous session's temporary files. Same-path Resume (hot rebuild
3244 // migration via AdoptHistory) keeps the generation.
3245 func (c *Controller) Resume(s *agent.Session, path string) {
3246 if c.NativeLegacySession() && strings.TrimSpace(path) != "" {
3247 if service := c.SessionService(); service != nil {
3248 ref, found, err := service.ExistingCanonicalForLegacy(path, s)
3249 if err != nil {
3250 c.failTurnEventLedger(err)
3251 return
3252 }
3253 if found {
3254 if _, err := c.OpenSession(context.Background(), ref); err != nil {
3255 c.failTurnEventLedger(err)
3256 }
3257 return
3258 }
3259 }
3260 }
3261 if c.sessionEngineEnabled() {
3262 if _, runtime, _ := c.v3Binding(); runtime != nil && strings.TrimSpace(path) == "" {
3263 c.restoreExecutorFromSessionEvents()
3264 return
3265 }
3266 if strings.TrimSpace(path) == "" {
3267 if _, err := c.BindFreshSession(context.Background(), ""); err != nil {
3268 c.failTurnEventLedger(err)
3269 }
3270 return
3271 }
3272 if _, statErr := os.Stat(path); os.IsNotExist(statErr) {
3273 if _, err := c.BindFreshSession(context.Background(), ""); err != nil {
3274 c.failTurnEventLedger(err)
3275 return
3276 }
3277 if s != nil && len(s.Snapshot()) > 0 {
3278 if err := c.replaceSessionEventProjection(context.Background(), "fresh-resume", s.Snapshot()); err != nil {
3279 c.failTurnEventLedger(err)
3280 } else {
3281 c.restoreExecutorFromSessionEvents()
3282 }
3283 }
3284 return
3285 }
3286 if err := c.ResumeNativeSession(s, path); err != nil {
3287 slog.Warn("controller: resume native legacy session", "path", path, "err", err)
3288 c.failTurnEventLedger(err)
3289 }
3290 return
3291 }
3292 // See snapshotMu: the swap must not interleave with an in-flight save.
3293 // recoverInterruptedTurn and maybeColdResumePrune snapshot on their own,
3294 // so they stay outside the locked section (snapshotMu is not reentrant).
3295 prevPath := c.SessionPath()
3296 c.snapshotMu.Lock()
3297 if c.executor != nil {
3298 c.executor.SetSession(s)
3299 }
3300 c.mu.Lock()
3301 c.sessionPath = path
3302 c.guardianPath = guardian.PathFor(path)
3303 c.mu.Unlock()
3304 c.bindExecutorProjection(path, true)
3305 c.ResetPlannerSession()
3306 c.setActiveJobSession(path)
3307 c.rebindCheckpoints(path)
3308 if err := c.seedSessionEventsFromExecutor("legacy-resume"); err != nil {
3309 slog.Warn("controller: import resumed transcript into v3", "err", err)
3310 }
3311 if err := c.importLegacyResumeOverPlaceholder(s); err != nil {
3312 slog.Warn("controller: replace placeholder v3 history from legacy resume", "err", err)
3313 }
3314 if err := c.adoptResumeSystemPrompt(s); err != nil {
3315 slog.Warn("controller: record refreshed system prompt in v3", "err", err)
3316 }
3317 // A host-managed replacement is deliberately unbound until its final lease
3318 // handoff. Keep the exact carried transcript on that private candidate; the
3319 // successful BindSessionWriteAuthority call publishes it to the shared v3
3320 // projection. Restoring the currently-active projection here would erase the
3321 // carried tail before the candidate had any chance to become authoritative.
3322 if c.sessionEventCommitAllowed() {
3323 c.restoreExecutorFromSessionEvents()
3324 }
3325 migPath, migData, migrated, legacy := c.goals.restoreFromState(path)
3326 if !c.restorePendingLegacyGoal(legacy) && migrated {
3327 c.persistGoalState(migPath, migData, true)
3328 }
3329 if c.executor != nil {
3330 c.executor.RestoreDeliveryCheckpoint(c.goals.deliveryState())
3331 }
3332 c.loadGuardianSession()
3333 c.loadRecoveryState(path)
3334 if shouldRotateSessionTempOnResume(prevPath, path) {
3335 c.rotateSessionTemp()
3336 }
3337 c.snapshotMu.Unlock()
3338 c.rebindInbox()
3339 c.recoverCheckpointTransactions()
3340 c.recoverInterruptedTurn(path)
3341 c.maybeColdResumePrune(path)
3342 // session.load: Resume has no failure channel, so the session_policy
3343 // strategy is advisory this stage — a required-class failure is surfaced
3344 // as a warning and the load stands. The event still carries the final
3345 // (possibly owner-adjusted) phase payload.
3346 if err := c.extensionSessionPhase(context.Background(), extension.PointSessionLoad, dispatch.PhaseLoad, path); err != nil {
3347 c.extensionWarn("session policy failed at session.load", err)
3348 }
3349 }
3350
3351 func shouldRotateSessionTempOnResume(prevPath, nextPath string) bool {
3352 prevPath = strings.TrimSpace(prevPath)
3353 nextPath = strings.TrimSpace(nextPath)
3354 if prevPath == "" || nextPath == "" {
3355 return false
3356 }
3357 return filepath.Clean(prevPath) != filepath.Clean(nextPath)
3358 }
3359
3360 func (c *Controller) loadGuardianSession() {
3361 if c.guardianSess == nil {
3362 return
3363 }
3364 c.guardianSess.Reset()
3365 path := c.guardianPath
3366 if path == "" {
3367 return
3368 }
3369 if err := c.guardianSess.Load(path); err != nil && !os.IsNotExist(err) {
3370 slog.Warn("controller: load guardian session", "err", err)
3371 }
3372 }
3373
3374 // ResetPlannerSession clears the planner's conversation history so the next
3375 // plan starts fresh. In dual-model (Plan+Execute) mode, this prevents stale
3376 // planner output from a previous session or tab from contaminating the current
3377 // executor's handoff. Safe to call on a single-model controller (no-op).
3378 func (c *Controller) ResetPlannerSession() {
3379 runner, ok := c.runner.(plannerSessionResetter)
3380 if ok {
3381 runner.ResetPlannerSession()
3382 }
3383 }
3384
3385 // cacheColdAfter resolves how long the active provider keeps a prompt prefix
3386 // cached. A session idle longer than this resumes against a cold cache, so a
3387 // history rewrite at that moment costs no extra cache misses — it only shrinks
3388 // the full-price first request. The TTL is vendor-aware: DeepSeek/unknown
3389 // 24h (legacy default deliberately preserved), DashScope 5m, Anthropic 5m.
3390 // Users can override per-provider
3391 // with cache_ttl_minutes in config.toml.
3392 func (c *Controller) cacheColdAfter() time.Duration {
3393 if c.testCacheColdAfter != 0 {
3394 if c.testCacheColdAfter == -1 {
3395 return 0
3396 }
3397 return c.testCacheColdAfter
3398 }
3399 // 查询路径只读:LoadForRootReadOnly 不触发配置迁移写盘(评审 #7168
3400 // 第 4 点);失败时保守回退 24h(DeepSeek/未知 vendor 默认),避免
3401 // 把 cache TTL 过期误当成历史改写信号(resume 只记录 warm/cold/unknown)。
3402 cfg, err := config.LoadForRootReadOnly(c.workspaceRoot)
3403 if err != nil {
3404 return 24 * time.Hour
3405 }
3406 ref := c.selection.ref
3407 if ref == "" {
3408 ref = cfg.DefaultModel
3409 }
3410 entry, ok := cfg.ResolveModel(ref)
3411 if !ok {
3412 return 24 * time.Hour
3413 }
3414 return entry.EffectiveCacheTTL()
3415 }
3416
3417 // Snapshot writes the executor's conversation to the active session file. No-op
3418 // when the executor is absent or the session has never been used (no user
3419 // interaction). Returns errNoSessionPath when there IS content but no resolved
3420 // path, so a misconfigured deployment surfaces instead of dropping data.
3421 // Called after every turn so a crash loses at most one in-flight prompt.
3422 func (c *Controller) Snapshot() error {
3423 return c.snapshot(false, false, false)
3424 }
3425
3426 // SnapshotForShutdown performs the final session snapshot. Only when the
3427 // compatibility file lock stays held for the full bounded wait does it fall
3428 // back: a schema-2 log takes the unsaved tail without the lock, a schema-1
3429 // session gets a distinct recovery branch. Other snapshot errors retain their
3430 // normal behavior and remain visible to the caller.
3431 func (c *Controller) SnapshotForShutdown() error {
3432 return c.snapshot(false, false, true)
3433 }
3434
3435 // SnapshotActivity writes the active conversation and marks the session as
3436 // recently active. Use it only after a real user/model turn changes the
3437 // transcript; switch/close snapshots should call Snapshot so they do not reorder
3438 // recent-session pickers.
3439 func (c *Controller) SnapshotActivity() error {
3440 return c.snapshot(true, false, false)
3441 }
3442
3443 // SnapshotRewrite persists an intentional history rewrite, such as rewind or
3444 // manual compaction. Ordinary autosave paths should use Snapshot so stale
3445 // controllers cannot overwrite a newer transcript.
3446 func (c *Controller) SnapshotRewrite() error {
3447 return c.snapshot(false, true, false)
3448 }
3449
3450 func (c *Controller) snapshot(markActivity, forceRewrite, shutdownRecovery bool) error {
3451 _, err := c.snapshotWithDurability(markActivity, forceRewrite, shutdownRecovery)
3452 return err
3453 }
3454
3455 // snapshotWithDurability reports whether the canonical transcript reached disk
3456 // even when a later sidecar update failed. Callers that guard a crash marker
3457 // need this distinction: a metadata error must not make a complete transcript
3458 // look like an in-memory-only turn.
3459 func (c *Controller) snapshotWithDurability(markActivity, forceRewrite, shutdownRecovery bool) (bool, error) {
3460 c.snapshotMu.Lock()
3461 defer c.snapshotMu.Unlock()
3462
3463 c.mu.Lock()
3464 path := c.sessionPath
3465 modelRef := c.selection.ref
3466 c.mu.Unlock()
3467 if c.executor == nil {
3468 return false, nil
3469 }
3470 s := c.executor.Session()
3471 if !s.HasContent() {
3472 // Nothing to persist yet (e.g. a fresh session with only a system
3473 // prompt) — staying quiet here is correct, not a data-loss path.
3474 return false, nil
3475 }
3476 if !s.HasSystemMessage() {
3477 // The session has user/assistant/tool messages but no leading system
3478 // prompt. Persisting it would create a session file that, when
3479 // reloaded, has no agent-identity contract — the model falls back to
3480 // its training-data defaults, giving wrong answers to identity
3481 // queries ("who are you?"). Log the anomaly so the root cause
3482 // (typically an empty sysPrompt reaching NewSession) can be
3483 // diagnosed, then refuse to write a corrupted transcript.
3484 slog.Warn("controller: refusing to snapshot session with content but no system message",
3485 "label", c.Label(), "session_dir", c.SessionDir(), "message_count", len(s.Snapshot()))
3486 return false, nil
3487 }
3488 if c.sessionEngineEnabled() {
3489 if _, runtime, _ := c.v3Binding(); runtime == nil || c.sessionEventStore() == nil {
3490 return false, errors.New("exclusive v3 controller has no bound session runtime")
3491 }
3492 // Export/switch/shutdown callers ask for durability explicitly. The
3493 // transcript, Goal, Plan, Todo and runtime facts already live in the one
3494 // event sequence; no legacy transcript or business sidecar is refreshed.
3495 if _, err := c.flushSessionEvents(context.Background()); err != nil {
3496 return false, err
3497 }
3498 return true, nil
3499 }
3500 if path == "" {
3501 // There IS content but nowhere to write it: this silently dropped whole
3502 // bot conversations (#4414). Surface it loudly instead of returning nil
3503 // so the missing session path can be diagnosed and fixed at the source.
3504 slog.Warn("controller: session has content but no session path; conversation will not be persisted",
3505 "label", c.Label(), "session_dir", c.SessionDir())
3506 return false, errNoSessionPath
3507 }
3508 // Snapshot is an explicit persistence request (session switch, export or
3509 // shutdown), so it is also a v3 semantic checkpoint. Live events normally
3510 // remain eligible for the 200 ms write-behind window; callers asking for a
3511 // snapshot must receive a durable receipt before the rebuildable legacy
3512 // transcript cache is refreshed.
3513 if _, err := c.flushSessionEvents(context.Background()); err != nil {
3514 return false, err
3515 }
3516 // session.save: the session_policy owner rules on the impending save; a
3517 // failure (required-class) vetoes the write. The event goes out after a
3518 // successful save carrying the final payload. The early no-content and
3519 // no-path returns above are not saves and stay unobserved. Conflict
3520 // recovery below may rewrite the path; the phase payload reports the path
3521 // the save targeted.
3522 savePayload, strategyErr := c.extensionSessionStrategy(context.Background(), extension.PointSessionSave, dispatch.PhaseSave, path)
3523 if strategyErr != nil {
3524 return false, strategyErr
3525 }
3526 err, forceRewrite := persistSessionSnapshot(s, path, forceRewrite)
3527 if authoritySaveError(err) {
3528 // Missing/stale authority must not enter diverged/recovery. Frontends
3529 // rebind the lease or surface the typed error.
3530 return false, err
3531 }
3532 if err != nil {
3533 if shutdownRecovery && errors.Is(err, agent.ErrSessionFileLockHeld) {
3534 recoveredPath, recoverErr := c.recoverShutdownSave(s, path, err, forceRewrite)
3535 if recoverErr != nil {
3536 return false, recoverErr
3537 }
3538 path = recoveredPath
3539 s = c.executor.Session()
3540 err = nil
3541 }
3542 }
3543 if err != nil {
3544 if errors.Is(err, agent.ErrSessionExternallyRemoved) {
3545 recoveredPath, recoverErr := c.recoverExternallyRemovedSession(path, err)
3546 if recoverErr != nil {
3547 return false, recoverErr
3548 }
3549 path = recoveredPath
3550 s = c.executor.Session()
3551 err = nil
3552 }
3553 }
3554 if err != nil {
3555 if !errors.Is(err, agent.ErrSessionSnapshotConflict) {
3556 return false, err
3557 }
3558 recoveredPath, outcome, recoverErr := c.recoverSnapshotConflict(path, err, forceRewrite)
3559 if recoverErr != nil {
3560 if shutdownRecovery && errors.Is(recoverErr, agent.ErrSessionFileLockHeld) {
3561 recoveredPath, recoverErr = c.recoverShutdownSave(s, path, recoverErr, forceRewrite)
3562 if recoverErr != nil {
3563 return false, recoverErr
3564 }
3565 path = recoveredPath
3566 s = c.executor.Session()
3567 } else {
3568 return false, recoverErr
3569 }
3570 } else {
3571 if outcome == conflictDropped {
3572 return false, nil
3573 }
3574 // Whatever recovery did — adopted the disk transcript, isolated the
3575 // depth-capped copy, or forked — the rewrite baseline lives on
3576 // the session object and was advanced by the save that succeeded, so
3577 // there is nothing to re-anchor here.
3578 path = recoveredPath
3579 s = c.executor.Session()
3580 }
3581 }
3582 // Persist guardian session so the prefix cache stays warm after restart.
3583 if gp := c.guardianPath; c.guardianSess != nil && gp != "" {
3584 if gerr := c.guardianSess.Save(gp); gerr != nil {
3585 slog.Warn("controller: guardian snapshot", "err", gerr)
3586 }
3587 }
3588 c.emitHeadEvents()
3589 transcriptDurable := true
3590 // Persist recovery gate state so unresolved checkpoints survive restart.
3591 c.saveRecoveryState(path)
3592 // Record the listing-only sidecar fields (model, preview, user-turn count)
3593 // straight from the in-memory conversation, so the sidebar and resume picker
3594 // never have to decode the whole .jsonl just to show them. markActivity bumps
3595 // UpdatedAt exactly like the previous TouchBranchMeta did; false preserves it
3596 // like SetBranchModelPreserveUpdated. The single write subsumes the old
3597 // EnsureBranchMeta / SetBranchModel / TouchBranchMeta sequence.
3598 preview, turns := agent.SessionPreviewFromMessages(s.Snapshot())
3599 if err := updateSessionModelProjection(s, path, modelRef, c.selection.identity, preview, turns, markActivity); err != nil && !listingDeferredAfterUnlockedAppend(s, path, err) {
3600 return transcriptDurable, err
3601 }
3602 c.extensionSessionPayloadEvent(extension.PointSessionSave, savePayload)
3603 return transcriptDurable, nil
3604 }
3605
3606 func updateSessionModelProjection(s *agent.Session, path, modelRef, identity, preview string, turns int, markActivity bool) error {
3607 persisted, ok := s.PersistedState(path)
3608 if !ok {
3609 return fmt.Errorf("session persistence baseline missing after save")
3610 }
3611 var err error
3612 if s.WriteAuthorityRequired() {
3613 _, err = agent.UpdateOwnedSessionListingProjectionIfCurrent(path, modelRef, identity, preview, turns, markActivity, persisted, s.WriteAuthority())
3614 } else {
3615 _, err = agent.UpdateSessionListingProjectionIfCurrent(path, modelRef, identity, preview, turns, markActivity, persisted)
3616 }
3617 return err
3618 }
3619
3620 func (c *Controller) recoverExternallyRemovedSession(path string, saveErr error) (string, error) {
3621 if c.executor == nil || strings.TrimSpace(path) == "" {
3622 return "", saveErr
3623 }
3624 const reason = "session removed while open"
3625 req := SessionRecoveryRequest{OriginalPath: path, Reason: reason, Mode: "external-removal"}
3626 meta := agent.BranchMeta{}
3627 if c.sessionRecoveryMeta != nil {
3628 meta = c.sessionRecoveryMeta(req)
3629 }
3630 info, err := c.executor.Session().SaveConflictRecoveryBranch(agent.RecoveryBranchOptions{
3631 OriginalPath: path,
3632 Reason: reason,
3633 BranchMeta: meta,
3634 })
3635 if err != nil {
3636 return "", fmt.Errorf("preserve externally removed session: %w", err)
3637 }
3638 if err := c.commitRecoveredSession(path, reason, info); err != nil {
3639 return "", err
3640 }
3641 appendSnapshotConflictDiagnostic(path, "external-removal", "moved_to_stable_recovery", saveErr, info.Path, info.Existing)
3642 slog.Warn("controller: active session was removed externally; moved runtime to stable recovery path",
3643 "path", path, "recovery", info.Path, "existing", info.Existing)
3644 c.sink.Emit(sessionRecoveryNotice(event.NoticeCodeSessionRecoveryForked,
3645 "the open session file was removed outside Reasonix; your active conversation was preserved as one recovery copy"))
3646 return info.Path, nil
3647 }
3648
3649 // snapshotConflictLogAttrs flattens a snapshot-conflict error into slog attrs.
3650 // Field reports of #6069-class "session changed on disk" spam are only
3651 // diagnosable when the logs say which trigger fired and what the revision
3652 // ledger looked like, so every recoverSnapshotConflict outcome logs these.
3653 func snapshotConflictLogAttrs(saveErr error, path, mode string) []any {
3654 attrs := []any{"path", path, "mode", mode}
3655 var conflict *agent.SessionSnapshotConflictError
3656 if errors.As(saveErr, &conflict) && conflict != nil {
3657 attrs = append(attrs,
3658 "kind", string(conflict.Kind),
3659 "disk_messages", conflict.ExistingMessages,
3660 "snapshot_messages", conflict.SnapshotMessages,
3661 "base_revision", conflict.BaseRevision,
3662 "disk_revision", conflict.DiskRevision,
3663 )
3664 }
3665 return attrs
3666 }
3667
3668 // conflictOutcome is recoverSnapshotConflict's declared result. Callers act
3669 // on it directly instead of re-deriving what happened from path or session
3670 // pointer comparisons — the misclassification that broke the depth-cap
3671 // rewrite baseline (#6120) hid in exactly that inference.
3672 type conflictOutcome int
3673
3674 const (
3675 // conflictDropped: nothing was recovered and the disk transcript could
3676 // not be adopted; this snapshot was deliberately dropped.
3677 conflictDropped conflictOutcome = iota
3678 // conflictAdoptedDisk: the executor session object was replaced by the
3679 // newer disk transcript; adoptDiskSession already reset its baselines.
3680 conflictAdoptedDisk
3681 // conflictForkedBranch: the same in-memory session moved to a freshly
3682 // forked recovery branch path.
3683 conflictForkedBranch
3684 )
3685
3686 func sessionRecoveryNotice(code, text string) event.Event {
3687 return event.Event{
3688 Kind: event.Notice,
3689 Level: event.LevelWarn,
3690 Audience: event.NoticeAudienceOperator,
3691 Code: code,
3692 Text: text,
3693 }
3694 }
3695
3696 func (c *Controller) recoverSnapshotConflict(path string, saveErr error, forceRewrite bool) (string, conflictOutcome, error) {
3697 if c.executor == nil || strings.TrimSpace(path) == "" {
3698 return "", conflictDropped, saveErr
3699 }
3700 mode := "snapshot"
3701 if forceRewrite {
3702 mode = "rewrite"
3703 }
3704 logAttrs := snapshotConflictLogAttrs(saveErr, path, mode)
3705 if kind, ok := agent.SnapshotConflictKind(saveErr); ok && kind == agent.SessionSnapshotConflictStalePrefix {
3706 if c.adoptDiskSession(path) {
3707 appendSnapshotConflictDiagnostic(path, mode, "adopted_newer_disk_transcript", saveErr, "", false)
3708 slog.Warn("controller: snapshot conflict; adopted newer disk transcript", logAttrs...)
3709 c.sink.Emit(sessionRecoveryNotice(event.NoticeCodeSessionRecoveryAdopted,
3710 "session changed on disk; adopted the newer transcript"))
3711 return path, conflictAdoptedDisk, nil
3712 }
3713 }
3714 reason := "snapshot conflict"
3715 if forceRewrite {
3716 reason = "rewrite conflict"
3717 }
3718 req := SessionRecoveryRequest{OriginalPath: path, Reason: reason, Mode: mode}
3719 meta := agent.BranchMeta{}
3720 if c.sessionRecoveryMeta != nil {
3721 meta = c.sessionRecoveryMeta(req)
3722 }
3723 baseRevision, diskRevision := snapshotConflictRevisions(saveErr)
3724 info, err := c.executor.Session().SaveRecoveryBranch(agent.RecoveryBranchOptions{
3725 OriginalPath: path,
3726 Reason: reason,
3727 BranchMeta: meta,
3728 BaseRevision: baseRevision,
3729 DiskRevision: diskRevision,
3730 })
3731 if err != nil {
3732 if errors.Is(err, agent.ErrSessionRecoveryNotNeeded) {
3733 if c.adoptDiskSession(path) {
3734 appendSnapshotConflictDiagnostic(path, mode, "recovery_not_needed_adopted_disk_transcript", saveErr, "", false)
3735 slog.Warn("controller: snapshot conflict; recovery not needed, adopted disk transcript", logAttrs...)
3736 c.sink.Emit(sessionRecoveryNotice(event.NoticeCodeSessionRecoveryAdoptedCovered,
3737 "session changed on disk; adopted the newer transcript (local changes already covered)"))
3738 return path, conflictAdoptedDisk, nil
3739 }
3740 // Nothing was recovered AND the disk transcript could not be
3741 // adopted: the snapshot is silently dropped. Leave a trace so
3742 // "my last turns vanished" reports can be tied to this path.
3743 appendSnapshotConflictDiagnostic(path, mode, "recovery_not_needed_adopt_failed", saveErr, "", false)
3744 slog.Warn("controller: snapshot conflict; recovery not needed but disk transcript could not be adopted", logAttrs...)
3745 return "", conflictDropped, nil
3746 }
3747 return "", conflictDropped, fmt.Errorf("recover stale session snapshot: %w", err)
3748 }
3749 if err := c.commitRecoveredSession(path, reason, info); err != nil {
3750 return "", conflictDropped, err
3751 }
3752 appendSnapshotConflictDiagnostic(path, mode, "forked_recovery_branch", saveErr, info.Path, info.Existing)
3753 slog.Warn("controller: snapshot conflict; forked recovery branch",
3754 append(logAttrs, "recovery", info.Path, "existing", info.Existing)...)
3755 c.sink.Emit(sessionRecoveryNotice(event.NoticeCodeSessionRecoveryForked,
3756 "session changed on disk; unsaved local transcript was saved as a conflict copy"))
3757 return info.Path, conflictForkedBranch, nil
3758 }
3759
3760 func (c *Controller) recoverShutdownSnapshot(path string, saveErr error) (string, error) {
3761 if c.executor == nil || strings.TrimSpace(path) == "" {
3762 return "", saveErr
3763 }
3764 const reason = "shutdown session file lock timeout"
3765 req := SessionRecoveryRequest{OriginalPath: path, Reason: reason, Mode: "shutdown"}
3766 meta := agent.BranchMeta{}
3767 if c.sessionRecoveryMeta != nil {
3768 meta = c.sessionRecoveryMeta(req)
3769 }
3770 info, err := c.executor.Session().SaveShutdownRecoveryBranch(agent.RecoveryBranchOptions{
3771 OriginalPath: path,
3772 Reason: reason,
3773 BranchMeta: meta,
3774 })
3775 if err != nil {
3776 return "", fmt.Errorf("save shutdown recovery branch: %w", err)
3777 }
3778 if err := c.commitRecoveredSession(path, reason, info); err != nil {
3779 return "", err
3780 }
3781 appendSnapshotConflictDiagnostic(path, "shutdown", "forked_file_lock_recovery", saveErr, info.Path, info.Existing)
3782 slog.Warn("controller: shutdown snapshot lock timed out; forked recovery branch",
3783 "path", path, "recovery", info.Path, "existing", info.Existing)
3784 c.sink.Emit(sessionRecoveryNotice(event.NoticeCodeSessionShutdownRecoveryForked,
3785 "session file stayed busy during shutdown; unsaved transcript was saved as a recovery copy"))
3786 return info.Path, nil
3787 }
3788
3789 func (c *Controller) commitRecoveredSession(originalPath, reason string, info agent.RecoveryBranchInfo) error {
3790 commit := &sessionRecoveryCommit{}
3791 recoveryInfo := SessionRecoveryInfo{
3792 OriginalPath: originalPath,
3793 RecoveryPath: info.Path,
3794 Existing: info.Existing,
3795 Reason: reason,
3796 BaseRevision: info.Meta.BaseRevision,
3797 DiskRevision: info.Meta.DiskRevision,
3798 Meta: info.Meta,
3799 commit: commit,
3800 }
3801 if onSessionRecovered := c.sessionRecoveredHandler(); onSessionRecovered != nil {
3802 if err := onSessionRecovered(recoveryInfo); err != nil {
3803 return fmt.Errorf("commit recovered session: %w", err)
3804 }
3805 }
3806 c.mu.Lock()
3807 c.sessionPath = info.Path
3808 c.guardianPath = guardian.PathFor(info.Path)
3809 c.mu.Unlock()
3810 // Recovery branch is a new lineage path. Load an inherited projection
3811 // sidecar when present so the model view stays compressed across the fork.
3812 c.bindExecutorProjection(info.Path, true)
3813 c.setActiveJobSession(info.Path)
3814 c.rebindCheckpoints(info.Path)
3815 c.transplantInFlightTurnMarker(originalPath, info.Path)
3816 commit.publish()
3817 return nil
3818 }
3819
3820 func (c *Controller) adoptDiskSession(path string) bool {
3821 loaded, err := agent.LoadSession(path)
3822 if err != nil || loaded == nil {
3823 return false
3824 }
3825 c.executor.SetSession(loaded)
3826 c.bindExecutorProjection(path, true)
3827 c.ResetPlannerSession()
3828 c.rebindCheckpoints(path)
3829 c.setActiveJobSession(path)
3830 return true
3831 }
3832
3833 func (c *Controller) messageCount() int {
3834 if c.executor == nil {
3835 return 0
3836 }
3837 return c.executor.Session().Len()
3838 }
3839
3840 // stripTurnMessagesAfter truncates the executor's session to keep only messages
3841 // before the given index, discarding an incomplete synthetic turn (the synthetic
3842 // user prompt plus every assistant/tool message that followed).
3843 func (c *Controller) stripTurnMessagesAfter(idx int) {
3844 if c.executor == nil {
3845 return
3846 }
3847 msgs := c.executor.Session().Snapshot()
3848 if len(msgs) <= idx {
3849 // Compaction may have removed the entire synthetic workset. The
3850 // explicit turn identities still need retraction from display history.
3851 c.replaceSessionAfterCancel(msgs)
3852 return
3853 }
3854 c.replaceSessionAfterCancel(msgs[:idx])
3855 }
3856
3857 // stripInterruptedSyntheticTurnMessagesAfter relocates a synthetic turn after
3858 // an in-turn compaction has rewritten the pre-turn message index, then drops
3859 // that whole controller-created turn.
3860 func (c *Controller) stripInterruptedSyntheticTurnMessagesAfter(idx int) {
3861 if c.executor == nil {
3862 return
3863 }
3864 msgs := c.executor.Session().Snapshot()
3865 startedAt := c.inFlightTurnStartedAt()
3866 if start, ok := resolveInterruptedTurnStart(msgs, idx, false, startedAt, provider.Message{}); ok {
3867 idx = start
3868 }
3869 c.stripTurnMessagesAfter(idx)
3870 }
3871
3872 // stripCancelledVisibleTurnMessagesAfterWithFallback preserves the real user
3873 // prompt and fully paired tool rounds from a cancelled visible turn. Unsafe
3874 // assistant/tool fragments are retained as provider-excluded display history.
3875 // It also covers coordinator
3876 // cancellation before the executor has appended the visible user message. The
3877 // orchestrator owns that input, so it supplies the exact message rather than
3878 // letting cancellation infer the current turn from older transcript history.
3879 func (c *Controller) stripCancelledVisibleTurnMessagesAfterWithFallback(idx int, fallback provider.Message) {
3880 c.stripCancelledVisibleTurnMessagesAfterWithFallbackAt(idx, fallback, c.inFlightTurnStartedAt())
3881 }
3882
3883 func (c *Controller) stripCancelledVisibleTurnMessagesAfterWithFallbackAt(idx int, fallback provider.Message, startedAt time.Time) {
3884 if c.executor == nil {
3885 return
3886 }
3887 before := c.executor.Session().Snapshot()
3888 next := planCancelledMessages(before, idx, fallback, startedAt, c.executor.CanReplayAssistantMessage, c.ledgerTailEvidence())
3889 if next != nil {
3890 c.replaceSessionAfterCancelFrom(before, next)
3891 }
3892 }
3893
3894 func (c *Controller) inFlightTurnStartedAt() time.Time {
3895 path := c.SessionPath()
3896 if path == "" {
3897 return time.Time{}
3898 }
3899 meta, ok, err := agent.LoadBranchMeta(path)
3900 if err != nil || !ok || meta.InFlightTurn == nil {
3901 return time.Time{}
3902 }
3903 return meta.InFlightTurn.StartedAt
3904 }
3905
3906 // resolveInterruptedTurnStart turns the pre-run array index into a stable
3907 // boundary after compaction. New user messages carry a creation timestamp set
3908 // after the marker, and graceful cleanup also has the exact composed prompt as
3909 // a fallback. We only fall back to the legacy index when it still points at a
3910 // plausible turn-start user message, keeping recovery data-safe for older
3911 // sidecars without timestamps.
3912 func resolveInterruptedTurnStart(msgs []provider.Message, idx int, preserveUser bool, startedAt time.Time, fallback provider.Message) (int, bool) {
3913 fallbackContent := ""
3914 if agent.IsUserAuthoredTurnMessage(fallback) {
3915 fallbackContent = StripComposePrefixes(fallback.Content)
3916 }
3917 matchesKind := func(m provider.Message) bool {
3918 if m.Role != provider.RoleUser || agent.IsPinnedContextRevision(m) {
3919 return false
3920 }
3921 if preserveUser {
3922 if !agent.IsUserAuthoredTurnMessage(m) {
3923 return false
3924 }
3925 if fallbackContent != "" && StripComposePrefixes(m.Content) != fallbackContent {
3926 return false
3927 }
3928 }
3929 return true
3930 }
3931 startedMillis := startedAt.UnixMilli()
3932 if !startedAt.IsZero() {
3933 for i, m := range msgs {
3934 if matchesKind(m) && m.CreatedAt >= startedMillis {
3935 return i, true
3936 }
3937 }
3938 }
3939 // Tests/headless runners may not persist an in-flight sidecar. The exact
3940 // graceful fallback still distinguishes the current visible turn; search
3941 // backward so a repeated prompt selects the newest occurrence.
3942 if fallbackContent != "" {
3943 for i, msg := range slices.Backward(msgs) {
3944 if matchesKind(msg) {
3945 return i, true
3946 }
3947 }
3948 }
3949 if idx >= 0 && idx < len(msgs) && matchesKind(msgs[idx]) {
3950 return idx, true
3951 }
3952 return 0, false
3953 }
3954
3955 func (c *Controller) hasInterruptedDisplayAfter(idx int, fallback provider.Message) bool {
3956 if c.executor == nil {
3957 return false
3958 }
3959 msgs := c.executor.Session().Snapshot()
3960 if start, ok := resolveInterruptedTurnStart(msgs, idx, true, c.inFlightTurnStartedAt(), fallback); ok {
3961 idx = start
3962 }
3963 idx = max(0, min(idx, len(msgs)))
3964 for _, m := range msgs[idx:] {
3965 if m.LocalOnly && m.InterruptedTurn != nil {
3966 return true
3967 }
3968 }
3969 return false
3970 }
3971
3972 func completeToolTurnEnd(msgs []provider.Message, i int) (int, bool) {
3973 if i < 0 || i >= len(msgs) {
3974 return i, false
3975 }
3976 m := msgs[i]
3977 if m.LocalOnly || m.Role != provider.RoleAssistant || len(m.ToolCalls) == 0 {
3978 return i, false
3979 }
3980 end := i + 1
3981 for end < len(msgs) && msgs[end].Role == provider.RoleTool && !msgs[end].LocalOnly {
3982 end++
3983 }
3984 results := msgs[i+1 : end]
3985 if len(results) != len(m.ToolCalls) {
3986 return i, false
3987 }
3988 for k, call := range m.ToolCalls {
3989 if strings.TrimSpace(call.Name) == "" || (call.Arguments != "" && !json.Valid([]byte(call.Arguments))) {
3990 return i, false
3991 }
3992 if results[k].ToolCallID != call.ID || results[k].Name != call.Name {
3993 return i, false
3994 }
3995 }
3996 return end, true
3997 }
3998
3999 func interruptedToolSummary(call provider.ToolCall) provider.InterruptedToolSummary {
4000 summary := provider.InterruptedToolSummary{
4001 ID: call.ID, Name: strings.TrimSpace(call.Name), Added: call.Added, Removed: call.Removed,
4002 }
4003 addFile := func(path string) {
4004 path = strings.TrimSpace(path)
4005 if path == "" || path == "/dev/null" || len(summary.Files) >= 8 {
4006 return
4007 }
4008 if slices.Contains(summary.Files, path) {
4009 return
4010 }
4011 summary.Files = append(summary.Files, path)
4012 }
4013 var args map[string]any
4014 if json.Unmarshal([]byte(call.Arguments), &args) == nil {
4015 for _, key := range []string{"path", "file", "file_path", "filename"} {
4016 if value, ok := args[key].(string); ok && strings.TrimSpace(value) != "" {
4017 addFile(value)
4018 }
4019 }
4020 }
4021 for line := range strings.SplitSeq(call.Diff, "\n") {
4022 line = strings.TrimSpace(line)
4023 switch {
4024 case strings.HasPrefix(line, "+++ b/"):
4025 addFile(strings.TrimPrefix(line, "+++ b/"))
4026 case strings.HasPrefix(line, "--- a/"):
4027 addFile(strings.TrimPrefix(line, "--- a/"))
4028 case strings.HasPrefix(line, "*** Update File: "):
4029 addFile(strings.TrimPrefix(line, "*** Update File: "))
4030 case strings.HasPrefix(line, "*** Add File: "):
4031 addFile(strings.TrimPrefix(line, "*** Add File: "))
4032 case strings.HasPrefix(line, "*** Delete File: "):
4033 addFile(strings.TrimPrefix(line, "*** Delete File: "))
4034 }
4035 }
4036 return summary
4037 }
4038
4039 func (c *Controller) replaceSessionAfterCancel(msgs []provider.Message) {
4040 if c.executor == nil {
4041 return
4042 }
4043 c.replaceSessionAfterCancelFromScoped(c.executor.Session().Snapshot(), msgs, true)
4044 }
4045
4046 func (c *Controller) replaceLegacySessionAfterCancelLocked(msgs []provider.Message) {
4047 // The whole cleanup is a save/recovery handoff like snapshot's: hold
4048 // snapshotMu from the in-memory truncation onward. Truncating outside the
4049 // lock would let an in-flight save capture the shortened transcript, read
4050 // the longer partial autosave on disk as a stale-prefix conflict, and
4051 // adopt it back into the executor — silently undoing the cancel cleanup
4052 // before the flush below could persist it.
4053 c.executor.Session().Replace(append([]provider.Message(nil), msgs...))
4054 // The mid-turn autosave may have already written a partial transcript to
4055 // disk. snapshotActivityIfChanged skips the write when messageCount()
4056 // returns to startMessages, so flush the cleaned transcript here. SaveRewrite
4057 // still checks that this controller owns the current on-disk baseline before
4058 // overwriting it, and also covers the edge case where the strip leaves only a
4059 // system message (HasContent() == false). The path is read under the lock so
4060 // an in-flight recovery retarget cannot leave it stale.
4061 c.mu.Lock()
4062 path := c.sessionPath
4063 c.mu.Unlock()
4064 if path != "" {
4065 if err := c.executor.Session().SaveRewrite(path); err != nil {
4066 if errors.Is(err, agent.ErrSessionSnapshotConflict) {
4067 if _, outcome, recoverErr := c.recoverSnapshotConflict(path, err, true); recoverErr != nil {
4068 slog.Warn("controller: post-cancel transcript recovery", "err", recoverErr)
4069 } else if outcome == conflictDropped {
4070 slog.Warn("controller: post-cancel transcript dropped after conflict", "path", path)
4071 }
4072 } else {
4073 slog.Warn("controller: post-cancel transcript flush", "err", err)
4074 }
4075 }
4076 }
4077 }
4078
4079 func (c *Controller) snapshotActivityIfChanged(startMessages int) (bool, error) {
4080 if c.messageCount() <= startMessages {
4081 return true, nil
4082 }
4083 return c.snapshotWithDurability(true, false, false)
4084 }
4085
4086 // SetSessionPath rebinds auto-save without changing the current session
4087 // preference. Callers creating a genuinely fresh conversation should use
4088 // SetFreshSessionPath; callers resuming history should use Resume.
4089 func (c *Controller) SetSessionPath(p string) {
4090 if c.sessionEngineEnabled() {
4091 // Path-based execution rebinding is intentionally unavailable. Hosts use
4092 // ContinueLegacySession for imports or publish an exact SessionRuntime.
4093 slog.Warn("controller: ignored legacy path rebind for exclusive v3 session", "path", p)
4094 return
4095 }
4096 c.setSessionPath(p, false)
4097 }
4098
4099 // SetFreshSessionPath binds a path that is known to belong to a newly-created
4100 // session and samples the configured new-session recovery default.
4101 func (c *Controller) SetFreshSessionPath(p string) {
4102 if service, _, exclusive := c.v3Binding(); exclusive && service != nil {
4103 if _, err := c.BindFreshSession(context.Background(), ""); err != nil {
4104 c.failTurnEventLedger(err)
4105 }
4106 return
4107 }
4108 c.setSessionPath(p, true)
4109 }
4110
4111 func (c *Controller) setSessionPath(p string, fresh bool) {
4112 defer c.refreshRuntimeState(event.Event{})
4113 // See snapshotMu: the swap must not interleave with an in-flight save.
4114 c.snapshotMu.Lock()
4115 c.mu.Lock()
4116 c.sessionPath = p
4117 c.guardianPath = guardian.PathFor(p)
4118 c.mu.Unlock()
4119 // Fresh paths clear projection; rebinds keep/load the target sidecar.
4120 c.bindExecutorProjection(p, !fresh)
4121 c.setActiveJobSession(p)
4122 c.rebindCheckpoints(p)
4123 // A path binding is the ownership boundary for the typed session log. Seed
4124 // the exact current transcript before any subsequent runtime or UI event can
4125 // make an otherwise empty v3 projection look authoritative. This covers the
4126 // initial EnsureSessionPath path as well as compatibility callers that bind
4127 // an already-loaded transcript without going through Resume.
4128 if err := c.seedSessionEventsFromExecutor("session-path-bind"); err != nil {
4129 slog.Warn("controller: seed bound session events", "path", p, "err", err)
4130 c.failTurnEventLedger(err)
4131 }
4132 if fresh {
4133 c.resetRecoveryForNewSession(p)
4134 // A newly-created conversation must not share the previous logical
4135 // session's temporary files (e.g. after EnsureSessionPath on a
4136 // controller that already ran commands).
4137 c.rotateSessionTemp()
4138 } else {
4139 c.loadRecoveryState(p)
4140 }
4141 c.snapshotMu.Unlock()
4142 c.rebindInbox()
4143 if !fresh {
4144 c.recoverCheckpointTransactions()
4145 }
4146 }
4147
4148 func (c *Controller) setActiveJobSession(sessionPath string) {
4149 if c.jobs != nil {
4150 // A candidate borrows the registry without changing its active owner.
4151 // Same-session replacement preserves the already bound artifact path.
4152 if c.background.scope != nil && c.jobs.ReplacementInProgress() {
4153 return
4154 }
4155 c.jobs.SetActiveSessionPath(agent.BranchID(sessionPath), sessionPath)
4156 }
4157 }
4158
4159 // SessionDir reports the directory new session files land in ("" disables
4160 // persistence), so the caller can decide whether to mint a path.
4161 func (c *Controller) SessionDir() string { return c.sessionDir }
4162
4163 // SessionPath reports the file the current conversation auto-saves to ("" when
4164 // persistence is disabled), so a history view can mark the active session.
4165 func (c *Controller) SessionPath() string {
4166 c.mu.Lock()
4167 defer c.mu.Unlock()
4168 return c.sessionPath
4169 }
4170
4171 // SessionRef returns the immutable v3 execution identity. The legacy path is
4172 // intentionally absent from this contract and may only remain as an import or
4173 // display locator while hosts complete their catalog transition.
4174 func (c *Controller) SessionRef() (session.SessionRef, bool) {
4175 _, runtime, _ := c.v3Binding()
4176 if runtime == nil {
4177 return session.SessionRef{}, false
4178 }
4179 return runtime.Ref(), true
4180 }
4181
4182 func (c *Controller) parentSessionID() string {
4183 if ref, ok := c.SessionRef(); ok {
4184 return ref.SessionID
4185 }
4186 return agent.BranchID(c.SessionPath())
4187 }
4188
4189 // History returns the executor's current message log (for repopulating a
4190 // resumed frontend's view).
4191 func (c *Controller) History() []provider.Message {
4192 if c.executor == nil {
4193 return nil
4194 }
4195 if snapshot, ok := c.sessionEventSnapshot(); ok && snapshot.EventSequence > 0 {
4196 if c.sessionEngineEnabled() {
4197 // Controller.History is the provider workset compatibility API used by
4198 // rebuilds and model switches. Durable UI history is deliberately
4199 // separate and flows through SessionService.Query/TranscriptReplay.
4200 return append([]provider.Message(nil), snapshot.Projection.ModelMessages...)
4201 }
4202 projected := snapshot.Projection.ModelMessages
4203 if store := c.sessionEventStore(); store != nil {
4204 // The retired path-bound compatibility store still owns its complete
4205 // persisted transcript. Its provider projection intentionally strips
4206 // host-origin metadata, so explicit history reads use the stored view.
4207 projected = store.Snapshot().Projection.Messages
4208 }
4209 current := c.executor.Session().Snapshot()
4210 // Compatibility runners and an interrupted pre-v3 caller can still leave
4211 // a live, unsaved message tail in the legacy Session cache. Keep that tail
4212 // visible without inferring or persisting business state from it. Normal
4213 // agent execution records exact messages first, so this branch disappears
4214 // as the remaining compatibility callers are retired.
4215 if len(current) > len(projected) && (len(projected) == 0 || reflect.DeepEqual(current[:len(projected)], projected)) {
4216 return current
4217 }
4218 return append([]provider.Message(nil), projected...)
4219 }
4220 return c.executor.Session().Snapshot() // copy — a turn may be appending concurrently
4221 }
4222
4223 // SessionHasUnsavedChanges tells desktop history whether it may safely refresh
4224 // an idle view from the durable WAL. A failed or contended save can leave the
4225 // controller with a newer in-memory transcript; replacing that view from disk
4226 // would hide the user's latest turn until the next retry.
4227 func (c *Controller) SessionHasUnsavedChanges() bool {
4228 if c == nil || c.executor == nil {
4229 return false
4230 }
4231 if snapshot, ok := c.sessionEventSnapshot(); ok && snapshot.EventSequence > 0 {
4232 // The v3 projection is the authoritative transcript. A mismatch identifies
4233 // obsolete compatibility code that bypassed the explicit event recorder;
4234 // checkpoints deliberately do not infer or mirror that unrecorded tail.
4235 current := c.executor.Session().Snapshot()
4236 return snapshot.EventSequence > snapshot.DurableSequence ||
4237 snapshot.PersistenceStatus == "failed" ||
4238 snapshot.PersistenceStatus == "uncertain" ||
4239 !reflect.DeepEqual(snapshot.Projection.ModelMessages, current)
4240 }
4241 return c.executor.Session().HasUnsavedChanges(c.SessionPath())
4242 }
4243
4244 // HistoryLen returns the number of messages in the live log.
4245 func (c *Controller) HistoryLen() int {
4246 if c.executor == nil {
4247 return 0
4248 }
4249 return c.executor.Session().Len()
4250 }
4251
4252 // HistoryWindow returns a copy of the messages in [start, end) of the live
4253 // log. Paging frontends use it to convert a display window without copying
4254 // the whole history.
4255 func (c *Controller) HistoryWindow(start, end int) []provider.Message {
4256 if c.executor == nil {
4257 return []provider.Message{}
4258 }
4259 return c.executor.Session().MessageRange(start, end)
4260 }
4261
4262 // SessionPersistedState exposes the session's persistence baseline for the
4263 // controller's current session path, so a paging frontend can validate a
4264 // display-index sidecar against the live session.
4265 func (c *Controller) SessionPersistedState() (agent.PersistedState, bool) {
4266 if c.executor == nil {
4267 return agent.PersistedState{}, false
4268 }
4269 return c.executor.Session().PersistedState(c.SessionPath())
4270 }
4271
4272 // ContextSnapshot returns (usedTokens, contextWindow) for the gauge. usedTokens
4273 // is what the next request will send, measured the way the compaction trigger
4274 // measures it, so the gauge and the trigger can never disagree. Both zero means
4275 // no data yet — a gauge hides itself.
4276 func (c *Controller) ContextSnapshot() (int, int) {
4277 if c.executor == nil {
4278 return 0, 0
4279 }
4280 return c.executor.ContextUsedTokens(), c.executor.ContextWindow()
4281 }
4282
4283 // CompactRatio returns the auto-compaction threshold as a fraction of the window
4284 // (0 when the executor is unset). The status line shows headroom against it.
4285 func (c *Controller) CompactRatio() float64 {
4286 if c.executor == nil {
4287 return 0
4288 }
4289 return c.executor.CompactRatio()
4290 }
4291
4292 // LastUsage returns the most recent turn's token telemetry (nil before the first
4293 // turn), so frontends can derive the prompt cache-hit rate for the status line.
4294 func (c *Controller) LastUsage() *provider.Usage {
4295 if c.executor == nil {
4296 return nil
4297 }
4298 return c.executor.LastUsage()
4299 }
4300
4301 // SessionCache returns cumulative cache hit/miss prompt tokens for the session,
4302 // so a frontend can render the aggregate (session-wide) cache-hit rate — steadier
4303 // than the single-turn rate and unaffected by compaction.
4304 func (c *Controller) SessionCache() (hit, miss int) {
4305 if c.executor == nil {
4306 return 0, 0
4307 }
4308 return c.executor.SessionCache()
4309 }
4310
4311 // Todos returns the committed event projection used by every frontend. The
4312 // executor's mutable copy is only a tool-loop convenience and cannot override
4313 // a failed or superseded durable commit.
4314 func (c *Controller) Todos() []evidence.TodoItem {
4315 var todos []event.Todo
4316 if ledger := c.turnEventLedger(); ledger != nil {
4317 todos, _ = ledger.TodoState()
4318 } else {
4319 todos, _ = c.volatileTodoState()
4320 }
4321 out := make([]evidence.TodoItem, len(todos))
4322 for i, todo := range todos {
4323 out[i] = evidence.TodoItem{Content: todo.Content, Status: todo.Status}
4324 }
4325 return out
4326 }
4327
4328 // Balance queries the active provider's wallet balance, or (nil, nil) when the
4329 // provider declares no balance_url — so a caller treats "not configured" and
4330 // "fetched" the same and just omits the readout when nil.
4331 func (c *Controller) Balance(ctx context.Context) (*billing.Balance, error) {
4332 if strings.TrimSpace(c.balanceURL) == "" {
4333 return nil, nil
4334 }
4335 ctx, cancel := context.WithTimeout(ctx, 12*time.Second)
4336 defer cancel()
4337 return billing.FetchWithClient(ctx, c.balanceClient, c.balanceURL, c.balanceKey)
4338 }
4339
4340 // Host returns the running MCP host (nil when no plugins), for frontends that
4341 // list servers / resolve MCP prompts.
4342 func (c *Controller) Host() *plugin.Host { return c.mcp.hostRef() }
4343
4344 // Commands returns the loaded custom slash commands.
4345 func (c *Controller) Commands() []command.Command {
4346 if p := c.commands.Load(); p != nil {
4347 return *p
4348 }
4349 return nil
4350 }
4351
4352 // ReloadCommands rescans all command directories and hot-swaps the slash_command
4353 // tool and the internal command slice — no MCP restart, no hook rerun.
4354 func (c *Controller) ReloadCommands(ctx context.Context) error {
4355 select {
4356 case <-ctx.Done():
4357 return ctx.Err()
4358 default:
4359 }
4360 cmds, loadErr := command.LoadRoots(config.CommandRootsForRoot(c.workspaceRoot)...)
4361 var cmdSkills []skill.Skill
4362 if !c.disableImplicitSkillInvocation {
4363 cmdSkills = c.SlashSkills()
4364 }
4365
4366 entries := make([]command.SlashEntry, 0, len(cmdSkills)+len(cmds))
4367 for _, sk := range cmdSkills {
4368
4369 entries = append(entries, command.SlashEntry{
4370 Name: sk.SlashName(),
4371 Description: sk.Description,
4372 Render: func(args []string) string { return c.skills.renderInvocation(sk, strings.Join(args, " ")) },
4373 })
4374 }
4375 for _, cmd := range cmds {
4376 if cmd.Hidden {
4377 continue
4378 }
4379
4380 entries = append(entries, command.SlashEntry{
4381 Name: cmd.Name,
4382 Description: cmd.Description,
4383 ArgHint: cmd.ArgHint,
4384 Render: func(args []string) string { return cmd.Render(args) },
4385 })
4386 }
4387 c.mcp.registerTool(command.NewSlashCommandTool(entries))
4388 cmdSlice := cmds
4389 c.commands.Store(&cmdSlice)
4390 return loadErr
4391 }
4392
4393 // Skills returns the discoverable skills (for the slash menu and `/skills`).
4394 // When a live Store is available, scan it on demand so skills installed during
4395 // this session appear without rewriting the cache-stable system prompt.
4396 // Executor returns the underlying agent when present (nil for pure runners).
4397 func (c *Controller) Executor() *agent.Agent {
4398 if c == nil {
4399 return nil
4400 }
4401 return c.executor
4402 }
4403
4404 func (c *Controller) Skills() []skill.Skill {
4405 return c.skills.list()
4406 }
4407
4408 // ImplicitSkillInvocationEnabled reports whether skills are exposed to the
4409 // model for automatic discovery and invocation. Explicit /skill handling is
4410 // independent of this model-facing capability.
4411 func (c *Controller) ImplicitSkillInvocationEnabled() bool {
4412 return c != nil && !c.disableImplicitSkillInvocation
4413 }
4414
4415 // SlashSkills returns the user-visible skill directory. Plugin skills use
4416 // package-qualified names while Skills keeps bare model/run_skill identifiers.
4417 func (c *Controller) SlashSkills() []skill.Skill {
4418 return c.skills.slashList()
4419 }
4420
4421 // AllSkills returns every discoverable skill, including disabled ones, for
4422 // management surfaces that need to re-enable a hidden skill.
4423 func (c *Controller) AllSkills() []skill.Skill {
4424 return c.skills.listAll()
4425 }
4426
4427 // LoadSkill reads the selected skill body without expanding every catalog
4428 // candidate. It is intentionally outside the shared capability interface:
4429 // management surfaces may opt into body loading, while search/list stay
4430 // metadata-only and cache-stable.
4431 func (c *Controller) LoadSkill(name string) (skill.Skill, bool) {
4432 return c.skills.load(name)
4433 }
4434
4435 // DisabledSkills returns all discoverable skills that are disabled in config.
4436 func (c *Controller) DisabledSkills() []skill.Skill {
4437 cfg, err := config.Load()
4438 if err != nil {
4439 return nil
4440 }
4441 var out []skill.Skill
4442 for _, sk := range c.AllSkills() {
4443 if cfg.IsSkillDisabled(sk.Name) {
4444 out = append(out, sk)
4445 }
4446 }
4447 return out
4448 }
4449
4450 // SkillEnabled reports whether a discoverable skill is enabled.
4451 func (c *Controller) SkillEnabled(name string) bool {
4452 cfg, err := config.Load()
4453 if err != nil {
4454 return true
4455 }
4456 return !cfg.IsSkillDisabled(name)
4457 }
4458
4459 // SetSkillEnabled persists a skill enable/disable preference. The caller should
4460 // rebuild the controller for the prompt/tool registry to reflect it immediately.
4461 func (c *Controller) SetSkillEnabled(name string, enabled bool) error {
4462 found := false
4463 for _, sk := range c.AllSkills() {
4464 if config.SkillNameKey(sk.Name) == config.SkillNameKey(name) {
4465 name = sk.Name
4466 found = true
4467 break
4468 }
4469 }
4470 if !found {
4471 return fmt.Errorf("unknown skill: %s", name)
4472 }
4473 // Serialize the load-modify-save against other in-process user-config
4474 // editors so concurrent writers (bot mapping persistence, desktop
4475 // settings) don't drop this toggle or lose their own fields.
4476 unlock := config.LockUserConfigEdits()
4477 defer unlock()
4478 cfg := config.LoadForEdit(config.UserConfigPath())
4479 if err := cfg.SetSkillEnabled(name, enabled); err != nil {
4480 return err
4481 }
4482 return cfg.SaveTo(config.UserConfigPath())
4483 }
4484
4485 // CreateSkill writes a new skill file at the given scope and returns its
4486 // path. Skills()/AllSkills()/RunSkill() read the live store on demand, so the
4487 // new skill is usable (by name) immediately with no rebuild and appears in the
4488 // next real user turn's live session-context catalog. Rebuilds remain necessary
4489 // when the tool registry or enabled-skill configuration changes.
4490 func (c *Controller) CreateSkill(name string, scope skill.Scope, content string) (string, error) {
4491 w := c.skills.writer()
4492 if w == nil {
4493 return "", fmt.Errorf("no writable skill store in this session")
4494 }
4495 return w.CreateWithContent(name, scope, content)
4496 }
4497
4498 // UpdateSkill overwrites an existing user-authored skill file in place. See
4499 // skill.Store.UpdateContent for the builtin-refusal and scope-match rules.
4500 func (c *Controller) UpdateSkill(name string, scope skill.Scope, content string) error {
4501 w := c.skills.writer()
4502 if w == nil {
4503 return fmt.Errorf("no writable skill store in this session")
4504 }
4505 return w.UpdateContent(name, scope, content)
4506 }
4507
4508 // DeleteSkill removes a user-authored skill file at the given scope. See
4509 // skill.Store.Delete for the builtin-refusal and scope-match rules.
4510 func (c *Controller) DeleteSkill(name string, scope skill.Scope) error {
4511 w := c.skills.writer()
4512 if w == nil {
4513 return fmt.Errorf("no writable skill store in this session")
4514 }
4515 return w.Delete(name, scope)
4516 }
4517
4518 // HookRunner returns the session's hook runner (nil-safe; may hold zero hooks),
4519 // so a frontend can list the active hooks via `/hooks`.
4520 func (c *Controller) HookRunner() *hook.Runner { return c.hooks }
4521
4522 // AddMCPServer connects an MCP server live and persists it to the user-global
4523 // config. Its tools are registered immediately and become available on the next
4524 // turn (the agent reads the registry per turn). The raw entry — ${VARS} intact —
4525 // is what's written to disk; the live connection uses the expanded form. Returns
4526 // the number of tools the server exposed. Persistence is transactional: a config
4527 // or activation failure removes the just-connected client so the live registry
4528 // never claims an install that will disappear after restart.
4529 func (c *Controller) AddMCPServer(e config.PluginEntry) (int, error) {
4530 // AddMCPServer is an explicit user action. Mark the live entry with the same
4531 // provenance it will receive when the saved user config is loaded next time,
4532 // so /mcp add is add-and-use in the current session too.
4533 e.Source = config.MCPSourceUserConfig
4534 if effective, loadErr := config.LoadForRootReadOnly(c.workspaceRoot); loadErr != nil {
4535 return 0, loadErr
4536 } else {
4537 for _, configured := range effective.Plugins {
4538 if configured.Name != e.Name {
4539 continue
4540 }
4541 if configured.Source != config.MCPSourceUserConfig && configured.Source != config.MCPSourceLegacyUser {
4542 return 0, fmt.Errorf("MCP server %q is already configured by %s; edit or remove that declaration before installing a global server with the same name", e.Name, configured.Source)
4543 }
4544 break
4545 }
4546 }
4547 n, err := c.connectMCPServer(e)
4548 if err != nil {
4549 return 0, err
4550 }
4551 if _, err := config.InstallUserPluginForRoot(c.workspaceRoot, e, true); err != nil {
4552 c.DisconnectMCPServer(e.Name)
4553 return 0, fmt.Errorf("saving MCP server config: %w", err)
4554 }
4555 return n, nil
4556 }
4557
4558 // ConnectMCPServer connects an MCP server entry for this session without writing
4559 // it to config. Desktop owns config placement so it can keep user-level settings
4560 // out of project reasonix.toml while preserving the CLI AddMCPServer semantics.
4561 func (c *Controller) ConnectMCPServer(e config.PluginEntry) (int, error) {
4562 return c.connectMCPServer(e)
4563 }
4564
4565 // RegisterMCPServerOnDemand restores a configured server's cached provider
4566 // surface without forcing a handshake. It is the durable-enable counterpart to
4567 // ConnectMCPServer, which remains the explicit install/retry operation.
4568 func (c *Controller) RegisterMCPServerOnDemand(e config.PluginEntry) (int, error) {
4569 spec := c.mcpSpec(e)
4570 n, err := c.mcp.registerSpecOnDemand(spec)
4571 if err == nil && c.capabilityRuntime != nil {
4572 c.capabilityRuntime.UpsertServer(e, spec, true)
4573 }
4574 return n, err
4575 }
4576
4577 // connectMCPServer expands an entry's ${VARS}, applies the known-server
4578 // overrides scoped to the workspace, and connects it live via the mcp manager.
4579 func (c *Controller) connectMCPServer(e config.PluginEntry) (int, error) {
4580 spec := c.mcpSpec(e)
4581 n, err := c.mcp.connectSpec(spec)
4582 if err == nil && c.capabilityRuntime != nil {
4583 c.capabilityRuntime.UpsertServer(e, spec, true)
4584 }
4585 return n, err
4586 }
4587
4588 func (c *Controller) mcpSpec(e config.PluginEntry) plugin.Spec {
4589 exp := e.ExpandedPlugin()
4590 configSource := strings.TrimSpace(string(exp.Source))
4591 spec := plugin.ApplyKnownOverrides(plugin.Spec{
4592 Name: exp.Name,
4593 Type: exp.Type,
4594 Command: exp.Command,
4595 Args: exp.Args,
4596 Env: exp.Env,
4597 URL: exp.URL,
4598 Headers: exp.Headers,
4599 StartupTimeout: controllerMCPTimeout(exp.StartupTimeoutSeconds),
4600 DefaultCallTimeout: c.mcpDefaultCallTimeout,
4601 CallTimeout: controllerMCPTimeout(exp.CallTimeoutSeconds),
4602 ToolTimeouts: controllerMCPToolTimeouts(exp.ToolTimeoutSeconds),
4603 WorkspaceRoot: c.WorkspaceRoot(),
4604 ConfigSource: configSource,
4605 Authorized: exp.Source.UserAuthorized(),
4606 // Explicit user installs and reconnects run as trusted host processes.
4607 ProcessMode: plugin.MCPProcessHost,
4608 }, c.WorkspaceRoot())
4609 if exp.Source.ProjectScoped() && strings.TrimSpace(spec.Dir) == "" {
4610 spec.Dir = c.WorkspaceRoot()
4611 }
4612 if c.mcpConfigureSpec != nil {
4613 c.mcpConfigureSpec(&spec)
4614 if spec.ProcessMode == "" {
4615 spec.ProcessMode = plugin.MCPProcessHost
4616 }
4617 }
4618 return spec
4619 }
4620
4621 // syncCapabilityRuntimeFromConfig restores one server's authoritative runtime
4622 // entry after a transactional disconnect/rollback. enabledOverride is used for
4623 // a session-only disconnect; nil re-resolves the durable activation state.
4624 func (c *Controller) syncCapabilityRuntimeFromConfig(name string, enabledOverride *bool) {
4625 if c == nil || c.capabilityRuntime == nil {
4626 return
4627 }
4628 name = strings.TrimSpace(name)
4629 cfg, err := config.LoadForRoot(c.workspaceRoot)
4630 if err != nil {
4631 // The caller revokes first. A config read failure must not re-enable a
4632 // potentially stale spec or shared-Host client.
4633 return
4634 }
4635 for _, entry := range cfg.Plugins {
4636 if strings.TrimSpace(entry.Name) != name {
4637 continue
4638 }
4639 enabled := config.DeclaredDefaultOn(entry)
4640 if enabledOverride != nil {
4641 enabled = *enabledOverride
4642 } else if resolved, resolveErr := config.DefaultMCPActivationStore().IsEnabled(entry, c.workspaceRoot); resolveErr == nil {
4643 enabled = resolved
4644 }
4645 c.capabilityRuntime.UpsertServer(entry, c.mcpSpec(entry), enabled)
4646 return
4647 }
4648 c.capabilityRuntime.RemoveServer(name)
4649 }
4650
4651 func controllerMCPTimeout(seconds int) time.Duration {
4652 if seconds <= 0 {
4653 return 0
4654 }
4655 return time.Duration(seconds) * time.Second
4656 }
4657
4658 func controllerMCPToolTimeouts(values map[string]int) map[string]time.Duration {
4659 if len(values) == 0 {
4660 return nil
4661 }
4662 out := make(map[string]time.Duration, len(values))
4663 for name, seconds := range values {
4664 if name = strings.TrimSpace(name); name != "" && seconds > 0 {
4665 out[name] = time.Duration(seconds) * time.Second
4666 }
4667 }
4668 if len(out) == 0 {
4669 return nil
4670 }
4671 return out
4672 }
4673
4674 // ImportMCPEntries persists selected MCP entries and attempts to connect them
4675 // live. A connection failure does not roll back the config import: the user can
4676 // fix local dependencies and reconnect in a later session.
4677 func (c *Controller) ImportMCPEntries(entries []config.PluginEntry) (total, added, updated, connected, failed, skipped int, err error) {
4678 total, added, updated, err = config.ImportCCSwitchMCPEntries(entries)
4679 if err != nil {
4680 return 0, 0, 0, 0, 0, 0, err
4681 }
4682 effectiveCfg, loadErr := config.LoadForRoot(c.workspaceRoot)
4683 if loadErr != nil {
4684 return 0, 0, 0, 0, 0, 0, loadErr
4685 }
4686 effective := make(map[string]config.PluginEntry, len(effectiveCfg.Plugins))
4687 for _, entry := range effectiveCfg.Plugins {
4688 effective[entry.Name] = entry
4689 }
4690 for _, imported := range entries {
4691 e, ok := effective[imported.Name]
4692 if !ok || e.Source != config.MCPSourceUserConfig {
4693 // A project declaration with the same name remains effective. The
4694 // imported global entry is saved as its lower-priority fallback.
4695 skipped++
4696 continue
4697 }
4698 if c.mcp.hasServer(e.Name) {
4699 if c.capabilityRuntime != nil {
4700 // Import updates may intentionally keep an existing live client, but
4701 // future proxy reconnects must use the newly persisted spec.
4702 c.capabilityRuntime.UpsertServer(e, c.mcpSpec(e), true)
4703 }
4704 skipped++
4705 continue
4706 }
4707 if _, err := c.AddMCPServer(e); err != nil {
4708 failed++
4709 continue
4710 }
4711 connected++
4712 }
4713 return total, added, updated, connected, failed, skipped, nil
4714 }
4715
4716 func (c *Controller) ConfiguredMCPNames() []string {
4717 cfg, err := config.LoadForRootReadOnly(c.workspaceRoot)
4718 if err != nil {
4719 return nil
4720 }
4721 names := make([]string, 0, len(cfg.Plugins))
4722 for _, p := range cfg.Plugins {
4723 names = append(names, p.Name)
4724 }
4725 return names
4726 }
4727
4728 func (c *Controller) DisconnectedMCPNames() []string {
4729 cfg, err := config.LoadForRootReadOnly(c.workspaceRoot)
4730 if err != nil {
4731 return nil
4732 }
4733 connected := map[string]bool{}
4734 for _, name := range c.mcp.serverNames() {
4735 connected[name] = true
4736 }
4737 var names []string
4738 for _, p := range cfg.Plugins {
4739 if !connected[p.Name] {
4740 names = append(names, p.Name)
4741 }
4742 }
4743 return names
4744 }
4745
4746 // ConnectConfiguredMCPServer starts a configured server at the user's request;
4747 // for a project-declared one that request is the approval, and is recorded.
4748 func (c *Controller) ConnectConfiguredMCPServer(name string) (int, error) {
4749 p, err := c.configuredMCPServer(name)
4750 if err != nil {
4751 return 0, err
4752 }
4753 if err := config.RecordExplicitStart(p, c.workspaceRoot); err != nil {
4754 return 0, err
4755 }
4756 return c.connectMCPServer(p)
4757 }
4758
4759 func (c *Controller) configuredMCPServer(name string) (config.PluginEntry, error) {
4760 cfg, err := config.LoadForRoot(c.workspaceRoot)
4761 if err != nil {
4762 return config.PluginEntry{}, err
4763 }
4764 for _, p := range cfg.Plugins {
4765 if p.Name == name {
4766 return p, nil
4767 }
4768 }
4769 return config.PluginEntry{}, fmt.Errorf("no configured MCP server named %q", name)
4770 }
4771
4772 // RemoveMCPServer removes writable config before disconnecting the live server.
4773 // A persistence failure must not produce a false-successful session-only removal.
4774 // MCPs contributed by installed plugin packages cannot be removed independently.
4775 func (c *Controller) RemoveMCPServer(name string) (disconnected bool, err error) {
4776 cfg, lerr := config.LoadForRoot(c.workspaceRoot)
4777 if lerr != nil {
4778 return false, lerr
4779 }
4780 if owner, ok := cfg.PluginPackageOwner(name); ok {
4781 return false, fmt.Errorf("MCP server %q is managed by plugin %q; disable or remove the plugin instead", name, owner)
4782 }
4783 entry, removed, _, rerr := config.RemovePluginFromEffectiveSourceForRoot(c.workspaceRoot, name)
4784 if rerr != nil {
4785 return false, rerr
4786 }
4787 if !removed {
4788 return false, fmt.Errorf("no removable MCP server named %q", name)
4789 }
4790 _ = config.DefaultMCPActivationStore().ClearServer(entry, c.workspaceRoot)
4791 removedState := reconcileRemovedMCPState(c.workspaceRoot, name)
4792 if c.capabilityRuntime != nil {
4793 // Revoke before touching the shared Host so an overlapping resolver cannot
4794 // reuse a sibling tab's still-connected client.
4795 c.capabilityRuntime.RemoveServer(name)
4796 }
4797 disconnected = c.mcp.disconnect(name)
4798 if !disconnected {
4799 c.mcp.removeToolPrefix(name)
4800 }
4801 // A lower-priority same-name declaration may now be effective. Restore its
4802 // cached/on-demand surface without starting a process; otherwise ensure the
4803 // removed name stays absent.
4804 if removedState.fallbackFound {
4805 enabled := config.MCPServerEnabled(removedState.fallback, c.workspaceRoot)
4806 if enabled {
4807 _, _ = c.RegisterMCPServerOnDemand(removedState.fallback)
4808 } else {
4809 c.syncCapabilityRuntimeFromConfig(name, &enabled)
4810 }
4811 } else {
4812 c.syncCapabilityRuntimeFromConfig(name, nil)
4813 }
4814 return disconnected, removedState.cleanupErr
4815 }
4816
4817 // DisconnectMCPServer disconnects a live server for this session without touching
4818 // config — the connector toggle's "off". Its tools vanish next turn; it reconnects
4819 // on the next session start, or now via ConnectConfiguredMCPServer (the "on").
4820 // Reports whether a live server was actually disconnected.
4821 func (c *Controller) DisconnectMCPServer(name string) bool {
4822 if c.capabilityRuntime != nil {
4823 c.capabilityRuntime.SetServerEnabled(name, false)
4824 }
4825 disconnected := c.mcp.disconnect(name)
4826 removedPlaceholder := 0
4827 if !disconnected {
4828 removedPlaceholder = c.mcp.removeToolPrefix(name)
4829 }
4830 // Keep configured servers discoverable as disabled, but forget runtime-only
4831 // or rolled-back installs that no longer exist in configuration.
4832 disabled := false
4833 c.syncCapabilityRuntimeFromConfig(name, &disabled)
4834 return disconnected || removedPlaceholder > 0
4835 }
4836
4837 // UnregisterMCPServerTools hides a shared MCP server from this controller only.
4838 // The desktop shared-host path uses this for per-tab connector toggles: the
4839 // shared client stays alive for sibling tabs, while this session's registry drops
4840 // the server's provider-visible tools before the next turn.
4841 func (c *Controller) UnregisterMCPServerTools(name string) bool {
4842 if c.capabilityRuntime != nil {
4843 c.capabilityRuntime.SetServerEnabled(name, false)
4844 }
4845 return c.mcp.suspendToolPrefix(name)
4846 }
4847
4848 // Label returns the human-readable model label, e.g. "deepseek-flash".
4849 func (c *Controller) Label() string { return c.label }
4850
4851 // ModelRef returns the canonical provider/model reference for the session.
4852 func (c *Controller) ModelRef() string { return c.selection.ref }
4853
4854 // ModelSelectionIdentity is frozen with the provider assembled for this runtime.
4855 func (c *Controller) ModelSelectionIdentity() string { return c.selection.identity }
4856
4857 // WorkspaceRoot returns the workspace root for this controller's session
4858 // (the directory that file-writers and @-references are scoped to).
4859 // Empty means no scoping is in effect.
4860 func (c *Controller) WorkspaceRoot() string { return c.workspaceRoot }
4861
4862 // WorkspaceRepo is the workspace's git identity as the session resolved it when
4863 // it opened. Host git reads the workspace through it rather than rediscovering
4864 // the repository from files the session's own commands may have written.
4865 func (c *Controller) WorkspaceRepo() gitcmd.Repo { return c.workspaceRepo }
4866
4867 func (c *Controller) imageInputEnabled() bool {
4868 if c.frozenImageInput != nil {
4869 return *c.frozenImageInput
4870 }
4871 ref := c.selection.ref
4872 cfg, err := config.LoadForRoot(c.workspaceRoot)
4873 if err == nil && ref == "" {
4874 ref = cfg.DefaultModel
4875 }
4876 if err != nil || ref == "" {
4877 return false
4878 }
4879 entry, ok := cfg.ResolveModel(ref)
4880 if !ok {
4881 return false
4882 }
4883 if c.modelCapabilityResolver != nil {
4884 return c.modelCapabilityResolver(entry).State == config.CapabilitySupported
4885 }
4886 return config.EffectiveVision(entry)
4887 }
4888
4889 // ImageInputEnabled reports whether the current model accepts direct image
4890 // inputs, so frontends can gate image-only UX before a turn starts.
4891 func (c *Controller) ImageInputEnabled() bool { return c.imageInputEnabled() }
4892
4893 // ImageInputSnapshot avoids configuration reads on the Desktop metadata path.
4894 // Legacy/custom controllers without a frozen boot snapshot use the existing
4895 // background metadata fallback instead.
4896 func (c *Controller) ImageInputSnapshot() (enabled, fallback, available bool) {
4897 if c == nil || c.frozenImageInput == nil {
4898 return false, false, false
4899 }
4900 return *c.frozenImageInput, c.visionModel != "", true
4901 }
4902
4903 // ImageCapabilityChanged lets desktop refresh an idle runtime before admission.
4904 func (c *Controller) ImageCapabilityChanged() bool {
4905 return c.imageCapabilityChanged != nil && c.imageCapabilityChanged()
4906 }
4907
4908 // SessionAuthorizations snapshots this controller's same-session tool
4909 // grants ("Allow for this session") and Plan-mode read-only command trust,
4910 // for carrying into a replacement controller across a rebuild — see
4911 // RestoreSessionAuthorizations.
4912 func (c *Controller) SessionAuthorizations() SessionAuthorizations {
4913 auth := c.approval.snapshotSessionAuthorizations()
4914 if c.writeAccess.roots != nil {
4915 auth.WriteRoots = c.writeAccess.roots.SessionRoots()
4916 if auth.WriteRoots == nil {
4917 auth.WriteRoots = []string{}
4918 }
4919 }
4920 return auth
4921 }
4922
4923 // ReleaseResources stops plugin subprocesses and releases resources without
4924 // firing SessionEnd. Use it only when replacing the controller for the same
4925 // logical session.
4926 func (c *Controller) ReleaseResources() {
4927 c.close(false, closeJobsWithGrace)
4928 }
4929
4930 // Close stops plugin subprocesses and releases resources. A session that ever
4931 // started fires SessionEnd so a teardown hook runs.
4932 func (c *Controller) Close() {
4933 c.recordLifecycle("close", "controller_close", "", 0, "")
4934 c.close(true, closeJobsWithGrace)
4935 }
4936
4937 // Closed is signalled after teardown has released all Controller-owned stores.
4938 // Close itself only requests teardown when a turn is still finalizing.
4939 func (c *Controller) Closed() <-chan struct{} { return c.closeFinalized }
4940
4941 // CloseAfterDestroy releases controller resources after the caller has already
4942 // begun session-specific job teardown. It avoids a second synchronous job grace
4943 // wait while still cancelling the manager root and reaping temporary artifacts
4944 // once every job goroutine finally exits.
4945 func (c *Controller) CloseAfterDestroy() {
4946 c.close(true, closeJobsAsync)
4947 }
4948
4949 type closeJobsMode int
4950
4951 const (
4952 closeJobsWithGrace closeJobsMode = iota
4953 closeJobsAsync
4954 )
4955
4956 func (c *Controller) close(fireSessionEnd bool, jobsMode closeJobsMode) {
4957 defer c.refreshRuntimeState(event.Event{})
4958 // Desktop tab lifecycles can race a rebind/model-switch/close on the same
4959 // controller; make teardown idempotent so a duplicate Close cannot re-fire
4960 // SessionEnd hooks or re-run cleanup. The first caller's jobsMode wins.
4961 c.closeOnce.Do(func() {
4962 c.mu.Lock()
4963 cancel := c.turns.cancel
4964 maintenance := c.maintenance
4965 done := c.turns.done
4966 // A phase marker alone is not a live turn: recovery may retain one after
4967 // cancel/done ownership has gone. Only a live body or terminal fanout
4968 // defers final resource release.
4969 maintenanceActive := maintenance != nil && !maintenance.safeToRelease
4970 maintenanceNeedsCancel := maintenanceActive && maintenance.activity != "finalizing"
4971 turnActive := done != nil || c.finalizingLocked() || c.turns.recoveryFanout || maintenanceActive
4972 // Seal turn admission and drop anything already parked: a parked turn
4973 // must not start against a controller that is being torn down, and
4974 // without the closed flag a submit landing after this critical
4975 // section (while a running turn's TurnDone delivery is still in
4976 // flight) would park again and start after teardown.
4977 c.closed = true
4978 c.closeFireSessionEnd = fireSessionEnd
4979 c.closeJobsMode = jobsMode
4980 c.turns.pending = nil
4981 c.turns.wake = false
4982 if cancel != nil {
4983 c.turns.cancelRequested = true
4984 if c.turns.phase == session.RuntimeRunning {
4985 c.turns.phase = session.RuntimeCancelling
4986 c.noteExecutionLocked(session.RuntimeCancelling, "cancelling")
4987 }
4988 } else {
4989 c.turns.cancelRequested = false
4990 }
4991 if !turnActive {
4992 if maintenance != nil {
4993 c.maintenance = nil
4994 close(maintenance.done)
4995 }
4996 c.turns.phase = session.RuntimeClosed
4997 c.turns.finishingBound.end()
4998 c.turns.finishingBound.endIdle()
4999 }
5000 c.mu.Unlock()
5001 if maintenanceNeedsCancel {
5002 maintenance.cancel()
5003 }
5004 if cancel != nil {
5005 // Signal the owned turn before prompt bookkeeping or callbacks. A
5006 // stalled registry/adapter must never delay Stop during shutdown.
5007 cancel()
5008 c.startCancellationWatchdog(done)
5009 c.promptOwner.CancelAll()
5010 c.approval.clearAll()
5011 } else {
5012 c.promptOwner.Clear()
5013 }
5014 if c.goalDriverControl.cancel != nil {
5015 c.goalDriverControl.cancel()
5016 }
5017 if !turnActive {
5018 c.finalizeControllerClose()
5019 }
5020 })
5021 }
5022
5023 // finalizeControllerClose releases stores and process resources only after an
5024 // active turn has published its terminal boundary. Closing the ledger or the
5025 // session binding earlier makes the final TurnDone impossible to accept.
5026 func (c *Controller) finalizeControllerClose() {
5027 c.closeFinalizeOnce.Do(func() {
5028 if c.closeFinalized != nil {
5029 defer close(c.closeFinalized)
5030 }
5031 c.mu.Lock()
5032 started := c.startedOnce
5033 fireSessionEnd := c.closeFireSessionEnd && !c.background.retired
5034 jobsMode := c.closeJobsMode
5035 c.mu.Unlock()
5036 // Goal-driver workers may be inside the pre-admission durability
5037 // checkpoint. Join them before closing the v3 writer so teardown cannot
5038 // race a late Flush or recreate files under a test/session directory.
5039 c.goalDriverWG.Wait()
5040 // Join sidecar creation and queue scans without waiting for the
5041 // dispatcher itself: host admission may retire its own controller.
5042 c.inbox.scanMu.Lock()
5043 c.inbox.mu.Lock()
5044 c.inbox.closed = true
5045 if c.inbox.store != nil {
5046 c.inbox.store.Close()
5047 c.inbox.store = nil
5048 }
5049 if c.inbox.tempLease != nil {
5050 c.inbox.tempLease.Release()
5051 c.inbox.tempLease = nil
5052 }
5053 c.inbox.mu.Unlock()
5054 c.inbox.scanMu.Unlock()
5055 if fireSessionEnd && started {
5056 c.hooks.SessionEnd(context.Background(), "other")
5057 c.extensionSessionEvent(extension.PointSessionEnd, dispatch.PhaseEnd, c.SessionPath())
5058 }
5059 if c.background.scope != nil {
5060 c.background.scope.Release(jobsMode == closeJobsAsync)
5061 } else if c.jobs != nil {
5062 switch jobsMode {
5063 case closeJobsAsync:
5064 c.jobs.CloseAsync()
5065 default:
5066 c.jobs.Close() // cancel any still-running background jobs
5067 }
5068 }
5069 c.turnEvents.mu.RLock()
5070 projection := c.turnEvents.projection
5071 c.turnEvents.mu.RUnlock()
5072 if projection != nil {
5073 projection.CloseFollowers()
5074 }
5075 if ledger := c.turnEventLedger(); ledger != nil {
5076 if err := ledger.Close(); err != nil {
5077 slog.Warn("controller: close turn event ledger", "err", err)
5078 }
5079 }
5080 c.turnEvents.commitMu.Lock()
5081 if pending := c.turnEvents.pendingExecutionCommit; pending != nil {
5082 pending.Release()
5083 c.turnEvents.pendingExecutionCommit = nil
5084 }
5085 c.turnEvents.commitMu.Unlock()
5086 service, runtime, exclusive := c.v3Binding()
5087 if exclusive && runtime != nil {
5088 c.releaseSessionRuntimeBinding(service)
5089 } else if v3 := c.sessionEventStore(); v3 != nil {
5090 c.turnEvents.mu.RLock()
5091 release := c.turnEvents.v3Release
5092 c.turnEvents.mu.RUnlock()
5093 var err error
5094 if release != nil {
5095 err = release(context.Background())
5096 } else {
5097 err = v3.Close(context.Background())
5098 }
5099 if err != nil {
5100 slog.Warn("controller: flush and close v3 session", "err", err)
5101 }
5102 }
5103 if c.cleanup != nil {
5104 c.cleanup()
5105 }
5106 // Drop the Controller owner reference last so background job leases
5107 // that outlive close still pin retired generations until they exit.
5108 if c.sessionTemp != nil {
5109 c.sessionTemp.Release()
5110 }
5111 if c.persistentShell != nil {
5112 c.persistentShell.Release()
5113 }
5114 c.finishBackgroundReplacement(false)
5115 })
5116 }
5117
5118 // SessionTemp returns the logical-session private temporary directory manager.
5119 // Hot rebuilds pass this to the replacement Controller so the directory survives
5120 // model/settings swaps. Nil only when the Controller was constructed without one
5121 // (should not happen after New).
5122 func (c *Controller) SessionTemp() *sessiontemp.Manager {
5123 if c == nil {
5124 return nil
5125 }
5126 return c.sessionTemp
5127 }
5128
5129 // rotateSessionTemp advances the private temporary generation so a new logical
5130 // session cannot see the previous session's temporary files. In-flight command
5131 // leases keep the old generation alive until they release.
5132 func (c *Controller) rotateSessionTemp() {
5133 if c == nil {
5134 return
5135 }
5136 if c.sessionTemp != nil {
5137 c.sessionTemp.Rotate()
5138 }
5139 if c.persistentShell != nil {
5140 c.persistentShell.Rotate()
5141 }
5142 }
5143
5144 // PersistentShell returns the session-scoped PTY manager. Hot rebuilds pass
5145 // this to the replacement Controller so shell state survives model/settings
5146 // swaps. Nil only when the Controller was constructed without one.
5147 func (c *Controller) PersistentShell() *persistentshell.Manager {
5148 if c == nil {
5149 return nil
5150 }
5151 return c.persistentShell
5152 }
5153
5154 // Jobs returns the still-running background jobs for the status bar (nil when
5155 // background jobs are disabled).
5156 func (c *Controller) Jobs() []jobs.View {
5157 if c.jobs == nil {
5158 return nil
5159 }
5160 return c.jobs.RunningForSession(c.parentSessionID())
5161 }
5162
5163 // KillJob cancels a running background job by ID.
5164 func (c *Controller) KillJob(id string) bool {
5165 if c.jobs == nil {
5166 return false
5167 }
5168 return c.jobs.Kill(id)
5169 }
5170
5171 // TaskRuntimeOwnerID identifies the recorder that admitted this runtime's jobs.
5172 func (c *Controller) TaskRuntimeOwnerID() string {
5173 if c.jobs == nil {
5174 return ""
5175 }
5176 return c.jobs.TaskRuntimeOwnerID()
5177 }
5178
5179 // CancelJob stops one background job owned by this controller's session.
5180 func (c *Controller) CancelJob(id string) bool {
5181 if c.jobs == nil {
5182 return false
5183 }
5184 return c.jobs.KillForSession(c.parentSessionID(), id)
5185 }
5186
5187 // WorkspaceLeaseState reports the process-local held and waiting lease scope.
5188 // Canonical keys are used only for matching another local controller and are
5189 // never copied into user-facing payloads.
5190 func (c *Controller) WorkspaceLeaseState() workspacelease.State {
5191 return c.workspaceLease.State()
5192 }
5193
5194 // WorkspaceLeaseHeldKeys is the lock-domain identity this controller currently
5195 // holds. Empty when no write lease is held.
5196 func (c *Controller) WorkspaceLeaseHeldKeys() []string {
5197 return c.workspaceLease.HeldKeys()
5198 }
5199
5200 // SetToolApprovalMode changes the runtime approval posture for permission-gated
5201 // tools. It does not answer business asks or plan approval. Sub-agents (task,
5202 // writer-capable skill sub-agents, the planner) have no UI to prompt through,
5203 // so this also pushes the mode to the shared headless gate they read from —
5204 // without it, a mode switch (Shift+Tab) would only rebuild the parent
5205 // executor's gate and leave sub-agents pinned to whatever mode was active
5206 // when the session booted.
5207 func (c *Controller) SetToolApprovalMode(mode string) {
5208 c.ApplyToolApprovalMode(mode)
5209 }
5210
5211 // ApplyToolApprovalMode updates the preset without answering an older prompt.
5212 // A real mode change invalidates the active turn and its request IDs, so an
5213 // approval created under an earlier permission revision can never authorize a
5214 // call under the new revision.
5215 func (c *Controller) ApplyToolApprovalMode(mode string) []string {
5216 c.permissionMu.Lock()
5217 defer c.permissionMu.Unlock()
5218 return c.applyToolApprovalModeLocked(mode)
5219 }
5220
5221 func (c *Controller) applyToolApprovalModeLocked(mode string) []string {
5222 mode = normalizeToolApprovalMode(mode)
5223 c.promptResolveMu.Lock()
5224 previousMode := c.approval.mode()
5225 if previousMode == mode {
5226 c.promptResolveMu.Unlock()
5227 return nil
5228 }
5229 c.permissionStateMu.Lock()
5230 pending := c.approval.setMode(mode)
5231 if c.subagentGate != nil {
5232 c.subagentGate.Update(mode)
5233 }
5234 c.refreshInteractiveGate()
5235 // Publish the revision only after every enforcement owner has adopted the
5236 // new mode. promptResolveMu makes this one transaction with approval commit.
5237 c.permissionRevision.Add(1)
5238 c.permissionStateMu.Unlock()
5239 drained := make([]string, 0, len(pending))
5240 for _, p := range pending {
5241 p.reply <- approvalReply{allow: true}
5242 drained = append(drained, p.id)
5243 }
5244 // A change that does not widen the preset invalidates the active turn so an
5245 // older approval never authorizes work the new snapshot forbids. A widening
5246 // revokes nothing the turn holds, so the turn keeps running. Avoid the idle
5247 // Cancel path: it stops an active Goal, and the preset is a separate axis.
5248 upgrade := permissionPresetRank(mode) > permissionPresetRank(previousMode)
5249 turnID, cancelled := "", false
5250 if !upgrade && c.Running() {
5251 turnID, cancelled = c.cancelTurnLocked()
5252 }
5253 c.promptResolveMu.Unlock()
5254 if cancelled {
5255 c.finishCancel(turnID, true)
5256 }
5257 // Processes admitted under a broader preset may outlive their spawning
5258 // turn. Only a downgrade must terminate them; an upgrade does not revoke
5259 // any capability they already held.
5260 if permissionPresetRank(mode) < permissionPresetRank(previousMode) && !c.isBackgroundCandidate() {
5261 for _, job := range c.Jobs() {
5262 c.CancelJob(job.ID)
5263 }
5264 }
5265 c.refreshRuntimeState(event.Event{})
5266 return drained
5267 }
5268
5269 func permissionPresetRank(mode string) int {
5270 switch normalizeToolApprovalMode(mode) {
5271 case ToolApprovalDangerFullAccess:
5272 return 2
5273 case ToolApprovalWorkspaceWrite:
5274 return 1
5275 default:
5276 return 0
5277 }
5278 }
5279
5280 func (c *Controller) ToolApprovalMode() string {
5281 return c.approval.mode()
5282 }
5283
5284 // SetAutoApproveTools is the legacy boolean compatibility binding. Both values
5285 // migrate to the safe workspace preset; full access requires an explicit preset.
5286 func (c *Controller) SetAutoApproveTools(on bool) {
5287 _ = on
5288 c.SetToolApprovalMode(ToolApprovalWorkspaceWrite)
5289 }
5290
5291 // SetBypass is the legacy name for SetAutoApproveTools. Keep it for existing
5292 // desktop/serve bindings and CLI code that still uses the bypass wording.
5293 func (c *Controller) SetBypass(on bool) {
5294 c.SetAutoApproveTools(on)
5295 }
5296
5297 // SetMode is the legacy combined Plan/permission binding. New callers use
5298 // ApplyComposerProfile with an explicit permission preset.
5299 func (c *Controller) SetMode(plan, autoApproveTools bool) {
5300 c.ApplyMode(plan, autoApproveTools)
5301 }
5302
5303 // ApplyMode is the legacy SetMode variant that reports invalidated prompt IDs.
5304 func (c *Controller) ApplyMode(plan, autoApproveTools bool) []string {
5305 c.applyPlanMode(plan)
5306 _ = autoApproveTools
5307 return c.ApplyToolApprovalMode(ToolApprovalWorkspaceWrite)
5308 }
5309
5310 // ApplyComposerProfile publishes the collaboration, approval, and Goal axes as
5311 // one controller operation. The only fallible mutation (durable Goal state)
5312 // commits first, so a persistence failure leaves Plan and approval unchanged.
5313 // Serve serializes this call with turn admission and controller replacement.
5314 func (c *Controller) ApplyComposerProfile(plan bool, toolApprovalMode, goal string) ([]string, error) {
5315 c.permissionMu.Lock()
5316 defer c.permissionMu.Unlock()
5317 return c.applyComposerProfileLocked(plan, toolApprovalMode, goal)
5318 }
5319
5320 // ApplyComposerProfileAt is the revision-checked protocol entry point used by
5321 // desktop and remote clients. It keeps Goal, Plan and permission changes under
5322 // one controller mutation boundary.
5323 func (c *Controller) ApplyComposerProfileAt(plan bool, toolApprovalMode, goal string, expectedPermissionRevision uint64) ([]string, error) {
5324 c.permissionMu.Lock()
5325 defer c.permissionMu.Unlock()
5326 if current := c.permissionRevision.Load(); current != expectedPermissionRevision {
5327 return nil, fmt.Errorf("permission revision changed: have %d, expected %d", current, expectedPermissionRevision)
5328 }
5329 return c.applyComposerProfileLocked(plan, toolApprovalMode, goal)
5330 }
5331
5332 // ApplyComposerProfileAtDurable is the canonical-session protocol entry point.
5333 // The explicit preset reaches durable storage before it changes enforcement.
5334 func (c *Controller) ApplyComposerProfileAtDurable(ctx context.Context, plan bool, toolApprovalMode, goal string, expectedPermissionRevision uint64) ([]string, error) {
5335 c.permissionMu.Lock()
5336 defer c.permissionMu.Unlock()
5337 if current := c.permissionRevision.Load(); current != expectedPermissionRevision {
5338 return nil, fmt.Errorf("permission revision changed: have %d, expected %d", current, expectedPermissionRevision)
5339 }
5340 raw := strings.ToLower(strings.TrimSpace(toolApprovalMode))
5341 if raw != "ask" && raw != "auto" && raw != "yolo" && !permissionpreset.Valid(raw) {
5342 return nil, fmt.Errorf("permission preset must be read-only, workspace-write, or danger-full-access")
5343 }
5344 raw = normalizeToolApprovalMode(raw)
5345 capabilities := platformPermissionCapabilities()
5346 if !slices.Contains(capabilities.SupportedPresets, raw) {
5347 return nil, fmt.Errorf("permission preset %q is unavailable: %s", raw, capabilities.UnavailableReason)
5348 }
5349 _, runtime, exclusive := c.v3Binding()
5350 if !exclusive || runtime == nil {
5351 return nil, session.ErrSessionNotRunning
5352 }
5353 goal = strings.TrimSpace(goal)
5354 if strings.TrimSpace(c.Goal()) != goal {
5355 if err := c.SetGoalDurable(goal); err != nil {
5356 return nil, fmt.Errorf("persist goal state: %w", err)
5357 }
5358 }
5359 if err := c.persistSessionPermissionPreset(ctx, runtime, raw); err != nil {
5360 return c.applyFailedPermissionDowngradeLocked(raw, err)
5361 }
5362 if goal != "" {
5363 plan = false
5364 }
5365 c.applyPlanMode(plan)
5366 return c.applyDurablePermissionPresetLocked(raw), nil
5367 }
5368
5369 func (c *Controller) applyComposerProfileLocked(plan bool, toolApprovalMode, goal string) ([]string, error) {
5370 if raw := strings.ToLower(strings.TrimSpace(toolApprovalMode)); raw != "ask" && raw != "auto" && raw != "yolo" &&
5371 raw != ToolApprovalReadOnly && raw != ToolApprovalWorkspaceWrite && raw != ToolApprovalDangerFullAccess {
5372 return nil, fmt.Errorf("permission preset must be read-only, workspace-write, or danger-full-access")
5373 }
5374 toolApprovalMode = normalizeToolApprovalMode(toolApprovalMode)
5375 goal = strings.TrimSpace(goal)
5376 if strings.TrimSpace(c.Goal()) != goal {
5377 if err := c.SetGoalDurable(goal); err != nil {
5378 return nil, fmt.Errorf("persist goal state: %w", err)
5379 }
5380 }
5381 if goal != "" {
5382 plan = false
5383 }
5384 c.applyPlanMode(plan)
5385 return c.applyToolApprovalModeLocked(toolApprovalMode), nil
5386 }
5387
5388 // AutoApproveTools is the legacy status query for explicit full access.
5389 func (c *Controller) AutoApproveTools() bool {
5390 return c.ToolApprovalMode() == ToolApprovalDangerFullAccess
5391 }
5392
5393 // Bypass is the legacy name for AutoApproveTools.
5394 func (c *Controller) Bypass() bool {
5395 return c.AutoApproveTools()
5396 }
5397
5398 // memory
5399 //
5400 // The memory snapshot, pending standing-doc notes, and write serialization
5401 // live in c.memory (a memoryManager) behind its own locks, off c.mu — so a
5402 // memory-panel save never stalls an approval or status poll. These methods are
5403 // the SessionAPI surface; each is a thin delegation. See memory.go.
5404
5405 // QuickAdd appends a one-line note to the doc-memory file for scope (project
5406 // REASONIX.md by default) — the write side of "#<note>". Returns the file written.
5407 func (c *Controller) QuickAdd(scope memory.Scope, note string) (string, error) {
5408 return c.memory.quickAdd(scope, note)
5409 }
5410
5411 // SaveDoc overwrites a recognized memory doc with body — the save side of the
5412 // desktop panel's in-place editor. Returns the file written.
5413 func (c *Controller) SaveDoc(path, body string) (string, error) {
5414 return c.memory.saveDoc(path, body)
5415 }
5416
5417 // SaveMemory writes an active auto-memory fact and refreshes the in-session
5418 // snapshot. It is the explicit user-confirmed counterpart to the model-owned
5419 // remember tool, used by management surfaces that preview a candidate first.
5420 func (c *Controller) SaveMemory(m memory.Memory) (string, error) {
5421 return c.memory.saveMemory(m)
5422 }
5423
5424 // ForgetMemory removes a saved auto-memory by name — the panel/TUI forget action,
5425 // the manual counterpart to the model's `forget` tool.
5426 func (c *Controller) ForgetMemory(name string) error {
5427 return c.memory.forget(name)
5428 }
5429
5430 // QueueMemory implements memory.Queue: model remember/forget tool results are
5431 // already visible in the current loop, so this refreshes the background snapshot
5432 // that will be published in session-context on the next real user turn.
5433 func (c *Controller) QueueMemory(note string) {
5434 c.memory.queue(note)
5435 }
5436
5437 // ClaimAutoMemoryWrite consumes the one-shot create-only authorization issued
5438 // by gateApprover for a low-risk project fact.
5439 func (c *Controller) ClaimAutoMemoryWrite(args json.RawMessage) bool {
5440 return c.memory.claimAutoRemember(args)
5441 }
5442
5443 func (c *Controller) MemoryRevisions(ref string) []memory.Memory {
5444 return c.memory.revisions(ref)
5445 }
5446
5447 // RestoreMemory restores an older active-memory revision as a new audited
5448 // revision and applies it to the next user turn.
5449 func (c *Controller) RestoreMemory(ref string, revision int) (memory.Memory, error) {
5450 return c.memory.restore(ref, revision)
5451 }
5452
5453 // RestoreArchivedMemory recovers an archived fact as a new audited revision and
5454 // applies it to the next user turn.
5455 func (c *Controller) RestoreArchivedMemory(archivePath string) (memory.Memory, error) {
5456 return c.memory.restoreArchived(archivePath)
5457 }
5458
5459 // Memory returns the loaded memory snapshot (nil when memory is disabled), for
5460 // frontends that surface a memory panel or the /memory command. The returned
5461 // *Set is immutable — mutations go through QuickAdd / SaveDoc.
5462 func (c *Controller) Memory() *memory.Set {
5463 return c.memory.current()
5464 }
5465
5466 // approval bridge (agent gate → events)
5467
5468 // gateApprover adapts the Controller to permission.Approver. It is distinct
5469 // from the public Approve command (different signature, different direction).
5470 type gateApprover struct{ c *Controller }
5471
5472 func (g gateApprover) Approve(ctx context.Context, tool, subject string, args json.RawMessage) (bool, bool, error) {
5473 allow, remember, _, err := g.ApproveWithReason(ctx, tool, subject, args)
5474 return allow, remember, err
5475 }
5476
5477 func (g gateApprover) ApproveWithReason(ctx context.Context, tool, subject string, args json.RawMessage) (bool, bool, string, error) {
5478 return g.approveWithPolicyReason(ctx, tool, subject, args, "")
5479 }
5480
5481 func (g gateApprover) ApproveWithPolicyReason(ctx context.Context, tool, subject string, args json.RawMessage, policyReason string) (bool, bool, string, error) {
5482 return g.approveWithPolicyReason(ctx, tool, subject, args, policyReason)
5483 }
5484
5485 func (g gateApprover) approveWithPolicyReason(ctx context.Context, tool, subject string, args json.RawMessage, policyReason string) (bool, bool, string, error) {
5486 if tool == memoryRememberTool && g.c.allowLowRiskRemember(args) {
5487 return true, false, "", nil
5488 }
5489 subject = approvalDisplaySubject(tool, subject, args)
5490 // Check pre-approval before any prompt or Guardian review. OS sandboxing is
5491 // the authority for nested shell syntax, so pipes, command substitutions and
5492 // inline interpreters follow the same permission decision as other commands.
5493 if g.c.approval.preApproved(tool, subject, args) {
5494 return true, false, "", nil
5495 }
5496 allow, remember, err := g.c.requestApprovalWithReason(ctx, tool, subject, args, policyReason)
5497 return allow, remember, "", err
5498 }
5499
5500 type planModeReadOnlyTrustApprover struct{ c *Controller }
5501
5502 type sandboxEscapeApprover struct{ c *Controller }
5503
5504 func (s sandboxEscapeApprover) ApproveSandboxEscape(ctx context.Context, req sandbox.EscapeRequest) (bool, string, error) {
5505 subject := sandboxEscapeApprovalSubject(req.Command)
5506 reason := sandboxEscapeApprovalReason(req.Reason)
5507 reply, err := s.c.requestFreshApprovalDecision(ctx, SandboxEscapeApprovalTool, subject, req.Args, reason)
5508 if err != nil {
5509 return false, "approval aborted", err
5510 }
5511 if !reply.allow {
5512 return false, i18n.M.SandboxEscapeDeclined, nil
5513 }
5514 if reply.session {
5515 s.c.permissionStateMu.Lock()
5516 s.c.approval.grantSession(SandboxEscapeApprovalTool, subject)
5517 s.c.permissionStateMu.Unlock()
5518 }
5519 return true, "", nil
5520 }
5521
5522 func (s sandboxEscapeApprover) SandboxEscapeSessionAllowed(_ context.Context, req sandbox.EscapeRequest) bool {
5523 return s.c.approval.preApprovedForDecision(SandboxEscapeApprovalTool, sandboxEscapeApprovalSubject(req.Command), nil, true)
5524 }
5525
5526 func sandboxEscapeApprovalSubject(command string) string {
5527 subject := strings.TrimSpace(command)
5528 if subject == "" {
5529 return i18n.M.SandboxEscapeSubjectFallback
5530 }
5531 return i18n.M.SandboxEscapeSubjectPrefix + subject
5532 }
5533
5534 func sandboxEscapeApprovalReason(reason string) string {
5535 reason = strings.TrimSpace(reason)
5536 if reason == "" {
5537 return i18n.M.SandboxEscapeRuntimeReason
5538 }
5539 return reason
5540 }
5541
5542 // managedConfigWriteApprover routes a file tool's Reasonix-managed config write
5543 // through the fresh-human approval prompt (see ManagedConfigWriteApprovalTool).
5544 // A session grant is tool-wide (mirroring sandbox_escape): one "allow for this
5545 // session" covers the rest of the repair flow across the handful of managed
5546 // config files without re-prompting on every incremental edit.
5547 type managedConfigWriteApprover struct{ c *Controller }
5548
5549 func (m managedConfigWriteApprover) ApproveManagedConfigWrite(ctx context.Context, req tool.ConfigWriteRequest) (bool, string, error) {
5550 subject := managedConfigWriteApprovalSubject(req.Path)
5551 args, _ := json.Marshal(map[string]string{"path": req.Path})
5552 reply, err := m.c.requestFreshApprovalDecision(ctx, ManagedConfigWriteApprovalTool, subject, args, i18n.M.ConfigWriteReason)
5553 if err != nil {
5554 return false, "approval aborted", err
5555 }
5556 if !reply.allow {
5557 return false, i18n.M.ConfigWriteDeclined, nil
5558 }
5559 if reply.session {
5560 m.c.permissionStateMu.Lock()
5561 m.c.approval.grantSession(ManagedConfigWriteApprovalTool, subject)
5562 m.c.permissionStateMu.Unlock()
5563 }
5564 return true, "", nil
5565 }
5566
5567 func (m managedConfigWriteApprover) ManagedConfigWriteSessionAllowed(_ context.Context, req tool.ConfigWriteRequest) bool {
5568 return m.c.approval.preApprovedForDecision(ManagedConfigWriteApprovalTool, managedConfigWriteApprovalSubject(req.Path), nil, true)
5569 }
5570
5571 func managedConfigWriteApprovalSubject(path string) string {
5572 return i18n.M.ConfigWriteSubjectPrefix + strings.TrimSpace(path)
5573 }
5574
5575 func (p planModeReadOnlyTrustApprover) CheckPlanModeReadOnlyTrust(ctx context.Context, req agent.PlanModeReadOnlyTrustRequest) (bool, string, error) {
5576 prefix := normalizePlanModeReadOnlyCommandPrefix(req.Prefix)
5577 if prefix == "" {
5578 return false, "missing plan-mode read-only command prefix", nil
5579 }
5580 return p.checkBashReadOnlyCommandTrust(ctx, req, prefix)
5581 }
5582
5583 func (p planModeReadOnlyTrustApprover) checkBashReadOnlyCommandTrust(ctx context.Context, req agent.PlanModeReadOnlyTrustRequest, prefix string) (bool, string, error) {
5584 if p.c.approval.planModeReadOnlyCommandTrusted(prefix) {
5585 return true, "", nil
5586 }
5587 command := strings.TrimSpace(req.Command)
5588 if command == "" {
5589 command = strings.TrimSpace(string(req.Args))
5590 }
5591 subject := fmt.Sprintf(i18n.M.PlanModeBashTrustSubjectFmt, prefix, command)
5592 reason := i18n.M.PlanModeBashTrustReason
5593 reply, err := p.c.requestFreshApprovalDecision(ctx, agent.PlanModeReadOnlyCommandApprovalTool, subject, req.Args, reason)
5594 if err != nil {
5595 return false, "approval aborted", err
5596 }
5597 if !reply.allow {
5598 return false, i18n.M.PlanModeBashTrustDeclined, nil
5599 }
5600 if reply.session {
5601 p.c.permissionStateMu.Lock()
5602 p.c.approval.grantPlanModeReadOnlyCommand(prefix)
5603 p.c.permissionStateMu.Unlock()
5604 }
5605 if reply.persist && p.c.onRememberPlanModeReadOnlyCommand != nil {
5606 p.c.emitPlanModeReadOnlyCommandTrustResult(p.c.onRememberPlanModeReadOnlyCommand(prefix))
5607 p.c.permissionStateMu.Lock()
5608 p.c.approval.grantPlanModeReadOnlyCommand(prefix)
5609 p.c.permissionStateMu.Unlock()
5610 }
5611 return true, "", nil
5612 }
5613
5614 func approvalDisplaySubject(tool, subject string, args json.RawMessage) string {
5615 switch tool {
5616 case memoryRememberTool:
5617 return rememberApprovalSubject(subject, args)
5618 case memoryForgetTool:
5619 return forgetApprovalSubject(subject, args)
5620 case "move_file":
5621 return moveApprovalSubject(subject, args)
5622 default:
5623 return subject
5624 }
5625 }
5626
5627 func moveApprovalSubject(fallback string, args json.RawMessage) string {
5628 if len(args) == 0 {
5629 return fallback
5630 }
5631 var in struct {
5632 SourcePath string `json:"source_path"`
5633 DestinationPath string `json:"destination_path"`
5634 }
5635 if err := json.Unmarshal(args, &in); err != nil {
5636 return fallback
5637 }
5638 if in.SourcePath == "" || in.DestinationPath == "" {
5639 return fallback
5640 }
5641 return in.SourcePath + " -> " + in.DestinationPath
5642 }
5643
5644 func rememberApprovalSubject(fallback string, args json.RawMessage) string {
5645 if len(args) == 0 {
5646 return fallback
5647 }
5648 var in struct {
5649 Name string `json:"name"`
5650 Title string `json:"title"`
5651 Description string `json:"description"`
5652 Type string `json:"type"`
5653 Body string `json:"body"`
5654 }
5655 if err := json.Unmarshal(args, &in); err != nil {
5656 return fallback
5657 }
5658 name := approvalCompactText(firstNonEmpty(in.Name, in.Title))
5659 desc := approvalTruncate(approvalCompactText(in.Description), 180)
5660 body := approvalTruncate(approvalCompactText(in.Body), 240)
5661 typ := string(memory.NormalizeType(in.Type))
5662
5663 var b strings.Builder
5664 b.WriteString(i18n.M.MemoryApprovalSaveUpdate)
5665 baseLen := b.Len()
5666 if name != "" {
5667 fmt.Fprintf(&b, " %q", name)
5668 }
5669 if typ != "" {
5670 fmt.Fprintf(&b, " [%s]", typ)
5671 }
5672 if desc != "" {
5673 b.WriteString(": ")
5674 b.WriteString(desc)
5675 }
5676 if body != "" {
5677 if desc == "" {
5678 b.WriteString(": ")
5679 } else {
5680 b.WriteString(" | ")
5681 }
5682 b.WriteString(i18n.M.MemoryApprovalBodyLabel)
5683 b.WriteString(": ")
5684 b.WriteString(body)
5685 }
5686 if b.Len() == baseLen && fallback != "" {
5687 return fallback
5688 }
5689 return b.String()
5690 }
5691
5692 func forgetApprovalSubject(fallback string, args json.RawMessage) string {
5693 if len(args) == 0 {
5694 return fallback
5695 }
5696 var in struct {
5697 Name string `json:"name"`
5698 }
5699 if err := json.Unmarshal(args, &in); err != nil {
5700 return fallback
5701 }
5702 name := approvalCompactText(in.Name)
5703 if name == "" {
5704 return fallback
5705 }
5706 return fmt.Sprintf(i18n.M.MemoryApprovalArchiveFmt, name)
5707 }
5708
5709 func firstNonEmpty(values ...string) string {
5710 for _, value := range values {
5711 if strings.TrimSpace(value) != "" {
5712 return value
5713 }
5714 }
5715 return ""
5716 }
5717
5718 func approvalCompactText(s string) string {
5719 return strings.Join(strings.Fields(s), " ")
5720 }
5721
5722 func approvalTruncate(s string, maxRunes int) string {
5723 if maxRunes <= 0 {
5724 return ""
5725 }
5726 runes := []rune(s)
5727 if len(runes) <= maxRunes {
5728 return s
5729 }
5730 return string(runes[:maxRunes]) + "..."
5731 }
5732
5733 func (c *Controller) sessionMessageCount() int {
5734 if c.executor == nil {
5735 return 0
5736 }
5737 return c.executor.Session().Len()
5738 }
5739
5740 // parseRewind parses the arguments after "/rewind". The user may provide:
5741 //
5742 // /rewind → latest checkpoint, both
5743 // /rewind <turn> → that turn, both
5744 // /rewind <turn> <scope> → that turn, code|conversation|both
5745 //
5746 // If no turn is given, the latest checkpoint is used. If no scope is given, Both is assumed.
5747 func parseRewind(args string, cps []checkpoint.Meta) (int, RewindScope, error) {
5748 fields := strings.Fields(args)
5749 if len(fields) == 0 {
5750 if len(cps) == 0 {
5751 return 0, RewindBoth, fmt.Errorf("no checkpoints available")
5752 }
5753 return cps[len(cps)-1].Turn, RewindBoth, nil
5754 }
5755 turn, err := strconv.Atoi(fields[0])
5756 if err != nil {
5757 return 0, RewindBoth, fmt.Errorf("invalid turn: %w", err)
5758 }
5759 scope := RewindBoth
5760 if len(fields) >= 2 {
5761 switch strings.ToLower(fields[1]) {
5762 case "code":
5763 scope = RewindCode
5764 case "conversation":
5765 scope = RewindConversation
5766 case "both":
5767 scope = RewindBoth
5768 default:
5769 return 0, RewindBoth, fmt.Errorf("unknown scope %q", fields[1])
5770 }
5771 }
5772 return turn, scope, nil
5773 }
5774
5775 // requestApproval emits an ApprovalRequest and blocks until Approve(ID, …)
5776 // answers or ctx is cancelled. A prior session grant (or a bypass posture) for
5777 // the same approval scope short-circuits. Each prompt waits independently;
5778 // this method keeps the I/O (events, hooks, remember) out of the registry.
5779 func (c *Controller) requestApproval(ctx context.Context, tool, subject string, args json.RawMessage) (bool, bool, error) {
5780 return c.requestApprovalWithReason(ctx, tool, subject, args, "")
5781 }
5782
5783 func (c *Controller) requestApprovalWithReason(ctx context.Context, tool, subject string, args json.RawMessage, reason string) (bool, bool, error) {
5784 return c.requestApprovalWithReasonOptions(ctx, tool, subject, args, reason, approvalDecisionOptions{})
5785 }
5786
5787 func (c *Controller) requestApprovalWithReasonOptions(ctx context.Context, tool, subject string, args json.RawMessage, reason string, opts approvalDecisionOptions) (bool, bool, error) {
5788 r, err := c.requestApprovalDecisionWithOptions(ctx, tool, subject, args, reason, opts)
5789 if err != nil {
5790 return false, false, err
5791 }
5792 // Plan approvals are one-shot — never persist a session grant for them, or
5793 // every future plan would auto-approve.
5794 if r.allow && r.session && !requiresFreshApprovalTool(tool) {
5795 c.permissionStateMu.Lock()
5796 c.approval.grantSession(tool, subject)
5797 c.permissionStateMu.Unlock()
5798 }
5799 if r.allow && r.persist && !requiresFreshApprovalTool(tool) && c.onRemember != nil {
5800 c.emitRememberResult(c.onRemember(permission.RememberRuleForScope(tool, subject)))
5801 }
5802 return r.allow, false, nil
5803 }
5804
5805 func (c *Controller) requestFreshApprovalDecision(ctx context.Context, tool, subject string, args json.RawMessage, reason string) (approvalReply, error) {
5806 return c.requestApprovalDecisionWithOptions(ctx, tool, subject, args, reason, approvalDecisionOptions{fresh: true})
5807 }
5808
5809 type approvalDecisionOptions struct {
5810 // fresh marks a user trust/business decision rather than an ordinary tool
5811 // permission. It may reuse an explicit session grant, but YOLO/auto approval
5812 // must not answer or drain the prompt.
5813 fresh bool
5814 // requireHuman marks an ordinary tool approval that Auto, an approved-plan
5815 // window, Guardian, or an allowing hook must not answer. Unlike fresh it
5816 // retains the ordinary four-choice UI and YOLO remains an explicit bypass.
5817 requireHuman bool
5818 }
5819
5820 func (c *Controller) requestApprovalDecisionWithOptions(ctx context.Context, tool, subject string, args json.RawMessage, reason string, opts approvalDecisionOptions) (approvalReply, error) {
5821 // YOLO/full access and the just-approved-plan execution window auto-allow
5822 // approval-gated tools without prompting. Plan approval is a user decision,
5823 // not a tool permission, so it deliberately stays interactive.
5824 if c.approval.preApprovedForDecisionOptions(tool, subject, args, opts.fresh, opts.requireHuman) {
5825 return approvalReply{allow: true}, nil
5826 }
5827
5828 // Claude's PermissionRequest contract answers the dialog on the plugin's
5829 // behalf (auto-allow/auto-deny) instead of merely observing it, so a
5830 // decision here must preempt the prompt rather than just notify — this
5831 // runs synchronously and before the dialog is shown. Native Reasonix
5832 // PermissionRequest hooks stay advisory-only (see claudePermissionBlocking).
5833 //
5834 // Hook auto-allow cannot replace a fresh-human decision. Interactive
5835 // YOLO may already have skipped remember/forget via preApproved. Deny
5836 // still applies universally — refusing is always safe to honor.
5837 if hookSubject, hookArgs, ok := permissionRequestHookPayload(tool, subject, args); ok {
5838 if decision, _ := c.hooks.PermissionRequest(ctx, tool, hookSubject, hookArgs); decision != nil {
5839 switch {
5840 case !*decision:
5841 return approvalReply{}, nil
5842 case !opts.fresh && !opts.requireHuman && !requiresFreshApprovalTool(tool):
5843 return approvalReply{allow: true}, nil
5844 }
5845 // An "allow" opinion on a fresh-human-required decision is
5846 // ignored; fall through to the normal interactive prompt.
5847 }
5848 }
5849
5850 c.approval.promptEmitMu.Lock()
5851 // Re-check at the publication boundary: a session grant may have landed
5852 // after the initial policy check.
5853 if c.approval.preApprovedForDecisionOptions(tool, subject, args, opts.fresh, opts.requireHuman) {
5854 c.approval.promptEmitMu.Unlock()
5855 return approvalReply{allow: true}, nil
5856 }
5857 var id string
5858 var reply chan approvalReply
5859 kind := ""
5860 if opts.fresh || opts.requireHuman || tool == planApprovalTool {
5861 if tool == planApprovalTool {
5862 kind = "plan"
5863 }
5864 id, reply = c.approval.registerDecisionKindWithInput(tool, subject, reason, args, opts.fresh, opts.requireHuman, kind, nil)
5865 ownerKind := PromptApproval
5866 if kind == "plan" {
5867 ownerKind = PromptPlan
5868 }
5869 c.registerOwnedPrompt(id, ownerKind)
5870 } else {
5871 id, reply = c.approval.registerWithInput(tool, subject, reason, args)
5872 c.registerOwnedPrompt(id, PromptApproval)
5873 }
5874
5875 if err := event.EmitChecked(c.sink, c.approvalRequestEvent(event.Approval{ID: id, Tool: tool, Subject: subject, Reason: reason, RawInput: append(json.RawMessage(nil), args...), Fresh: opts.fresh, Kind: kind})); err != nil {
5876 c.approval.promptEmitMu.Unlock()
5877 c.cancelOwnedPrompt(id)
5878 return approvalReply{}, fmt.Errorf("persist approval request: %w", err)
5879 }
5880 c.approval.promptEmitMu.Unlock()
5881 // The agent now needs the user's attention; a Notification hook can ping an
5882 // external channel (desktop notice, phone) while the run blocks on the reply.
5883 go c.hooks.Notification(ctx, approvalNotificationText(tool, subject), "permission_prompt")
5884
5885 waitCtx, cancelWait := c.approval.waitContext(ctx)
5886 defer cancelWait()
5887
5888 select {
5889 case r := <-reply:
5890 return r, nil
5891 case <-waitCtx.Done():
5892 c.cancelOwnedPrompt(id)
5893 return approvalReply{}, waitCtx.Err()
5894 }
5895 }
5896
5897 func (c *Controller) approvalRequestEvent(approval event.Approval) event.Event {
5898 approval.Generation = c.runtimeGeneration
5899 approval.PermissionRevision = c.permissionRevision.Load()
5900 if approval.TurnID == "" {
5901 approval.TurnID, _, _, _ = c.turnEventRuntimeStatus()
5902 }
5903 _, runtimeEpoch := c.promptIdentitySnapshot()
5904 if identity := c.bindOwnedPromptRouting(approval.ID, approval.TurnID, runtimeEpoch); identity.TurnID != "" {
5905 approval.TurnID = identity.TurnID
5906 }
5907 return event.Event{Kind: event.ApprovalRequest, TurnID: approval.TurnID, ItemID: approval.ID, Approval: approval}
5908 }
5909
5910 func (c *Controller) emitRememberResult(r RememberResult) {
5911 if r.Err != nil {
5912 c.sink.Emit(event.Event{
5913 Kind: event.Notice,
5914 Level: event.LevelWarn,
5915 Text: fmt.Sprintf(i18n.M.PermissionSaveFailedFmt, r.Rule, r.Err),
5916 })
5917 return
5918 }
5919 switch {
5920 case r.Saved:
5921 c.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelInfo, Text: fmt.Sprintf(i18n.M.PermissionSavedFmt, r.Path, r.Rule)})
5922 case strings.TrimSpace(r.CoveredBy) != "":
5923 c.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelInfo, Text: fmt.Sprintf(i18n.M.PermissionAlreadyAllowedFmt, r.Path, r.CoveredBy)})
5924 }
5925 }
5926
5927 func (c *Controller) emitPlanModeReadOnlyCommandTrustResult(r PlanModeReadOnlyCommandTrustResult) {
5928 prefix := strings.TrimSpace(r.Prefix)
5929 if r.Err != nil {
5930 c.sink.Emit(event.Event{
5931 Kind: event.Notice,
5932 Level: event.LevelWarn,
5933 Text: fmt.Sprintf(i18n.M.PlanModeReadOnlyCommandTrustFailedFmt, prefix, r.Err),
5934 })
5935 return
5936 }
5937 switch {
5938 case r.Saved:
5939 c.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelInfo, Text: fmt.Sprintf(i18n.M.PlanModeReadOnlyCommandTrustSavedFmt, r.Path, prefix)})
5940 case strings.TrimSpace(r.CoveredBy) != "":
5941 c.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelInfo, Text: fmt.Sprintf(i18n.M.PlanModeReadOnlyCommandTrustAlreadyFmt, r.Path, r.CoveredBy)})
5942 }
5943 }
5944
5944 lines GO