95ae900a58
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.
268 lines
7.7 KiB
Go
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
|
|
}
|