返回 DeepSeek-Reasonix
service.go
根目录 / internal / acp / service.go
1 package acp
2
3 import (
4 "context"
5 "crypto/rand"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "io"
10 "log/slog"
11 "maps"
12 "os"
13 "path/filepath"
14 "sort"
15 "strings"
16 "sync"
17 "time"
18
19 "reasonix/internal/agent"
20 "reasonix/internal/config"
21 "reasonix/internal/control"
22 "reasonix/internal/event"
23 "reasonix/internal/extension/uihub"
24 "reasonix/internal/fileutil"
25 fileencoding "reasonix/internal/fileutil/encoding"
26 "reasonix/internal/jobs"
27 "reasonix/internal/persistentshell"
28 "reasonix/internal/plugin"
29 "reasonix/internal/provider"
30 "reasonix/internal/sessioninbox"
31 "reasonix/internal/sessiontemp"
32 "reasonix/internal/store"
33 "reasonix/internal/tool/builtin"
34 )
35
36 // SessionParams is everything a Factory needs to assemble one ACP session's
37 // controller. Sink is owned by this package (an updateSink bound to the session
38 // id) and must be wired into the controller's event sink; the controller's
39 // interactive approval (see control.Controller.EnableInteractiveApproval) then
40 // routes "ask" decisions back through that sink as ApprovalRequest events, which
41 // the sink forwards to the client over session/request_permission.
42 //
43 // Cwd roots the session's file tools and bash (built via builtin.Workspace).
44 // Model, EffortOverride, and RuntimeProfile are optional session-local selectors
45 // from ACP config options. MCPServers are the MCP servers the client asked the
46 // agent to connect for this session. The path hooks keep service bookkeeping
47 // aligned; factories must wire both into the controller they build.
48 type SessionParams struct {
49 BackgroundScope *jobs.SessionBackgroundScope
50 SessionTemp *sessiontemp.Manager
51 PersistentShell *persistentshell.Manager
52 // MCPInteractions enables interactive MCP only after explicit client negotiation.
53 MCPInteractions bool
54 Cwd string
55 MCPServers []plugin.Spec
56 Sink event.Sink
57 Model string
58 EffortOverride *string
59 RuntimeProfile string
60 // NativeLegacySession asks the factory for a path-backed controller. It is
61 // set only when load/resume resolves an existing legacy transcript.
62 NativeLegacySession bool
63 OnSessionRecovered func(control.SessionRecoveryInfo) error
64 OnSessionTransition func(control.SessionTransitionInfo) error
65 // FileOverlay and Terminal are non-nil when the client advertised the
66 // matching capability at initialize: file tools then see unsaved editor
67 // buffers, and foreground bash can run in a client-owned terminal.
68 // Factories thread them into the controller's tool assembly.
69 FileOverlay builtin.FileOverlay
70 Terminal builtin.TerminalRunner
71 }
72
73 // Factory builds the per-session controller. The composition root (the cli's
74 // `reasonix acp` command) implements it by reusing setup()'s assembly: a
75 // Provider for Model, a tool Registry rooted at Cwd via builtin.Workspace, a
76 // per-session MCP host from MCPServers, the event Sink, all wired into a
77 // control.Controller. The returned controller owns its own cleanup (Close stops
78 // MCP subprocesses), so the service calls ctrl.Close() on teardown.
79 type Factory interface {
80 NewSession(ctx context.Context, p SessionParams) (*control.Controller, error)
81 }
82
83 // SessionConfigStateParams asks the Factory for normalized session config
84 // selectors. Empty Model and RuntimeProfile use configured defaults. Nil
85 // EffortOverride means provider config wins; a non-nil empty string means
86 // provider default for this session.
87 type SessionConfigStateParams struct {
88 Cwd string
89 Model string
90 EffortOverride *string
91 RuntimeProfile string
92 }
93
94 // SessionConfigState is the complete ACP-visible config state for a session.
95 type SessionConfigState struct {
96 Model string
97 EffortOverride *string
98 RuntimeProfile string
99 Models *SessionModelState
100 ConfigOptions []SessionConfigOption
101 }
102
103 // SessionConfigStateProvider lets a Factory expose model, effort, and work-mode
104 // selectors without making the ACP transport depend on a concrete config backend.
105 type SessionConfigStateProvider interface {
106 SessionConfigState(ctx context.Context, p SessionConfigStateParams) (SessionConfigState, error)
107 }
108
109 // SessionDirProvider lets a Factory expose the persistent session directory
110 // without forcing session/list to build a controller first.
111 type SessionDirProvider interface {
112 SessionDir() string
113 }
114
115 // SessionRebuilder lets a Factory rebuild a session's controller via
116 // boot.Rebuild: the replacement is built with the same boot.Options NewSession
117 // would use, and the session state (history, approval grants, goal/recovery,
118 // lifecycle) migrates off old inside the boot layer. The caller keeps the
119 // swap/close ordering. Factories that do not implement it leave
120 // _reasonix.io/session/reloadExtensions reporting unavailable.
121 type SessionRebuilder interface {
122 RebuildSession(ctx context.Context, p SessionParams, old *control.Controller) (*control.Controller, error)
123 }
124
125 // AgentInfo identifies this agent to clients in the initialize reply.
126 type AgentInfo struct {
127 Name string
128 Version string
129 }
130
131 // Serve runs an ACP agent on r/w (stdin/stdout in production) until the input
132 // ends or ctx is cancelled. It owns the JSON-RPC connection and the session
133 // registry; the Factory supplies the kernel wiring. This is the single entry
134 // point the `reasonix acp` command calls.
135 //
136 // stdout is the JSON-RPC channel: callers must keep all other output (logs,
137 // diagnostics) off w and on stderr, or the wire corrupts.
138 func Serve(ctx context.Context, r io.Reader, w io.Writer, factory Factory, info AgentInfo) error {
139 conn := NewConn(r, w)
140 svc := &service{
141 conn: conn,
142 factory: factory,
143 info: info,
144 sessions: make(map[string]*acpSession),
145 }
146 conn.Handle("initialize", svc.initialize)
147 conn.Handle("authenticate", svc.authenticate)
148 conn.Handle("session/new", svc.sessionNew)
149 conn.Handle("session/load", svc.sessionLoad)
150 conn.Handle("session/resume", svc.sessionResume)
151 conn.Handle("session/prompt", svc.sessionPrompt)
152 conn.Handle(sessionSteerMethod, svc.sessionSteer)
153 conn.Handle(sessionInboxEnqueueMethod, svc.sessionInboxEnqueue)
154 conn.Handle(sessionInboxListMethod, svc.sessionInboxList)
155 conn.Handle(sessionInboxGetMethod, svc.sessionInboxGet)
156 conn.Handle(sessionInboxUpdateMethod, svc.sessionInboxUpdate)
157 conn.Handle(sessionInboxDeleteMethod, svc.sessionInboxDelete)
158 conn.Handle(sessionInboxMoveMethod, svc.sessionInboxMove)
159 conn.Handle(sessionInboxPauseMethod, svc.sessionInboxSetPaused)
160 conn.Handle(sessionInboxRetryMethod, svc.sessionInboxRetry)
161 conn.Handle(sessionInboxRefreshMethod, svc.sessionInboxRefresh)
162 conn.Handle(sessionReloadExtensionsMethod, svc.sessionReloadExtensions)
163 conn.Handle(sessionStatusMethod, svc.sessionStatus)
164 conn.Handle("session/set_config_option", svc.sessionSetConfigOption)
165 conn.Handle("session/set_model", svc.sessionSetModel)
166 conn.Handle("session/set_mode", svc.sessionSetMode)
167 conn.Handle("session/close", svc.sessionClose)
168 conn.Handle("session/list", svc.sessionList)
169 conn.Handle("session/delete", svc.sessionDelete)
170 conn.HandleNotify("session/cancel", svc.sessionCancel)
171 defer svc.closeAll()
172 return conn.Serve(ctx)
173 }
174
175 // service holds the connection-wide ACP state: the factory, agent identity, and
176 // the live session registry.
177 type service struct {
178 conn *Conn
179 factory Factory
180 info AgentInfo
181
182 mu sync.Mutex
183 sessions map[string]*acpSession
184 // clientCaps is what the client offered at initialize (fs proxy, host
185 // terminals). Zero until initialize arrives; sessions opened later bind a
186 // clientIO built from it.
187 clientCaps ClientCapabilities
188 }
189
190 // afterResponse wraps a result with work that must run after the transport has
191 // successfully written that result. Session-opening notifications use this so a
192 // client can register the returned session before receiving its first update.
193 type afterResponse struct {
194 result any
195 after func()
196 }
197
198 func (r afterResponse) Response() any { return r.result }
199
200 func (r afterResponse) AfterResponse() {
201 if r.after != nil {
202 r.after()
203 }
204 }
205
206 func (s *service) setClientCapabilities(caps ClientCapabilities) {
207 s.mu.Lock()
208 s.clientCaps = caps
209 s.mu.Unlock()
210 }
211
212 func (s *service) clientCapabilities() ClientCapabilities {
213 s.mu.Lock()
214 defer s.mu.Unlock()
215 return s.clientCaps
216 }
217
218 // extensionSurfaceSupported reports whether the connected client advertised
219 // reasonix.extensionSurface support in its initialize handshake.
220 func (s *service) extensionSurfaceSupported() bool {
221 return clientExtensionSurfaceSupported(s.clientCapabilities())
222 }
223
224 // clientExtensionSurfaceSupported tolerantly parses the client's vendor
225 // capability block: _meta["reasonix.io"]["extensionSurface"]["supported"] must
226 // be an explicit true. Absent keys, wrong shapes, or a malformed block all
227 // mean unsupported — the sink then sends only the text fallback.
228 func clientExtensionSurfaceSupported(caps ClientCapabilities) bool {
229 vendor, ok := caps.Meta["reasonix.io"].(map[string]any)
230 if !ok {
231 return false
232 }
233 capability, ok := vendor["extensionSurface"].(map[string]any)
234 if !ok {
235 return false
236 }
237 supported, _ := capability["supported"].(bool)
238 return supported
239 }
240
241 // bindClientIO fills SessionParams' overlay/terminal fields from the client's
242 // declared capabilities. The nil checks keep absent capabilities as nil
243 // interface fields (a typed-nil *clientIO must never reach the interface).
244 func (s *service) bindClientIO(p *SessionParams, sessionID string) {
245 p.MCPInteractions = clientMCPInteractionSupported(s.clientCapabilities())
246 io := newClientIO(s.conn, sessionID, s.clientCapabilities())
247 if !io.hasAny() {
248 return
249 }
250 if fo := io.fileOverlay(); fo != nil {
251 p.FileOverlay = fo
252 }
253 if tr := io.terminalRunner(); tr != nil {
254 p.Terminal = tr
255 }
256 }
257
258 // acpController is the slice of the controller's driving port the ACP transport
259 // drives: session lifecycle + persistence, turn execution, interactive approval,
260 // and the capability surface (commands/skills/MCP prompts). ACP never touches
261 // goals, checkpoints, or memory, so it depends on those sub-ports only — not the
262 // concrete *control.Controller.
263 type acpController interface {
264 control.Lifecycle
265 control.TurnControl
266 RunTurnWithRaw(ctx context.Context, input, raw, invokedSkill string) error
267 RunSkillWithName(input string) (sent, name string, found bool)
268 RunFinalReadinessRecoveryWithAdmission(ctx context.Context, input string, onAdmitted func()) error
269 TrySteer(text string) bool
270 control.Approvals
271 control.Capabilities
272 control.SessionPersistence
273 // Goals backs ACP's normal/plan/goal collaboration-mode surface.
274 control.Goals
275 }
276
277 // acpSession is one open session: its controller, the on-disk transcript path
278 // (empty when persistence is off), and the cancel func of the in-flight turn
279 // (nil when idle) so session/cancel can abort it.
280 type acpSession struct {
281 id string
282 ctrl acpController
283 sink *updateSink
284 transcript string
285 cwd string
286 mcpServers []plugin.Spec
287 model string
288 // nil means use config; non-nil empty string means provider default.
289 effortOverride *string
290 runtimeProfile string
291 toolApprovalMode string
292 // runtimeState is the effective planner/sandbox posture captured after CLI
293 // hard overrides. status snapshots never reconstruct it from user config.
294 runtimeState SessionRuntimeState
295 status *statusTelemetry
296 // modeID is the ACP collaboration mode last reported to the client (normal |
297 // plan | goal). Goal draft mode turns the next user prompt into the goal.
298 // Both are guarded by mu; controller-side completion/plan exit is reconciled
299 // after each turn through current_mode_update.
300 modeID string
301 goalDraftMode bool
302 // pendingConfig queues config deltas requested while a turn or rebuild is
303 // in flight, holding at most one entry per axis: a later request replaces
304 // only its own axis (last-write-wins per axis), so a model change and a
305 // work-mode change queued back to back during one turn both survive to the
306 // drain instead of the second overwriting the first.
307 pendingConfig []sessionConfigDelta
308 // pendingReload coalesces _reasonix.io/session/reloadExtensions requests
309 // made while a turn or a rebuild is in flight; the finishTurn /
310 // post-maintenance drains run it once the session is idle.
311 pendingReload bool
312 title string
313 createdAt time.Time
314 updatedAt time.Time
315
316 mu sync.Mutex
317 // stateChangeMu serializes controller rebuilds with collaboration/approval
318 // changes so a swap cannot overwrite a newer user selection.
319 stateChangeMu sync.Mutex
320 cancel context.CancelFunc
321 done chan struct{}
322 running bool
323 deleted bool
324 // lease is the session lease guarding transcript against other runtimes
325 // (a desktop window, the CLI) for the life of this session. Held from
326 // session/new / session/load and released on close/delete/teardown.
327 // Config rebuilds keep the same transcript; when a snapshot conflict
328 // retargets the controller to a recovery branch, sessionRecoveredHandler
329 // moves transcript and this lease to the recovery file at commit time.
330 lease *agent.SessionLease
331 // retiredLeases tracks outgoing leases whose Release must run after the
332 // authority-guarded save that triggered a recovery callback returns. Any
333 // ACP operation that exposes a completed Snapshot waits for these channels,
334 // so callers never observe the old transcript as still owned after the
335 // handoff has completed.
336 retiredLeases []<-chan struct{}
337 // maintenanceDone is non-nil while session-owned maintenance, such as an
338 // idle config rebuild, is in flight outside mu.
339 maintenanceDone chan struct{}
340 }
341
342 func (s *acpSession) begin(ctx context.Context) (context.Context, context.CancelFunc, bool) {
343 runCtx, cancel := context.WithCancel(ctx)
344 // Prompt admission and config-axis changes share this lock. TryLock keeps
345 // ACP admission non-blocking while closing the idle-check/use window in an
346 // in-place role switch.
347 if !s.stateChangeMu.TryLock() {
348 cancel()
349 return nil, nil, false
350 }
351 defer s.stateChangeMu.Unlock()
352 s.mu.Lock()
353 // A queued pendingConfig blocks new turns so a prompt never runs on the
354 // outgoing config. The turn or maintenance that queued it applies it from
355 // its defer, so no new turn is needed to drain the queue.
356 if s.running || s.deleted || s.maintenanceDone != nil || len(s.pendingConfig) > 0 {
357 s.mu.Unlock()
358 cancel()
359 return nil, nil, false
360 }
361 s.running = true
362 s.cancel = cancel
363 s.done = make(chan struct{})
364 s.mu.Unlock()
365 return runCtx, cancel, true
366 }
367
368 func (s *acpSession) finish() {
369 s.mu.Lock()
370 done := s.done
371 s.running = false
372 s.cancel = nil
373 s.done = nil
374 s.mu.Unlock()
375 if done != nil {
376 close(done)
377 }
378 }
379
380 func (s *acpSession) abort() {
381 s.mu.Lock()
382 c := s.cancel
383 s.mu.Unlock()
384 if c != nil {
385 c()
386 }
387 }
388
389 func (s *acpSession) abortAndWait() {
390 s.mu.Lock()
391 c := s.cancel
392 done := s.done
393 maintenanceDone := s.maintenanceDone
394 s.mu.Unlock()
395 if c != nil {
396 c()
397 }
398 if done != nil {
399 <-done
400 }
401 if maintenanceDone != nil {
402 <-maintenanceDone
403 }
404 }
405
406 func (s *acpSession) deleteAndWait() {
407 s.mu.Lock()
408 s.deleted = true
409 c := s.cancel
410 done := s.done
411 maintenanceDone := s.maintenanceDone
412 s.mu.Unlock()
413 if c != nil {
414 c()
415 }
416 if done != nil {
417 <-done
418 }
419 if maintenanceDone != nil {
420 <-maintenanceDone
421 }
422 }
423
424 func (s *acpSession) finishMaintenance(done chan struct{}) {
425 if done == nil {
426 return
427 }
428 closeDone := false
429 s.mu.Lock()
430 if s.maintenanceDone == done {
431 s.maintenanceDone = nil
432 closeDone = true
433 }
434 s.mu.Unlock()
435 if closeDone {
436 close(done)
437 }
438 }
439
440 // swapModeID records the mode reported to the client and returns the previous
441 // value, so callers can emit current_mode_update only on change.
442 func (s *acpSession) swapModeID(id string) (old string) {
443 s.mu.Lock()
444 old = s.modeID
445 s.modeID = id
446 s.mu.Unlock()
447 return old
448 }
449
450 // currentModeID returns the mode last reported to the client.
451 func (s *acpSession) currentModeID() string {
452 s.mu.Lock()
453 defer s.mu.Unlock()
454 if s.modeID == "" {
455 return sessionModeNormal
456 }
457 return s.modeID
458 }
459
460 func (s *acpSession) setGoalDraftMode(on bool) {
461 s.mu.Lock()
462 s.goalDraftMode = on
463 s.mu.Unlock()
464 }
465
466 func (s *acpSession) takeGoalDraftMode() bool {
467 s.mu.Lock()
468 on := s.goalDraftMode
469 s.goalDraftMode = false
470 s.mu.Unlock()
471 return on
472 }
473
474 func (s *acpSession) isGoalDraftMode() bool {
475 s.mu.Lock()
476 defer s.mu.Unlock()
477 return s.goalDraftMode
478 }
479
480 func (s *acpSession) setToolApprovalMode(mode string) {
481 s.mu.Lock()
482 s.toolApprovalMode = normalizeACPToolApprovalMode(mode)
483 s.mu.Unlock()
484 }
485
486 func (s *acpSession) swapToolApprovalMode(mode string) (old string) {
487 mode = normalizeACPToolApprovalMode(mode)
488 s.mu.Lock()
489 old = normalizeACPToolApprovalMode(s.toolApprovalMode)
490 s.toolApprovalMode = mode
491 s.mu.Unlock()
492 return old
493 }
494
495 func (s *acpSession) saveMetaIfPresent() {
496 s.mu.Lock()
497 path := s.transcript
498 meta := s.metaLocked()
499 s.mu.Unlock()
500 if path != "" && sessionFileExists(path) {
501 _ = saveACPMeta(path, meta)
502 }
503 }
504
505 // currentCtrl returns the session's controller under mu. rebuildSession swaps
506 // ctrl while holding mu, so any read of the field outside mu races with a
507 // concurrent config rebuild; always go through this accessor unless mu is
508 // already held.
509 func (s *acpSession) currentCtrl() acpController {
510 s.mu.Lock()
511 defer s.mu.Unlock()
512 return s.ctrl
513 }
514
515 // releaseSessionLease drops the session's transcript lease, if any. Idempotent.
516 func (s *acpSession) releaseSessionLease() {
517 s.mu.Lock()
518 lease := s.lease
519 s.lease = nil
520 s.mu.Unlock()
521 if lease != nil {
522 lease.Release()
523 }
524 s.waitForRetiredSessionLeases()
525 }
526
527 // retireSessionLease defers Release until the authority-guarded save that
528 // invoked a recovery callback can return. Releasing synchronously inside that
529 // callback would wait on the very save executing the callback and deadlock.
530 func (s *acpSession) retireSessionLease(lease *agent.SessionLease) {
531 if lease == nil {
532 return
533 }
534 done := make(chan struct{})
535 s.mu.Lock()
536 s.retiredLeases = append(s.retiredLeases, done)
537 s.mu.Unlock()
538 go func() {
539 lease.Release()
540 close(done)
541 }()
542 }
543
544 func (s *acpSession) waitForRetiredSessionLeases() {
545 s.mu.Lock()
546 retired := append([]<-chan struct{}(nil), s.retiredLeases...)
547 s.mu.Unlock()
548 for _, done := range retired {
549 <-done
550 }
551 s.mu.Lock()
552 pending := s.retiredLeases[:0]
553 for _, done := range s.retiredLeases {
554 select {
555 case <-done:
556 default:
557 pending = append(pending, done)
558 }
559 }
560 s.retiredLeases = pending
561 s.mu.Unlock()
562 }
563
564 // sessionLeaseBindError maps a lease-acquisition failure to the protocol
565 // error the client sees: a held session names its holder with the shared CLI
566 // wording; anything else is an internal error.
567 func sessionLeaseBindError(method string, err error) *RPCError {
568 if errors.Is(err, agent.ErrSessionLeaseHeld) {
569 return &RPCError{
570 Code: ErrInvalidRequest,
571 Message: method + ": " + control.SessionInUseMessage(err) + "; " + control.SessionLeaseCloseHint,
572 }
573 }
574 return &RPCError{Code: ErrInternal, Message: method + ": session lease: " + err.Error()}
575 }
576
577 // initialize advertises the agent's capability set: persisted load plus ACP v1
578 // list/resume/close/delete lifecycle helpers, prompts carrying inline resource
579 // text (embeddedContext) but not image/audio, and stdio / Streamable HTTP MCP
580 // (no legacy sse).
581 func (s *service) initialize(_ context.Context, raw json.RawMessage) (any, error) {
582 var p InitializeParams
583 if len(raw) > 0 && json.Unmarshal(raw, &p) == nil {
584 s.setClientCapabilities(p.ClientCapabilities)
585 }
586 return InitializeResult{
587 ProtocolVersion: ProtocolVersion,
588 AgentCapabilities: AgentCapabilities{
589 LoadSession: true,
590 SessionCapabilities: SessionCapabilities{
591 List: &EmptyCapability{},
592 Resume: &EmptyCapability{},
593 Close: &EmptyCapability{},
594 Delete: &EmptyCapability{},
595 },
596 PromptCapabilities: PromptCapabilities{
597 Image: false,
598 Audio: false,
599 EmbeddedContext: true,
600 },
601 MCPCapabilities: MCPCapabilities{HTTP: true, SSE: false},
602 Meta: map[string]any{
603 "reasonix.io": ReasonixExtensionCapabilities{
604 MCPInteraction: &MCPInteractionCapability{Supported: true, SchemaVersion: 1, Method: mcpInteractionMethod},
605 SessionSteer: &SessionSteerCapability{Method: sessionSteerMethod},
606 SessionInbox: &SessionInboxCapability{
607 SchemaVersion: sessionInboxSchemaVersion,
608 Methods: map[string]string{
609 "enqueue": sessionInboxEnqueueMethod,
610 "list": sessionInboxListMethod,
611 "get": sessionInboxGetMethod,
612 "update": sessionInboxUpdateMethod,
613 "delete": sessionInboxDeleteMethod,
614 "move": sessionInboxMoveMethod,
615 "setPaused": sessionInboxPauseMethod,
616 "retry": sessionInboxRetryMethod,
617 "refresh": sessionInboxRefreshMethod,
618 },
619 },
620 SessionReloadExtensions: &SessionReloadExtensionsCapability{Method: sessionReloadExtensionsMethod},
621 ExtensionSurface: &ExtensionSurfaceCapability{Supported: true, SchemaVersion: reasonixExtensionSurfaceSchemaVersion},
622 },
623 sessionStatusMethod: ReasonixSchemaCapability{SchemaVersion: reasonixStatusSchemaVersion},
624 sessionStatusUpdateMethod: ReasonixSchemaCapability{SchemaVersion: reasonixStatusSchemaVersion},
625 },
626 },
627 AgentInfo: Implementation{Name: s.info.Name, Version: s.info.Version},
628 AuthMethods: []AuthMethod{reasonixSetupAuthMethod()},
629 }, nil
630 }
631
632 func reasonixSetupAuthMethod() AuthMethod {
633 return AuthMethod{
634 ID: "reasonix-setup",
635 Name: "Reasonix setup",
636 Description: "Configure Reasonix providers and credentials in a terminal",
637 Type: "terminal",
638 Args: []string{"setup"},
639 }
640 }
641
642 func (s *service) authenticate(_ context.Context, raw json.RawMessage) (any, error) {
643 var p AuthenticateParams
644 if err := json.Unmarshal(raw, &p); err != nil {
645 return nil, &RPCError{Code: ErrInvalidParams, Message: "authenticate: " + err.Error()}
646 }
647 if strings.TrimSpace(p.MethodID) != reasonixSetupAuthMethod().ID {
648 return nil, &RPCError{Code: ErrInvalidParams, Message: "authenticate: unknown methodId " + p.MethodID}
649 }
650 return AuthenticateResult{}, nil
651 }
652
653 // sessionNew opens a session: it mints an id, builds the session's sink bound to
654 // that id, asks the Factory to assemble the controller, switches the controller
655 // to interactive approval (so tool gates surface as ApprovalRequest events the
656 // sink forwards), and registers it.
657 func (s *service) sessionNew(ctx context.Context, raw json.RawMessage) (any, error) {
658 var p SessionNewParams
659 if len(raw) > 0 {
660 if err := json.Unmarshal(raw, &p); err != nil {
661 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/new: " + err.Error()}
662 }
663 }
664 cwd, err := s.resolveSessionCwd(p.Cwd, "")
665 if err != nil {
666 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/new: " + err.Error()}
667 }
668 mcpServers, err := mcpSpecs(p.MCPServers, cwd)
669 if err != nil {
670 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/new: " + err.Error()}
671 }
672 cfgState, err := s.sessionConfigState(ctx, SessionConfigStateParams{Cwd: cwd})
673 if err != nil {
674 return nil, &RPCError{Code: ErrInternal, Message: "session/new: " + err.Error()}
675 }
676 cfgState = withToolApprovalConfig(cfgState, control.ToolApprovalWorkspaceWrite)
677 runtimeState, err := s.sessionRuntimeState(ctx, SessionRuntimeStateParams{
678 Cwd: cwd, Model: cfgState.Model, RuntimeProfile: cfgState.RuntimeProfile,
679 })
680 if err != nil {
681 return nil, &RPCError{Code: ErrInternal, Message: "session/new: " + err.Error()}
682 }
683
684 id, err := newSessionID()
685 if err != nil {
686 return nil, &RPCError{Code: ErrInternal, Message: "session/new: " + err.Error()}
687 }
688
689 sink := newUpdateSink(s.conn, id)
690 sink.bindCwd(cwd)
691 sink.bindExtensionSurface(s.extensionSurfaceSupported())
692 sessionParams := SessionParams{
693 Cwd: cwd,
694 MCPServers: mcpServers,
695 Sink: sink,
696 Model: cfgState.Model,
697 EffortOverride: cloneStringPtr(cfgState.EffortOverride),
698 RuntimeProfile: cfgState.RuntimeProfile,
699 }
700 s.bindSessionClients(id, &sessionParams)
701 ctrl, err := s.factory.NewSession(ctx, sessionParams)
702 if err != nil {
703 return nil, &RPCError{Code: ErrInternal, Message: "session/new: " + err.Error()}
704 }
705 ctrl.EnableInteractiveApproval()
706 // The session metadata and advertised selector both start in workspace-write.
707 // Apply the same preset to the controller before admitting the first turn so
708 // a later same-value reconciliation cannot look like a permission change and
709 // cancel work that was already admitted under the advertised boundary.
710 ctrl.SetToolApprovalMode(control.ToolApprovalWorkspaceWrite)
711 sink.bindControllerPrompts(ctrl, sessionParams.MCPInteractions)
712
713 now := time.Now().UTC()
714 sess := &acpSession{
715 id: id,
716 ctrl: ctrl,
717 sink: sink,
718 cwd: cwd,
719 mcpServers: clonePluginSpecs(mcpServers),
720 model: cfgState.Model,
721 effortOverride: cloneStringPtr(cfgState.EffortOverride),
722 runtimeProfile: cfgState.RuntimeProfile,
723 toolApprovalMode: control.ToolApprovalWorkspaceWrite,
724 runtimeState: runtimeState,
725 status: newStatusTelemetry(),
726 modeID: sessionModeNormal,
727 createdAt: now,
728 updatedAt: now,
729 }
730 s.bindStatusEvents(sess)
731 // Exclusive v3 sessions bind the ACP id directly to the immutable storage
732 // identity. They never manufacture an id.jsonl transcript or acquire its
733 // legacy lease. Older factories retain the isolated compatibility path.
734 if ctrl.UsesExclusiveSession() {
735 if _, err := ctrl.BindFreshSession(ctx, id); err != nil {
736 ctrl.Close()
737 return nil, &RPCError{Code: ErrInternal, Message: "session/new: " + err.Error()}
738 }
739 } else if dir := ctrl.SessionDir(); dir != "" {
740 sess.transcript = transcriptPath(dir, id)
741 lease, err := agent.TryAcquireSessionLease(sess.transcript)
742 if err != nil {
743 ctrl.Close()
744 return nil, sessionLeaseBindError("session/new", err)
745 }
746 sess.lease = lease
747 ctrl.SetFreshSessionPath(sess.transcript)
748 if err := bindACPWriteAuthorityOrClose(ctrl, lease); err != nil {
749 sess.lease = nil
750 return nil, sessionLeaseBindError("session/new", err)
751 }
752 }
753
754 s.mu.Lock()
755 s.sessions[id] = sess
756 s.mu.Unlock()
757
758 // Fold in the live controller's extension catalog so plugin/... models
759 // are discoverable from the very first session/new result.
760 cfgState = enrichStateWithExtensionModels(cfgState, ctrl.ProviderCatalog())
761 return afterResponse{
762 result: SessionNewResult{
763 SessionID: id,
764 Models: cfgState.Models,
765 Modes: sessionModesState(sessionModeNormal),
766 ConfigOptions: cfgState.ConfigOptions,
767 },
768 after: func() { s.sendAvailableCommands(sess) },
769 }, nil
770 }
771
772 // Session modes exposed over ACP describe how the agent advances the task.
773 // Tool approval and runtime profile are independent config options. The legacy
774 // default/auto ids remain accepted for clients that used the old mixed axis.
775 const (
776 sessionModeNormal = "normal"
777 sessionModePlan = "plan"
778 sessionModeGoal = "goal"
779 sessionModeLegacyDefault = "default"
780 sessionModeLegacyAuto = "auto"
781 )
782
783 func sessionModesState(current string) *SessionModeState {
784 return &SessionModeState{
785 CurrentModeID: current,
786 AvailableModes: []SessionMode{
787 {ID: sessionModeNormal, Name: "Normal", Description: "Work directly and pause when user input is required"},
788 {ID: sessionModePlan, Name: "Plan", Description: "Research and propose a plan before making changes"},
789 {ID: sessionModeGoal, Name: "Goal", Description: "Keep advancing the next prompt as a goal until complete or blocked"},
790 },
791 }
792 }
793
794 // sessionSetMode switches the session's operating mode and confirms it with a
795 // current_mode_update, per the ACP session-mode contract.
796 func (s *service) sessionSetMode(ctx context.Context, raw json.RawMessage) (any, error) {
797 var p SessionSetModeParams
798 if err := json.Unmarshal(raw, &p); err != nil {
799 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/set_mode: " + err.Error()}
800 }
801 sess := s.session(p.SessionID)
802 if sess == nil {
803 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/set_mode: unknown session " + p.SessionID}
804 }
805 sess.stateChangeMu.Lock()
806 defer sess.stateChangeMu.Unlock()
807 ctrl := sess.currentCtrl()
808 nextMode, legacyApproval, rpcErr := applyACPSessionMode(ctrl, p.ModeID)
809 if rpcErr != nil {
810 return nil, rpcErr
811 }
812 // Entering Goal mode only arms a draft when no lifecycle exists. A restored,
813 // blocked, paused, or disarmed Goal must retain its complete objective so the
814 // user's next prompt can authorize recovery instead of replacing it.
815 sess.setGoalDraftMode(selectedGoalDraftMode(nextMode, ctrl.Goal()))
816 if legacyApproval != "" {
817 ctrl.SetToolApprovalMode(legacyApproval)
818 sess.setToolApprovalMode(legacyApproval)
819 if cfgState, err := s.configStateForSession(ctx, sess); err == nil {
820 sess.sink.send(configOptionUpdate{SessionUpdate: "config_option_update", ConfigOptions: cfgState.ConfigOptions})
821 }
822 }
823 if sess.swapModeID(nextMode) != nextMode {
824 sess.sink.send(currentModeUpdate{SessionUpdate: "current_mode_update", CurrentModeID: nextMode})
825 }
826 sess.saveMetaIfPresent()
827 return SessionSetModeResult{}, nil
828 }
829
830 // emitModeDrift reports controller-side mode flips (plan mode auto-exits when
831 // a plan is approved, a config rebuild resets switches) as current_mode_update
832 // so the client's mode picker stays truthful.
833 func (s *service) emitModeDrift(sess *acpSession) {
834 // Hold stateChangeMu across the controller read and the session-state swap:
835 // a session/set_mode completing between them (it holds this lock) would
836 // otherwise be read back as drift, roll the session's modeID and metadata
837 // back to the pre-selection value, and make the next rebuild re-apply that
838 // stale mode to the replacement controller.
839 sess.stateChangeMu.Lock()
840 defer sess.stateChangeMu.Unlock()
841 ctrl := sess.currentCtrl()
842 current := sessionModeNormal
843 switch {
844 case ctrl.PlanMode():
845 current = sessionModePlan
846 case ctrl.GoalStatus() == control.GoalStatusRunning || sess.isGoalDraftMode():
847 current = sessionModeGoal
848 }
849 if sess.swapModeID(current) != current {
850 sess.sink.send(currentModeUpdate{SessionUpdate: "current_mode_update", CurrentModeID: current})
851 sess.saveMetaIfPresent()
852 }
853 }
854
855 func (s *service) emitToolApprovalDrift(ctx context.Context, sess *acpSession) {
856 // Same contract as emitModeDrift: serialize with switchSessionToolApproval
857 // and rebuilds so a user selection landing between the controller read and
858 // the swap below is never reverted.
859 sess.stateChangeMu.Lock()
860 defer sess.stateChangeMu.Unlock()
861 current := normalizeACPToolApprovalMode(sess.currentCtrl().ToolApprovalMode())
862 if sess.swapToolApprovalMode(current) == current {
863 return
864 }
865 if cfgState, err := s.configStateForSession(ctx, sess); err == nil {
866 sess.sink.send(configOptionUpdate{SessionUpdate: "config_option_update", ConfigOptions: cfgState.ConfigOptions})
867 }
868 sess.saveMetaIfPresent()
869 }
870
871 // sessionLoad resumes a previously-saved session by id: it builds a controller
872 // (rooted at the requested cwd), seeds it from the on-disk transcript, replays
873 // the conversation to the client as session/update notifications, and registers
874 // it for subsequent prompts. A session already live in this process is replayed
875 // from memory without rebuilding.
876 func (s *service) sessionLoad(ctx context.Context, raw json.RawMessage) (any, error) {
877 var p SessionLoadParams
878 if err := json.Unmarshal(raw, &p); err != nil {
879 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/load: " + err.Error()}
880 }
881 cfgState, err := s.openExistingSession(ctx, "session/load", p.SessionID, p.Cwd, p.MCPServers, true)
882 if err != nil {
883 return nil, err
884 }
885 return afterResponse{
886 result: SessionLoadResult{Models: cfgState.Models, Modes: s.sessionModesFor(p.SessionID), ConfigOptions: cfgState.ConfigOptions},
887 after: func() { s.sendSessionProjection(s.session(p.SessionID)) },
888 }, nil
889 }
890
891 // sessionModesFor reports the modes state for a just-opened session. A live
892 // session keeps its current normal/plan/goal selection, so load/resume must not
893 // reset a reconnecting client's mode picker to normal.
894 func (s *service) sessionModesFor(id string) *SessionModeState {
895 if sess := s.session(id); sess != nil {
896 return sessionModesState(sess.currentModeID())
897 }
898 return sessionModesState(sessionModeNormal)
899 }
900
901 // sessionResume restores a previously-saved session without replaying its
902 // conversation history to the client.
903 func (s *service) sessionResume(ctx context.Context, raw json.RawMessage) (any, error) {
904 var p SessionResumeParams
905 if err := json.Unmarshal(raw, &p); err != nil {
906 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/resume: " + err.Error()}
907 }
908 cfgState, err := s.openExistingSession(ctx, "session/resume", p.SessionID, p.Cwd, p.MCPServers, false)
909 if err != nil {
910 return nil, err
911 }
912 return afterResponse{
913 result: SessionResumeResult{Models: cfgState.Models, Modes: s.sessionModesFor(p.SessionID), ConfigOptions: cfgState.ConfigOptions},
914 after: func() { s.sendSessionProjection(s.session(p.SessionID)) },
915 }, nil
916 }
917
918 // sendSessionProjection publishes current host-owned state after load/resume.
919 // History replay is presentation data and may contain legacy todo tool cards;
920 // the committed runtime snapshot is the only source of the current ACP plan.
921 func (s *service) sendSessionProjection(sess *acpSession) {
922 if sess == nil {
923 return
924 }
925 s.sendAvailableCommands(sess)
926 reader, ok := sess.currentCtrl().(control.RuntimeStateReader)
927 if !ok {
928 return
929 }
930 snapshot := reader.RuntimeStateSnapshot()
931 sess.sink.send(planUpdate{SessionUpdate: "plan", Entries: planEntriesFromTodos(snapshot.Todos)})
932 }
933
934 func (s *service) openExistingSession(ctx context.Context, method, id, cwdParam string, servers []MCPServerSpec, replay bool) (SessionConfigState, error) {
935 if err := validateSessionID(method, id); err != nil {
936 return SessionConfigState{}, err
937 }
938 cwd, err := s.resolveSessionCwd(cwdParam, id)
939 if err != nil {
940 return SessionConfigState{}, &RPCError{Code: ErrInvalidParams, Message: method + ": " + err.Error()}
941 }
942 mcpServers, err := mcpSpecs(servers, cwd)
943 if err != nil {
944 return SessionConfigState{}, &RPCError{Code: ErrInvalidParams, Message: method + ": " + err.Error()}
945 }
946
947 if sess := s.session(id); sess != nil {
948 if agent.IsCleanupPending(sess.transcript) {
949 return SessionConfigState{}, &RPCError{Code: ErrInvalidParams, Message: method + ": unknown session " + id}
950 }
951 if replay {
952 ctrl := sess.currentCtrl()
953 replaySink := newUpdateSink(s.conn, id)
954 replaySink.bindCwd(sess.cwd)
955 replaySink.replay(ctrl.History())
956 }
957 cfgState, err := s.configStateForSession(ctx, sess)
958 if err != nil {
959 return SessionConfigState{}, &RPCError{Code: ErrInternal, Message: method + ": " + err.Error()}
960 }
961 return cfgState, nil
962 }
963
964 var saved acpSessionMeta
965 persistedPath := ""
966 if dir := s.sessionDir(); dir != "" {
967 persistedPath = resolveTranscriptPath(dir, id)
968 if agent.IsCleanupPending(persistedPath) {
969 return SessionConfigState{}, &RPCError{Code: ErrInvalidParams, Message: method + ": unknown session " + id}
970 }
971 meta, _, metaErr := loadACPMeta(persistedPath)
972 if metaErr != nil {
973 return SessionConfigState{}, &RPCError{Code: ErrInternal, Message: method + ": " + metaErr.Error()}
974 }
975 saved = meta
976 }
977 cfgParams := SessionConfigStateParams{
978 Cwd: cwd,
979 Model: saved.Model,
980 EffortOverride: cloneStringPtr(saved.EffortOverride),
981 RuntimeProfile: saved.RuntimeProfile,
982 }
983 cfgState, err := s.sessionConfigState(ctx, cfgParams)
984 if err != nil && (strings.TrimSpace(saved.Model) != "" || saved.EffortOverride != nil || strings.TrimSpace(saved.RuntimeProfile) != "") {
985 cfgState, err = s.sessionConfigState(ctx, SessionConfigStateParams{Cwd: cwd})
986 }
987 if err != nil {
988 return SessionConfigState{}, &RPCError{Code: ErrInternal, Message: method + ": " + err.Error()}
989 }
990 runtimeState, err := s.sessionRuntimeState(ctx, SessionRuntimeStateParams{
991 Cwd: cwd, Model: cfgState.Model, RuntimeProfile: cfgState.RuntimeProfile,
992 })
993 if err != nil {
994 return SessionConfigState{}, &RPCError{Code: ErrInternal, Message: method + ": " + err.Error()}
995 }
996
997 sink := newUpdateSink(s.conn, id)
998 sink.bindCwd(cwd)
999 sink.bindExtensionSurface(s.extensionSurfaceSupported())
1000 sessionParams := SessionParams{
1001 Cwd: cwd,
1002 MCPServers: mcpServers,
1003 Sink: sink,
1004 Model: cfgState.Model,
1005 EffortOverride: cloneStringPtr(cfgState.EffortOverride),
1006 RuntimeProfile: cfgState.RuntimeProfile,
1007 NativeLegacySession: existingTranscript(persistedPath),
1008 }
1009 s.bindSessionClients(id, &sessionParams)
1010 ctrl, err := s.factory.NewSession(ctx, sessionParams)
1011 if err != nil {
1012 return SessionConfigState{}, &RPCError{Code: ErrInternal, Message: method + ": " + err.Error()}
1013 }
1014 ctrl.EnableInteractiveApproval()
1015 sink.bindControllerPrompts(ctrl, sessionParams.MCPInteractions)
1016
1017 path := ""
1018 var lease *agent.SessionLease
1019 canonicalExists, statErr := canonicalSessionExists(ctx, ctrl, id)
1020 if statErr != nil {
1021 ctrl.Close()
1022 return SessionConfigState{}, &RPCError{Code: ErrInternal, Message: method + ": " + statErr.Error()}
1023 }
1024 if canonicalExists {
1025 if err := openCanonicalSession(ctx, ctrl, id, method); err != nil {
1026 ctrl.Close()
1027 return SessionConfigState{}, err
1028 }
1029 } else {
1030 dir := ctrl.SessionDir()
1031 if dir == "" {
1032 ctrl.Close()
1033 return SessionConfigState{}, &RPCError{Code: ErrInternal, Message: method + ": persistence is disabled"}
1034 }
1035 path = resolveTranscriptPath(dir, id)
1036 if path != persistedPath && agent.IsCleanupPending(path) {
1037 ctrl.Close()
1038 return SessionConfigState{}, &RPCError{Code: ErrInvalidParams, Message: method + ": unknown session " + id}
1039 }
1040 // Legacy sessions keep the path lease until their one-time migration path
1041 // is selected by an explicit legacy client.
1042 var leaseErr error
1043 lease, leaseErr = agent.TryAcquireSessionLease(path)
1044 if leaseErr != nil {
1045 ctrl.Close()
1046 return SessionConfigState{}, sessionLeaseBindError(method, leaseErr)
1047 }
1048 loaded, loadErr := agent.LoadSession(path)
1049 if loadErr != nil {
1050 lease.Release()
1051 ctrl.Close()
1052 return SessionConfigState{}, &RPCError{Code: ErrInvalidParams, Message: method + ": unknown session " + id}
1053 }
1054 if err := resumeACPControllerForWrite(ctrl, loaded, path, lease); err != nil {
1055 return SessionConfigState{}, sessionLeaseBindError(method, err)
1056 }
1057 }
1058 toolApprovalMode := normalizeACPToolApprovalMode(saved.ToolApprovalMode)
1059 if strings.TrimSpace(saved.ToolApprovalMode) == "" {
1060 toolApprovalMode = control.ToolApprovalWorkspaceWrite
1061 }
1062 ctrl.SetToolApprovalMode(toolApprovalMode)
1063 modeID, goalDraftMode := applyLoadedACPMode(ctrl, saved.CollaborationMode)
1064
1065 meta := metadataForLoadedSession(path, id, cwd, ctrl.History())
1066 meta.Model = cfgState.Model
1067 meta.EffortOverride = cloneStringPtr(cfgState.EffortOverride)
1068 meta.RuntimeProfile = cfgState.RuntimeProfile
1069 meta.ToolApprovalMode = toolApprovalMode
1070 meta.CollaborationMode = modeID
1071 cfgState = withToolApprovalConfig(cfgState, toolApprovalMode)
1072 sess := &acpSession{
1073 id: id,
1074 ctrl: ctrl,
1075 sink: sink,
1076 transcript: path,
1077 cwd: meta.Cwd,
1078 mcpServers: clonePluginSpecs(mcpServers),
1079 model: cfgState.Model,
1080 effortOverride: cloneStringPtr(cfgState.EffortOverride),
1081 runtimeProfile: cfgState.RuntimeProfile,
1082 toolApprovalMode: toolApprovalMode,
1083 runtimeState: runtimeState,
1084 status: restoreStatusTelemetry(saved.Status),
1085 modeID: modeID,
1086 goalDraftMode: goalDraftMode,
1087 title: meta.Title,
1088 createdAt: meta.CreatedAt,
1089 updatedAt: meta.UpdatedAt,
1090 lease: lease,
1091 }
1092 s.bindStatusEvents(sess)
1093 if path != "" {
1094 if err := saveACPMeta(path, sess.meta()); err != nil {
1095 sess.releaseSessionLease()
1096 ctrl.Close()
1097 return SessionConfigState{}, &RPCError{Code: ErrInternal, Message: method + ": " + err.Error()}
1098 }
1099 }
1100 s.mu.Lock()
1101 s.sessions[id] = sess
1102 s.mu.Unlock()
1103
1104 if replay {
1105 sink.replay(ctrl.History())
1106 }
1107 return enrichStateWithExtensionModels(cfgState, ctrl.ProviderCatalog()), nil
1108 }
1109
1110 // transcriptPath is where a session's transcript lives — keyed by id so
1111 // session/load can recover it. Distinct from the cli's timestamp-labelled
1112 // chat/run session files (those are addressed by a picker, not by id).
1113 func transcriptPath(dir, id string) string {
1114 return filepath.Join(dir, id+".jsonl")
1115 }
1116
1117 // resolveTranscriptPath returns the transcript file session id currently
1118 // lives in. That is the id-keyed path by default; after a snapshot recovery
1119 // moved the live session onto a recovery branch, the id-keyed sidecar carries
1120 // an ActiveTranscript redirect (written by sessionRecoveredHandler) that
1121 // load/resume/delete/meta lookups must follow, or a restart silently reopens
1122 // the pre-recovery transcript. The redirect is a basename, must stay inside
1123 // dir, and its target must exist and claim the same session id; anything else
1124 // falls back to the id-keyed path.
1125 func resolveTranscriptPath(dir, id string) string {
1126 path := transcriptPath(dir, id)
1127 meta, ok, err := loadACPMeta(path)
1128 if err != nil || !ok {
1129 return path
1130 }
1131 active := strings.TrimSpace(meta.ActiveTranscript)
1132 if active == "" || active == filepath.Base(path) {
1133 return path
1134 }
1135 if filepath.Base(active) != active {
1136 return path
1137 }
1138 resolved := filepath.Join(dir, active)
1139 if !sessionFileExists(resolved) {
1140 return path
1141 }
1142 targetMeta, ok, err := loadACPMeta(resolved)
1143 if err != nil || !ok || targetMeta.SessionID != id {
1144 return path
1145 }
1146 return resolved
1147 }
1148
1149 // sessionPrompt runs one turn. It flattens the prompt blocks to text and runs the
1150 // session's controller synchronously under a per-turn cancelable context (so
1151 // session/cancel can stop it), then reports why the turn ended. The controller
1152 // streams the turn's events to the session's sink as it runs.
1153 func (s *service) sessionPrompt(ctx context.Context, raw json.RawMessage) (any, error) {
1154 var p SessionPromptParams
1155 if err := json.Unmarshal(raw, &p); err != nil {
1156 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/prompt: " + err.Error()}
1157 }
1158 sess := s.session(p.SessionID)
1159 if sess == nil {
1160 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/prompt: unknown session " + p.SessionID}
1161 }
1162 text := FlattenPrompt(p.Prompt)
1163 rawText := text
1164 invokedSkill := ""
1165 if text == "" && p.Action != control.ProtocolRecoveryAction {
1166 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/prompt: empty prompt"}
1167 }
1168 protocolRecovery := p.Action == control.ProtocolRecoveryAction
1169 if protocolRecovery && p.RecoveryID == "" {
1170 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/prompt: missing recoveryId"}
1171 }
1172 recovery := p.Action == control.FinalReadinessRecoveryAction
1173 if p.Action != "" && !recovery && !protocolRecovery {
1174 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/prompt: unsupported action " + p.Action}
1175 }
1176 if p.Action == "" {
1177 if id, guidance, ok := control.ParseProtocolRecoveryCommand(text); ok {
1178 protocolRecovery = true
1179 p.RecoveryID = id
1180 text = guidance
1181 } else if prompt, ok := control.ParseFinalReadinessRecoveryCommand(text); ok {
1182 recovery = true
1183 text = prompt
1184 } else {
1185 text, invokedSkill = s.resolveSlashPrompt(ctx, sess, text)
1186 }
1187 }
1188
1189 runCtx, cancel, ok := sess.begin(ctx)
1190 if !ok {
1191 return nil, &RPCError{Code: ErrInvalidRequest, Message: "session/prompt: session already has an active prompt"}
1192 }
1193 defer func() {
1194 sess.sink.clearTurnContext()
1195 s.finishTurn(ctx, sess)
1196 cancel()
1197 }()
1198 statusStarted := false
1199 if rpcErr := prepareACPGoalPrompt(sess, text); rpcErr != nil {
1200 return nil, rpcErr
1201 }
1202 beginTurn := func() {
1203 if sess.status == nil {
1204 sess.status = newStatusTelemetry()
1205 }
1206 sess.status.beginTurn()
1207 s.publishStatus(sess, "phase")
1208 sess.sink.setTurnContext(runCtx)
1209 statusStarted = true
1210 }
1211 var runErr error
1212 if protocolRecovery {
1213 if runner, ok := sess.ctrl.(interface {
1214 RunProtocolRecoveryWithAdmission(context.Context, string, string, func()) error
1215 }); ok {
1216 runErr = runner.RunProtocolRecoveryWithAdmission(runCtx, p.RecoveryID, text, beginTurn)
1217 } else {
1218 return nil, &RPCError{Code: ErrInvalidRequest, Message: "protocol recovery is unsupported by this controller"}
1219 }
1220 } else if recovery {
1221 runErr = sess.ctrl.RunFinalReadinessRecoveryWithAdmission(runCtx, text, beginTurn)
1222 } else {
1223 beginTurn()
1224 runErr = sess.ctrl.RunTurnWithRaw(runCtx, text, rawText, invokedSkill)
1225 }
1226 if errors.Is(runErr, agent.ErrProtocolRecoveryUnavailable) && !statusStarted {
1227 return nil, &RPCError{Code: ErrInvalidRequest, Message: "session/prompt: protocol recovery is unavailable or stale"}
1228 }
1229 if errors.Is(runErr, control.ErrNoFinalReadinessRecovery) && !statusStarted {
1230 return nil, &RPCError{
1231 Code: ErrInvalidRequest,
1232 Message: "session/prompt: no pending final-readiness check to continue",
1233 }
1234 }
1235 runErr = drainACPInbox(runCtx, sess.ctrl, runErr)
1236 cancelled := runCtx.Err() != nil
1237
1238 statusEvent := sess.status.finishTurn(
1239 runErr,
1240 cancelled,
1241 sess.currentCtrl().GoalStatus(),
1242 finalAssistantSummary(sess.currentCtrl()),
1243 )
1244 s.publishStatus(sess, statusEvent)
1245 // Persist after status finalization (best-effort) so reconnect recovers both
1246 // the transcript and the same sequence/usage/outcome snapshot.
1247 sess.persistAfterTurn(text)
1248
1249 stop, warning, promptErr := promptStopReason(runErr, cancelled, p.SessionID)
1250 if promptErr != nil {
1251 return nil, promptErr
1252 }
1253 if warning != "" {
1254 // The TUI keeps completed work for deliberate run boundaries; mirror
1255 // that for ACP and tell clients how the successful turn ended.
1256 sess.sink.Emit(promptPauseNotice(runErr, warning))
1257 }
1258 res := SessionPromptResult{StopReason: stop}
1259 if sess.transcript != "" {
1260 res.TranscriptPath = &sess.transcript
1261 }
1262 return res, nil
1263 }
1264
1265 // finalReadinessNotice is the warning text ACP clients receive when a completed
1266 // turn's final-readiness gate stays unsatisfied; the TUI shows the same gaps in
1267 // its recovery card.
1268 func finalReadinessNotice(e *agent.FinalReadinessError) string {
1269 const maxNoticeBytes = 2_048
1270 const fallback = "final-answer readiness gate not satisfied"
1271 if e == nil {
1272 return fallback
1273 }
1274 if reason := strings.TrimSpace(e.Reason); reason != "" {
1275 return clipStatusCredentialText(fallback+": "+reason, maxNoticeBytes)
1276 }
1277 return clipStatusError(e, maxNoticeBytes)
1278 }
1279
1280 // sessionSteer durably persists guidance then attempts mid-turn admission.
1281 // Parameter/session errors remain RPC errors; busy rejection returns a
1282 // disposition so clients can keep the durable follow-up.
1283 func (s *service) sessionSteer(_ context.Context, raw json.RawMessage) (any, error) {
1284 var p SessionSteerParams
1285 if err := json.Unmarshal(raw, &p); err != nil {
1286 return nil, &RPCError{Code: ErrInvalidParams, Message: sessionSteerMethod + ": " + err.Error()}
1287 }
1288 sess := s.session(p.SessionID)
1289 if sess == nil {
1290 return nil, &RPCError{Code: ErrInvalidParams, Message: sessionSteerMethod + ": unknown session " + p.SessionID}
1291 }
1292 text := FlattenPrompt(p.Prompt)
1293 if text == "" {
1294 return nil, &RPCError{Code: ErrInvalidParams, Message: sessionSteerMethod + ": empty prompt"}
1295 }
1296 ctrl := sess.currentCtrl()
1297 if api, ok := ctrl.(control.SessionAPI); ok {
1298 if ensurer, ok := any(api).(interface{ EnsureSessionPath() }); ok {
1299 ensurer.EnsureSessionPath()
1300 }
1301 // Durable path when the session has a transcript path; ephemeral
1302 // test controllers without persistence fall back to TrySteer.
1303 if api.SessionPath() != "" {
1304 rec, err := api.TryEnqueueAndSteer(control.InboxRequest{
1305 Intent: sessioninbox.IntentSteer,
1306 Display: text,
1307 Raw: text,
1308 Submit: text,
1309 Source: "acp",
1310 })
1311 if err != nil {
1312 return nil, &RPCError{Code: ErrInvalidRequest, Message: sessionSteerMethod + ": " + err.Error()}
1313 }
1314 return SessionSteerResult{ItemID: rec.ItemID, Disposition: string(rec.Disposition)}, nil
1315 }
1316 }
1317 // Compatibility for older controller stubs / pathless sessions.
1318 if !ctrl.TrySteer(text) {
1319 return nil, &RPCError{Code: ErrInvalidRequest, Message: sessionSteerMethod + ": session has no active prompt"}
1320 }
1321 return SessionSteerResult{Disposition: "steer_accepted"}, nil
1322 }
1323
1324 // sessionReloadExtensions rebuilds a session's agent runtime in place —
1325 // tools, skills, commands, hooks, MCP servers, and providers are re-discovered
1326 // — while the session (transcript, approval grants, goal and recovery state)
1327 // carries over via boot.Rebuild. It follows the same contract as a config
1328 // switch: a turn or rebuild in flight coalesces exactly one queued reload,
1329 // drained when the session goes idle; a failure keeps the old controller fully
1330 // usable; the old controller's resources are released only after the swap.
1331 func (s *service) sessionReloadExtensions(ctx context.Context, raw json.RawMessage) (any, error) {
1332 var p SessionReloadExtensionsParams
1333 if err := json.Unmarshal(raw, &p); err != nil {
1334 return nil, &RPCError{Code: ErrInvalidParams, Message: sessionReloadExtensionsMethod + ": " + err.Error()}
1335 }
1336 sess := s.session(p.SessionID)
1337 if sess == nil {
1338 return nil, &RPCError{Code: ErrInvalidParams, Message: sessionReloadExtensionsMethod + ": unknown session " + p.SessionID}
1339 }
1340 return s.reloadSessionExtensions(ctx, sess)
1341 }
1342
1343 func (s *service) reloadSessionExtensions(ctx context.Context, sess *acpSession) (any, error) {
1344 rebuilder, ok := s.factory.(SessionRebuilder)
1345 if !ok {
1346 return nil, &RPCError{Code: ErrInvalidRequest, Message: sessionReloadExtensionsMethod + ": runtime reload is unavailable in this session"}
1347 }
1348 if !sess.stateChangeMu.TryLock() {
1349 // A config switch or reload is in maintenance: coalesce one reload
1350 // behind it; the maintenance owner's post-maintenance drain runs it
1351 // (mirrors the pendingConfig queue contract in rebuildSession).
1352 sess.mu.Lock()
1353 if sess.maintenanceDone != nil && !sess.deleted {
1354 sess.pendingReload = true
1355 sess.mu.Unlock()
1356 return SessionReloadExtensionsResult{Queued: true}, nil
1357 }
1358 sess.mu.Unlock()
1359 sess.stateChangeMu.Lock()
1360 }
1361 didMaintenance := false
1362 res, err := s.reloadSessionExtensionsLocked(ctx, sess, rebuilder, &didMaintenance)
1363 sess.stateChangeMu.Unlock()
1364 if didMaintenance {
1365 s.reportPendingSessionConfigError(ctx, sess, s.applyPendingSessionConfig(ctx, sess), "after maintenance")
1366 s.drainPendingReload(ctx, sess)
1367 }
1368 return res, err
1369 }
1370
1371 // reloadSessionExtensionsLocked is reloadSessionExtensions' body; callers hold
1372 // stateChangeMu. The busy/queue checks and the publish/close ordering mirror
1373 // rebuildSessionLocked, but the build itself goes through the factory's
1374 // boot.Rebuild path instead of NewSession + manual migration.
1375 func (s *service) reloadSessionExtensionsLocked(ctx context.Context, sess *acpSession, rebuilder SessionRebuilder, didMaintenance *bool) (any, error) {
1376 sess.mu.Lock()
1377 if sess.deleted {
1378 sess.mu.Unlock()
1379 return nil, &RPCError{Code: ErrInvalidRequest, Message: sessionReloadExtensionsMethod + ": session is deleted"}
1380 }
1381 status := sess.ctrl.RuntimeStatus()
1382 if status.PendingPrompt {
1383 sess.mu.Unlock()
1384 return nil, sessionConfigActiveWorkError("answer pending prompts before reloading the runtime")
1385 }
1386 if !sess.running && !status.Running && status.BackgroundJobs > 0 {
1387 sess.mu.Unlock()
1388 return nil, sessionConfigActiveWorkError("stop background jobs before reloading the runtime")
1389 }
1390 if sess.running || status.Running || sess.maintenanceDone != nil {
1391 // Busy: coalesce exactly one reload; finishTurn (or the maintenance
1392 // owner's post-maintenance drain) runs it once the session is idle.
1393 sess.pendingReload = true
1394 sess.mu.Unlock()
1395 return SessionReloadExtensionsResult{Queued: true}, nil
1396 }
1397 // Claim the queued reload and raise maintenance in the same critical
1398 // section (mirrors rebuildSessionLocked): begin must never observe an
1399 // idle session between the two.
1400 sess.pendingReload = false
1401 cur := sess.ctrl
1402 sink := sess.sink
1403 mcpServers := clonePluginSpecs(sess.mcpServers)
1404 cwd := sess.cwd
1405 model := sess.model
1406 effortOverride := cloneStringPtr(sess.effortOverride)
1407 runtimeProfile := sess.runtimeProfile
1408 maintenanceDone := make(chan struct{})
1409 sess.maintenanceDone = maintenanceDone
1410 *didMaintenance = true
1411 sess.mu.Unlock()
1412 defer func() {
1413 sess.finishMaintenance(maintenanceDone)
1414 }()
1415
1416 if err := snapshotACPController(sess, cur); err != nil {
1417 return nil, &RPCError{Code: ErrInternal, Message: sessionReloadExtensionsMethod + ": snapshot before reload: " + err.Error()}
1418 }
1419 // Read the path only after Snapshot: a conflict can retarget cur to a
1420 // recovery branch, and boot.Rebuild binds the replacement to whatever
1421 // cur reports now (see rebuildSessionLocked). SessionPath is
1422 // controller-locked, so reading it off sess.mu is safe.
1423 prevPath := cur.SessionPath()
1424 old, ok := cur.(*control.Controller)
1425 if !ok {
1426 return nil, &RPCError{Code: ErrInternal, Message: sessionReloadExtensionsMethod + ": session controller does not support rebuild"}
1427 }
1428 rebuildParams := SessionParams{
1429 Cwd: cwd,
1430 MCPServers: mcpServers,
1431 Sink: sink,
1432 Model: model,
1433 EffortOverride: effortOverride,
1434 RuntimeProfile: runtimeProfile,
1435 NativeLegacySession: old.NativeLegacySession(),
1436 }
1437 s.bindSessionClients(sess.id, &rebuildParams)
1438 newCtrl, err := rebuilder.RebuildSession(ctx, rebuildParams, old)
1439 if err != nil {
1440 return nil, &RPCError{Code: ErrInternal, Message: sessionReloadExtensionsMethod + ": " + err.Error()}
1441 }
1442 newCtrl.EnableInteractiveApproval()
1443 // Config on disk may have changed the effective planner/sandbox posture;
1444 // recompute the status snapshot from the same resolved inputs.
1445 runtimeState, err := s.sessionRuntimeState(ctx, SessionRuntimeStateParams{
1446 Cwd: cwd, Model: model, RuntimeProfile: runtimeProfile,
1447 })
1448 if err != nil {
1449 newCtrl.ReleaseResources()
1450 return nil, &RPCError{Code: ErrInternal, Message: sessionReloadExtensionsMethod + ": runtime state: " + err.Error()}
1451 }
1452 // Persist before publishing the replacement. If this fails, the outgoing
1453 // controller and transcript still agree and remain fully usable (mirrors
1454 // the config switch).
1455 if err := s.prepareACPReplacementAuthority(sess, newCtrl, cur, prevPath, "snapshot after reload"); err != nil {
1456 newCtrl.ReleaseResources()
1457 return nil, &RPCError{Code: ErrInternal, Message: sessionReloadExtensionsMethod + ": " + err.Error()}
1458 }
1459
1460 sess.mu.Lock()
1461 if sess.deleted {
1462 sess.mu.Unlock()
1463 newCtrl.ReleaseResources()
1464 return nil, &RPCError{Code: ErrInvalidRequest, Message: sessionReloadExtensionsMethod + ": session is deleted"}
1465 }
1466 if sess.ctrl != cur {
1467 sess.mu.Unlock()
1468 newCtrl.ReleaseResources()
1469 return nil, sessionConfigActiveWorkError("session changed while reloading; retry")
1470 }
1471 oldCtrl, _ := cur.(*control.Controller)
1472 if err := control.ActivateControllerReplacement(oldCtrl, newCtrl); err != nil {
1473 sess.mu.Unlock()
1474 newCtrl.ReleaseResources()
1475 return nil, &RPCError{Code: ErrInternal, Message: sessionReloadExtensionsMethod + ": activate replacement: " + err.Error()}
1476 }
1477 sess.ctrl = newCtrl
1478 sess.runtimeState = runtimeState
1479 if sess.transcript != "" && sessionFileExists(sess.transcript) {
1480 _ = saveACPMeta(sess.transcript, sess.metaLocked())
1481 }
1482 sess.mu.Unlock()
1483 newCtrl.ActivateGoalDriverAfterRebuild()
1484 sink.bindControllerPrompts(newCtrl, rebuildParams.MCPInteractions)
1485
1486 // Release the outgoing controller only after the swap published the
1487 // replacement. ReleaseResources (not Close): the session logically
1488 // continues, so SessionEnd hooks must not fire — mirrors the config
1489 // switch.
1490 cur.ReleaseResources()
1491 // Clients see refreshed plugin commands without waiting for the next turn.
1492 s.sendAvailableCommands(sess)
1493 return SessionReloadExtensionsResult{}, nil
1494 }
1495
1496 // drainPendingReload runs the coalesced reloadExtensions request once the
1497 // session is idle. Called from finishTurn and after a config switch's or a
1498 // reload's own maintenance completes; callers must NOT hold stateChangeMu
1499 // (the reload re-acquires it).
1500 func (s *service) drainPendingReload(ctx context.Context, sess *acpSession) {
1501 if _, ok := s.factory.(SessionRebuilder); !ok {
1502 return
1503 }
1504 sess.mu.Lock()
1505 if !sess.pendingReload || sess.deleted || sess.running || sess.maintenanceDone != nil || len(sess.pendingConfig) > 0 {
1506 sess.mu.Unlock()
1507 return
1508 }
1509 sess.mu.Unlock()
1510 if _, err := s.reloadSessionExtensions(ctx, sess); err != nil {
1511 s.reportPendingSessionConfigError(ctx, sess, err, "after queued reload")
1512 }
1513 }
1514
1515 // finishTurn reconciles controller-side drift and drains any config switch
1516 // queued during the turn. Drift must be reconciled before finish() exposes
1517 // the session as idle: a concurrent config switch races on sess.running, and
1518 // if it wins that race while modeID/toolApprovalMode are still stale (a
1519 // slash command or plan/goal completion changed them inside the turn), it
1520 // rebuilds the replacement controller from the outgoing state instead of the
1521 // one this turn actually ended in.
1522 func (s *service) finishTurn(ctx context.Context, sess *acpSession) {
1523 s.emitModeDrift(sess)
1524 s.emitToolApprovalDrift(ctx, sess)
1525 sess.finish()
1526 s.reportPendingSessionConfigError(ctx, sess, s.applyPendingSessionConfig(ctx, sess), "after turn")
1527 // A reloadExtensions request queued during the turn runs now that the
1528 // session may be idle; the drain re-checks busy state.
1529 s.drainPendingReload(ctx, sess)
1530 // Re-check after a rebuild in case the replacement normalized state.
1531 s.emitModeDrift(sess)
1532 s.emitToolApprovalDrift(ctx, sess)
1533 }
1534
1535 // sessionSetConfigOption applies ACP's generic session-level selectors for
1536 // model, reasoning effort, work mode, and tool approval.
1537 func (s *service) sessionSetConfigOption(ctx context.Context, raw json.RawMessage) (any, error) {
1538 var p SetSessionConfigOptionParams
1539 if err := json.Unmarshal(raw, &p); err != nil {
1540 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/set_config_option: " + err.Error()}
1541 }
1542 sess := s.session(p.SessionID)
1543 if sess == nil {
1544 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/set_config_option: unknown session " + p.SessionID}
1545 }
1546 // Retired execution-mode IDs remain accepted for old clients, but they are
1547 // no longer advertised and never change the live controller.
1548 if id := normalizeConfigID(p.ConfigID); id == "work_mode" || id == "agent_preset" || id == "quality_floor" {
1549 if err := validateDeprecatedModeValue(p.Value); err != nil {
1550 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/set_config_option: " + err.Error()}
1551 }
1552 cfgState, err := s.configStateForSession(ctx, sess)
1553 if err != nil {
1554 return nil, &RPCError{Code: ErrInternal, Message: "session/set_config_option: " + err.Error()}
1555 }
1556 return SetSessionConfigOptionResult{
1557 ConfigOptions: cfgState.ConfigOptions,
1558 DeprecatedNotice: "Execution modes have been retired; this setting is accepted for compatibility and uses standard execution.",
1559 }, nil
1560 }
1561 cfgState, err := s.configStateForSession(ctx, sess)
1562 if err != nil {
1563 return nil, &RPCError{Code: ErrInternal, Message: "session/set_config_option: " + err.Error()}
1564 }
1565 option, ok := findConfigOption(cfgState.ConfigOptions, p.ConfigID)
1566 if !ok {
1567 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/set_config_option: unknown config option " + p.ConfigID}
1568 }
1569 if !configOptionHasValue(option, p.Value) {
1570 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/set_config_option: invalid value " + p.Value + " for " + option.ID}
1571 }
1572
1573 var next SessionConfigState
1574 switch configOptionCategory(option) {
1575 case "model":
1576 next, err = s.switchSessionModel(ctx, sess, p.Value)
1577 case "thought_level":
1578 next, err = s.switchSessionEffort(ctx, sess, p.Value)
1579 case "tool_approval":
1580 next, err = s.switchSessionToolApproval(ctx, sess, p.Value)
1581 default:
1582 err = &RPCError{Code: ErrInvalidParams, Message: "session/set_config_option: unsupported config option " + option.ID}
1583 }
1584 if err != nil {
1585 return nil, err
1586 }
1587 return SetSessionConfigOptionResult{ConfigOptions: next.ConfigOptions}, nil
1588 }
1589
1590 // sessionSetModel keeps older ACP clients working while configOptions becomes
1591 // the preferred model selector.
1592 func (s *service) sessionSetModel(ctx context.Context, raw json.RawMessage) (any, error) {
1593 var p SetSessionModelParams
1594 if err := json.Unmarshal(raw, &p); err != nil {
1595 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/set_model: " + err.Error()}
1596 }
1597 sess := s.session(p.SessionID)
1598 if sess == nil {
1599 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/set_model: unknown session " + p.SessionID}
1600 }
1601 if _, err := s.switchSessionModel(ctx, sess, p.ModelID); err != nil {
1602 return nil, err
1603 }
1604 return SetSessionModelResult{}, nil
1605 }
1606
1607 // resolveSessionConfigDeltas resolves deltas against the session's current
1608 // config baseline. Calling this fresh at every apply — instead of reusing a
1609 // snapshot taken when a change was first requested — is what keeps a queued
1610 // delta for one axis from clobbering another axis that rebuilt in between.
1611 func (s *service) resolveSessionConfigDeltas(ctx context.Context, sess *acpSession, deltas []sessionConfigDelta) (SessionConfigState, error) {
1612 params := sess.configStateParams()
1613 for _, delta := range deltas {
1614 delta.applyTo(&params)
1615 }
1616 cfgState, err := s.sessionConfigState(ctx, params)
1617 if err != nil {
1618 return SessionConfigState{}, err
1619 }
1620 return withToolApprovalConfig(cfgState, sess.currentToolApprovalMode()), nil
1621 }
1622
1623 func (s *service) switchSessionModel(ctx context.Context, sess *acpSession, modelID string) (SessionConfigState, error) {
1624 deltas := []sessionConfigDelta{{axis: "model", model: modelID}}
1625 return s.switchSessionConfig(ctx, sess, deltas)
1626 }
1627
1628 func (s *service) switchSessionEffort(ctx context.Context, sess *acpSession, effort string) (SessionConfigState, error) {
1629 level := strings.TrimSpace(effort)
1630 if level == "auto" {
1631 level = ""
1632 }
1633 deltas := []sessionConfigDelta{{axis: "thought_level", effortOverride: &level}}
1634 return s.switchSessionConfig(ctx, sess, deltas)
1635 }
1636
1637 func (s *service) switchSessionConfig(ctx context.Context, sess *acpSession, deltas []sessionConfigDelta) (SessionConfigState, error) {
1638 resolve := func() (SessionConfigState, error) {
1639 cfgState, err := s.resolveSessionConfigDeltas(ctx, sess, deltas)
1640 if err != nil {
1641 method := "session/set_config_option"
1642 if len(deltas) == 1 && deltas[0].axis == "model" {
1643 method = "session/set_model"
1644 }
1645 return SessionConfigState{}, &RPCError{Code: ErrInvalidParams, Message: method + ": " + err.Error()}
1646 }
1647 if len(deltas) == 1 && deltas[0].axis == "model" && cfgState.Model == "" {
1648 return SessionConfigState{}, &RPCError{Code: ErrInvalidRequest, Message: "session/set_model: model switching is unavailable in this session"}
1649 }
1650 return cfgState, nil
1651 }
1652
1653 if !sess.stateChangeMu.TryLock() {
1654 // Preserve the non-blocking queue contract while a rebuild is already in
1655 // maintenance. Resolve once for validation and the immediate client update;
1656 // the drain resolves the queued deltas again against live state.
1657 cfgState, err := resolve()
1658 if err != nil {
1659 return SessionConfigState{}, err
1660 }
1661 sess.mu.Lock()
1662 if sess.maintenanceDone != nil && !sess.deleted {
1663 for _, delta := range deltas {
1664 sess.pendingConfig = mergePendingConfig(sess.pendingConfig, delta)
1665 }
1666 sess.mu.Unlock()
1667 sess.sink.send(configOptionUpdate{SessionUpdate: "config_option_update", ConfigOptions: cfgState.ConfigOptions})
1668 return cfgState, nil
1669 }
1670 sess.mu.Unlock()
1671 sess.stateChangeMu.Lock()
1672 }
1673
1674 // Always resolve inside the serialization domain. Even a successful TryLock
1675 // can follow a concurrent rebuild that completed after this request began.
1676 cfgState, err := resolve()
1677 if err != nil {
1678 sess.stateChangeMu.Unlock()
1679 return SessionConfigState{}, err
1680 }
1681 didMaintenance := false
1682 err = s.rebuildSessionLocked(ctx, sess, cfgState, deltas, &didMaintenance)
1683 sess.stateChangeMu.Unlock()
1684 if didMaintenance {
1685 pendingErr := s.applyPendingSessionConfig(ctx, sess)
1686 s.reportPendingSessionConfigError(ctx, sess, pendingErr, "after maintenance")
1687 // A reloadExtensions request queued behind this maintenance runs next.
1688 s.drainPendingReload(ctx, sess)
1689 // The pending drain completes before this request returns. Refresh the RPC
1690 // result so an older response cannot overwrite the newer config_option_update
1691 // with the pre-drain full snapshot on the client.
1692 if current, stateErr := s.configStateForSession(ctx, sess); stateErr == nil {
1693 cfgState = current
1694 }
1695 }
1696 if err != nil {
1697 return SessionConfigState{}, err
1698 }
1699 return cfgState, nil
1700 }
1701
1702 func (s *service) switchSessionToolApproval(ctx context.Context, sess *acpSession, mode string) (SessionConfigState, error) {
1703 sess.stateChangeMu.Lock()
1704 defer sess.stateChangeMu.Unlock()
1705 mode = normalizeACPToolApprovalMode(mode)
1706 ctrl := sess.currentCtrl()
1707 ctrl.SetToolApprovalMode(mode)
1708 sess.setToolApprovalMode(mode)
1709 sess.saveMetaIfPresent()
1710 cfgState, err := s.configStateForSession(ctx, sess)
1711 if err != nil {
1712 return SessionConfigState{}, &RPCError{Code: ErrInternal, Message: "session/set_config_option: " + err.Error()}
1713 }
1714 sess.sink.send(configOptionUpdate{SessionUpdate: "config_option_update", ConfigOptions: cfgState.ConfigOptions})
1715 return cfgState, nil
1716 }
1717
1718 func (s *service) rebuildSession(ctx context.Context, sess *acpSession, cfgState SessionConfigState, deltas []sessionConfigDelta) error {
1719 if !sess.stateChangeMu.TryLock() {
1720 // Preserve the existing queue contract: a config change arriving during
1721 // a controller build returns immediately and is applied after that
1722 // build. The queue keeps one delta per axis (last-write-wins within an
1723 // axis), so changes queued for different axes never clobber each other.
1724 // Collaboration/approval changes do not use this queue; they wait for
1725 // the swap and then update the replacement controller.
1726 sess.mu.Lock()
1727 if sess.maintenanceDone != nil && !sess.deleted {
1728 for _, delta := range deltas {
1729 sess.pendingConfig = mergePendingConfig(sess.pendingConfig, delta)
1730 }
1731 sess.mu.Unlock()
1732 sess.sink.send(configOptionUpdate{SessionUpdate: "config_option_update", ConfigOptions: cfgState.ConfigOptions})
1733 return nil
1734 }
1735 sess.mu.Unlock()
1736 sess.stateChangeMu.Lock()
1737 }
1738 didMaintenance := false
1739 err := s.rebuildSessionLocked(ctx, sess, cfgState, deltas, &didMaintenance)
1740 sess.stateChangeMu.Unlock()
1741 if didMaintenance {
1742 pendingErr := s.applyPendingSessionConfig(ctx, sess)
1743 s.reportPendingSessionConfigError(ctx, sess, pendingErr, "after maintenance")
1744 // A reloadExtensions request queued behind this maintenance runs next.
1745 s.drainPendingReload(ctx, sess)
1746 }
1747 return err
1748 }
1749
1750 func (s *service) rebuildSessionLocked(ctx context.Context, sess *acpSession, cfgState SessionConfigState, deltas []sessionConfigDelta, didMaintenance *bool) (retErr error) {
1751 sess.mu.Lock()
1752 if sess.deleted {
1753 sess.mu.Unlock()
1754 return &RPCError{Code: ErrInvalidRequest, Message: "session config: session is deleted"}
1755 }
1756 status := sess.ctrl.RuntimeStatus()
1757 if status.PendingPrompt {
1758 sess.mu.Unlock()
1759 return sessionConfigActiveWorkError("answer pending prompts before switching config")
1760 }
1761 modelOnly := modelOnlyConfigDeltas(deltas)
1762 if !sess.running && !status.Running && status.BackgroundJobs > 0 && configBackgroundBlocked(sess.ctrl, modelOnly) {
1763 sess.mu.Unlock()
1764 return sessionConfigActiveWorkError("stop background jobs before switching config")
1765 }
1766 if sess.running || status.Running || sess.maintenanceDone != nil {
1767 for _, delta := range deltas {
1768 sess.pendingConfig = mergePendingConfig(sess.pendingConfig, delta)
1769 }
1770 sess.mu.Unlock()
1771 sess.sink.send(configOptionUpdate{SessionUpdate: "config_option_update", ConfigOptions: cfgState.ConfigOptions})
1772 return nil
1773 }
1774 // Claim this rebuild's axes from the queue in the same critical section
1775 // that raises maintenanceDone below: begin must never observe an idle
1776 // session between the two. Axes queued by other requests stay queued and
1777 // are drained by the post-maintenance apply.
1778 sess.pendingConfig = removePendingAxes(sess.pendingConfig, deltas)
1779
1780 cur := sess.ctrl
1781 sink := sess.sink
1782 mcpServers := clonePluginSpecs(sess.mcpServers)
1783 cwd := sess.cwd
1784 modeID := normalizeACPCollaborationMode(sess.modeID)
1785 goalDraftMode := sess.goalDraftMode
1786 toolApprovalMode := normalizeACPToolApprovalMode(sess.toolApprovalMode)
1787 if strings.TrimSpace(cfgState.RuntimeProfile) == "" {
1788 cfgState.RuntimeProfile = sess.runtimeProfile
1789 }
1790 maintenanceDone := make(chan struct{})
1791 sess.maintenanceDone = maintenanceDone
1792 *didMaintenance = true
1793 sess.mu.Unlock()
1794 defer func() {
1795 sess.finishMaintenance(maintenanceDone)
1796 }()
1797
1798 if err := snapshotACPController(sess, cur); err != nil {
1799 return &RPCError{Code: ErrInternal, Message: "session config: snapshot before switch: " + err.Error()}
1800 }
1801 // Capture the adopt path and history only after Snapshot: a snapshot
1802 // conflict can retarget cur to a recovery branch (or adopt the newer disk
1803 // transcript), and a pre-snapshot capture would bind the rebuilt controller
1804 // back to the original file, re-conflicting on every later save. When that
1805 // recovery fired, sessionRecoveredHandler already moved sess.transcript
1806 // and the session lease to the recovery file, so prevPath, the session
1807 // bookkeeping, and the controller agree on one path here.
1808 // SessionPath is controller-locked, so reading it off sess.mu is safe.
1809 prevPath := cur.SessionPath()
1810 carried := cur.History()
1811 carriedGoal := ""
1812 if cur.GoalStatus() == control.GoalStatusRunning {
1813 carriedGoal = cur.Goal()
1814 }
1815
1816 rebuildParams := SessionParams{
1817 Cwd: cwd,
1818 MCPServers: mcpServers,
1819 Sink: sink,
1820 Model: cfgState.Model,
1821 EffortOverride: cloneStringPtr(cfgState.EffortOverride),
1822 RuntimeProfile: cfgState.RuntimeProfile,
1823 NativeLegacySession: prevPath != "",
1824 }
1825 s.bindSessionClients(sess.id, &rebuildParams)
1826 newCtrl, err := s.buildConfigReplacement(ctx, cur, rebuildParams, modelOnly)
1827 if err != nil {
1828 return &RPCError{Code: ErrInternal, Message: "session config: " + err.Error()}
1829 }
1830 newCtrl.EnableInteractiveApproval()
1831 runtimeState, err := s.sessionRuntimeState(ctx, SessionRuntimeStateParams{
1832 Cwd: cwd, Model: cfgState.Model, RuntimeProfile: cfgState.RuntimeProfile,
1833 })
1834 if err != nil {
1835 newCtrl.ReleaseResources()
1836 return &RPCError{Code: ErrInternal, Message: "session config: runtime state: " + err.Error()}
1837 }
1838 // The freshly built controller's own leading system message carries the
1839 // target profile's contract (see boot/token_profile.go); AdoptHistory below
1840 // replaces the whole history with carried, so splice that message in first
1841 // or the model keeps seeing the outgoing profile's contract after every
1842 // switch.
1843 if fresh := newCtrl.History(); len(fresh) > 0 && fresh[0].Role == provider.RoleSystem {
1844 if len(carried) > 0 && carried[0].Role == provider.RoleSystem {
1845 carried[0] = fresh[0]
1846 } else {
1847 carried = append([]provider.Message{fresh[0]}, carried...)
1848 }
1849 }
1850 newCtrl.AdoptHistory(carried, prevPath)
1851 // Re-apply all three independent session axes. A controller rebuild must not
1852 // turn Plan into tool approval, drop a running Goal, or reset Ask/Auto/Yolo.
1853 newCtrl.SetToolApprovalMode(toolApprovalMode)
1854 switch modeID {
1855 case sessionModePlan:
1856 newCtrl.SetPlanMode(true)
1857 case sessionModeGoal:
1858 newCtrl.SetPlanMode(false)
1859 if carriedGoal != "" {
1860 newCtrl.SetGoal(carriedGoal)
1861 }
1862 default:
1863 newCtrl.SetPlanMode(false)
1864 }
1865 // InheritLifecycleFrom wires two concrete controllers' turn/hook state; it's a
1866 // construction concern, not part of the driving port. cur is always the
1867 // *control.Controller the factory built for this session, so this is safe.
1868 if rpcErr := inheritACPControllerLifecycle(newCtrl, cur); rpcErr != nil {
1869 newCtrl.ReleaseResources()
1870 return rpcErr
1871 }
1872 // Persist before publishing the replacement. If this fails, the outgoing
1873 // controller and transcript still agree and remain fully usable; publishing
1874 // first would report a successful switch whose refreshed profile contract
1875 // disappears on restart. AdoptHistory preserves the loaded CAS baseline, so
1876 // this compatible leading-system rewrite is safe to snapshot here.
1877 if err := s.prepareACPReplacementAuthority(sess, newCtrl, cur, prevPath, "snapshot after switch"); err != nil {
1878 newCtrl.ReleaseResources()
1879 return &RPCError{Code: ErrInternal, Message: "session config: " + err.Error()}
1880 }
1881
1882 sess.mu.Lock()
1883 if sess.deleted {
1884 sess.mu.Unlock()
1885 newCtrl.ReleaseResources()
1886 return &RPCError{Code: ErrInvalidRequest, Message: "session config: session is deleted"}
1887 }
1888 if sess.ctrl != cur {
1889 sess.mu.Unlock()
1890 newCtrl.ReleaseResources()
1891 return sessionConfigActiveWorkError("session changed while switching config; retry")
1892 }
1893 oldCtrl, _ := cur.(*control.Controller)
1894 if err := control.ActivateControllerReplacement(oldCtrl, newCtrl); err != nil {
1895 sess.mu.Unlock()
1896 newCtrl.ReleaseResources()
1897 return &RPCError{Code: ErrInternal, Message: "session config: activate replacement: " + err.Error()}
1898 }
1899 sess.ctrl = newCtrl
1900 sess.model = cfgState.Model
1901 sess.effortOverride = cloneStringPtr(cfgState.EffortOverride)
1902 sess.runtimeProfile = cfgState.RuntimeProfile
1903 sess.toolApprovalMode = toolApprovalMode
1904 sess.runtimeState = runtimeState
1905 sess.modeID = modeID
1906 sess.goalDraftMode = goalDraftMode
1907 if sess.transcript != "" && sessionFileExists(sess.transcript) {
1908 _ = saveACPMeta(sess.transcript, sess.metaLocked())
1909 }
1910 sess.mu.Unlock()
1911 newCtrl.ActivateGoalDriverAfterRebuild()
1912 sink.bindControllerPrompts(newCtrl, rebuildParams.MCPInteractions)
1913
1914 cur.ReleaseResources()
1915 s.sendAvailableCommands(sess)
1916 sink.send(configOptionUpdate{SessionUpdate: "config_option_update", ConfigOptions: cfgState.ConfigOptions})
1917 return nil
1918 }
1919
1920 func (s *service) applyPendingSessionConfig(ctx context.Context, sess *acpSession) error {
1921 var firstErr error
1922 for {
1923 if s.session(sess.id) != sess {
1924 return firstErr
1925 }
1926 // Claim the queue in the same serialization domain as explicit config
1927 // switches. Without this lock, a newer same-axis request can rebuild after
1928 // the clone below but before this apply starts, then the stale cloned delta
1929 // queues behind it and wins last instead of preserving request order.
1930 sess.stateChangeMu.Lock()
1931 didMaintenance := false
1932 sess.mu.Lock()
1933 if sess.deleted || len(sess.pendingConfig) == 0 {
1934 sess.mu.Unlock()
1935 sess.stateChangeMu.Unlock()
1936 return firstErr
1937 }
1938 deltas := clonePendingConfig(sess.pendingConfig)
1939 // Keep pendingConfig set while rebuilding: begin refuses new turns until
1940 // rebuildSession claims it together with raising maintenanceDone, so no
1941 // promptable instant is visible in between.
1942 sess.mu.Unlock()
1943
1944 // Re-resolve against the session's current state rather than reusing
1945 // whatever baseline existed when each delta queued: another axis may have
1946 // finished rebuilding in the meantime, and replaying its old value here
1947 // would silently roll it back. All queued axes resolve into one state so a
1948 // single rebuild applies them together.
1949 cfgState, err := s.resolveSessionConfigDeltas(ctx, sess, deltas)
1950 if err != nil {
1951 sess.mu.Lock()
1952 if !sess.deleted && !sess.running && sess.maintenanceDone == nil {
1953 sess.pendingConfig = removePendingAxes(sess.pendingConfig, deltas)
1954 }
1955 sess.mu.Unlock()
1956 sess.stateChangeMu.Unlock()
1957 if firstErr != nil {
1958 s.reportPendingSessionConfigError(ctx, sess, err, "after failed maintenance")
1959 return firstErr
1960 }
1961 return err
1962 }
1963
1964 err = s.rebuildSessionLocked(ctx, sess, cfgState, deltas, &didMaintenance)
1965 if err != nil && !didMaintenance {
1966 // Once this attempt failed nothing in flight is left to retry the
1967 // claimed axes, and begin refuses new turns while any are queued — drop
1968 // them so the session stays promptable. Once maintenance started, those
1969 // axes were already removed; anything queued now is a newer request and
1970 // must survive this failure.
1971 sess.mu.Lock()
1972 if !sess.deleted && !sess.running && sess.maintenanceDone == nil {
1973 sess.pendingConfig = removePendingAxes(sess.pendingConfig, deltas)
1974 }
1975 sess.mu.Unlock()
1976 }
1977 sess.stateChangeMu.Unlock()
1978
1979 if err != nil {
1980 if firstErr == nil {
1981 firstErr = err
1982 } else {
1983 s.reportPendingSessionConfigError(ctx, sess, err, "after failed maintenance")
1984 }
1985 if !didMaintenance {
1986 return firstErr
1987 }
1988 }
1989 if !didMaintenance {
1990 return firstErr
1991 }
1992 // Requests can queue while NewSession/Snapshot runs. Iterate even when this
1993 // rebuild failed so their already-successful RPCs cannot leave the session
1994 // blocked. A loop keeps sustained config traffic from growing the call stack.
1995 }
1996 }
1997
1998 func (s *service) reportPendingSessionConfigError(ctx context.Context, sess *acpSession, err error, when string) {
1999 if err == nil || sess == nil || sess.sink == nil {
2000 return
2001 }
2002 sess.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Text: "session config switch failed " + when + ": " + err.Error()})
2003 // A queued request already announced its desired config to the client. Every
2004 // apply failure leaves the outgoing controller/config active, so always send
2005 // the live state back; otherwise snapshot/build/resolve failures leave the
2006 // picker claiming a switch that never happened.
2007 if current, stateErr := s.configStateForSession(ctx, sess); stateErr == nil {
2008 sess.sink.send(configOptionUpdate{SessionUpdate: "config_option_update", ConfigOptions: current.ConfigOptions})
2009 }
2010 }
2011
2012 type activeSessionConfigWorkError struct {
2013 *RPCError
2014 }
2015
2016 func (e *activeSessionConfigWorkError) Unwrap() error {
2017 return e.RPCError
2018 }
2019
2020 func sessionConfigActiveWorkError(message string) error {
2021 return &activeSessionConfigWorkError{
2022 RPCError: &RPCError{Code: ErrInvalidRequest, Message: "session config: " + message},
2023 }
2024 }
2025
2026 // sessionClose releases an active session. Unknown sessions are accepted as a
2027 // no-op because closing is an idempotent resource cleanup request.
2028 func (s *service) sessionClose(_ context.Context, raw json.RawMessage) (any, error) {
2029 var p SessionCloseParams
2030 if err := json.Unmarshal(raw, &p); err != nil {
2031 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/close: " + err.Error()}
2032 }
2033 if err := validateSessionID("session/close", p.SessionID); err != nil {
2034 return nil, err
2035 }
2036 if sess := s.takeSession(p.SessionID); sess != nil {
2037 sess.abortAndWait()
2038 sess.ctrl.Close()
2039 sess.releaseSessionLease()
2040 }
2041 return SessionCloseResult{}, nil
2042 }
2043
2044 // sessionList returns ACP sessions known to this process or persisted as ACP
2045 // sidecars. It deliberately ignores ordinary CLI timestamp sessions.
2046 func (s *service) sessionList(_ context.Context, raw json.RawMessage) (any, error) {
2047 var p SessionListParams
2048 if len(raw) > 0 {
2049 if err := json.Unmarshal(raw, &p); err != nil {
2050 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/list: " + err.Error()}
2051 }
2052 }
2053 filterCwd := strings.TrimSpace(p.Cwd)
2054 if filterCwd != "" && !filepath.IsAbs(filterCwd) {
2055 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/list: cwd must be an absolute path"}
2056 }
2057 if strings.TrimSpace(p.Cursor) != "" {
2058 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/list: unsupported cursor"}
2059 }
2060
2061 byID := map[string]SessionInfo{}
2062 if dir := s.sessionDir(); dir != "" {
2063 metas, err := listACPMetas(dir)
2064 if err != nil {
2065 return nil, &RPCError{Code: ErrInternal, Message: "session/list: " + err.Error()}
2066 }
2067 // A recovered session has two sidecars claiming the same id: the
2068 // active recovery transcript's own meta and the id-keyed redirect.
2069 // Reduce to one representative per id before filtering, so the entry
2070 // shown never carries the stale pre-recovery title/timestamps.
2071 best := map[string]acpSessionMeta{}
2072 for _, meta := range metas {
2073 cur, ok := best[meta.SessionID]
2074 if !ok || listMetaBeats(meta, cur) {
2075 best[meta.SessionID] = meta
2076 }
2077 }
2078 for _, meta := range best {
2079 info := meta.info(nil)
2080 if sessionInfoMatchesCwd(info, filterCwd) {
2081 byID[info.SessionID] = info
2082 }
2083 }
2084 }
2085 for _, sess := range s.liveSessions() {
2086 info := sess.info()
2087 if sessionInfoMatchesCwd(info, filterCwd) {
2088 byID[info.SessionID] = info
2089 }
2090 }
2091
2092 sessions := make([]SessionInfo, 0, len(byID))
2093 for _, info := range byID {
2094 sessions = append(sessions, info)
2095 }
2096 sort.Slice(sessions, func(i, j int) bool {
2097 ti := parseSessionUpdatedAt(sessions[i].UpdatedAt)
2098 tj := parseSessionUpdatedAt(sessions[j].UpdatedAt)
2099 if ti.Equal(tj) {
2100 return sessions[i].SessionID < sessions[j].SessionID
2101 }
2102 return ti.After(tj)
2103 })
2104 return SessionListResult{Sessions: sessions}, nil
2105 }
2106
2107 // sessionDelete removes a session from future list results. Deleting a missing
2108 // session succeeds silently, matching ACP's idempotent delete guidance.
2109 func (s *service) sessionDelete(_ context.Context, raw json.RawMessage) (any, error) {
2110 var p SessionDeleteParams
2111 if err := json.Unmarshal(raw, &p); err != nil {
2112 return nil, &RPCError{Code: ErrInvalidParams, Message: "session/delete: " + err.Error()}
2113 }
2114 if err := validateSessionID("session/delete", p.SessionID); err != nil {
2115 return nil, err
2116 }
2117
2118 path := ""
2119 var destroy control.SessionDestroyHandle
2120 var delayed bool
2121 if sess := s.takeSession(p.SessionID); sess != nil {
2122 sess.deleteAndWait()
2123 // The session is going away; drop its lease before removing files so
2124 // the lease sidecars retire with the release (they are not in
2125 // SessionSidecarFiles and would otherwise linger).
2126 sess.releaseSessionLease()
2127 path = sess.transcript
2128 destroy = sess.ctrl.BeginDestroySession(path)
2129 if result := destroy.Wait(); result.HasTimedOut() {
2130 if err := agent.MarkCleanupPending(path, "delete"); err != nil {
2131 go delayedDeleteSessionFiles(path, destroy)
2132 sess.ctrl.CloseAfterDestroy()
2133 return nil, &RPCError{Code: ErrInternal, Message: "session/delete: " + err.Error()}
2134 }
2135 go delayedDeleteSessionFiles(path, destroy)
2136 delayed = true
2137 }
2138 sess.ctrl.CloseAfterDestroy()
2139 }
2140 if path == "" {
2141 if dir := s.sessionDir(); dir != "" {
2142 path = resolveTranscriptPath(dir, p.SessionID)
2143 }
2144 }
2145 if path != "" && !delayed {
2146 if err := deleteSessionFiles(path); err != nil {
2147 return nil, &RPCError{Code: ErrInternal, Message: "session/delete: " + err.Error()}
2148 }
2149 if destroy.Finish != nil {
2150 destroy.Finish()
2151 }
2152 }
2153 // A recovered session lives in two files: the recovery transcript (deleted
2154 // above) and the id-keyed original holding the redirect. Remove the twin
2155 // too, or it resurfaces in session/list as a ghost that delete-by-id can
2156 // never reach again.
2157 if dir := s.sessionDir(); dir != "" {
2158 if idPath := transcriptPath(dir, p.SessionID); idPath != path {
2159 if err := deleteSessionFiles(idPath); err != nil {
2160 return nil, &RPCError{Code: ErrInternal, Message: "session/delete: " + err.Error()}
2161 }
2162 }
2163 }
2164 return SessionDeleteResult{}, nil
2165 }
2166
2167 // sessionCancel aborts a session's in-flight turn, if any. It is a notification:
2168 // no reply, and an unknown session is silently ignored.
2169 func (s *service) sessionCancel(_ context.Context, raw json.RawMessage) {
2170 var p SessionCancelParams
2171 if err := json.Unmarshal(raw, &p); err != nil {
2172 return
2173 }
2174 if sess := s.session(p.SessionID); sess != nil {
2175 sess.abort()
2176 }
2177 }
2178
2179 func (s *service) session(id string) *acpSession {
2180 s.mu.Lock()
2181 defer s.mu.Unlock()
2182 return s.sessions[id]
2183 }
2184
2185 func (s *service) takeSession(id string) *acpSession {
2186 s.mu.Lock()
2187 defer s.mu.Unlock()
2188 sess := s.sessions[id]
2189 delete(s.sessions, id)
2190 return sess
2191 }
2192
2193 func (s *service) liveSessions() []*acpSession {
2194 s.mu.Lock()
2195 defer s.mu.Unlock()
2196 out := make([]*acpSession, 0, len(s.sessions))
2197 for _, sess := range s.sessions {
2198 out = append(out, sess)
2199 }
2200 return out
2201 }
2202
2203 func (s *service) sessionDir() string {
2204 if p, ok := s.factory.(SessionDirProvider); ok {
2205 if dir := strings.TrimSpace(p.SessionDir()); dir != "" {
2206 return dir
2207 }
2208 }
2209 s.mu.Lock()
2210 defer s.mu.Unlock()
2211 for _, sess := range s.sessions {
2212 if dir := sess.currentCtrl().SessionDir(); dir != "" {
2213 return dir
2214 }
2215 }
2216 return ""
2217 }
2218
2219 func (s *service) sessionConfigState(ctx context.Context, p SessionConfigStateParams) (SessionConfigState, error) {
2220 if provider, ok := s.factory.(SessionConfigStateProvider); ok {
2221 state, err := provider.SessionConfigState(ctx, p)
2222 if err != nil {
2223 return SessionConfigState{}, err
2224 }
2225 return withoutQualityFloorConfig(state), nil
2226 }
2227 return withoutQualityFloorConfig(SessionConfigState{}), nil
2228 }
2229
2230 func (s *service) configStateForSession(ctx context.Context, sess *acpSession) (SessionConfigState, error) {
2231 state, err := s.sessionConfigState(ctx, sess.configStateParams())
2232 if err != nil {
2233 return SessionConfigState{}, err
2234 }
2235 // Fold in the live controller's extension catalog so plugin/... models
2236 // are discoverable on every config-state read, not only when current.
2237 state = enrichStateWithExtensionModels(state, sess.currentCtrl().ProviderCatalog())
2238 state = withToolApprovalConfig(state, sess.currentToolApprovalMode())
2239 return withoutQualityFloorConfig(state), nil
2240 }
2241
2242 func (s *acpSession) configStateParams() SessionConfigStateParams {
2243 s.mu.Lock()
2244 defer s.mu.Unlock()
2245 return SessionConfigStateParams{
2246 Cwd: s.cwd,
2247 Model: s.model,
2248 EffortOverride: cloneStringPtr(s.effortOverride),
2249 RuntimeProfile: s.runtimeProfile,
2250 }
2251 }
2252
2253 func (s *acpSession) currentToolApprovalMode() string {
2254 s.mu.Lock()
2255 defer s.mu.Unlock()
2256 return normalizeACPToolApprovalMode(s.toolApprovalMode)
2257 }
2258
2259 func normalizeACPToolApprovalMode(mode string) string {
2260 return config.NormalizeToolApprovalMode(mode)
2261 }
2262
2263 func normalizeACPCollaborationMode(mode string) string {
2264 switch strings.ToLower(strings.TrimSpace(mode)) {
2265 case sessionModePlan:
2266 return sessionModePlan
2267 case sessionModeGoal:
2268 return sessionModeGoal
2269 default:
2270 return sessionModeNormal
2271 }
2272 }
2273
2274 func withToolApprovalConfig(state SessionConfigState, mode string) SessionConfigState {
2275 mode = normalizeACPToolApprovalMode(mode)
2276 option := SessionConfigOption{
2277 ID: "tool_approval",
2278 Name: "Permissions",
2279 Category: "tool_approval",
2280 Type: "select",
2281 CurrentValue: mode,
2282 Options: []SessionConfigSelectOption{
2283 {Value: control.ToolApprovalReadOnly, Name: "Read only", Description: "Read files; ask before writes and external side effects"},
2284 {Value: control.ToolApprovalWorkspaceWrite, Name: "Workspace access", Description: "Write inside the workspace and private session temp directory"},
2285 {Value: control.ToolApprovalDangerFullAccess, Name: "Full access", Description: "Skip ordinary prompts while explicit deny rules remain active"},
2286 },
2287 }
2288 for i := range state.ConfigOptions {
2289 if normalizeConfigID(state.ConfigOptions[i].ID) == option.ID {
2290 state.ConfigOptions[i] = option
2291 return state
2292 }
2293 }
2294 state.ConfigOptions = append(state.ConfigOptions, option)
2295 return state
2296 }
2297
2298 func findConfigOption(options []SessionConfigOption, id string) (SessionConfigOption, bool) {
2299 id = normalizeConfigID(id)
2300 for _, opt := range options {
2301 if normalizeConfigID(opt.ID) == id {
2302 return opt, true
2303 }
2304 }
2305 return SessionConfigOption{}, false
2306 }
2307
2308 // validateDeprecatedModeValue accepts the historical execution-mode vocabulary
2309 // so one-version-old clients get a precise error for garbage values while
2310 // well-formed values succeed as no-ops.
2311 func validateDeprecatedModeValue(value string) error {
2312 switch strings.ToLower(strings.TrimSpace(value)) {
2313 case "", "light", "economy", "eco", "lite", "save", "saving", "low", "minimal",
2314 "standard", "normal", "balanced", "full", "delivery", "deliver", "quality":
2315 return nil
2316 }
2317 return fmt.Errorf("invalid value %q for quality floor (accepted: standard, delivery; legacy light folds to standard)", value)
2318 }
2319
2320 func normalizeConfigID(id string) string {
2321 switch strings.TrimSpace(id) {
2322 case "models":
2323 return "model"
2324 case "reasoning_effort", "thought_level":
2325 return "effort"
2326 case "profile", "runtime_profile", "token_mode":
2327 return "work_mode"
2328 case "approval", "approval_mode", "tool_approval_mode":
2329 return "tool_approval"
2330 default:
2331 return strings.TrimSpace(id)
2332 }
2333 }
2334
2335 func configOptionHasValue(option SessionConfigOption, value string) bool {
2336 for _, opt := range option.Options {
2337 if opt.Value == value {
2338 return true
2339 }
2340 }
2341 return false
2342 }
2343
2344 func configOptionCategory(option SessionConfigOption) string {
2345 if option.Category != "" {
2346 return option.Category
2347 }
2348 switch normalizeConfigID(option.ID) {
2349 case "model":
2350 return "model"
2351 case "effort":
2352 return "thought_level"
2353 case "work_mode":
2354 return "work_mode"
2355 case "tool_approval":
2356 return "tool_approval"
2357 default:
2358 return ""
2359 }
2360 }
2361
2362 func cloneStringPtr(p *string) *string {
2363 if p == nil {
2364 return nil
2365 }
2366 cp := *p
2367 return &cp
2368 }
2369
2370 func clonePluginSpecs(in []plugin.Spec) []plugin.Spec {
2371 if len(in) == 0 {
2372 return nil
2373 }
2374 out := make([]plugin.Spec, len(in))
2375 copy(out, in)
2376 return out
2377 }
2378
2379 func (s *service) resolveSessionCwd(cwd, sessionID string) (string, error) {
2380 cwd = strings.TrimSpace(cwd)
2381 if cwd != "" {
2382 if !filepath.IsAbs(cwd) {
2383 return "", fmt.Errorf("cwd must be an absolute path")
2384 }
2385 return filepath.Clean(cwd), nil
2386 }
2387 if sessionID != "" {
2388 if meta, ok := s.loadMeta(sessionID); ok && meta.Cwd != "" {
2389 if !filepath.IsAbs(meta.Cwd) {
2390 return "", fmt.Errorf("stored cwd must be an absolute path")
2391 }
2392 return filepath.Clean(meta.Cwd), nil
2393 }
2394 }
2395 wd, err := os.Getwd()
2396 if err != nil {
2397 return "", fmt.Errorf("resolve cwd: %w", err)
2398 }
2399 return wd, nil
2400 }
2401
2402 func (s *service) loadMeta(id string) (acpSessionMeta, bool) {
2403 dir := s.sessionDir()
2404 if dir == "" {
2405 return acpSessionMeta{}, false
2406 }
2407 meta, ok, err := loadACPMeta(resolveTranscriptPath(dir, id))
2408 if err != nil {
2409 return acpSessionMeta{}, false
2410 }
2411 return meta, ok
2412 }
2413
2414 // closeAll tears down every open session (aborting any in-flight turn and
2415 // stopping its MCP subprocesses) when the connection ends.
2416 func (s *service) closeAll() {
2417 s.mu.Lock()
2418 sessions := s.sessions
2419 s.sessions = make(map[string]*acpSession)
2420 s.mu.Unlock()
2421 for _, sess := range sessions {
2422 sess.abortAndWait()
2423 sess.currentCtrl().Close()
2424 sess.releaseSessionLease()
2425 }
2426 }
2427
2428 func (s *acpSession) persistAfterTurn(prompt string) {
2429 s.mu.Lock()
2430 if s.deleted {
2431 s.mu.Unlock()
2432 return
2433 }
2434 ctrl := s.ctrl
2435 s.mu.Unlock()
2436
2437 _ = snapshotACPController(s, ctrl)
2438
2439 s.mu.Lock()
2440 defer s.mu.Unlock()
2441 if s.deleted || s.ctrl != ctrl {
2442 return
2443 }
2444 if s.title == "" {
2445 s.title = previewTitle(prompt)
2446 }
2447 s.updatedAt = time.Now().UTC()
2448 if s.createdAt.IsZero() {
2449 s.createdAt = s.updatedAt
2450 }
2451 if s.transcript != "" && sessionFileExists(s.transcript) {
2452 _ = saveACPMeta(s.transcript, s.metaLocked())
2453 }
2454 }
2455
2456 func (s *acpSession) meta() acpSessionMeta {
2457 s.mu.Lock()
2458 defer s.mu.Unlock()
2459 return s.metaLocked()
2460 }
2461
2462 func (s *acpSession) metaLocked() acpSessionMeta {
2463 return acpSessionMeta{
2464 SessionID: s.id,
2465 Cwd: s.cwd,
2466 Model: s.model,
2467 EffortOverride: cloneStringPtr(s.effortOverride),
2468 RuntimeProfile: s.runtimeProfile,
2469 ToolApprovalMode: normalizeACPToolApprovalMode(s.toolApprovalMode),
2470 CollaborationMode: normalizeACPCollaborationMode(s.modeID),
2471 Title: s.title,
2472 CreatedAt: s.createdAt,
2473 UpdatedAt: s.updatedAt,
2474 Status: s.status.persisted(),
2475 }
2476 }
2477
2478 func (s *acpSession) info() SessionInfo {
2479 meta := s.meta()
2480 ctrl := s.currentCtrl()
2481 extra := map[string]any{}
2482 if n := len(ctrl.History()); n > 0 {
2483 extra["messageCount"] = n
2484 }
2485 if len(extra) == 0 {
2486 extra = nil
2487 }
2488 return meta.info(extra)
2489 }
2490
2491 func (s *service) sendAvailableCommands(sess *acpSession) {
2492 if sess == nil {
2493 return
2494 }
2495 ctrl := sess.currentCtrl()
2496 if ctrl == nil {
2497 return
2498 }
2499 cmds := availableCommandsFor(ctrl)
2500 if len(cmds) == 0 {
2501 return
2502 }
2503 sess.sink.send(availableCommandsUpdate{
2504 SessionUpdate: "available_commands_update",
2505 AvailableCommands: cmds,
2506 })
2507 }
2508
2509 func availableCommandsFor(ctrl acpController) []AvailableCommand {
2510 if ctrl == nil {
2511 return nil
2512 }
2513 byName := map[string]AvailableCommand{}
2514 for _, cmd := range ctrl.Commands() {
2515 if cmd.Hidden {
2516 continue
2517 }
2518 name := strings.TrimSpace(cmd.Name)
2519 if name == "" {
2520 continue
2521 }
2522 desc := strings.TrimSpace(cmd.Description)
2523 if desc == "" {
2524 desc = "Run the " + name + " command"
2525 }
2526 ac := AvailableCommand{Name: name, Description: desc}
2527 if hint := strings.TrimSpace(cmd.ArgHint); hint != "" {
2528 ac.Input = &AvailableCommandInput{Hint: hint}
2529 }
2530 byName[name] = ac
2531 }
2532 for _, sk := range ctrl.SlashSkills() {
2533 name := strings.TrimSpace(sk.SlashName())
2534 if name == "" {
2535 continue
2536 }
2537 if _, exists := byName[name]; exists {
2538 continue
2539 }
2540 desc := strings.TrimSpace(sk.Description)
2541 if desc == "" {
2542 desc = "Run the " + name + " skill"
2543 }
2544 byName[name] = AvailableCommand{
2545 Name: name,
2546 Description: desc,
2547 Input: &AvailableCommandInput{Hint: "instructions"},
2548 }
2549 }
2550 if host := ctrl.Host(); host != nil {
2551 for _, prompt := range host.Prompts() {
2552 name := strings.TrimSpace(prompt.Name)
2553 if name == "" {
2554 continue
2555 }
2556 desc := strings.TrimSpace(prompt.Description)
2557 if desc == "" {
2558 desc = "Run the " + name + " MCP prompt"
2559 }
2560 ac := AvailableCommand{Name: name, Description: desc}
2561 if len(prompt.Args) > 0 {
2562 ac.Input = &AvailableCommandInput{Hint: "arguments"}
2563 }
2564 byName[name] = ac
2565 }
2566 }
2567 // Extension actions surface as "<plugin>:<action>" commands so ACP clients
2568 // can discover them in the slash menu alongside commands/skills/prompts.
2569 for _, action := range ctrl.ExtensionActions() {
2570 name := strings.TrimPrefix(strings.TrimSpace(action.Slash), "/")
2571 if name == "" {
2572 continue
2573 }
2574 if _, exists := byName[name]; exists {
2575 continue
2576 }
2577 desc := strings.TrimSpace(action.Label)
2578 if desc == "" {
2579 desc = "Run the " + name + " extension action"
2580 }
2581 byName[name] = AvailableCommand{
2582 Name: name,
2583 Description: desc,
2584 Input: &AvailableCommandInput{Hint: "arguments"},
2585 }
2586 }
2587 out := make([]AvailableCommand, 0, len(byName))
2588 for _, cmd := range byName {
2589 out = append(out, cmd)
2590 }
2591 sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name })
2592 return out
2593 }
2594
2595 func (s *service) resolveSlashPrompt(ctx context.Context, sess *acpSession, text string) (string, string) {
2596 line := strings.TrimSpace(text)
2597 if sess == nil || !strings.HasPrefix(line, "/") {
2598 return text, ""
2599 }
2600 ctrl := sess.currentCtrl()
2601 if ctrl == nil {
2602 return text, ""
2603 }
2604 if sent, ok := ctrl.CustomCommand(line); ok {
2605 return sent, ""
2606 }
2607 if sent, name, ok := ctrl.RunSkillWithName(line); ok {
2608 return sent, name
2609 }
2610 if sent, ok, err := ctrl.MCPPrompt(ctx, line); err == nil && ok {
2611 return sent, ""
2612 }
2613 if sent, ok := invokeExtensionAction(ctx, ctrl, line); ok {
2614 return sent, ""
2615 }
2616 return text, ""
2617 }
2618
2619 // invokeExtensionAction resolves a "/<plugin>:<action> args…" line against the
2620 // handshake-declared extension actions and invokes it — the last resolution
2621 // step in resolveSlashPrompt, after custom commands, skills, and MCP prompts.
2622 // The extension's result message becomes the prompt text. A parse miss, an
2623 // undeclared action, an invocation error, or an empty result all leave the
2624 // line untouched (ok=false), matching how unknown slash commands fall through.
2625 func invokeExtensionAction(ctx context.Context, ctrl acpController, line string) (string, bool) {
2626 fields := strings.Fields(line)
2627 if len(fields) == 0 {
2628 return "", false
2629 }
2630 pluginID, actionID, ok := uihub.ParseSlashName(fields[0])
2631 if !ok {
2632 return "", false
2633 }
2634 declared := false
2635 for _, action := range ctrl.ExtensionActions() {
2636 if action.PluginID == pluginID && action.ActionID == actionID {
2637 declared = true
2638 break
2639 }
2640 }
2641 if !declared {
2642 return "", false
2643 }
2644 message, err := ctrl.InvokeExtensionAction(ctx, fields[0], control.ParseExtensionActionArgs(fields[1:]))
2645 if err != nil || strings.TrimSpace(message) == "" {
2646 return "", false
2647 }
2648 return message, true
2649 }
2650
2651 type acpSessionMeta struct {
2652 SessionID string `json:"sessionId"`
2653 Cwd string `json:"cwd"`
2654 Model string `json:"model,omitempty"`
2655 EffortOverride *string `json:"effortOverride,omitempty"`
2656 RuntimeProfile string `json:"runtimeProfile,omitempty"`
2657 ToolApprovalMode string `json:"toolApprovalMode,omitempty"`
2658 CollaborationMode string `json:"collaborationMode,omitempty"`
2659 Title string `json:"title,omitempty"`
2660 CreatedAt time.Time `json:"createdAt"`
2661 UpdatedAt time.Time `json:"updatedAt"`
2662 Status *persistedStatusTelemetry `json:"status,omitempty"`
2663 // ActiveTranscript, when set on the id-keyed sidecar, is the basename of
2664 // the transcript this session currently lives in: a snapshot recovery
2665 // moved the live session onto a recovery branch and left this redirect
2666 // behind so restart-time lookups (resolveTranscriptPath) follow the
2667 // session instead of reopening the pre-recovery file.
2668 ActiveTranscript string `json:"activeTranscript,omitempty"`
2669 }
2670
2671 func (m acpSessionMeta) info(extra map[string]any) SessionInfo {
2672 updatedAt := ""
2673 if !m.UpdatedAt.IsZero() {
2674 updatedAt = m.UpdatedAt.Format(time.RFC3339Nano)
2675 }
2676 return SessionInfo{
2677 SessionID: m.SessionID,
2678 Cwd: m.Cwd,
2679 Title: m.Title,
2680 UpdatedAt: updatedAt,
2681 Meta: extra,
2682 }
2683 }
2684
2685 func metadataForLoadedSession(path, id, cwd string, history []provider.Message) acpSessionMeta {
2686 now := time.Now().UTC()
2687 meta, ok, err := loadACPMeta(path)
2688 if err != nil || !ok {
2689 meta = acpSessionMeta{
2690 SessionID: id,
2691 Cwd: cwd,
2692 Title: titleFromHistory(history),
2693 CreatedAt: now,
2694 UpdatedAt: now,
2695 }
2696 if info, statErr := os.Stat(path); statErr == nil {
2697 meta.CreatedAt = info.ModTime().UTC()
2698 meta.UpdatedAt = info.ModTime().UTC()
2699 }
2700 }
2701 if meta.SessionID == "" {
2702 meta.SessionID = id
2703 }
2704 if cwd != "" {
2705 meta.Cwd = cwd
2706 }
2707 if meta.Title == "" {
2708 meta.Title = titleFromHistory(history)
2709 }
2710 if meta.CreatedAt.IsZero() {
2711 meta.CreatedAt = now
2712 }
2713 if meta.UpdatedAt.IsZero() {
2714 meta.UpdatedAt = meta.CreatedAt
2715 }
2716 return meta
2717 }
2718
2719 func loadACPMeta(sessionPath string) (acpSessionMeta, bool, error) {
2720 path := acpMetaPath(sessionPath)
2721 if path == "" {
2722 return acpSessionMeta{}, false, nil
2723 }
2724 b, err := fileencoding.ReadFileUTF8(path)
2725 if err != nil {
2726 if os.IsNotExist(err) {
2727 return acpSessionMeta{}, false, nil
2728 }
2729 return acpSessionMeta{}, false, err
2730 }
2731 var meta acpSessionMeta
2732 if err := json.Unmarshal(b, &meta); err != nil {
2733 return acpSessionMeta{}, false, fmt.Errorf("decode ACP session metadata %s: %w", path, err)
2734 }
2735 return meta, true, nil
2736 }
2737
2738 func saveACPMeta(sessionPath string, meta acpSessionMeta) error {
2739 path := acpMetaPath(sessionPath)
2740 if path == "" {
2741 return nil
2742 }
2743 now := time.Now().UTC()
2744 if meta.SessionID == "" {
2745 meta.SessionID = sessionIDFromTranscript(sessionPath)
2746 }
2747 if meta.CreatedAt.IsZero() {
2748 meta.CreatedAt = now
2749 }
2750 if meta.UpdatedAt.IsZero() {
2751 meta.UpdatedAt = meta.CreatedAt
2752 }
2753 if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
2754 return err
2755 }
2756 b, err := json.MarshalIndent(meta, "", " ")
2757 if err != nil {
2758 return err
2759 }
2760 b = append(b, '\n')
2761 tmp, err := os.CreateTemp(filepath.Dir(path), ".acp-session.*.tmp")
2762 if err != nil {
2763 return err
2764 }
2765 tmpPath := tmp.Name()
2766 if _, err := tmp.Write(b); err != nil {
2767 tmp.Close()
2768 os.Remove(tmpPath)
2769 return err
2770 }
2771 if err := tmp.Close(); err != nil {
2772 os.Remove(tmpPath)
2773 return err
2774 }
2775 return fileutil.ReplaceFile(tmpPath, path)
2776 }
2777
2778 func listACPMetas(dir string) ([]acpSessionMeta, error) {
2779 entries, err := os.ReadDir(dir)
2780 if err != nil {
2781 if os.IsNotExist(err) {
2782 return nil, nil
2783 }
2784 return nil, err
2785 }
2786 out := []acpSessionMeta{}
2787 for _, e := range entries {
2788 if e.IsDir() || !strings.HasSuffix(e.Name(), ".acp.json") {
2789 continue
2790 }
2791 id := strings.TrimSuffix(e.Name(), ".acp.json")
2792 sessionPath := transcriptPath(dir, id)
2793 if agent.IsCleanupPending(sessionPath) {
2794 continue
2795 }
2796 if !sessionFileExists(sessionPath) {
2797 continue
2798 }
2799 meta, ok, err := loadACPMeta(sessionPath)
2800 if err != nil || !ok {
2801 continue
2802 }
2803 if meta.SessionID == "" {
2804 meta.SessionID = id
2805 }
2806 if meta.Cwd == "" {
2807 continue
2808 }
2809 out = append(out, meta)
2810 }
2811 return out, nil
2812 }
2813
2814 func sessionFileExists(path string) bool {
2815 info, err := os.Stat(path)
2816 return err == nil && !info.IsDir()
2817 }
2818
2819 func acpMetaPath(sessionPath string) string {
2820 if sessionPath == "" {
2821 return ""
2822 }
2823 return strings.TrimSuffix(sessionPath, filepath.Ext(sessionPath)) + ".acp.json"
2824 }
2825
2826 func sessionIDFromTranscript(path string) string {
2827 base := filepath.Base(path)
2828 if ext := filepath.Ext(base); ext != "" {
2829 base = strings.TrimSuffix(base, ext)
2830 }
2831 return base
2832 }
2833
2834 // listMetaBeats reports whether a should represent its session id in
2835 // session/list over b. A meta without an ActiveTranscript redirect is the
2836 // session's live transcript and always beats a redirect sidecar; between two
2837 // of the same kind the later UpdatedAt wins.
2838 func listMetaBeats(a, b acpSessionMeta) bool {
2839 aRedirect := strings.TrimSpace(a.ActiveTranscript) != ""
2840 bRedirect := strings.TrimSpace(b.ActiveTranscript) != ""
2841 if aRedirect != bRedirect {
2842 return !aRedirect
2843 }
2844 return a.UpdatedAt.After(b.UpdatedAt)
2845 }
2846
2847 func sessionInfoMatchesCwd(info SessionInfo, filter string) bool {
2848 if filter == "" {
2849 return true
2850 }
2851 return filepath.Clean(info.Cwd) == filepath.Clean(filter)
2852 }
2853
2854 func titleFromHistory(history []provider.Message) string {
2855 for _, m := range history {
2856 if agent.IsUserAuthoredTurnMessage(m) {
2857 if title := previewTitle(m.Content); title != "" {
2858 return title
2859 }
2860 }
2861 }
2862 return ""
2863 }
2864
2865 func previewTitle(text string) string {
2866 text = strings.Join(strings.Fields(text), " ")
2867 if len([]rune(text)) <= 80 {
2868 return text
2869 }
2870 runes := []rune(text)
2871 return string(runes[:77]) + "..."
2872 }
2873
2874 func validateSessionID(method, id string) error {
2875 trimmed := strings.TrimSpace(id)
2876 if trimmed == "" {
2877 return &RPCError{Code: ErrInvalidParams, Message: method + ": missing sessionId"}
2878 }
2879 if trimmed != id || trimmed == "." || trimmed == ".." || !isSafeSessionID(trimmed) {
2880 return &RPCError{Code: ErrInvalidParams, Message: method + ": invalid sessionId"}
2881 }
2882 return nil
2883 }
2884
2885 func isSafeSessionID(id string) bool {
2886 for _, r := range id {
2887 if r >= 'a' && r <= 'z' {
2888 continue
2889 }
2890 if r >= 'A' && r <= 'Z' {
2891 continue
2892 }
2893 if r >= '0' && r <= '9' {
2894 continue
2895 }
2896 if r == '-' || r == '_' || r == '.' {
2897 continue
2898 }
2899 return false
2900 }
2901 return true
2902 }
2903
2904 func parseSessionUpdatedAt(s string) time.Time {
2905 t, err := time.Parse(time.RFC3339Nano, s)
2906 if err != nil {
2907 return time.Time{}
2908 }
2909 return t
2910 }
2911
2912 func deleteSessionFiles(sessionPath string) error {
2913 paths := []string{
2914 sessionPath,
2915 acpMetaPath(sessionPath),
2916 }
2917 paths = append(paths, store.SessionSidecarFiles(sessionPath)...)
2918 for _, path := range paths {
2919 if path == "" {
2920 continue
2921 }
2922 if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
2923 return err
2924 }
2925 }
2926 if dir := checkpointPath(sessionPath); dir != "" {
2927 if err := os.RemoveAll(dir); err != nil && !os.IsNotExist(err) {
2928 return err
2929 }
2930 }
2931 if err := agent.DeleteSubagentsByParent(filepath.Dir(sessionPath), agent.BranchID(sessionPath)); err != nil {
2932 return err
2933 }
2934 if err := jobs.RemoveArtifacts(sessionPath); err != nil {
2935 return err
2936 }
2937 return agent.ClearCleanupPending(sessionPath)
2938 }
2939
2940 // ReconcileCleanupPending retries delayed ACP session cleanup left by a previous
2941 // process, including ACP's own metadata sidecar.
2942 func ReconcileCleanupPending(dir string) error {
2943 return agent.ReconcileCleanupPending(dir, func(item agent.CleanupPendingInfo) error {
2944 return deleteSessionFiles(item.SessionPath)
2945 })
2946 }
2947
2948 func delayedDeleteSessionFiles(sessionPath string, destroy control.SessionDestroyHandle) {
2949 if destroy.WaitAll != nil {
2950 destroy.WaitAll()
2951 }
2952 if err := deleteSessionFiles(sessionPath); err != nil {
2953 slog.Warn("acp: delayed session delete failed", "path", sessionPath, "err", err)
2954 }
2955 if destroy.Finish != nil {
2956 destroy.Finish()
2957 }
2958 }
2959
2960 func checkpointPath(sessionPath string) string {
2961 return store.SessionCheckpointDir(sessionPath)
2962 }
2963
2964 // mcpSpecs converts ACP MCP server declarations to plugin.Spec.
2965 func mcpSpecs(in []MCPServerSpec, cwd string) ([]plugin.Spec, error) {
2966 if len(in) == 0 {
2967 return nil, nil
2968 }
2969 out := make([]plugin.Spec, 0, len(in))
2970 for _, m := range in {
2971 typ := strings.ToLower(strings.TrimSpace(m.Type))
2972 if typ == "" {
2973 typ = "stdio"
2974 }
2975 if strings.TrimSpace(m.Name) == "" {
2976 return nil, fmt.Errorf("MCP server name is required")
2977 }
2978 switch typ {
2979 case "stdio":
2980 if strings.TrimSpace(m.Command) == "" {
2981 return nil, fmt.Errorf("MCP server %q command is required", m.Name)
2982 }
2983 case "http", "streamable-http", "streamable_http", "sse":
2984 if strings.TrimSpace(m.URL) == "" {
2985 return nil, fmt.Errorf("MCP server %q url is required", m.Name)
2986 }
2987 if typ != "sse" {
2988 typ = "http"
2989 }
2990 default:
2991 return nil, fmt.Errorf("MCP server %q uses unsupported transport %q", m.Name, m.Type)
2992 }
2993 out = append(out, plugin.Spec{
2994 Name: strings.TrimSpace(m.Name),
2995 Type: typ,
2996 Command: strings.TrimSpace(m.Command),
2997 Args: append([]string(nil), m.Args...),
2998 Env: mapString(m.Env),
2999 URL: strings.TrimSpace(m.URL),
3000 Headers: mapString(m.Headers),
3001 Dir: cwd,
3002 WorkspaceRoot: cwd,
3003 })
3004 }
3005 return out, nil
3006 }
3007
3008 func mapString(in map[string]string) map[string]string {
3009 if len(in) == 0 {
3010 return nil
3011 }
3012 out := make(map[string]string, len(in))
3013 maps.Copy(out, in)
3014 return out
3015 }
3016
3017 // newSessionID returns a random RFC 4122 v4 UUID string used to address a session.
3018 func newSessionID() (string, error) {
3019 var b [16]byte
3020 if _, err := rand.Read(b[:]); err != nil {
3021 return "", err
3022 }
3023 b[6] = (b[6] & 0x0f) | 0x40
3024 b[8] = (b[8] & 0x3f) | 0x80
3025 return fmt.Sprintf("%x-%x-%x-%x-%x", b[0:4], b[4:6], b[6:8], b[8:10], b[10:16]), nil
3026 }
3027
3027 lines GO