Files
Maven/internal/mcp/client.go
kami 95ae900a58 Talk MCP: a client for external tool servers (#251)
docs/plans/06-mcp-support.md asks for the host direction — Maven connects OUT
to MCP servers and consumes what they offer. This is the client half: the
protocol, the transports, the connection manager, the config block. Nothing is
wired into a turn yet, and nothing here exposes Maven's own capabilities to an
outside caller.

internal/mcp:
  - hand-rolled JSON-RPC 2.0 (the wire format is four fields, and the repo
    vendors its deps, so a library would cost more than it saves);
  - two transports: a stdio subprocess on this box, and streamable HTTP, which
    accepts a plain JSON reply or an SSE frame because servers disagree about
    which they send;
  - Client: initialize handshake, tools/list, tools/call, resources/list,
    resources/read. Text content only — everything downstream is a sentence;
  - Manager: lazy dial, per-server failure that never blocks boot or the other
    servers, backoff reconnect, Status for a web surface, graceful Close;
  - the allowlist encoding: a discovered tool becomes the store row
    "vikunja_list_tasks" with cmd ["mcp","vikunja","list_tasks"], scope
    "mcp:vikunja". No new column, no migration, and ProposeTool, EnableTool,
    the act matcher and the confirm turn all keep working untouched.

Constraints held, in code rather than in prose:
  - OFF unless configured, and a server is dark until "enabled": true.
  - A url server goes through internal/webfetch, so the SSRF guard, the size
    cap, the redirect cap and the per-host rate limit apply. Reaching loopback
    needs allow_private on THAT server, and each server gets its own fetcher so
    one loopback exemption cannot become a hole for a public endpoint.
  - readOnlyHint decides destructive: no hint means "assume it mutates", which
    will route the call through the existing confirm turn. Guessing wrong in
    that direction only costs a question.
  - The catalogue stays small on purpose — allow_tools, and max_tools=12 per
    server. The resident model is a 1.7B with a 4096-token context; a tool name
    it half-remembers is a wrong act.
  - Only the tool name and the router's arguments are sent. There is no API
    here through which a note, a fact or the persona block could travel.

webfetch grows Post (JSON-RPC cannot be a GET) and surfaces response headers
for Mcp-Session-Id. It shares Get's guards exactly: a body buys a caller
nothing, a POST to the LAN is refused for the same reason a GET is.

Verified against the real Vikunja MCP server on homesrv
(http://localhost:9100/mcp): handshake, three discovered tools with update_task
correctly NOT read-only, a live list_projects call, a tool excluded by
allow_tools refused, and the same server refused outright once allow_private
was dropped. Tests cover both transports (the stdio one against a real
subprocess), SSE and JSON framing, session echo, reconnect, and the config
validation.
2026-08-01 04:22:52 +04:00

268 lines
7.7 KiB
Go

package mcp
import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"sync"
"sync/atomic"
)
// Errors callers distinguish.
var (
// ErrClosed — the transport is gone (subprocess died, client closed).
ErrClosed = errors.New("mcp: connection is closed")
// ErrNotInitialized — a call was made before the initialize handshake.
ErrNotInitialized = errors.New("mcp: not initialized")
// ErrToolFailed — the server ran the tool and reported an error result.
ErrToolFailed = errors.New("mcp: tool reported an error")
)
// Tool is one tool a server offers, in the form Maven cares about.
//
// ReadOnly comes from the server's own readOnlyHint annotation and decides
// whether the allowlist row is marked destructive: no hint, or a false one,
// means "assume it mutates", which routes the call through the confirm turn.
// Guessing wrong in that direction only costs a question.
type Tool struct {
Server string
Name string
Description string
InputSchema json.RawMessage
ReadOnly bool
}
// Resource is one resource a server offers. Contents are fetched separately —
// listing is cheap, reading is not.
type Resource struct {
Server string
URI string
Name string
MIMEType string
}
// ServerInfo is what came back from the handshake.
type ServerInfo struct {
Name string `json:"name"`
Version string `json:"version"`
ProtocolVersion string `json:"-"`
}
// Client is one connected MCP server. Safe for concurrent use.
type Client struct {
name string
tr transport
next atomic.Int64
mu sync.Mutex
info ServerInfo
ready bool
}
// newClient wraps a transport. Callers use Dial* in manager.go.
func newClient(name string, tr transport) *Client {
return &Client{name: name, tr: tr}
}
// Name — the local name of this server (the config key, not the server's own).
func (c *Client) Name() string { return c.name }
// Info — what the server said about itself during the handshake.
func (c *Client) Info() ServerInfo {
c.mu.Lock()
defer c.mu.Unlock()
return c.info
}
// Initialize performs the MCP handshake and sends notifications/initialized.
// Capabilities we declare are empty on purpose: Maven consumes, she does not
// offer sampling or roots back to the server.
func (c *Client) Initialize(ctx context.Context) error {
var out struct {
ProtocolVersion string `json:"protocolVersion"`
ServerInfo ServerInfo `json:"serverInfo"`
}
err := c.call(ctx, "initialize", map[string]any{
"protocolVersion": ProtocolVersion,
"capabilities": map[string]any{},
"clientInfo": map[string]any{"name": "maven", "version": "1.0"},
}, &out)
if err != nil {
return err
}
if strings.TrimSpace(out.ProtocolVersion) == "" {
return fmt.Errorf("mcp: %s: handshake returned no protocol version", c.name)
}
out.ServerInfo.ProtocolVersion = out.ProtocolVersion
c.mu.Lock()
c.info, c.ready = out.ServerInfo, true
c.mu.Unlock()
// Best effort: a stateless HTTP server may not care, and a failure here is
// not worth dropping a working connection over.
_ = c.tr.Notify(ctx, "notifications/initialized", map[string]any{})
return nil
}
// ListTools discovers the server's tools.
func (c *Client) ListTools(ctx context.Context) ([]Tool, error) {
if !c.initialized() {
return nil, ErrNotInitialized
}
var out struct {
Tools []struct {
Name string `json:"name"`
Description string `json:"description"`
InputSchema json.RawMessage `json:"inputSchema"`
Annotations *struct {
ReadOnlyHint bool `json:"readOnlyHint"`
} `json:"annotations"`
} `json:"tools"`
}
if err := c.call(ctx, "tools/list", map[string]any{}, &out); err != nil {
return nil, err
}
tools := make([]Tool, 0, len(out.Tools))
for _, t := range out.Tools {
if strings.TrimSpace(t.Name) == "" {
continue
}
tools = append(tools, Tool{
Server: c.name,
Name: t.Name,
Description: strings.TrimSpace(t.Description),
InputSchema: t.InputSchema,
ReadOnly: t.Annotations != nil && t.Annotations.ReadOnlyHint,
})
}
return tools, nil
}
// CallTool runs one tool and returns its text content, joined by newlines.
// Non-text content (images, blobs) is dropped: everything downstream of here
// is a spoken or written sentence.
//
// args is exactly what the router produced. Nothing else — no history, no
// notes, no persona — is in scope here, by construction.
func (c *Client) CallTool(ctx context.Context, name string, args map[string]any) (string, error) {
if !c.initialized() {
return "", ErrNotInitialized
}
if args == nil {
args = map[string]any{}
}
var out struct {
IsError bool `json:"isError"`
Content []struct {
Type string `json:"type"`
Text string `json:"text"`
} `json:"content"`
}
if err := c.call(ctx, "tools/call", map[string]any{"name": name, "arguments": args}, &out); err != nil {
return "", err
}
var parts []string
for _, ct := range out.Content {
if ct.Type == "text" && strings.TrimSpace(ct.Text) != "" {
parts = append(parts, strings.TrimSpace(ct.Text))
}
}
text := strings.Join(parts, "\n")
if out.IsError {
return text, fmt.Errorf("%w: %s/%s: %s", ErrToolFailed, c.name, name, text)
}
return text, nil
}
// ListResources discovers the server's resources. A server without the
// resources capability answers with an error; that is not fatal, the caller
// gets an empty list.
func (c *Client) ListResources(ctx context.Context) ([]Resource, error) {
if !c.initialized() {
return nil, ErrNotInitialized
}
var out struct {
Resources []struct {
URI string `json:"uri"`
Name string `json:"name"`
MIMEType string `json:"mimeType"`
} `json:"resources"`
}
if err := c.call(ctx, "resources/list", map[string]any{}, &out); err != nil {
return nil, err
}
res := make([]Resource, 0, len(out.Resources))
for _, r := range out.Resources {
if strings.TrimSpace(r.URI) == "" {
continue
}
res = append(res, Resource{Server: c.name, URI: r.URI, Name: r.Name, MIMEType: r.MIMEType})
}
return res, nil
}
// ReadResource returns a resource's text contents, joined by newlines. This is
// the RAG-hint path: the text can be pasted into a router or phraser prompt.
func (c *Client) ReadResource(ctx context.Context, uri string) (string, error) {
if !c.initialized() {
return "", ErrNotInitialized
}
var out struct {
Contents []struct {
Text string `json:"text"`
} `json:"contents"`
}
if err := c.call(ctx, "resources/read", map[string]any{"uri": uri}, &out); err != nil {
return "", err
}
var parts []string
for _, ct := range out.Contents {
if strings.TrimSpace(ct.Text) != "" {
parts = append(parts, strings.TrimSpace(ct.Text))
}
}
return strings.Join(parts, "\n"), nil
}
// Close drops the connection.
func (c *Client) Close() error {
c.mu.Lock()
c.ready = false
c.mu.Unlock()
return c.tr.Close()
}
func (c *Client) initialized() bool {
c.mu.Lock()
defer c.mu.Unlock()
return c.ready
}
// alive reports whether the underlying transport can still carry a call. HTTP
// is stateless, so it is always alive; a dead subprocess is not.
func (c *Client) alive() bool {
if s, ok := c.tr.(*stdioTransport); ok {
return s.alive()
}
return true
}
func (c *Client) call(ctx context.Context, method string, params any, out any) error {
req := &rpcRequest{JSONRPC: "2.0", ID: c.next.Add(1), Method: method, Params: params}
resp, err := c.tr.Call(ctx, req)
if err != nil {
return fmt.Errorf("mcp: %s: %s: %w", c.name, method, err)
}
if resp.Error != nil {
return fmt.Errorf("mcp: %s: %s: %w", c.name, method, resp.Error)
}
if out == nil || len(resp.Result) == 0 {
return nil
}
if err := json.Unmarshal(resp.Result, out); err != nil {
return fmt.Errorf("mcp: %s: %s: decode result: %w", c.name, method, err)
}
return nil
}