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.
This commit is contained in:
@@ -0,0 +1,149 @@
|
||||
package mcp
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// maxLine bounds one JSON-RPC frame from a subprocess. A tool result bigger
|
||||
// than this is a misbehaving server, not something to buffer.
|
||||
const maxLine = 1 << 20 // 1 MiB
|
||||
|
||||
// stdioTransport speaks newline-delimited JSON-RPC to a child process. This is
|
||||
// the local transport: the server runs on this box, under this user, and gets
|
||||
// no network guard because it never touches the network on our behalf.
|
||||
//
|
||||
// Args are argv, never a shell string — the same discipline internal/tool
|
||||
// keeps, for the same reason.
|
||||
type stdioTransport struct {
|
||||
mu sync.Mutex
|
||||
cmd *exec.Cmd
|
||||
in io.WriteCloser
|
||||
out *bufio.Reader
|
||||
dead bool
|
||||
}
|
||||
|
||||
func newStdioTransport(ctx context.Context, argv []string, env []string, dir string) (*stdioTransport, error) {
|
||||
if len(argv) == 0 {
|
||||
return nil, errors.New("mcp: stdio server needs a command")
|
||||
}
|
||||
cmd := exec.Command(argv[0], argv[1:]...)
|
||||
cmd.Dir = dir
|
||||
if len(env) > 0 {
|
||||
cmd.Env = append(os.Environ(), env...)
|
||||
}
|
||||
cmd.Stderr = os.Stderr
|
||||
in, err := cmd.StdinPipe()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("mcp: stdin pipe: %w", err)
|
||||
}
|
||||
out, err := cmd.StdoutPipe()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("mcp: stdout pipe: %w", err)
|
||||
}
|
||||
if err := cmd.Start(); err != nil {
|
||||
return nil, fmt.Errorf("mcp: start %q: %w", argv[0], err)
|
||||
}
|
||||
return &stdioTransport{cmd: cmd, in: in, out: bufio.NewReaderSize(out, 64<<10)}, nil
|
||||
}
|
||||
|
||||
func (t *stdioTransport) Call(ctx context.Context, req *rpcRequest) (*rpcResponse, error) {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
if t.dead {
|
||||
return nil, ErrClosed
|
||||
}
|
||||
if err := t.write(req); err != nil {
|
||||
t.dead = true
|
||||
return nil, err
|
||||
}
|
||||
// Read until the frame with our id turns up; anything else on the pipe is
|
||||
// a notification or a server-initiated request we do not answer.
|
||||
for {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
line, err := t.readLine()
|
||||
if err != nil {
|
||||
t.dead = true
|
||||
return nil, err
|
||||
}
|
||||
var resp rpcResponse
|
||||
if err := json.Unmarshal(line, &resp); err != nil {
|
||||
continue // not a response frame; ignore rather than break the turn
|
||||
}
|
||||
if resp.ID == nil || *resp.ID != req.ID {
|
||||
continue
|
||||
}
|
||||
return &resp, nil
|
||||
}
|
||||
}
|
||||
|
||||
func (t *stdioTransport) Notify(ctx context.Context, method string, params any) error {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
if t.dead {
|
||||
return ErrClosed
|
||||
}
|
||||
return t.write(&rpcRequest{JSONRPC: "2.0", Method: method, Params: params})
|
||||
}
|
||||
|
||||
func (t *stdioTransport) write(req *rpcRequest) error {
|
||||
req.JSONRPC = "2.0"
|
||||
raw, err := json.Marshal(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := t.in.Write(append(raw, '\n')); err != nil {
|
||||
return fmt.Errorf("mcp: write %s: %w", req.Method, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (t *stdioTransport) readLine() ([]byte, error) {
|
||||
for {
|
||||
line, err := t.out.ReadString('\n')
|
||||
if err != nil {
|
||||
if len(strings.TrimSpace(line)) == 0 {
|
||||
return nil, fmt.Errorf("mcp: read: %w", err)
|
||||
}
|
||||
return []byte(line), nil
|
||||
}
|
||||
if len(line) > maxLine {
|
||||
return nil, fmt.Errorf("mcp: frame exceeds %d bytes", maxLine)
|
||||
}
|
||||
if s := strings.TrimSpace(line); s != "" {
|
||||
return []byte(s), nil
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (t *stdioTransport) Close() error {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
t.dead = true
|
||||
if t.in != nil {
|
||||
_ = t.in.Close()
|
||||
}
|
||||
if t.cmd.Process != nil {
|
||||
_ = t.cmd.Process.Kill()
|
||||
_ = t.cmd.Wait()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// alive reports whether the transport can still carry a call. The manager uses
|
||||
// it to decide on a reconnect instead of retrying into a dead pipe.
|
||||
func (t *stdioTransport) alive() bool {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
return !t.dead
|
||||
}
|
||||
Reference in New Issue
Block a user