返回 DeepSeek-Reasonix
usecapability.go
根目录 / internal / agent / usecapability.go
1 package agent
2
3 import (
4 "context"
5 "encoding/json"
6 "fmt"
7 "maps"
8 "sort"
9 "strings"
10 "sync"
11 "time"
12
13 "reasonix/internal/capability"
14 "reasonix/internal/config"
15 "reasonix/internal/event"
16 "reasonix/internal/plugin"
17 "reasonix/internal/tool"
18 )
19
20 // MCPCapabilityRuntime is the session-shared MCP substrate: Host, boot specs,
21 // provider-visible registry for already-registered tools, schema cache catalog,
22 // and live connection snapshots. Each agent (executor, planner, task/fleet
23 // child) gets its own UseCapabilityTool frontend so ledger/audit never cross
24 // agent boundaries, while process connections remain on the shared Host.
25 type MCPCapabilityRuntime struct {
26 lifeCtx context.Context
27 host *plugin.Host
28 registry *tool.Registry
29 catalog func() capability.Catalog
30
31 // dispatchMu linearizes server enable/spec mutations against MCP process
32 // startup and tools/call. Calls may run concurrently under RLock; a disable,
33 // uninstall, or hot update waits for in-flight dispatch and invalidates every
34 // target that has not begun its final runtime-bound execution check.
35 dispatchMu sync.RWMutex
36 mu sync.RWMutex
37 servers map[string]mcpRuntimeServer
38 gates mcpServerGates
39 // shared connection observation across all frontends on this session.
40 state *mcpProxySharedState
41 frontendsMu sync.RWMutex
42 frontends map[*UseCapabilityTool]int
43 }
44
45 type mcpRuntimeServer struct {
46 entry config.PluginEntry
47 spec plugin.Spec
48 enabled bool
49 decision config.MCPDecision // why a disabled server is off
50 cached []plugin.CachedTool
51 cacheKeyOK bool
52 }
53
54 type mcpProxySharedState struct {
55 mu sync.Mutex
56 connected map[string]bool
57 liveTools map[string][]plugin.CachedTool
58 }
59
60 // hostProfile returns the session host's capability profile. A nil host (tests)
61 // resolves to core-v1.
62 func (r *MCPCapabilityRuntime) hostProfile() plugin.HostProfile {
63 if r == nil || r.host == nil {
64 return plugin.HostProfileCore
65 }
66 return r.host.Profile()
67 }
68
69 // hostProfileFor mirrors hostProfile for UseCapabilityTool, which holds its
70 // own host reference shared with the runtime.
71 func (t *UseCapabilityTool) hostProfileFor() plugin.HostProfile {
72 if t == nil || t.host == nil {
73 return plugin.HostProfileCore
74 }
75 return t.host.Profile()
76 }
77
78 // NewMCPCapabilityRuntime builds the session-shared MCP substrate. lifeCtx owns
79 // on-demand MCP child process lifetimes; specs must be the boot-converted specs.
80 func NewMCPCapabilityRuntime(lifeCtx context.Context, host *plugin.Host, specs []plugin.Spec, reg *tool.Registry, catalog func() capability.Catalog) *MCPCapabilityRuntime {
81 r := &MCPCapabilityRuntime{
82 lifeCtx: lifeCtx,
83 host: host,
84 registry: reg,
85 catalog: catalog,
86 servers: map[string]mcpRuntimeServer{},
87 state: &mcpProxySharedState{connected: map[string]bool{}},
88 frontends: map[*UseCapabilityTool]int{},
89 }
90 r.ConfigureServers(nil, specs, nil)
91 if host != nil {
92 host.SubscribeToolListChangesWithReplay(lifeCtx, r.applyToolListChange)
93 }
94 return r
95 }
96
97 // ConfigureServers replaces the runtime's configured MCP inventory. enabled is
98 // keyed by server name; nil keeps the standalone/test default that every spec is
99 // enabled. Boot passes the activation-resolved set so disabled servers are
100 // visible to discovery but cannot reuse a sibling tab's shared Host client.
101 func (r *MCPCapabilityRuntime) ConfigureServers(entries []config.PluginEntry, specs []plugin.Spec, enabled map[string]bool) {
102 if r == nil {
103 return
104 }
105 r.dispatchMu.Lock()
106 defer r.dispatchMu.Unlock()
107 byName := make(map[string]config.PluginEntry, len(entries))
108 for _, entry := range entries {
109 name := strings.TrimSpace(entry.Name)
110 if name != "" {
111 byName[name] = entry
112 }
113 }
114 next := make(map[string]mcpRuntimeServer, len(specs))
115 for _, raw := range specs {
116 spec := cloneMCPSpec(raw)
117 name := strings.TrimSpace(spec.Name)
118 if name == "" {
119 continue
120 }
121 entry, ok := byName[name]
122 if !ok {
123 entry = config.PluginEntry{Name: name}
124 }
125 isEnabled := true
126 if enabled != nil {
127 isEnabled = enabled[name]
128 }
129 decision := config.MCPDecisionOn
130 if ok {
131 decision = serverDecision(entry, spec, isEnabled)
132 }
133 cached, keyOK := cachedToolsForSpec(spec, r.hostProfile())
134 next[name] = mcpRuntimeServer{
135 entry: runtimePluginEntry(entry),
136 spec: spec,
137 enabled: isEnabled,
138 decision: decision,
139 cached: cached,
140 cacheKeyOK: keyOK,
141 }
142 }
143 r.mu.Lock()
144 previous := r.servers
145 r.servers = next
146 r.mu.Unlock()
147 r.syncRegistryInventory(previous, next)
148 }
149
150 // UpsertServer makes a hot-added or updated MCP spec authoritative for every
151 // frontend on this controller. Dynamic state stays host-local and never changes
152 // the provider-visible use_capability schema.
153 func (r *MCPCapabilityRuntime) UpsertServer(entry config.PluginEntry, raw plugin.Spec, enabled bool) {
154 if r == nil {
155 return
156 }
157 r.dispatchMu.Lock()
158 defer r.dispatchMu.Unlock()
159 spec := cloneMCPSpec(raw)
160 name := strings.TrimSpace(spec.Name)
161 if name == "" {
162 return
163 }
164 decision := serverDecision(entry, spec, enabled)
165 entry = runtimePluginEntry(entry)
166 if strings.TrimSpace(entry.Name) == "" {
167 entry.Name = name
168 }
169 cached, keyOK := cachedToolsForSpec(spec, r.hostProfile())
170 r.mu.Lock()
171 r.servers[name] = mcpRuntimeServer{
172 entry: entry,
173 spec: spec,
174 enabled: enabled,
175 decision: decision,
176 cached: cached,
177 cacheKeyOK: keyOK,
178 }
179 r.mu.Unlock()
180 r.setRegistryServerEnabled(name, enabled)
181 // Endpoint/tool metadata may have changed. Never route a stale live snapshot
182 // across an update; a connected client or the next call will repopulate it.
183 r.state.clearServer(name)
184 }
185
186 // SetServerEnabled revokes or restores this controller's right to use a server.
187 // It is intentionally independent from Host connectivity because desktop tabs
188 // may share one Host while keeping different enable states.
189 func (r *MCPCapabilityRuntime) SetServerEnabled(name string, enabled bool) bool {
190 if r == nil {
191 return false
192 }
193 r.dispatchMu.Lock()
194 defer r.dispatchMu.Unlock()
195 name = strings.TrimSpace(name)
196 r.mu.Lock()
197 server, ok := r.servers[name]
198 if ok {
199 server.enabled = enabled
200 server.decision = config.MCPDecisionOff
201 if enabled {
202 server.decision = config.MCPDecisionOn
203 }
204 r.servers[name] = server
205 }
206 r.mu.Unlock()
207 if ok {
208 r.setRegistryServerEnabled(name, enabled)
209 if !enabled {
210 r.state.clearServer(name)
211 }
212 }
213 return ok
214 }
215
216 // RemoveServer removes an uninstalled/runtime-only MCP from discovery and
217 // clears any live tool snapshot that could otherwise keep it routable.
218 func (r *MCPCapabilityRuntime) RemoveServer(name string) bool {
219 if r == nil {
220 return false
221 }
222 r.dispatchMu.Lock()
223 defer r.dispatchMu.Unlock()
224 name = strings.TrimSpace(name)
225 r.mu.Lock()
226 _, ok := r.servers[name]
227 delete(r.servers, name)
228 r.mu.Unlock()
229 r.setRegistryServerEnabled(name, false)
230 r.state.clearServer(name)
231 return ok
232 }
233
234 // CatalogState returns deterministic, privacy-minimal routing inputs for this
235 // controller. Configuration secrets are never copied into the transient route.
236 func (r *MCPCapabilityRuntime) CatalogState() (entries []config.PluginEntry, cached map[string][]plugin.CachedTool, keyOK map[string]bool, disabled map[string]bool) {
237 if r == nil {
238 return nil, nil, nil, nil
239 }
240 r.dispatchMu.RLock()
241 defer r.dispatchMu.RUnlock()
242 return r.catalogStateLocked()
243 }
244
245 // CapabilityCatalogState returns configuration and live proxy tools from one
246 // lifecycle generation. Callers must use this combined snapshot when building
247 // a route: taking the two halves separately can otherwise pair a just-updated
248 // spec with a stale pre-update live-tool directory.
249 func (r *MCPCapabilityRuntime) CapabilityCatalogState() (entries []config.PluginEntry, cached map[string][]plugin.CachedTool, keyOK map[string]bool, disabled map[string]bool, proxyTools map[string][]plugin.CachedTool) {
250 if r == nil {
251 return nil, nil, nil, nil, nil
252 }
253 r.dispatchMu.RLock()
254 defer r.dispatchMu.RUnlock()
255 entries, cached, keyOK, disabled = r.catalogStateLocked()
256 proxyTools = r.connectedProxyToolsLocked()
257 return entries, cached, keyOK, disabled, proxyTools
258 }
259
260 func (r *MCPCapabilityRuntime) catalogStateLocked() (entries []config.PluginEntry, cached map[string][]plugin.CachedTool, keyOK map[string]bool, disabled map[string]bool) {
261 r.mu.RLock()
262 names := make([]string, 0, len(r.servers))
263 for name := range r.servers {
264 names = append(names, name)
265 }
266 sort.Strings(names)
267 entries = make([]config.PluginEntry, 0, len(names))
268 cached = make(map[string][]plugin.CachedTool, len(names))
269 keyOK = make(map[string]bool, len(names))
270 disabled = make(map[string]bool)
271 for _, name := range names {
272 server := r.servers[name]
273 entries = append(entries, runtimePluginEntry(server.entry))
274 if len(server.cached) > 0 {
275 cached[name] = cloneCachedTools(server.cached)
276 keyOK[name] = server.cacheKeyOK
277 }
278 if !server.enabled {
279 disabled[name] = true
280 }
281 }
282 r.mu.RUnlock()
283 if len(cached) == 0 {
284 cached = nil
285 keyOK = nil
286 }
287 if len(disabled) == 0 {
288 disabled = nil
289 }
290 return entries, cached, keyOK, disabled
291 }
292
293 func (r *MCPCapabilityRuntime) configuredServers() []mcpRuntimeServer {
294 if r == nil {
295 return nil
296 }
297 r.mu.RLock()
298 names := make([]string, 0, len(r.servers))
299 for name := range r.servers {
300 names = append(names, name)
301 }
302 sort.Strings(names)
303 out := make([]mcpRuntimeServer, 0, len(names))
304 for _, name := range names {
305 server := r.servers[name]
306 server.spec = cloneMCPSpec(server.spec)
307 server.entry = runtimePluginEntry(server.entry)
308 server.cached = cloneCachedTools(server.cached)
309 out = append(out, server)
310 }
311 r.mu.RUnlock()
312 return out
313 }
314
315 func (r *MCPCapabilityRuntime) enabledSpec(server string) (plugin.Spec, bool) {
316 if r == nil {
317 return plugin.Spec{}, false
318 }
319 r.mu.RLock()
320 configured, ok := r.servers[strings.TrimSpace(server)]
321 r.mu.RUnlock()
322 if !ok || !configured.enabled {
323 return plugin.Spec{}, false
324 }
325 return cloneMCPSpec(configured.spec), true
326 }
327
328 func (r *MCPCapabilityRuntime) serverEnabled(server string) bool {
329 if r == nil {
330 return false
331 }
332 r.mu.RLock()
333 configured, ok := r.servers[strings.TrimSpace(server)]
334 r.mu.RUnlock()
335 return ok && configured.enabled
336 }
337
338 func cachedToolsForSpec(spec plugin.Spec, profile plugin.HostProfile) ([]plugin.CachedTool, bool) {
339 cached, keyOK := capability.LoadCachedToolsForSpecs([]plugin.Spec{spec}, profile)
340 return cloneCachedTools(cached[spec.Name]), keyOK[spec.Name]
341 }
342
343 func runtimePluginEntry(entry config.PluginEntry) config.PluginEntry {
344 out := config.PluginEntry{
345 Name: strings.TrimSpace(entry.Name),
346 Concurrency: strings.ToLower(strings.TrimSpace(entry.Concurrency)),
347 Source: entry.Source,
348 }
349 if entry.AutoStart != nil {
350 value := *entry.AutoStart
351 out.AutoStart = &value
352 }
353 return out
354 }
355
356 func cloneCachedTools(in []plugin.CachedTool) []plugin.CachedTool {
357 if len(in) == 0 {
358 return nil
359 }
360 out := make([]plugin.CachedTool, len(in))
361 copy(out, in)
362 for i := range out {
363 out[i].Schema = append(json.RawMessage(nil), in[i].Schema...)
364 }
365 return out
366 }
367
368 func cloneMCPSpec(in plugin.Spec) plugin.Spec {
369 out := in
370 out.Args = append([]string(nil), in.Args...)
371 out.LaunchArgs = append([]string(nil), in.LaunchArgs...)
372 out.LauncherIdentityArgs = append([]string(nil), in.LauncherIdentityArgs...)
373 out.Env = cloneStringMap(in.Env)
374 out.Headers = cloneStringMap(in.Headers)
375 if in.ToolTimeouts != nil {
376 out.ToolTimeouts = make(map[string]time.Duration, len(in.ToolTimeouts))
377 maps.Copy(out.ToolTimeouts, in.ToolTimeouts)
378 }
379 return out
380 }
381
382 func cloneStringMap(in map[string]string) map[string]string {
383 if in == nil {
384 return nil
385 }
386 out := make(map[string]string, len(in))
387 maps.Copy(out, in)
388 return out
389 }
390
391 // NewFrontend returns a per-agent use_capability instance. ledger/audit may be
392 // nil for ordinary sub-agents that do not run Delivery capability gates.
393 func (r *MCPCapabilityRuntime) NewFrontend(ledger *capability.Ledger, audit *capability.Audit) *UseCapabilityTool {
394 if r == nil {
395 return NewUseCapabilityTool(context.Background(), nil, nil, nil, ledger, audit, nil)
396 }
397 frontend := &UseCapabilityTool{
398 host: r.host,
399 lifeCtx: r.lifeCtx,
400 runtime: r,
401 registry: r.registry,
402 ledger: ledger,
403 audit: audit,
404 catalog: r.catalog,
405 state: r.state,
406 }
407 return frontend
408 }
409
410 func (r *MCPCapabilityRuntime) activateFrontend(frontend *UseCapabilityTool) func() {
411 if r == nil || frontend == nil {
412 return func() {}
413 }
414 r.frontendsMu.Lock()
415 if r.frontends == nil {
416 r.frontends = map[*UseCapabilityTool]int{}
417 }
418 r.frontends[frontend]++
419 r.frontendsMu.Unlock()
420 var once sync.Once
421 return func() {
422 once.Do(func() {
423 r.frontendsMu.Lock()
424 if r.frontends[frontend] <= 1 {
425 delete(r.frontends, frontend)
426 } else {
427 r.frontends[frontend]--
428 }
429 r.frontendsMu.Unlock()
430 })
431 }
432 }
433
434 func (r *MCPCapabilityRuntime) notifyToolListChanged(server string, tools []tool.Tool) {
435 if r == nil {
436 return
437 }
438 schemaBytes := 0
439 for _, target := range tools {
440 if target != nil {
441 schemaBytes += len(target.Schema())
442 }
443 }
444 r.frontendsMu.RLock()
445 frontends := make([]*UseCapabilityTool, 0, len(r.frontends))
446 for frontend := range r.frontends {
447 frontends = append(frontends, frontend)
448 }
449 r.frontendsMu.RUnlock()
450 for _, frontend := range frontends {
451 frontend.capabilityAudit().RecordMCPList("remote", "list_changed", 0, len(tools), schemaBytes)
452 frontend.observeMCPList(mcpListObservation{
453 Server: server, Source: "remote", Trigger: "list_changed",
454 ToolCount: len(tools), SchemaBytes: schemaBytes, NetworkCall: true,
455 })
456 }
457 }
458
459 // ConnectedProxyTools returns live tool metadata for servers connected through
460 // any frontend on this runtime, keyed by server name.
461 func (r *MCPCapabilityRuntime) ConnectedProxyTools() map[string][]plugin.CachedTool {
462 if r == nil || r.state == nil {
463 return nil
464 }
465 r.dispatchMu.RLock()
466 defer r.dispatchMu.RUnlock()
467 return r.connectedProxyToolsLocked()
468 }
469
470 func (r *MCPCapabilityRuntime) connectedProxyToolsLocked() map[string][]plugin.CachedTool {
471 live := r.state.snapshotLiveTools()
472 if len(live) == 0 {
473 return nil
474 }
475 r.mu.RLock()
476 for name := range live {
477 server, ok := r.servers[name]
478 if !ok || !server.enabled {
479 delete(live, name)
480 }
481 }
482 r.mu.RUnlock()
483 if len(live) == 0 {
484 return nil
485 }
486 return live
487 }
488
489 // UseCapabilityTool is the stable MCP capability proxy for Delivery, the
490 // two-model Planner, and task/fleet sub-agents. It lists, inspects, calls, or
491 // declines catalog capabilities without adding dynamic MCP tools to the
492 // provider-visible registry — subsequent calls keep using this stable schema.
493 // Multiple frontends may share one MCPCapabilityRuntime (Host + connection
494 // state) while keeping independent ledger/audit.
495 type UseCapabilityTool struct {
496 host *plugin.Host
497 // lifeCtx is the session-scoped context that owns on-demand MCP child
498 // processes (mirrors lazySpawn.ctx): a proxied server must outlive the tool
499 // call that started it and die with the session, not with a resolve-phase
500 // timeout. nil falls back to context.Background() for direct/test use.
501 lifeCtx context.Context
502 // specs are the boot-converted plugin specs (env expansion, workspace
503 // overrides and timeouts). The proxy never rebuilds
504 // specs from raw config entries — that would fork the conversion logic.
505 specs []plugin.Spec
506 runtime *MCPCapabilityRuntime
507 registry *tool.Registry // live registry for already-exposed MCP tools
508 ledger *capability.Ledger
509 audit *capability.Audit
510 catalog func() capability.Catalog
511 // toolResultSession is bound by Agent.New; clones leave it empty so parent transcripts cannot
512 // leak into planner or child frontends before their Agent binds them.
513 toolResultMu sync.RWMutex
514 toolResultSession func() *Session
515 mcpListMu sync.RWMutex
516 mcpListObserver func(mcpListObservation)
517 // state is session-shared connection observation when built via
518 // MCPCapabilityRuntime; nil falls back to a private map for tests.
519 state *mcpProxySharedState
520 }
521
522 // runtimeBoundMCPTool keeps the provider-visible MCP adapter unchanged while
523 // binding execution to the current controller runtime. The underlying Host may
524 // be shared by sibling tabs, so a server name alone must never authorize reuse.
525 type runtimeBoundMCPTool struct {
526 proxy *UseCapabilityTool
527 target tool.Tool
528 server string
529 authorized bool
530 }
531
532 func (b *runtimeBoundMCPTool) Name() string { return b.target.Name() }
533 func (b *runtimeBoundMCPTool) Description() string { return b.target.Description() }
534 func (b *runtimeBoundMCPTool) Schema() json.RawMessage { return b.target.Schema() }
535 func (b *runtimeBoundMCPTool) ReadOnly() bool { return b.target.ReadOnly() }
536 func (b *runtimeBoundMCPTool) MCPServerAuthorized() bool { return b.authorized }
537 func (b *runtimeBoundMCPTool) MCPServerName() string { return b.server }
538 func (b *runtimeBoundMCPTool) MCPRawToolName() string { return mcpRawToolName(b.target) }
539 func (b *runtimeBoundMCPTool) MCPDestructiveHint() bool { return mcpDestructiveHint(b.target) }
540 func (b *runtimeBoundMCPTool) MCPVisibleToolName() string {
541 if meta, ok := b.target.(tool.MCPVisibleMetadata); ok {
542 return meta.MCPVisibleToolName()
543 }
544 return b.MCPRawToolName()
545 }
546 func (b *runtimeBoundMCPTool) MCPPackageName() string {
547 if meta, ok := b.target.(tool.MCPPackageMetadata); ok {
548 return meta.MCPPackageName()
549 }
550 return ""
551 }
552
553 func (b *runtimeBoundMCPTool) Execute(ctx context.Context, args json.RawMessage) (string, error) {
554 var out string
555 err := b.proxy.withRuntimeBoundMCP(ctx, b.server, b.target, func() error {
556 var execErr error
557 out, execErr = b.target.Execute(ctx, args)
558 return execErr
559 })
560 return out, err
561 }
562
563 func (b *runtimeBoundMCPTool) ExecuteWithImages(ctx context.Context, args json.RawMessage) (string, []string, error) {
564 var out string
565 var images []string
566 err := b.proxy.withRuntimeBoundMCP(ctx, b.server, b.target, func() error {
567 if imageTool, ok := b.target.(tool.ImageTool); ok {
568 var execErr error
569 out, images, execErr = imageTool.ExecuteWithImages(ctx, args)
570 return execErr
571 }
572 var execErr error
573 out, execErr = b.target.Execute(ctx, args)
574 return execErr
575 })
576 return out, images, err
577 }
578
579 func mcpRawToolName(target tool.Tool) string {
580 if meta, ok := target.(tool.MCPMetadata); ok {
581 return meta.MCPRawToolName()
582 }
583 return ""
584 }
585
586 // NewUseCapabilityTool builds a standalone capability proxy (tests and simple
587 // boots). Prefer MCPCapabilityRuntime.NewFrontend when multiple agents share
588 // one session Host.
589 func NewUseCapabilityTool(lifeCtx context.Context, host *plugin.Host, specs []plugin.Spec, reg *tool.Registry, ledger *capability.Ledger, audit *capability.Audit, catalog func() capability.Catalog) *UseCapabilityTool {
590 return &UseCapabilityTool{
591 host: host,
592 lifeCtx: lifeCtx,
593 specs: append([]plugin.Spec(nil), specs...),
594 registry: reg,
595 ledger: ledger,
596 audit: audit,
597 catalog: catalog,
598 state: &mcpProxySharedState{connected: map[string]bool{}},
599 }
600 }
601
602 // CloneForAgent returns a new frontend sharing Host/specs/connection state but
603 // with independent ledger and audit (nil unless provided).
604 func (t *UseCapabilityTool) CloneForAgent(ledger *capability.Ledger, audit *capability.Audit) *UseCapabilityTool {
605 if t == nil {
606 return nil
607 }
608 state := t.state
609 if state == nil {
610 state = &mcpProxySharedState{connected: map[string]bool{}}
611 }
612 clone := &UseCapabilityTool{
613 host: t.host,
614 lifeCtx: t.lifeCtx,
615 specs: t.specs,
616 runtime: t.runtime,
617 registry: t.registry,
618 ledger: ledger,
619 audit: audit,
620 catalog: t.catalog,
621 state: state,
622 }
623 return clone
624 }
625
626 func (*UseCapabilityTool) Name() string { return tool.HostUseCapability }
627
628 func (*UseCapabilityTool) Description() string {
629 return "Fixed-schema capability proxy. Prefer search(query, limit<=8), then inspect one exact capability, then call it. list is a compact diagnostic inventory only. Supports stable ids such as tool:grep, skill:review, mcp-tool:server/tool, task:subagent, workflow:name, and web:/lsp:/session:/memory: namespaces. memory:remember saves facts (description+body required; activation=\"relevant\" on create; omit activation on update; \"pinned\" only if user asks); memory:forget(name); tool:memory(operation=search|read|list). decline records a reason for a prefer capability. Independent list/search/inspect calls are read-only and may be issued together. Calls keep the provider-visible schema fixed; real writers still pass permission, plan mode, sandbox, write-path, and workspace-lease checks."
630 }
631
632 func (*UseCapabilityTool) ReadOnly() bool { return true }
633
634 func (*UseCapabilityTool) Schema() json.RawMessage {
635 // Stable schema — must not change across turns or when MCP connects.
636 // capability_id is optional only for action=list; inspect/call/decline still
637 // require it at resolve time. One intentional prefix upgrade; thereafter
638 // install/connect churn does not change this schema.
639 return json.RawMessage(`{
640 "type":"object",
641 "properties":{
642 "action":{"type":"string","enum":["list","search","inspect","call","decline"],"description":"Use search for discovery, inspect one exact result, then call. list is diagnostic only."},
643 "capability_id":{"type":"string","description":"Capability id such as skill:review, mcp-server:github, or mcp-tool:github/search_issues. Not required for action=list."},
644 "query":{"type":"string","description":"Local catalog query required for action=search. No process or network is started."},
645 "limit":{"type":"integer","minimum":1,"maximum":100,"description":"Maximum results. Search defaults to 5 and allows at most 8; list defaults to 50 and allows at most 100."},
646 "cursor":{"type":"string","description":"Opaque cursor returned by action=list. It is valid only for the same catalog version."},
647 "arguments":{"type":"object","description":"Raw MCP tool arguments for action=call"},
648 "reason":{"type":"string","description":"Required non-empty reason when action=decline"}
649 },
650 "required":["action"],
651 "additionalProperties":false
652 }`)
653 }
654
655 // ResolveCall implements tool.CallResolver so the agent can run permission,
656 // hooks, and evidence against the real MCP target before execution.
657 func (t *UseCapabilityTool) ResolveCall(ctx context.Context, args json.RawMessage) (tool.ResolvedCall, error) {
658 p, action, id, err := parseUseCapabilityArgs(args)
659 if err != nil {
660 return tool.ResolvedCall{}, &capabilityInputError{err}
661 }
662 base := tool.ResolvedCall{
663 DisplayName: "use_capability",
664 ProxyAction: action,
665 CapabilityID: id,
666 Args: p.Arguments,
667 }
668 switch action {
669 case "list", "search", "inspect":
670 return t.resolveDiscovery(ctx, p, action, id, base)
671 case "decline":
672 if id == "" {
673 return tool.ResolvedCall{}, capabilityInputErrorf("capability_id is required for action=decline")
674 }
675 reason := strings.TrimSpace(p.Reason)
676 if reason == "" {
677 return tool.ResolvedCall{}, capabilityInputErrorf("reason is required for action=decline")
678 }
679 // Decline must not skip require. The mutation itself is delayed until the
680 // agent has applied its post-resolution host boundary.
681 if t.ledger != nil {
682 if e, ok := t.ledger.Get(id); ok && e.Policy == capability.AutoUseRequire {
683 return tool.ResolvedCall{}, fmt.Errorf("cannot decline a require capability %q", id)
684 }
685 }
686 base.SkipExecute = true
687 base.Result = fmt.Sprintf("declined capability %s: %s", id, reason)
688 base.ReadOnly = true
689 base.Commit = func() error {
690 if t.ledger != nil {
691 if err := t.ledger.MarkDeclined(id, reason); err != nil {
692 return err
693 }
694 }
695 if t.audit != nil {
696 t.audit.RecordDecline()
697 }
698 return nil
699 }
700 return base, nil
701 case "call":
702 if id == "" {
703 return tool.ResolvedCall{}, capabilityInputErrorf("capability_id is required for action=call")
704 }
705 if id == sessionToolResultCapabilityID {
706 return t.resolveSessionToolResult(p.Arguments, base)
707 }
708 return t.resolveCall(ctx, id, p.Arguments, base)
709 default:
710 return tool.ResolvedCall{}, capabilityInputErrorf("unknown action %q; use list, search, inspect, call, or decline", p.Action)
711 }
712 }
713
714 func (t *UseCapabilityTool) Execute(ctx context.Context, args json.RawMessage) (string, error) {
715 resolved, err := t.ResolveCall(ctx, args)
716 if err != nil {
717 return "", err
718 }
719 if resolved.SkipExecute {
720 if resolved.Commit != nil {
721 if err := resolved.Commit(); err != nil {
722 return "", err
723 }
724 }
725 if resolved.ProxyAction == "call" && !resolved.Unavailable {
726 if t.ledger != nil {
727 t.ledger.MarkSucceeded(resolved.CapabilityID)
728 }
729 if t.audit != nil {
730 t.audit.RecordMCPProxy(false, true, false)
731 }
732 }
733 return resolved.Result, nil
734 }
735 if resolved.Unavailable {
736 if t.ledger != nil {
737 t.ledger.MarkUnavailable(resolved.CapabilityID, resolved.UnavailableReason)
738 }
739 return "", fmt.Errorf("capability unavailable: %s", resolved.UnavailableReason)
740 }
741 if resolved.Target == nil {
742 return "", fmt.Errorf("no target tool resolved for %s", resolved.CapabilityID)
743 }
744 if t.ledger != nil {
745 t.ledger.MarkInvoked(resolved.CapabilityID)
746 }
747 if t.audit != nil {
748 t.audit.RecordMCPProxy(false, true, false)
749 }
750 out, err := resolved.Target.Execute(ctx, resolved.Args)
751 if err != nil {
752 if t.ledger != nil {
753 t.ledger.MarkFailed(resolved.CapabilityID, err.Error())
754 }
755 if t.audit != nil {
756 t.audit.RecordMCPProxy(false, true, true)
757 }
758 return out, err
759 }
760 if t.ledger != nil {
761 t.ledger.MarkSucceeded(resolved.CapabilityID)
762 }
763 return out, nil
764 }
765
766 func (t *UseCapabilityTool) resolveCall(ctx context.Context, id string, args json.RawMessage, base tool.ResolvedCall) (tool.ResolvedCall, error) {
767 // Registry-backed tools and skills share the unified proxy. Real writers
768 // still pass permission/plan/sandbox/lease checks via ResolvedCall.Target.
769 if name, ok := strings.CutPrefix(id, "tool:"); ok {
770 return t.resolveRegistryTool(ctx, strings.TrimSpace(name), id, args, base)
771 }
772 if name, ok := strings.CutPrefix(id, "skill:"); ok {
773 return t.resolveSkillCall(strings.TrimSpace(name), id, args, base)
774 }
775 if name, ok := strings.CutPrefix(id, "task:"); ok {
776 return t.resolveRegistryTool(ctx, taskToolName(name), id, args, base)
777 }
778 if name, ok := strings.CutPrefix(id, "workflow:"); ok {
779 return t.resolveRegistryTool(ctx, strings.TrimSpace(name), "tool:"+strings.TrimSpace(name), args, base)
780 }
781 for _, prefix := range []string{"web:", "lsp:", "session:", "memory:"} {
782 if rest, ok := strings.CutPrefix(id, prefix); ok {
783 return t.resolveRegistryTool(ctx, strings.TrimSpace(rest), id, args, base)
784 }
785 }
786 // Server-level call is the first-discovery path for servers with no
787 // schema cache: it resolves to a gated connect-and-list target so the
788 // model can learn tool names without inspect ever starting a process.
789 if server, ok := parseMCPServerCapabilityID(id); ok {
790 return t.resolveServerConnect(ctx, server, base)
791 }
792 server, raw, err := parseMCPCapabilityID(id)
793 if err != nil {
794 return tool.ResolvedCall{}, err
795 }
796 if !t.serverEnabled(server) {
797 return t.resolveUnavailable(base, id, plugin.ModelToolName(server, raw), t.serverUnavailableReason(server)), nil
798 }
799 var runtimeSpec plugin.Spec
800 releaseRuntime := func() {}
801 if t.runtime != nil {
802 spec, unlock, lockErr := t.lockAuthorizedRuntimeServer(ctx, server)
803 if lockErr != nil {
804 return t.resolveUnavailable(base, id, plugin.ModelToolName(server, raw), lockErr.Error()), nil
805 }
806 releaseRuntime = unlock
807 defer func() { releaseRuntime() }()
808 runtimeSpec = spec
809 }
810 // Prefer already-exposed registry tool (auto-started MCP). The model name
811 // MUST come from the plugin layer's canonical constructor: it appends a
812 // collision hash for sanitised raw names, and permission/hook rules are
813 // written against that executed name — a proxy-local normalization would
814 // let them silently miss.
815 modelName := plugin.ModelToolName(server, raw)
816 if t.registry != nil {
817 if tl, ok := t.registry.Get(modelName); ok {
818 if t.runtime != nil && !plugin.MCPToolMatchesSpec(tl, runtimeSpec) {
819 return t.resolveUnavailable(base, id, modelName, fmt.Sprintf("connected MCP server %q identity does not match the current runtime configuration", server)), nil
820 }
821 base.TargetName = modelName
822 base.Target = t.bindRuntimeMCP(runtimeSpec, tl)
823 base.ReadOnly = tl.ReadOnly()
824 if len(args) == 0 {
825 base.Args = json.RawMessage(`{}`)
826 } else {
827 base.Args = args
828 }
829 return base, nil
830 }
831 }
832 // Server already connected (auto-started, a previous proxy call, or a
833 // sibling tab sharing this host): resolving against live tools is
834 // side-effect-free. serverTools also refreshes the catalog snapshot so a
835 // cross-tab connect still yields routable mcp-tool entries here.
836 if t.host != nil && t.host.HasClient(server) {
837 var tools []tool.Tool
838 if t.runtime != nil {
839 tools, err = t.serverToolsForSpec(ctx, server, runtimeSpec)
840 } else {
841 tools, err = t.serverTools(ctx, server)
842 }
843 if err != nil {
844 return t.resolveUnavailable(base, id, modelName, err.Error()), nil
845 }
846 target := findMCPTool(tools, raw, modelName)
847 if target == nil {
848 return t.resolveUnavailable(base, id, modelName, fmt.Sprintf("MCP tool %q not found on server %q", raw, server)), nil
849 }
850 base.Target = t.bindRuntimeMCP(runtimeSpec, target)
851 base.TargetName = target.Name()
852 base.ReadOnly = target.ReadOnly()
853 if len(args) == 0 {
854 base.Args = json.RawMessage(`{}`)
855 } else {
856 base.Args = args
857 }
858 return base, nil
859 }
860 // Unconnected server: resolution must stay pure — no subprocess, no network.
861 // Return a deferred target that connects in Execute, after the permission
862 // gate and PreToolUse hooks have approved the real target name/arguments.
863 spec := runtimeSpec
864 if t.runtime == nil {
865 var ok bool
866 spec, ok = t.specFor(server)
867 if !ok {
868 return t.resolveUnavailable(base, id, modelName, mcpServerUnregisteredMessage(server)), nil
869 }
870 spec = plugin.ResolveStoredAuthorization(ctx, spec)
871 }
872 // Execute rechecks the immutable spec snapshot. Release dispatchMu because
873 // Boot's catalog callback takes its own runtime snapshot.
874 releaseRuntime()
875 releaseRuntime = func() {}
876 return t.resolveUnconnectedMCPCall(id, args, base, spec, server, raw, modelName), nil
877 }
878
879 // resolveSkillCall routes skill:<name> through run_skill / read_only_skill /
880 // read_skill when present, preserving the real skill tool name for evidence.
881 func (t *UseCapabilityTool) resolveSkillCall(skillName, id string, args json.RawMessage, base tool.ResolvedCall) (tool.ResolvedCall, error) {
882 skillName = strings.TrimSpace(skillName)
883 if skillName == "" {
884 return tool.ResolvedCall{}, fmt.Errorf("capability id %q is missing a skill name", id)
885 }
886 if t.registry == nil {
887 return t.resolveUnavailable(base, id, skillName, "tool registry is unavailable"), nil
888 }
889 // Prefer full run_skill; fall back to read-only variants when the session
890 // only exposes them (planner / plan mode).
891 for _, toolName := range []string{"run_skill", "read_only_skill", "read_skill"} {
892 if tl, ok := t.registry.Get(toolName); ok {
893 payload := normalizeSkillCapabilityArgs(skillName, args)
894 base.Target = tl
895 base.TargetName = toolName
896 base.ReadOnly = tl.ReadOnly()
897 base.Args = payload
898 return base, nil
899 }
900 }
901 return t.resolveUnavailable(base, id, skillName, fmt.Sprintf("skill tools are not available for %q", skillName)), nil
902 }
903
904 func taskToolName(name string) string {
905 name = strings.TrimSpace(name)
906 switch name {
907 case "subagent", "task":
908 return "task"
909 case "read_only_subagent", "read_only_task", "research":
910 return "read_only_task"
911 case "parallel", "parallel_tasks":
912 return "parallel_tasks"
913 case "fleet":
914 return "fleet"
915 default:
916 return name
917 }
918 }
919
920 // resolveUnavailable fills the host-proven unavailable shape shared by the
921 // side-effect-free resolution failures (missing config, unknown tool).
922 func (t *UseCapabilityTool) resolveUnavailable(base tool.ResolvedCall, id, modelName, reason string) tool.ResolvedCall {
923 base.Unavailable = true
924 base.UnavailableReason = reason
925 base.SkipExecute = true
926 base.Result = "capability unavailable: " + reason
927 base.TargetName = modelName
928 base.ReadOnly = false
929 base.Commit = func() error {
930 if t.ledger != nil {
931 t.ledger.MarkUnavailable(id, reason)
932 }
933 if t.audit != nil {
934 t.audit.RecordMCPProxy(false, true, true)
935 }
936 return nil
937 }
938 return base
939 }
940
941 // findMCPTool matches a server's tool list by raw MCP name or by the
942 // canonical namespaced model-visible name (plugin.ModelToolName).
943 func findMCPTool(tools []tool.Tool, raw, modelName string) tool.Tool {
944 for _, tl := range tools {
945 if m, ok := tl.(tool.MCPMetadata); ok && m.MCPRawToolName() == raw {
946 return tl
947 }
948 if tl.Name() == modelName {
949 return tl
950 }
951 }
952 return nil
953 }
954
955 // onDemandMCPTool defers MCP server startup to Execute so permission and hook
956 // gates always run before any subprocess or network side effect. Before the live
957 // handshake it remains write-capable until the resolved MCP tool is classified.
958 type onDemandMCPTool struct {
959 proxy *UseCapabilityTool
960 spec plugin.Spec
961 server string
962 raw string
963 modelName string
964 // destructive comes from the schema cache when available. A live promotion
965 // is detected in Execute so a retry re-enters the current Plan/read-only
966 // execution boundary.
967 destructive bool
968 readOnly bool
969 schema json.RawMessage
970 }
971
972 func (o *onDemandMCPTool) Name() string { return o.modelName }
973
974 func (o *onDemandMCPTool) Description() string {
975 return "on-demand MCP tool " + o.server + "/" + o.raw + " (connects when first used)"
976 }
977
978 func (o *onDemandMCPTool) Schema() json.RawMessage {
979 if len(o.schema) > 0 {
980 return o.schema
981 }
982 return json.RawMessage(`{"type":"object"}`)
983 }
984
985 func (o *onDemandMCPTool) ReadOnly() bool {
986 return o.readOnly
987 }
988
989 func (o *onDemandMCPTool) ReadOnlyExecutionHostMutation() bool { return true }
990
991 func (o *onDemandMCPTool) MCPServerAuthorized() bool {
992 // Spec.Authorized is the single runtime authorization result. Boot/install
993 // and ResolveStoredAuthorization set it; this path never invents trust.
994 return o.spec.ServerAuthorized()
995 }
996
997 func (o *onDemandMCPTool) ReadOnlyExecutionBlockReason() string {
998 return "connect this MCP capability from a parent session first"
999 }
1000
1001 // MCPServerName/MCPRawToolName expose the deferred target for audit and
1002 // diagnostics (tool.MCPMetadata).
1003 func (o *onDemandMCPTool) MCPServerName() string { return o.server }
1004 func (o *onDemandMCPTool) MCPRawToolName() string { return o.raw }
1005 func (o *onDemandMCPTool) MCPDestructiveHint() bool {
1006 return o.destructive
1007 }
1008
1009 func (o *onDemandMCPTool) Execute(ctx context.Context, args json.RawMessage) (string, error) {
1010 text, _, err := o.executeWithImages(ctx, args)
1011 return text, err
1012 }
1013
1014 // ExecuteWithImages preserves structured MCP image results on the first call,
1015 // when the deferred target must connect the server before dispatch. Keeping the
1016 // resolution and safety checks in executeWithImages ensures text-only and image
1017 // callers share the same authorization and runtime-identity boundary.
1018 func (o *onDemandMCPTool) ExecuteWithImages(ctx context.Context, args json.RawMessage) (string, []string, error) {
1019 return o.executeWithImages(ctx, args)
1020 }
1021
1022 func (o *onDemandMCPTool) executeWithImages(ctx context.Context, args json.RawMessage) (string, []string, error) {
1023 // Final runtime-bound authorization and identity check before any
1024 // process/network start. The read lock also linearizes this dispatch against
1025 // disable, uninstall, and same-name hot replacement.
1026 spec, unlock, err := o.proxy.lockAuthorizedRuntimeServer(ctx, o.server)
1027 if err != nil {
1028 msg := err.Error()
1029 if o.proxy.ledger != nil {
1030 o.proxy.ledger.MarkUnavailable("mcp-tool:"+o.server+"/"+o.raw, msg)
1031 }
1032 return "", nil, err
1033 }
1034 defer unlock()
1035 if !plugin.MCPRuntimeSpecMatches(spec, o.spec) {
1036 return "", nil, fmt.Errorf("MCP server %q runtime identity changed after resolution; retry so Reasonix can bind the current configuration", o.server)
1037 }
1038 tools, err := o.proxy.ensureServerToolsForSpec(ctx, o.server, spec)
1039 if err != nil {
1040 // Audit for the call path is recorded once by the agent loop
1041 // (noteCapabilityInvocation); only the ledger outcome lands here.
1042 if o.proxy.ledger != nil {
1043 o.proxy.ledger.MarkUnavailable("mcp-tool:"+o.server+"/"+o.raw, err.Error())
1044 }
1045 return "", nil, err
1046 }
1047 target := findMCPTool(tools, o.raw, o.modelName)
1048 if target == nil {
1049 msg := fmt.Sprintf("MCP tool %q not found on server %q", o.raw, o.server)
1050 if o.proxy.ledger != nil {
1051 o.proxy.ledger.MarkUnavailable("mcp-tool:"+o.server+"/"+o.raw, msg)
1052 }
1053 return "", nil, fmt.Errorf("%s", msg)
1054 }
1055 if !plugin.MCPToolMatchesSpec(target, spec) {
1056 return "", nil, fmt.Errorf("connected MCP server %q identity does not match the current runtime configuration; reconnect this server before retrying", o.server)
1057 }
1058 if _, err := plugin.ReconcileCachedToolSafety(o.server, o.raw, plugin.CachedToolSafety{
1059 ReadOnly: o.readOnly,
1060 Destructive: o.destructive,
1061 }, target); err != nil {
1062 return "", nil, err
1063 }
1064 // Planner non-destructive lane and reader lane: re-check live metadata
1065 // before tools/call even when Reconcile did not see a cache promotion.
1066 if tool.HasNonDestructiveMCPExecutionIntent(ctx) {
1067 if !mcpServerAuthorized(target) || mcpDestructiveHint(target) {
1068 return "", nil, fmt.Errorf("MCP server %q changed the authorization or destructive classification for tool %q; the call was blocked before dispatch — retry so Reasonix can re-apply the current Planner MCP safety boundary", o.server, o.raw)
1069 }
1070 }
1071 if blocked, msg := hostValidateBeforeDispatch(target, args, "mcp-tool:"+o.server+"/"+o.raw); blocked {
1072 return "", nil, fmt.Errorf("%s", msg)
1073 }
1074 if imageTool, ok := target.(tool.ImageTool); ok {
1075 return imageTool.ExecuteWithImages(ctx, args)
1076 }
1077 text, err := target.Execute(ctx, args)
1078 return text, nil, err
1079 }
1080
1081 func (t *UseCapabilityTool) ensureServerToolsForSpec(ctx context.Context, server string, spec plugin.Spec) ([]tool.Tool, error) {
1082 // Reuse shared host if already connected (including auto-started).
1083 if t.host.HasClient(server) {
1084 return t.serverToolsForSpec(ctx, server, spec)
1085 }
1086 // On-demand connect: the child and handshake belong to the session. This
1087 // tool call waits briefly, but a slow healthy server continues in the
1088 // background instead of being killed and restarted on every retry. Tools
1089 // stay off the main provider-visible registry.
1090 life := t.lifeCtx
1091 if life == nil {
1092 life = context.Background()
1093 }
1094 started := time.Now()
1095 result := t.host.EnsureConnectedInBackground(life, spec)
1096 waitBudget := plugin.DefaultStartupWaitBudget()
1097 timer := time.NewTimer(waitBudget)
1098 defer timer.Stop()
1099 var tools []tool.Tool
1100 var err error
1101 select {
1102 case connected := <-result:
1103 tools, err = connected.Tools, connected.Err
1104 case <-ctx.Done():
1105 return nil, ctx.Err()
1106 case <-timer.C:
1107 return nil, fmt.Errorf("MCP server %q is still initializing after %s; startup continues in background (limit %s) — retry on a later turn",
1108 server, waitBudget, spec.ResolvedStartupTimeout())
1109 }
1110 if err != nil {
1111 if plugin.IsServerAlreadyConnected(err) {
1112 return t.serverToolsForSpec(ctx, server, spec)
1113 }
1114 t.host.RecordFailure(spec, err)
1115 return nil, fmt.Errorf("connect %q: %w", server, err)
1116 }
1117 t.ensureState().markConnected(server)
1118 schemaBytes := 0
1119 for _, target := range tools {
1120 schemaBytes += len(target.Schema())
1121 }
1122 durationMs := time.Since(started).Milliseconds()
1123 t.capabilityAudit().RecordMCPList("remote", "connect", durationMs, len(tools), schemaBytes)
1124 t.observeMCPList(mcpListObservation{
1125 Server: server, Source: "remote", Trigger: "connect", DurationMs: durationMs,
1126 ToolCount: len(tools), SchemaBytes: schemaBytes, NetworkCall: true,
1127 })
1128 // Intentionally do NOT add tools to t.registry — provider schema stays stable.
1129 _ = tools
1130 return t.serverToolsForSpec(ctx, server, spec)
1131 }
1132
1133 // serverTools fetches the live tools for a connected server and refreshes the
1134 // shared catalog snapshot so mcp-tool entries stay routable once the server
1135 // is StatusReady (its tools are absent from the provider-visible registry).
1136 func (t *UseCapabilityTool) serverTools(ctx context.Context, server string) ([]tool.Tool, error) {
1137 spec, unlock, err := t.lockAuthorizedRuntimeServer(ctx, server)
1138 if err != nil {
1139 return nil, err
1140 }
1141 defer unlock()
1142 return t.serverToolsForSpec(ctx, server, spec)
1143 }
1144
1145 func (t *UseCapabilityTool) serverToolsForSpec(ctx context.Context, server string, spec plugin.Spec) ([]tool.Tool, error) {
1146 tools, err := t.host.ToolsForSpec(ctx, spec)
1147 if err != nil {
1148 return nil, err
1149 }
1150 t.ensureState().setLiveTools(server, snapshotMCPTools(tools))
1151 return tools, nil
1152 }
1153
1154 // ConnectedProxyTools returns raw tool metadata for servers connected through
1155 // any frontend sharing this proxy state, keyed by server name. Catalog builders
1156 // consume it so concrete mcp-tool capabilities survive an on-demand connect
1157 // without ever touching the provider-visible registry.
1158 func (t *UseCapabilityTool) ConnectedProxyTools() map[string][]plugin.CachedTool {
1159 if t == nil {
1160 return nil
1161 }
1162 return t.ensureState().snapshotLiveTools()
1163 }
1164
1165 func (t *UseCapabilityTool) ensureState() *mcpProxySharedState {
1166 if t.state == nil {
1167 t.state = &mcpProxySharedState{connected: map[string]bool{}}
1168 }
1169 return t.state
1170 }
1171
1172 // specFor looks up the boot-converted spec for server. The proxy deliberately
1173 // holds []plugin.Spec, not raw config entries: env expansion, workspace
1174 // overrides, call timeouts, and read-only tool names all live in the
1175 // shared conversion and must not be re-derived here.
1176 func (t *UseCapabilityTool) specFor(server string) (plugin.Spec, bool) {
1177 if t.runtime != nil {
1178 return t.runtime.enabledSpec(server)
1179 }
1180 for _, s := range t.specs {
1181 if s.Name == server {
1182 return s, true
1183 }
1184 }
1185 return plugin.Spec{}, false
1186 }
1187
1188 // lockAuthorizedRuntimeServer acquires the shared runtime dispatch read lock
1189 // and returns the current enabled, authorized spec. The caller must keep the
1190 // returned lock until the identity-bound Host operation or tools/call has
1191 // crossed its dispatch boundary; lifecycle mutations take the write lock.
1192 func (t *UseCapabilityTool) lockAuthorizedRuntimeServer(ctx context.Context, server string) (plugin.Spec, func(), error) {
1193 if t.runtime == nil {
1194 spec, ok := t.specFor(server)
1195 if !ok {
1196 return plugin.Spec{}, func() {}, mcpServerUnregisteredError(server)
1197 }
1198 spec = plugin.ResolveStoredAuthorization(ctx, spec)
1199 if !spec.ServerAuthorized() {
1200 return plugin.Spec{}, func() {}, fmt.Errorf("MCP server %q is not authorized; install it or complete project identity approval before connecting", server)
1201 }
1202 return spec, func() {}, nil
1203 }
1204
1205 t.runtime.dispatchMu.RLock()
1206 unlock := t.runtime.dispatchMu.RUnlock
1207 t.runtime.mu.RLock()
1208 configured, ok := t.runtime.servers[strings.TrimSpace(server)]
1209 t.runtime.mu.RUnlock()
1210 if !ok {
1211 unlock()
1212 return plugin.Spec{}, func() {}, mcpServerUnregisteredError(server)
1213 }
1214 if !configured.enabled {
1215 unlock()
1216 return plugin.Spec{}, func() {}, mcpServerOffError(server, configured.decision)
1217 }
1218 spec := plugin.ResolveStoredAuthorization(ctx, cloneMCPSpec(configured.spec))
1219 if !spec.ServerAuthorized() {
1220 unlock()
1221 return plugin.Spec{}, func() {}, fmt.Errorf("MCP server %q is not authorized; install it or complete project identity approval before connecting", server)
1222 }
1223 return spec, unlock, nil
1224 }
1225
1226 func (t *UseCapabilityTool) bindRuntimeMCP(spec plugin.Spec, target tool.Tool) tool.Tool {
1227 if t.runtime == nil || target == nil {
1228 return target
1229 }
1230 return &runtimeBoundMCPTool{
1231 proxy: t,
1232 target: target,
1233 server: spec.Name,
1234 authorized: spec.ServerAuthorized(),
1235 }
1236 }
1237
1238 func (t *UseCapabilityTool) withRuntimeBoundMCP(ctx context.Context, server string, target tool.Tool, execute func() error) error {
1239 if t.runtime == nil {
1240 return execute()
1241 }
1242 spec, unlock, err := t.lockAuthorizedRuntimeServer(ctx, server)
1243 if err != nil {
1244 return err
1245 }
1246 defer unlock()
1247 if !plugin.MCPToolMatchesSpec(target, spec) {
1248 return fmt.Errorf("connected MCP server %q identity does not match the current runtime configuration; reconnect this server before retrying", server)
1249 }
1250 return t.runtime.withServerGate(ctx, server, execute)
1251 }
1252
1253 func (t *UseCapabilityTool) serverEnabled(server string) bool {
1254 if t.runtime != nil {
1255 return t.runtime.serverEnabled(server)
1256 }
1257 // Standalone proxies predate the authoritative runtime and may resolve
1258 // already-registered MCP tools without carrying a duplicate spec slice.
1259 return true
1260 }
1261
1262 func (t *UseCapabilityTool) configuredServers() []mcpRuntimeServer {
1263 if t.runtime != nil {
1264 return t.runtime.configuredServers()
1265 }
1266 servers := make([]mcpRuntimeServer, 0, len(t.specs))
1267 seen := map[string]bool{}
1268 for _, raw := range t.specs {
1269 spec := cloneMCPSpec(raw)
1270 name := strings.TrimSpace(spec.Name)
1271 if name == "" || seen[name] {
1272 continue
1273 }
1274 seen[name] = true
1275 servers = append(servers, mcpRuntimeServer{
1276 entry: config.PluginEntry{Name: name},
1277 spec: spec,
1278 enabled: true,
1279 })
1280 }
1281 sort.Slice(servers, func(i, j int) bool { return servers[i].spec.Name < servers[j].spec.Name })
1282 return servers
1283 }
1284
1285 func (t *UseCapabilityTool) currentCatalog() capability.Catalog {
1286 if t.catalog != nil {
1287 return t.catalog()
1288 }
1289 return capability.Catalog{}
1290 }
1291
1292 // parseMCPServerCapabilityID extracts the server name from an mcp-server id.
1293 func parseMCPServerCapabilityID(id string) (string, bool) {
1294 if !strings.HasPrefix(id, "mcp-server:") {
1295 return "", false
1296 }
1297 name := strings.TrimSpace(strings.TrimPrefix(id, "mcp-server:"))
1298 return name, name != ""
1299 }
1300
1301 // resolveServerConnect resolves action=call on an mcp-server id. A connected
1302 // server lists its tools immediately (side-effect free); an unconnected one
1303 // resolves to a deferred connect target that runs only after the permission
1304 // gate and PreToolUse hooks approve it. Stored project authorization is applied
1305 // at resolve time so unauthorized project MCP never reaches process startup.
1306 func (t *UseCapabilityTool) resolveServerConnect(ctx context.Context, server string, base tool.ResolvedCall) (tool.ResolvedCall, error) {
1307 id := "mcp-server:" + server
1308 if !t.serverEnabled(server) {
1309 return t.resolveUnavailable(base, id, plugin.ToolPrefix(server), t.serverUnavailableReason(server)), nil
1310 }
1311 if t.host != nil && t.host.HasClient(server) {
1312 out, err := t.listServerTools(ctx, server)
1313 if err != nil {
1314 return t.resolveUnavailable(base, id, plugin.ToolPrefix(server), err.Error()), nil
1315 }
1316 base.SkipExecute = true
1317 base.HostCompleted = true
1318 base.Result = out
1319 base.ReadOnly = true
1320 return base, nil
1321 }
1322 spec, unlock, err := t.lockAuthorizedRuntimeServer(ctx, server)
1323 if err != nil {
1324 return t.resolveUnavailable(base, id, plugin.ToolPrefix(server), err.Error()), nil
1325 }
1326 unlock()
1327 connect := &onDemandMCPConnect{proxy: t, spec: spec, server: server}
1328 base.Target = connect
1329 // A dedicated exact identity names the connect for permission and hook
1330 // rules. It cannot collide with a real mcp__ tool, and rules do not need to
1331 // rely on unsupported tool-name glob matching.
1332 base.TargetName = connect.Name()
1333 // Connecting spawns a subprocess, so it is never a read-only fast path for
1334 // ordinary Plan/strict agents. PlannerMCPExecution may allow authorized
1335 // connects; unauthorized specs are blocked before process/network start.
1336 base.ReadOnly = false
1337 base.Args = json.RawMessage(`{}`)
1338 return base, nil
1339 }
1340
1341 // onDemandMCPConnect is the deferred first-discovery target: it connects the
1342 // server post-approval and returns the live tool directory.
1343 type onDemandMCPConnect struct {
1344 proxy *UseCapabilityTool
1345 spec plugin.Spec
1346 server string
1347 }
1348
1349 func (o *onDemandMCPConnect) Name() string { return plugin.MCPConnectPermissionName(o.server) }
1350
1351 func (o *onDemandMCPConnect) Description() string {
1352 return "connect MCP server " + o.server + " on demand and list its tools"
1353 }
1354
1355 func (o *onDemandMCPConnect) Schema() json.RawMessage { return json.RawMessage(`{"type":"object"}`) }
1356
1357 func (o *onDemandMCPConnect) ReadOnly() bool { return false }
1358
1359 // MCPLifecycleConnect marks this target as an MCP connect-and-list lifecycle
1360 // action for Planner authorization (not a remote tools/call).
1361 func (o *onDemandMCPConnect) MCPLifecycleConnect() bool { return true }
1362
1363 func (o *onDemandMCPConnect) MCPServerAuthorized() bool {
1364 return o.spec.ServerAuthorized()
1365 }
1366
1367 func (o *onDemandMCPConnect) MCPServerName() string { return o.server }
1368
1369 func (o *onDemandMCPConnect) ReadOnlyExecutionHostMutation() bool { return true }
1370
1371 func (o *onDemandMCPConnect) ReadOnlyExecutionBlockReason() string {
1372 if !o.spec.ServerAuthorized() {
1373 return "start an unauthorized MCP server (install it or complete project identity approval first)"
1374 }
1375 return "connect this MCP server from a parent session first"
1376 }
1377
1378 func (o *onDemandMCPConnect) Execute(ctx context.Context, _ json.RawMessage) (string, error) {
1379 // Zero process/network start when authorization, enable state, or exact
1380 // runtime identity changed after resolve.
1381 spec, unlock, err := o.proxy.lockAuthorizedRuntimeServer(ctx, o.server)
1382 if err != nil {
1383 msg := err.Error()
1384 if o.proxy.ledger != nil {
1385 o.proxy.ledger.MarkUnavailable("mcp-server:"+o.server, msg)
1386 }
1387 return "", err
1388 }
1389 defer unlock()
1390 if !plugin.MCPRuntimeSpecMatches(spec, o.spec) {
1391 return "", fmt.Errorf("MCP server %q runtime identity changed after resolution; retry so Reasonix can bind the current configuration", o.server)
1392 }
1393 if _, err := o.proxy.ensureServerToolsForSpec(ctx, o.server, spec); err != nil {
1394 if o.proxy.ledger != nil {
1395 o.proxy.ledger.MarkUnavailable("mcp-server:"+o.server, err.Error())
1396 }
1397 return "", err
1398 }
1399 return o.proxy.listServerToolsForSpec(ctx, o.server, spec)
1400 }
1401
1402 func parseMCPCapabilityID(id string) (server, raw string, err error) {
1403 id = strings.TrimSpace(id)
1404 switch {
1405 case strings.HasPrefix(id, "mcp-tool:"):
1406 rest := strings.TrimPrefix(id, "mcp-tool:")
1407 server, raw, ok := strings.Cut(rest, "/")
1408 if !ok || server == "" || raw == "" {
1409 return "", "", fmt.Errorf("invalid mcp-tool id %q; want mcp-tool:<server>/<tool>", id)
1410 }
1411 return server, raw, nil
1412 case strings.HasPrefix(id, "mcp-server:"):
1413 return "", "", fmt.Errorf("%q is a server id; call it directly to connect and list tools, or use mcp-tool:<server>/<tool>", id)
1414 default:
1415 return "", "", fmt.Errorf("action=call requires an mcp-tool capability id, got %q", id)
1416 }
1417 }
1418
1419 // Ensure UseCapabilityTool satisfies the tool contracts used by the agent.
1420 var (
1421 _ tool.Tool = (*UseCapabilityTool)(nil)
1422 _ tool.CallResolver = (*UseCapabilityTool)(nil)
1423 _ tool.BatchClassifier = (*UseCapabilityTool)(nil)
1424 )
1425
1426 // EmitProxyAudit is a helper for frontends: returns a notice describing the
1427 // proxy name and real target for user audit trails.
1428 func EmitProxyAudit(sink event.Sink, resolved tool.ResolvedCall, callIDs ...string) {
1429 if sink == nil || resolved.TargetName == "" {
1430 return
1431 }
1432 callID := ""
1433 if len(callIDs) > 0 {
1434 callID = callIDs[0]
1435 }
1436 detail, _ := json.Marshal(map[string]string{"callId": callID, "capabilityId": resolved.CapabilityID, "target": resolved.TargetName})
1437 sink.Emit(event.Event{
1438 Kind: event.Notice,
1439 Level: event.LevelInfo,
1440 Code: "capability_proxy_audit",
1441 Text: capabilityProxyNoticeText(resolved.DisplayName, resolved.TargetName),
1442 Detail: string(detail),
1443 })
1444 }
1445
1445 lines GO