Compare commits

...

3 Commits

Author SHA1 Message Date
kami d92349ca6e Store and describe images through a shared media intake (#252)
Vision needs a second model this box does not have, so the shipped half is
the part that works without one: an image arrives, is sniffed, is stored
content-addressed, and is prepared for inference. The describing half is
written and tested against a fake server, and refuses any endpoint that is
not on this box.

internal/media is the intake all three senses share — hearing and speaker
recognition store their audio in the same place under the same retention.
Blobs stay out of the sqlite store; only the derived text becomes a note,
and only when the caller asks. Retention is enforced by an hourly prune
loop rather than by a comment.

The plan's RemoteProvider step is refused: no cloud model, inference stays
on the box, and vision.NewLocal validates that at construction.
2026-08-01 04:53:07 +04:00
kami 8d5e357b57 Expose discovered MCP tools through the act allowlist (#251)
Second half of the MCP client: the tools the manager discovers become rows in
the existing act allowlist instead of a parallel capability system.

An MCP tool is encoded in the columns that already exist — cmd
["mcp",<server>,<tool>], scope mcp:<server> — so no migration, and
ProposeTool/EnableTool/DisableTool, tool.Matcher and the confirm turn need no
changes. One branch in Executor.Exec routes such a row to the manager instead
of exec, and "mcp" is never run as a binary.

Discovery only ever PROPOSES. destructive comes from the inverse of the MCP
readOnlyHint, so a tool that does not promise to be read-only inherits the
confirm turn, and enabling stays on /tools behind step-up.

Voice args are positional and MCP args are named, so CallPositional binds only
what it can defend: no required properties runs bare, and a read-only tool with
exactly one required string or number gets the tail. Everything else refuses
with ErrNeedsArgs rather than guessing. The read-only condition was learned
against the live Vikunja server: update_task requires only task_id and takes
the rest as optional, so one guessed argument blanked the fields it did not
mention. A partially-filled write destroys what it omits, so a mutating tool
never receives a guessed argument.

Also: a read-only mcp_servers IPC method and an "MCP servers" card on /tools
showing transport, target and state, with the trust level of a local target
spelled out. There is deliberately no call-a-tool IPC method and no run button,
so mutation keeps exactly one path.

Vikunja #251
2026-08-01 04:36:40 +04:00
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
44 changed files with 5415 additions and 36 deletions
+7
View File
@@ -5,6 +5,7 @@ import (
"errors"
"log"
"github.com/kami/maven/internal/mcp"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/tool"
)
@@ -52,6 +53,12 @@ func (h *reactiveHandler) actionAct(ctx context.Context, dec router.Decision) st
return "выполнить «" + phrase + "»? скажи «да» или «нет»."
case errors.Is(err, tool.ErrNotEnabled):
return h.proposeGap(ctx, dec)
case errors.Is(err, mcp.ErrNeedsArgs):
// An MCP tool that wants named arguments a spoken verb cannot
// supply. Guessing them would be a wrong act, so she says so
// instead — the tool is still runnable from the authed surface,
// where a human types them.
return "этому инструменту нужны аргументы, которые я из голоса не соберу — я не буду угадывать."
}
log.Printf("voice: tool %s: %v", dec.Slots.Fn, err)
if out != "" {
+19
View File
@@ -280,6 +280,9 @@ func run(args []string) error {
api := coreAPI.(*daemonAPI)
api.chatFn = voiceW.handler.handleText
}
if voiceW != nil && voiceW.mcp != nil {
coreAPI.(*daemonAPI).getMCPServers = voiceW.mcp.status
}
} else {
// locked mode: no real store yet, so there's no meaningful CoreAPI to
// serve. srv.Check below is the actual guard — every CoreAPI call is
@@ -330,6 +333,9 @@ func run(args []string) error {
if !locked {
wireMailIntake(srv, st, phr, cfg)
wireModelSwap(srv, phr, cfg)
// Vision + the media blob store (Vikunja #252). Both stay dark without a
// media block; MethodDescribeImage answers ErrUnknownMethod then.
wireVision(ctx, srv, st, embedderOf(voiceW), cfg)
}
// WrapKeyFn — wraps the env key with a passkey credential public key and
@@ -473,6 +479,7 @@ func run(args []string) error {
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
wireMailIntake(srv, st, phr, cfg)
wireModelSwap(srv, phr, cfg)
wireVision(ctx, srv, st, embedderOf(voiceW), cfg)
// Start voice server.
if voiceW != nil {
@@ -518,6 +525,11 @@ func run(args []string) error {
}()
}
// Keep MCP connections alive (nil unless configured).
if voiceW != nil && voiceW.mcp != nil {
go voiceW.mcp.run(ctx)
}
dl.unlock()
log.Printf("mavend: unlocked via passkey assertion")
return nil
@@ -577,6 +589,13 @@ func run(args []string) error {
crawlWkr.run(ctx)
}()
}
if voiceW != nil && voiceW.mcp != nil {
wg.Add(1)
go func() {
defer wg.Done()
voiceW.mcp.run(ctx)
}()
}
}
<-ctx.Done()
+154
View File
@@ -0,0 +1,154 @@
package main
import (
"context"
"fmt"
"log"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/mcp"
"github.com/kami/maven/internal/store"
"github.com/kami/maven/internal/webfetch"
)
// mcpRefreshInterval — how often the manager re-dials a server that is down.
// The manager applies its own backoff on top, so this being short is cheap.
const mcpRefreshInterval = time.Minute
// mcpWiring — the MCP client, when the `mcp` block configures at least one
// enabled server. nil ⇒ nothing was configured, nothing is connected, and an
// allowlist row that happens to look like an MCP row refuses to run.
//
// It lives on the voice wiring because MCP tools ARE acts: they run through
// tool.Executor, the enabled allowlist and the confirm turn, which only exist
// on the voice/chat path. No voice surface ⇒ nothing that could call a tool.
type mcpWiring struct {
mgr *mcp.Manager
st *store.Store
}
// wireMCP builds the manager, connects, and proposes what it found. It never
// fails the daemon: a server that is unreachable at boot is logged and retried,
// because Maven starting is not contingent on someone else's process.
func wireMCP(cfg *config.Config, st *store.Store) *mcpWiring {
servers := cfg.MCPServers()
if len(servers) == 0 {
return nil
}
limits := webfetch.Config{}
if cfg.MCP != nil {
limits.AllowHosts = cfg.MCP.AllowHosts
limits.DenyHosts = cfg.MCP.DenyHosts
limits.MaxBytes = cfg.MCP.MaxBytes
limits.Timeout = time.Duration(cfg.MCP.Timeout)
}
mgr, err := mcp.NewManager(mcp.WebfetchDoor(limits), servers)
if err != nil {
// Validation already ran in config.validate, so this is a programming
// error rather than a config one. Still not fatal: MCP off is a working
// Maven.
log.Printf("mcp: not wired: %v", err)
return nil
}
w := &mcpWiring{mgr: mgr, st: st}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
mgr.Connect(ctx)
w.propose(ctx)
return w
}
// propose writes a 'proposed' allowlist row for every discovered tool. It does
// NOT enable anything: a configured server is a place Maven may look, not a
// capability she has. Kami enables what he wants on /tools, behind step-up,
// which is the same gate a shell tool goes through.
//
// Re-running on every boot is idempotent — ProposeMCPTool never touches an
// existing row, so a tool he disabled stays disabled and one he enabled keeps
// the cmd he enabled it with.
func (w *mcpWiring) propose(ctx context.Context) {
if w == nil {
return
}
now := time.Now()
fresh := 0
for _, t := range w.mgr.Tools() {
name := mcp.LocalName(t.Server, t.Name)
// No readOnlyHint ⇒ assume it mutates ⇒ the confirm turn. Being wrong
// in this direction only costs a question.
destructive := !t.ReadOnly
provenance := fmt.Sprintf("mcp %s/%s", t.Server, t.Name)
if t.Description != "" {
provenance += ": " + t.Description
}
ok, err := w.st.ProposeMCPTool(ctx, name, mcp.Scope(t.Server),
mcp.Cmd(t.Server, t.Name), destructive, provenance, now)
if err != nil {
log.Printf("mcp: propose %s: %v", name, err)
continue
}
if ok {
fresh++
}
}
if fresh > 0 {
log.Printf("mcp: %d new tool proposal(s) waiting on /tools", fresh)
}
}
// run re-dials downed servers and picks up tools that appeared, until ctx is
// canceled.
func (w *mcpWiring) run(ctx context.Context) {
if w == nil {
return
}
t := time.NewTicker(mcpRefreshInterval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
w.mgr.Refresh(ctx)
w.propose(ctx)
}
}
}
// status maps the manager's view onto the wire type the web surface reads.
func (w *mcpWiring) status() []ipc.MCPServerStatus {
if w == nil {
return nil
}
in := w.mgr.Status()
out := make([]ipc.MCPServerStatus, 0, len(in))
for _, s := range in {
out = append(out, ipc.MCPServerStatus{
Name: s.Name,
Transport: s.Transport,
Target: s.Target,
Connected: s.Connected,
Server: s.Server,
Tools: s.Tools,
Err: s.Err,
})
}
return out
}
func (w *mcpWiring) close() {
if w == nil {
return
}
_ = w.mgr.Close()
}
// caller is the tool.MCPCaller the executor gets, or nil when MCP is off.
func (w *mcpWiring) caller() *mcp.Manager {
if w == nil {
return nil
}
return w.mgr
}
+77
View File
@@ -0,0 +1,77 @@
package main
import (
"context"
"strings"
"testing"
"github.com/kami/maven/internal/config"
)
func TestWireMCPOffWhenUnconfigured(t *testing.T) {
st := newTestStore(t)
for name, cfg := range map[string]*config.Config{
"no block": {},
"nothing enabled": {MCP: &config.MCPConfig{Servers: []config.MCPServerConfig{
{Name: "vikunja", URL: "http://192.168.1.104:9100/mcp"},
}}},
} {
t.Run(name, func(t *testing.T) {
if w := wireMCP(cfg, st); w != nil {
t.Fatal("MCP must be off unless a server is configured AND enabled")
}
})
}
// nil wiring must be safe to use everywhere it is reachable.
var w *mcpWiring
w.close()
w.propose(context.Background())
if w.status() != nil || w.caller() != nil {
t.Fatal("a nil wiring must report nothing")
}
}
// An unreachable server must not stop the daemon, must be reported as down, and
// must propose nothing.
func TestWireMCPUnreachableServerIsNotFatal(t *testing.T) {
st := newTestStore(t)
w := wireMCP(&config.Config{MCP: &config.MCPConfig{Servers: []config.MCPServerConfig{{
Name: "dead", Command: "/nonexistent/mcp-server", Enabled: true,
}}}}, st)
if w == nil {
t.Fatal("a configured server should still wire")
}
defer w.close()
st2 := w.status()
if len(st2) != 1 || st2[0].Connected || st2[0].Err == "" {
t.Fatalf("status = %+v", st2)
}
tools, err := st.ListTools(context.Background(), "")
if err != nil {
t.Fatal(err)
}
if len(tools) != 0 {
t.Fatalf("a server that never answered must propose nothing, got %+v", tools)
}
}
// A url server whose address is private is refused by webfetch unless that
// server sets allow_private. This is the guard the whole MCP path rides on, so
// it is asserted here too, at the wiring level.
func TestWireMCPPrivateURLRefusedWithoutAllowPrivate(t *testing.T) {
st := newTestStore(t)
w := wireMCP(&config.Config{MCP: &config.MCPConfig{Servers: []config.MCPServerConfig{{
Name: "lan", URL: "http://127.0.0.1:9100/mcp", Enabled: true,
}}}}, st)
if w == nil {
t.Fatal("should wire")
}
defer w.close()
s := w.status()[0]
if s.Connected {
t.Fatal("a loopback server must not connect without allow_private")
}
if !strings.Contains(s.Err, "private address") {
t.Fatalf("err = %q, want the private-address refusal", s.Err)
}
}
+11
View File
@@ -941,6 +941,7 @@ type daemonAPI struct {
getMorningStatus func(ctx context.Context) []ipc.MorningRoutineStatus
getDayPlan func(ctx context.Context) ipc.DayPlan
chatFn func(ctx context.Context, text string) string
getMCPServers func() []ipc.MCPServerStatus
}
func (d *daemonAPI) Chat(ctx context.Context, text string) (string, error) {
@@ -950,6 +951,16 @@ func (d *daemonAPI) Chat(ctx context.Context, text string) (string, error) {
return d.chatFn(ctx, text), nil
}
// MCPServers — the configured MCP servers and their health (Vikunja #251).
// Empty, not an error, when the mcp block is absent: "not configured" is the
// default state and the web surface renders it as such.
func (d *daemonAPI) MCPServers(ctx context.Context) ([]ipc.MCPServerStatus, error) {
if d.getMCPServers == nil {
return nil, nil
}
return d.getMCPServers(), nil
}
func (d *daemonAPI) TickTrace(ctx context.Context) (ipc.TickTrace, error) {
trace := d.getTrace()
if trace == nil {
+251
View File
@@ -0,0 +1,251 @@
// mavend/vision.go — core's half of image understanding (Vikunja #252,
// docs/plans/07-vision.md).
//
// The split: any surface that can receive a picture (mavweb upload, a Telegram
// photo through mavpoll, a path he names) hands the bytes to core over
// ipc.MethodDescribeImage. Core stores them content-addressed under
// media.dir, prepares a downscaled JPEG, and asks a local vision server what it
// is. The description comes back as words; nothing about the image is echoed.
//
// Off unless configured twice over: no `media` block ⇒ nowhere to keep the
// bytes, so the method does not exist; no `vision` block with enabled + a local
// endpoint ⇒ the store is wired but the describing half refuses, and the method
// still does not exist. A surface cannot make Maven look at pictures by merely
// sending one.
//
// Two things this file deliberately does not do:
//
// - No cloud vision call, ever. internal/vision refuses a non-private
// endpoint at construction; there is no config shape here that could reach
// an upstream API even if someone wanted one.
// - No automatic memory. SaveNote is opt-in per call. Glancing at a screenshot
// is not the same act as remembering it, and a 1.7B-class VLM's guess about
// a photo is not a fact worth carrying around.
package main
import (
"context"
"errors"
"fmt"
"log"
"path/filepath"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/media"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/store"
"github.com/kami/maven/internal/vision"
)
// prunePeriod — how often stored blobs are checked against media.retention.
// Hourly is far more often than needed for a 7-day retention and costs a
// directory walk over a handful of sidecars; the point is that the promise is
// kept by a loop that runs, not by an operator remembering a cron.
const prunePeriod = time.Hour
// mediaKeeper — the blob store plus the loop that enforces its retention. The
// two are one object because a store without the loop is a directory that grows
// forever, and shipping that would break the only interesting promise this
// capability makes.
type mediaKeeper struct {
store *media.Store
}
// openMediaStore builds the blob store from config, or returns nil when media is
// not configured. A relative dir resolves against StateDir, the same rule the db
// and socket paths follow.
func openMediaStore(cfg *config.Config) *mediaKeeper {
dir := cfg.Media.StoreDir()
if dir == "" {
return nil
}
if !filepath.IsAbs(dir) && cfg.StateDir != "" {
dir = filepath.Join(cfg.StateDir, dir)
}
st, err := media.Open(dir, cfg.Media.MaxBytes, time.Duration(cfg.Media.Retention))
if err != nil {
log.Printf("media: %v — image and audio intake disabled", err)
return nil
}
log.Printf("media: blob store at %s, retention %s", st.Dir(), st.Retention())
return &mediaKeeper{store: st}
}
// runPrune deletes over-retention blobs on a loop until ctx ends. It prunes once
// immediately, so a daemon restarted after a long downtime does not sit on a
// month of stale recordings until the first tick.
func (k *mediaKeeper) runPrune(ctx context.Context) {
prune := func() {
n, err := k.store.Prune()
if err != nil {
log.Printf("media: prune: %v", err)
return
}
if n > 0 {
log.Printf("media: pruned %d blob(s) older than %s", n, k.store.Retention())
}
}
prune()
t := time.NewTicker(prunePeriod)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
prune()
}
}
}
// visionIntake — one image at a time: store, prepare, describe, optionally note.
type visionIntake struct {
in *vision.Intake
st *store.Store
emb router.Embedder
now func() time.Time
}
// newVisionIntake returns nil when there is nothing to wire. keeper == nil means
// no media block, which disables the method outright; a missing or disabled
// vision block still wires the method, because storing an image and answering
// "I can't look at it yet" is more useful than pretending the surface does not
// exist — and it is exactly the state this box is in until a vision model is on
// disk.
func newVisionIntake(keeper *mediaKeeper, st *store.Store, emb router.Embedder, cfg *config.Config) *visionIntake {
if keeper == nil {
return nil
}
vc := cfg.Vision
maxDim := 0
var provider vision.Provider = vision.Disabled{}
if vc.LooksAtImages() {
p, err := vision.NewLocal(vision.Config{
Endpoint: vc.Endpoint,
Model: vc.Model,
Timeout: time.Duration(vc.Timeout),
MaxTokens: vc.MaxTokens,
Prompt: vc.Prompt,
})
if err != nil {
// A public endpoint, a hostname, a bad URL. Logged once here rather
// than failing every turn, and the store still works.
log.Printf("vision: %v — she can store images but not describe them", err)
} else {
provider = p
maxDim = vc.MaxDim
log.Printf("vision: enabled against %s", p.Endpoint())
}
} else {
log.Printf("vision: not configured — images are stored, not described")
}
return &visionIntake{
in: vision.NewIntake(keeper.store, provider, maxDim),
st: st,
emb: emb,
now: time.Now,
}
}
// describe handles one ipc.MethodDescribeImage call.
//
// A description failure is NOT an error out of this method when the bytes were
// stored: the caller gets the id and an empty description, which is honest ("it
// is kept, I cannot read it yet") and re-runnable. A failure to store, or bytes
// that are not an image at all, is an error — there is nothing to come back to.
func (v *visionIntake) describe(ctx context.Context, req ipc.DescribeImageReq) (ipc.DescribeImageResp, error) {
if len(req.Data) == 0 && req.ID == "" {
return ipc.DescribeImageResp{}, fmt.Errorf("describe image: neither data nor id")
}
var (
res vision.Result
err error
)
if req.ID != "" {
res, err = v.in.Rerun(ctx, req.ID, req.Question)
} else {
res, err = v.in.Accept(ctx, req.Data, sourceOrDefault(req.Source), req.Question)
}
if res.Blob.ID == "" {
// Nothing was stored: bad format, over the size cap, unwritable dir.
return ipc.DescribeImageResp{}, fmt.Errorf("describe image: %w", err)
}
resp := ipc.DescribeImageResp{
ID: res.Blob.ID,
Description: res.Description,
Width: res.Image.Width,
Height: res.Image.Height,
}
if err != nil {
// Bytes are safe, words are not available. The log names the blob and the
// reason; it never names what was in the picture.
if errors.Is(err, vision.ErrDisabled) {
log.Printf("vision: stored %s, no vision model configured", res.Blob)
} else {
log.Printf("vision: stored %s, describe failed: %v", res.Blob, err)
}
return resp, nil
}
if req.SaveNote {
id, werr := v.writeNote(ctx, res)
if werr != nil {
// The description is still returned: losing the note is worse as a
// silent failure than as a log line next to a successful answer.
log.Printf("vision: note write for %s failed: %v", res.Blob, werr)
} else {
resp.NoteID = id
}
}
log.Printf("vision: described %s (%dx%d)", res.Blob, res.Image.Width, res.Image.Height)
return resp, nil
}
// writeNote stores the description as an ordinary note so it is recallable. The
// note carries the blob id in its source, which is the only link back to the
// bytes — the note text is words about the picture, never the picture.
func (v *visionIntake) writeNote(ctx context.Context, res vision.Result) (int64, error) {
var vec []float32
if v.emb != nil {
// EmbedPassage, not Embed: a description is text being searched FOR, and
// the e5 embedder is asymmetric. Backwards here makes it unfindable by
// the question that should have matched it.
var err error
vec, err = router.EmbedPassage(ctx, v.emb, res.Description)
if err != nil {
return 0, fmt.Errorf("embed: %w", err)
}
}
source := "media:image:" + res.Blob.ID[:12]
return v.st.WriteNote(ctx, v.now(), res.Description, vec, source)
}
// sourceOrDefault labels a blob whose sender did not say where it came from.
func sourceOrDefault(s string) string {
if s == "" {
return "unknown"
}
return s
}
// wireVision installs the IPC hook and starts the retention loop, or leaves the
// hook nil so ipc.MethodDescribeImage reports ErrUnknownMethod. Called on both
// startup paths (unlocked boot and passkey unlock) so vision behaves the same
// either way.
func wireVision(ctx context.Context, srv *ipc.Server, st *store.Store, emb router.Embedder, cfg *config.Config) {
keeper := openMediaStore(cfg)
if keeper == nil {
return
}
go keeper.runPrune(ctx)
vi := newVisionIntake(keeper, st, emb, cfg)
if vi == nil {
return
}
srv.DescribeImageFn = vi.describe
}
+13
View File
@@ -40,6 +40,10 @@ type voiceWiring struct {
// mavsttd / mavttsd don't keep a stale conn into a restarting daemon.
sttClient *worker.Client
ttsClient *worker.Client
// mcp — the MCP client, nil unless the `mcp` block configures an enabled
// server (Vikunja #251). Its tools land in the same allowlist as every
// other act, so nothing else here has to know about it.
mcp *mcpWiring
}
// close releases the listener + worker conns. Safe to call on nil (when
@@ -60,6 +64,7 @@ func (w *voiceWiring) close() {
if w.ttsClient != nil {
_ = w.ttsClient.Close()
}
w.mcp.close()
}
// wireVoice builds the audio path from cfg + a CoreAPI + a router. Returns
@@ -131,6 +136,14 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
// daemon restart.
seedTools(coreAPI, cfg.Voice.Tools)
exec := tool.NewExecutor(coreAPI, time.Duration(cfg.Voice.ToolTimeout))
// MCP servers (Vikunja #251): discovery PROPOSES tools into the same
// allowlist, so an MCP tool is enabled by hand on /tools like any other and
// runs through the same confirm turn. Off unless the `mcp` block configures
// an enabled server.
w.mcp = wireMCP(cfg, dataStore)
if w.mcp != nil {
exec = exec.WithMCP(w.mcp.caller())
}
matcher := tool.NewMatcher(coreAPI)
// ----- weather provider (Open-Meteo when configured, Stub otherwise) -----
+54
View File
@@ -67,6 +67,14 @@ type fakeCore struct {
// for handleChatAPI tests
chatText string
chatErr error
// for the MCP section of /tools
mcpServers []ipc.MCPServerStatus
mcpErr error
}
func (f *fakeCore) MCPServers(context.Context) ([]ipc.MCPServerStatus, error) {
return f.mcpServers, f.mcpErr
}
func (f *fakeCore) Chat(_ context.Context, text string) (string, error) {
@@ -1118,3 +1126,49 @@ func TestHandleChatAPI_FailOpenByDefault(t *testing.T) {
t.Errorf("core.Chat text = %q, want %q", core.chatText, "привет")
}
}
// The MCP section renders the configured servers, and a proposal that already
// knows its cmd prefills the enable form so the argv is not retyped by hand.
func TestHandleTools_GET_MCPSection(t *testing.T) {
core := &fakeCore{
proposed: []ipc.Tool{{
Name: "vikunja_list_tasks", Scope: "mcp:vikunja",
Cmd: []string{"mcp", "vikunja", "list_tasks"}, Destructive: true,
Utterance: "mcp vikunja/list_tasks: List tasks in a project.",
}},
mcpServers: []ipc.MCPServerStatus{
{Name: "vikunja", Transport: "http", Target: "http://192.168.1.104:9100/mcp", Connected: true, Server: "vikunja 0.1.0", Tools: 4},
{Name: "files", Transport: "stdio", Target: "mcp-server-fs /srv", Err: "start: no such file"},
},
}
rr := httptest.NewRecorder()
handleTools(rr, httptest.NewRequest(http.MethodGet, "/tools", nil), core, nil, false)
if rr.Code != http.StatusOK {
t.Fatalf("status = %d", rr.Code)
}
body := rr.Body.String()
for _, want := range []string{
"MCP servers", "vikunja", "192.168.1.104:9100/mcp", "vikunja 0.1.0",
"files", "no such file",
`value="mcp vikunja list_tasks"`, // the enable form is prefilled
"checked", // and pre-marked destructive (no readOnlyHint)
} {
if !strings.Contains(body, want) {
t.Errorf("missing %q in /tools output", want)
}
}
}
// MCP off (or an older core that does not know the method) renders the section
// empty instead of breaking the page.
func TestHandleTools_GET_MCPUnavailable(t *testing.T) {
core := &fakeCore{mcpErr: ipc.ErrNotImplemented}
rr := httptest.NewRecorder()
handleTools(rr, httptest.NewRequest(http.MethodGet, "/tools", nil), core, nil, false)
if rr.Code != http.StatusOK {
t.Fatalf("status = %d, want 200", rr.Code)
}
if !strings.Contains(rr.Body.String(), "no MCP servers configured") {
t.Error("expected the empty-state copy")
}
}
+25 -4
View File
@@ -692,7 +692,7 @@ const toolsHTML = `{{template "shellTop" "tools"}}
{{if .Msg}}<div class="msg msg-ok">{{.Msg}}</div>{{end}}
<section class=card>
<h2 class=card-title>proposed <span class=badge>{{len .Proposed}}</span></h2>
{{if .Proposed}}<p class=hint>maven drafted these from acts she couldn't run. Fill the command (argv, space-separated) and enable.</p>
{{if .Proposed}}<p class=hint>maven drafted these from acts she couldn't run. Fill the command (argv, space-separated) and enable. A row in an <code>mcp:</code> scope came from an MCP server and already knows what it calls — check the command, then enable.</p>
<div class=scroll><table><tr><th>name</th><th>scope</th><th>from utterance</th><th>enable as</th></tr>
{{range .Proposed}}<tr>
<td><code>{{.Name}}</code></td><td><span class=badge>{{.Scope}}</span></td><td>{{.Utterance}}</td>
@@ -700,8 +700,8 @@ const toolsHTML = `{{template "shellTop" "tools"}}
<input type=hidden name=name value="{{.Name}}">
<input type=hidden name=scope value="{{.Scope}}">
<input type=hidden name=action value=enable>
<input type=text name=cmd class=input-wide placeholder="systemctl restart" required>
<label><input type=checkbox name=destructive> destructive</label>
<input type=text name=cmd class=input-wide placeholder="systemctl restart" value="{{join .Cmd " "}}" required>
<label><input type=checkbox name=destructive {{if .Destructive}}checked{{end}}> destructive</label>
<button class=btn>enable</button></form>
<form method=post action=/tools class=inline-form>
<input type=hidden name=name value="{{.Name}}">
@@ -730,6 +730,19 @@ const toolsHTML = `{{template "shellTop" "tools"}}
<div class=hint>enable proposed tools above, or ask maven to configure one</div>
</div>{{end}}
</section>
<section class=card>
<h2 class=card-title>MCP servers <span class=badge>{{len .MCP}}</span></h2>
{{if .MCP}}<p class=hint>servers she connects OUT to. Their tools appear above as proposals — a configured server is a place she may look, not a capability she has. A <code>stdio</code> target is a process on this box; an <code>http</code> one on a loopback or LAN address is inside the network, so treat its tools accordingly.</p>
<div class=scroll><table><tr><th>name</th><th>transport</th><th>target</th><th>state</th><th>tools</th></tr>
{{range .MCP}}<tr><td><code>{{.Name}}</code></td><td><span class=badge>{{.Transport}}</span></td><td><code>{{.Target}}</code></td>
<td>{{if .Connected}}connected{{if .Server}}{{.Server}}{{end}}{{else}}<span class=red>down</span>{{if .Err}}{{.Err}}{{end}}{{end}}</td>
<td>{{.Tools}}</td></tr>{{end}}</table></div>
{{else}}<div class=empty>
<svg class=icon width="20" height="20"><use href="/ethos-icons.svg#i-settings"/></svg>
<div>no MCP servers configured</div>
<div class=hint>add an <code>mcp.servers</code> block to mavend.json to let her use an external tool server</div>
</div>{{end}}
</section>
{{template "shellBottom"}}`
// routinesHTML — proposed routine review surface. One row per thing maven
@@ -1320,12 +1333,20 @@ func handleTools(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI, sessi
http.Error(w, "core read failed", http.StatusBadGateway)
return
}
// MCP is off by default and an older core may not know the method at all,
// so a failure here renders an empty section rather than breaking the page.
servers, err := core.MCPServers(ctx)
if err != nil {
log.Printf("tools: mcp servers: %v", err)
servers = nil
}
w.Header().Set("Content-Type", "text/html; charset=utf-8")
if err := toolsTmpl.Execute(w, struct {
Msg string
Proposed []ipc.Tool
Enabled []ipc.Tool
}{msg, proposed, enabled}); err != nil {
MCP []ipc.MCPServerStatus
}{msg, proposed, enabled, servers}); err != nil {
log.Printf("tools render: %v", err)
}
}
+14
View File
@@ -31,6 +31,20 @@
"cooldown": "24h"
},
"mcp": {
"timeout": "15s",
"servers": [
{
"name": "vikunja",
"url": "http://192.168.1.104:9100/mcp",
"allow_private": true,
"allow_tools": ["list_projects", "list_tasks", "get_task_details", "create_task"],
"max_tools": 6,
"enabled": false
}
]
},
"nexus": { "url": "http://nexus:9740" },
"praxis": { "url": "http://praxis:8989" },
"hexis": { "url": "http://hexis:9741" },
+101 -22
View File
@@ -1,27 +1,106 @@
# Plan: Vision — Image Understanding Capability
**Goal:** Maven can "see" — accept images (from mavweb upload, Telegram, or filesystem paths), run vision inference via a local or remote multimodal model, and answer questions about the image content or extract structured information.
**Goal:** Maven can "see" — accept images (from mavweb upload, Telegram, or filesystem paths), store them, run inference via a **local** multimodal model, and answer questions about the image content or extract text from it.
**Done when:**
- Vision model backend is configurable: local multimodal LLM (e.g., LLaVA, Qwen-VL via `llama-server` mmproj) or remote API
- `internal/vision/` package handles image preprocessing, model inference, result parsing
- Voice/text commands like "что на картинке?" or "прочитай текст с экрана" route to the vision handler
- Extracted information can be written as facts/notes through `ipc.CoreAPI`
- Telegram image messages are processed through the same pipeline
**Status (2026-08-01):** intake, storage, config seam and the provider are shipped. The
describing half is **BLOCKED on a model download** — see "What is blocked" below.
**Scope:**
- New `internal/vision/` package — image loader (Go stdlib `image` + `golang.org/x/image`), inference client
- New config block: `voice.vision` in `config.Config``{enabled, provider, model_path, mmproj_path, remote_url}`
- Router intent extension: new `IntentVision` or reuse `IntentQuery` with a vision flag
- Reuses `internal/llm.Client` for API-compatible backends (OpenAI-compatible vision API)
- Reuses `internal/ipc.CoreAPI` for writing extracted data
## What shipped
**Steps:**
1. Create `internal/vision/provider.go``Provider` interface with `Describe(image []byte, prompt string) (string, error)` and `ExtractText(image []byte) (string, error)`
2. Implement `LocalProvider` — spawns `llama-server` with mmproj, sends multimodal chat completion requests
3. Implement `RemoteProvider` — calls an OpenAI-compatible vision API endpoint, reuses `internal/llm.Client`
4. Create `internal/vision/processor.go` — image preprocessing (resize, format conversion to JPEG/PNG, base64 encoding)
5. Wire vision into `cmd/mavend/voice.go:reactiveHandler` — detect vision intent from router (new `IntentVision` or a `Slots.HasImage` flag)
6. Add IPC method `MethodDescribeImage` for programmatic access (mavweb upload, telegram bot)
7. Add vision config block to `config.Config` and wire in `cmd/mavend/main.go`
8. Test with a local multimodal model: send an image via mavweb, verify description and text extraction
| Piece | Where |
|---|---|
| Blob store (content-addressed, retention-pruned) | `internal/media/store.go` |
| Image decode / flatten / downscale / JPEG | `internal/media/image.go` |
| `Provider` seam + `Disabled` floor + `LocalProvider` | `internal/vision/vision.go` |
| Store-then-describe orchestration, re-runnable | `internal/vision/intake.go` |
| Config blocks `media` and `vision` | `internal/config/config.go` |
| IPC method `describe_image` (`AuthRead`) | `internal/ipc/{wire,api,client,server}.go`, `internal/auth/policy.go` |
| Daemon wiring + hourly retention prune | `cmd/mavend/vision.go` |
`internal/media` is deliberately shared: hearing (#253) and speaker recognition (#255) have
the same intake problem — a blob arrives, gets stored, gets described — and they store their
audio in the same place under the same retention.
## Design decisions worth knowing
**Store before describe.** `Intake.Accept` writes the blob to disk *first*, then asks the
model. If the model is missing or broken — which is this box's actual state — the answer is
"it's kept, I can't read it yet" with a content-addressed id, and `Intake.Rerun(id, question)`
describes it later. Nothing is lost to a missing model.
**No `RemoteProvider`.** The original step 3 called for "an OpenAI-compatible vision API
endpoint". Refused. The surviving hard constraint in CLAUDE.md after "never phones home" was
deprecated is *no cloud model, inference stays on the box*, and a photo of his flat is the
worst possible exception. `vision.NewLocal` therefore validates the endpoint at construction:
loopback, a private IP, or `localhost`. A hostname is refused too — it could resolve anywhere,
and resolving it would mean trusting DNS with his pictures.
**Blobs are not in the database.** The sqlite store is small, encrypted and read every tick;
a 40 MB blob has no business there. What lands in the database is the *text* the blob produced,
as an ordinary note (`source: media:image:<id-prefix>`), and only when the caller asks for it
(`save_note`). Glancing at a screenshot is not the same act as remembering it.
**Images are never search input and never embedded.** Only the derived description
participates in recall, and only after he can see it as a note.
**Retention is enforced by a loop, not by a promise.** `media.retention` defaults to 7 days
and `cmd/mavend` prunes hourly, starting at boot. A store that grows forever would be the real
failure mode of this capability.
**No webp.** The stdlib has no webp decoder and this repo takes no new dependencies (the box
is offline). `media.SniffImage` recognises webp well enough to refuse it *by name*, so the log
says "webp is not supported" instead of "not an image". Telegram sends webp for stickers; that
is a known gap, not a mystery.
**Text extraction is not a second method.** "прочитай текст с картинки" is a prompt. A VLM has
no separate OCR mode to select, and a second interface method would only duplicate the first.
## What is blocked, and on what
There is **no vision-capable gguf and no mmproj file on this box**. Checked 2026-08-01:
```
/mnt/hdd1/llms/{Bonsai,LFM2.5,llama3.2,ministral,nemotron3-nano,qwen3,qwen3.5}
```
— sixteen ggufs, all text-only, no `*mmproj*` anywhere. The resident Qwen3-1.7B is text-only
by construction, so vision needs a *second* model. The ≤1.7B ceiling in CLAUDE.md is about the
resident router/phraser, not about a second model loaded on demand — but iGPU VRAM still is,
so keep it small.
To unblock, download one pair to `/mnt/hdd1/llms/vision/` (bind-mounted to
`/opt/maven/models/llm`), a gguf **and** its mmproj:
- `Qwen2.5-VL-3B-Instruct` (Q4_K_M + `mmproj-F16.gguf`) — the safe default; reads Russian, and
its OCR is the best of this size class.
- `SmolVLM2-2.2B-Instruct` — smaller and faster, weaker at Cyrillic text in images.
- `moondream2` — smallest, English-only in practice. Do not bother, per the sub-500M lesson.
Then run a second llama-server on 8081 with `--mmproj`, point `vision.endpoint` at it, and
walk the QA steps on Vikunja #252.
## Config
```json
"media": { "dir": "media", "retention": "168h", "max_bytes": 67108864 },
"vision": {
"enabled": true,
"endpoint": "http://127.0.0.1:8081",
"model": "qwen2.5-vl-3b",
"max_dim": 896,
"max_tokens": 300,
"timeout": "90s"
}
```
Both absent by default. No `media` block ⇒ `describe_image` does not exist at all; a `media`
block with no `vision` block ⇒ images are stored and honestly not described.
## Still open
- **Router intent.** "что на картинке?" does not route anywhere yet. Adding an intent is
premature while nothing can answer it; the IPC method is the surface a Telegram photo or a
mavweb upload calls today.
- **Telegram photo path** in `mavpoll` (download the file, call `DescribeImage`).
- **mavweb upload page** and a `/media` listing so stored blobs are visible and deletable from
the authed surface.
+8
View File
@@ -87,6 +87,14 @@ func Requirement(m ipc.Method) Authority {
// reminder, or touch the tool allowlist, so a compromised mail reader can
// at worst put junk on a review page he clears in one click.
ipc.MethodIngestMail,
// Looking at one image (Vikunja #252). AuthRead because of what it can
// produce: words about a picture, and optionally a note. It cannot write
// a fact, set a reminder, or touch the tool allowlist. The invasive part
// of this capability is not the authority rung — it is that the bytes are
// kept on disk, which media.retention bounds, and that they never leave
// the box, which internal/vision enforces by refusing a non-private
// endpoint.
ipc.MethodDescribeImage,
// The read side of the model swap: which model is resident, which ones are
// allowlisted. It loads nothing and changes nothing.
ipc.MethodModelStatus:
+210
View File
@@ -18,10 +18,12 @@ import (
"fmt"
"os"
"path/filepath"
"strings"
"time"
"github.com/kami/maven/internal/delivery/ntfysink"
"github.com/kami/maven/internal/delivery/telegramsink"
"github.com/kami/maven/internal/mcp"
"github.com/kami/maven/internal/morning"
"github.com/kami/maven/internal/update"
"github.com/robfig/cron/v3"
@@ -191,6 +193,130 @@ type Config struct {
// discovers and executes capabilities through Hexis for ecosystem actions.
// nil ⇒ no capability-aware routing.
Hexis *HexisConfig `json:"hexis,omitempty"`
// Vision — image understanding (Vikunja #252). nil / absent ⇒ she cannot
// look at pictures at all: the intake refuses, and no vision server is
// contacted. See VisionConfig.
Vision *VisionConfig `json:"vision,omitempty"`
// Media — where images and captured audio are kept on disk, and for how
// long. nil / absent ⇒ no blob store is wired, which is what disables both
// vision intake and meeting capture regardless of their own blocks: nothing
// in this repo holds a recording only in memory. See MediaConfig.
Media *MediaConfig `json:"media,omitempty"`
// MCP — Model Context Protocol servers Maven connects OUT to (Vikunja
// #251). nil / absent / no enabled server ⇒ no connection is made and no
// tool is discovered, like every other capability that reaches outside the
// box. She is a client here, never a server: nothing exposes her own
// capabilities to an outside caller. See MCPConfig.
MCP *MCPConfig `json:"mcp,omitempty"`
}
// MCPConfig — the MCP client block. Servers are dark until one has
// `"enabled": true`, and a discovered tool is only ever PROPOSED: Kami enables
// it on /tools, on the authed surface, exactly as he would a shell tool. The
// voice path can never grant a capability to itself.
type MCPConfig struct {
// Servers — the configured servers. Each needs exactly one of command
// (a subprocess on this box) or url (a streamable-HTTP endpoint).
Servers []MCPServerConfig `json:"servers,omitempty"`
// Timeout — per-call budget for every server that does not set its own.
// 0 ⇒ mcp.DefaultTimeout (15s). A tool slower than this is not usable in a
// spoken turn.
Timeout Duration `json:"timeout,omitempty"`
// AllowHosts / DenyHosts — the host lists for the shared webfetch door that
// url servers go through. Deny wins. Private addresses are refused
// unconditionally unless the individual server sets allow_private.
AllowHosts []string `json:"allow_hosts,omitempty"`
DenyHosts []string `json:"deny_hosts,omitempty"`
// MaxBytes — cap on one JSON-RPC response. 0 ⇒ webfetch.DefaultMaxBytes.
MaxBytes int64 `json:"max_bytes,omitempty"`
}
// MCPServerConfig — one MCP server.
type MCPServerConfig struct {
// Name — the local handle. It prefixes every tool this server contributes
// ("vikunja" + "list_tasks" ⇒ the allowlist row "vikunja_list_tasks") and
// becomes the store scope "mcp:<name>", so its provenance is readable on
// /tools without opening the config.
Name string `json:"name"`
// Command / Args / Env / Dir — a stdio server: a child process of mavend,
// on this box, under this user. argv, never a shell string.
Command string `json:"command,omitempty"`
Args []string `json:"args,omitempty"`
Env []string `json:"env,omitempty"`
Dir string `json:"dir,omitempty"`
// URL — a streamable-HTTP endpoint. It is fetched through
// internal/webfetch, so the SSRF guard, the redirect cap, the size cap and
// the one-request-per-host-per-second limit all apply.
URL string `json:"url,omitempty"`
// AllowPrivate — let THIS server be a loopback or LAN address. The Vikunja
// server on homesrv is "http://localhost:9100/mcp", which is refused
// without this flag. Understand what it means before setting it: a local
// server is a DIFFERENT trust level from a public one. It is inside the
// network, it usually needs no credential, and it can change things that
// matter — so an argument the router got wrong lands somewhere real. Set it
// only for a server you run yourself, and prefer allow_tools with it.
AllowPrivate bool `json:"allow_private,omitempty"`
// AllowTools — when set, the ONLY remote tool names taken from this server.
// This is the knob that keeps the catalogue deliberate: the resident model
// is a 1.7B with a 4096-token context, and a tool name it half-remembers is
// a wrong act, so fewer and better-chosen beats complete.
AllowTools []string `json:"allow_tools,omitempty"`
// MaxTools — cap on this server's contribution. 0 ⇒ mcp.DefaultMaxTools (12).
MaxTools int `json:"max_tools,omitempty"`
// Timeout — per-call budget for this server. 0 ⇒ MCPConfig.Timeout.
Timeout Duration `json:"timeout,omitempty"`
// Enabled — false (the default) keeps a configured server described but
// dark, so a block can be written and reviewed before it is switched on.
Enabled bool `json:"enabled,omitempty"`
}
// MCPServers maps the config blocks onto the mcp package's own type. It lives
// here so config validation and daemon wiring cannot drift on the mapping.
// Returns nil when nothing is configured or nothing is enabled.
func (c *Config) MCPServers() []mcp.ServerConfig {
if c.MCP == nil {
return nil
}
out := make([]mcp.ServerConfig, 0, len(c.MCP.Servers))
for _, s := range c.MCP.Servers {
if !s.Enabled {
continue
}
timeout := time.Duration(s.Timeout)
if timeout <= 0 {
timeout = time.Duration(c.MCP.Timeout)
}
out = append(out, mcp.ServerConfig{
Name: s.Name,
Command: s.Command,
Args: s.Args,
Env: s.Env,
Dir: s.Dir,
URL: s.URL,
AllowPrivate: s.AllowPrivate,
AllowTools: s.AllowTools,
MaxTools: s.MaxTools,
Timeout: timeout,
Enabled: true,
})
}
if len(out) == 0 {
return nil
}
return out
}
// PraxisConfig — maven's connection to the Praxis attention service.
@@ -360,6 +486,78 @@ type VoiceConfig struct {
ToolTimeout Duration `json:"tool_timeout,omitempty"`
}
// MediaConfig — the on-disk blob store for images and captured audio
// (internal/media). It is shared by all three senses: vision intake, meeting
// capture, and speaker enrolment samples all write here.
//
// Absent ⇒ off, and off means Maven cannot accept an image or start a recording
// at all. That default is deliberate: a capability that keeps photos and audio of
// people on disk should require someone to have typed a path.
type MediaConfig struct {
// Dir — the blob store root, created 0700. Relative paths resolve against
// StateDir. Required; an empty dir means the store is not wired.
Dir string `json:"dir,omitempty"`
// Retention — how long a blob is kept before the tick prunes it. 0 ⇒
// media.DefaultRetention (7 days). This is the knob that stops recordings
// of people accumulating; raising it past a few weeks should need a reason.
Retention Duration `json:"retention,omitempty"`
// MaxBytes — per-blob cap. 0 ⇒ media.DefaultMaxBytes (64 MiB).
MaxBytes int64 `json:"max_bytes,omitempty"`
}
// StoreDir reports the configured blob directory, or "" when media is not
// wired. Safe on a nil receiver.
func (m *MediaConfig) StoreDir() string {
if m == nil {
return ""
}
return strings.TrimSpace(m.Dir)
}
// VisionConfig — the vision provider (internal/vision, docs/plans/07-vision.md).
//
// Absent, or enabled=false, ⇒ the daemon wires vision.Disabled and every attempt
// to look at an image answers that vision is not set up. There is no cloud
// option in this block on purpose: Endpoint must be a loopback or private
// address and internal/vision refuses anything else at startup, because
// inference stays on the box and a photo of his flat is the last thing to make
// an exception for.
type VisionConfig struct {
// Enabled — may she look at images. Default false.
Enabled bool `json:"enabled,omitempty"`
// Endpoint — base URL of a llama-server running a vision model with its
// mmproj, e.g. "http://127.0.0.1:8081". Loopback / private only.
Endpoint string `json:"endpoint,omitempty"`
// Model — model name sent in the request. llama-server ignores it.
Model string `json:"model,omitempty"`
// MaxDim — longest edge the image is scaled to before inference. 0 ⇒
// media.DefaultMaxDim (896).
MaxDim int `json:"max_dim,omitempty"`
// MaxTokens — cap on the description. 0 ⇒ vision.DefaultMaxTokens (300).
MaxTokens int `json:"max_tokens,omitempty"`
// Timeout — per-description budget. 0 ⇒ vision.DefaultTimeout (90s). A small
// VLM on an iGPU is slow; a tight timeout here just means no answer ever.
Timeout Duration `json:"timeout,omitempty"`
// Prompt — the default question when he only sent a picture. Empty ⇒
// vision.DefaultPrompt (Russian, "опиши что на изображении").
Prompt string `json:"prompt,omitempty"`
}
// LooksAtImages reports whether vision is configured well enough to try. Safe on
// a nil receiver, and false without an endpoint — enabled with nothing to talk
// to is a misconfiguration, not a capability.
func (v *VisionConfig) LooksAtImages() bool {
return v != nil && v.Enabled && strings.TrimSpace(v.Endpoint) != ""
}
// WeatherConfig configures the weather provider for voice queries.
type WeatherConfig struct {
Provider string `json:"provider,omitempty"` // "open-meteo" or "" → stub
@@ -783,6 +981,12 @@ func (c *Config) applyDefaults() {
c.Feeds = nil
}
// Same rule for MCP: a block with no server, or none enabled, is the same
// as no block at all. Normalising it to nil keeps "off" in one place.
if c.MCP != nil && len(c.MCPServers()) == 0 {
c.MCP = nil
}
// Same rule for the crawler: a block that neither answers on demand nor
// watches anything has nothing to do, so it is normalised to "off".
if c.Crawl != nil && !c.Crawl.OnDemand && len(c.Crawl.Watches) == 0 {
@@ -887,6 +1091,12 @@ func (c *Config) validate() error {
return fmt.Errorf("routine %q: bad cron %q: %w", r.Name, r.Cron, err)
}
}
// An MCP block with a typo (no name, both command and url, a bare hostname
// as the url) fails here, at startup, rather than at the first turn that
// needed the tool.
if err := mcp.Validate(c.MCPServers()); err != nil {
return err
}
if len(c.MorningRoutines) > 0 {
if err := morning.Validate(morningRoutinesFromConfig(c.MorningRoutines)); err != nil {
return err
+82
View File
@@ -0,0 +1,82 @@
package config
import (
"testing"
"time"
)
func TestMCPAbsentIsOff(t *testing.T) {
c, err := Load(writeConfig(t, `{}`))
if err != nil {
t.Fatal(err)
}
if c.MCP != nil {
t.Error("no mcp block ⇒ nil")
}
if got := c.MCPServers(); got != nil {
t.Errorf("MCPServers() = %+v, want nil", got)
}
}
// A described-but-not-enabled server must not be wired. This is how a block can
// sit in the config file, reviewed, before it is switched on.
func TestMCPDisabledServerIsOff(t *testing.T) {
c, err := Load(writeConfig(t, `{"mcp":{"servers":[
{"name":"vikunja","url":"http://localhost:9100/mcp","allow_private":true}]}}`))
if err != nil {
t.Fatal(err)
}
if c.MCP != nil {
t.Errorf("a block with nothing enabled must normalise to nil, got %+v", c.MCP)
}
if got := c.MCPServers(); len(got) != 0 {
t.Errorf("MCPServers() = %+v", got)
}
}
func TestMCPEnabledServerMapping(t *testing.T) {
c, err := Load(writeConfig(t, `{"mcp":{
"timeout":"5s",
"servers":[
{"name":"vikunja","url":"http://localhost:9100/mcp","allow_private":true,
"allow_tools":["list_tasks"],"max_tools":3,"enabled":true},
{"name":"files","command":"mcp-server-fs","args":["/srv"],"timeout":"1s","enabled":true},
{"name":"off","command":"nope"}
]}}`))
if err != nil {
t.Fatal(err)
}
got := c.MCPServers()
if len(got) != 2 {
t.Fatalf("servers = %+v", got)
}
if got[0].Name != "vikunja" || !got[0].AllowPrivate || got[0].MaxTools != 3 ||
len(got[0].AllowTools) != 1 || got[0].Timeout != 5*time.Second {
t.Errorf("vikunja mapped wrong: %+v", got[0])
}
if got[1].Command != "mcp-server-fs" || len(got[1].Args) != 1 || got[1].Timeout != time.Second {
t.Errorf("files mapped wrong: %+v", got[1])
}
// allow_private is per server and must not leak to the other one.
if got[1].AllowPrivate {
t.Error("allow_private leaked between servers")
}
}
func TestMCPBadServerFailsAtStartup(t *testing.T) {
cases := map[string]string{
"no name": `{"mcp":{"servers":[{"command":"x","enabled":true}]}}`,
"both": `{"mcp":{"servers":[{"name":"a","command":"x","url":"http://a.test","enabled":true}]}}`,
"neither": `{"mcp":{"servers":[{"name":"a","enabled":true}]}}`,
"bad scheme": `{"mcp":{"servers":[{"name":"a","url":"unix:///run/x.sock","enabled":true}]}}`,
"duplicate": `{"mcp":{"servers":[{"name":"a","command":"x","enabled":true},{"name":"a","command":"y","enabled":true}]}}`,
"spacey name": `{"mcp":{"servers":[{"name":"a b","command":"x","enabled":true}]}}`,
}
for name, body := range cases {
t.Run(name, func(t *testing.T) {
if _, err := Load(writeConfig(t, body)); err == nil {
t.Fatal("want a startup error")
}
})
}
}
+96
View File
@@ -0,0 +1,96 @@
package config
import (
"encoding/json"
"testing"
"time"
)
// Absent blocks must read as off on a nil receiver: the daemon calls these
// helpers before it knows whether the operator configured anything.
func TestSensesOffByDefault(t *testing.T) {
var cfg Config
if cfg.Media.StoreDir() != "" {
t.Error("media store dir is set with no media block")
}
if cfg.Vision.LooksAtImages() {
t.Error("vision is on with no vision block")
}
}
// enabled with nothing to talk to is a misconfiguration, not a capability.
func TestVisionNeedsBothEnabledAndEndpoint(t *testing.T) {
cases := []struct {
name string
v *VisionConfig
want bool
}{
{"absent", nil, false},
{"endpoint but not enabled", &VisionConfig{Endpoint: "http://127.0.0.1:8081"}, false},
{"enabled but no endpoint", &VisionConfig{Enabled: true}, false},
{"enabled, blank endpoint", &VisionConfig{Enabled: true, Endpoint: " "}, false},
{"both", &VisionConfig{Enabled: true, Endpoint: "http://127.0.0.1:8081"}, true},
}
for _, c := range cases {
if got := c.v.LooksAtImages(); got != c.want {
t.Errorf("%s: LooksAtImages() = %v, want %v", c.name, got, c.want)
}
}
}
func TestSensesBlocksParseFromJSON(t *testing.T) {
raw := `{
"db_path": "/tmp/x.db",
"socket_path": "/tmp/x.sock",
"media": {"dir": "media", "retention": "48h", "max_bytes": 1048576},
"vision": {
"enabled": true,
"endpoint": "http://127.0.0.1:8081",
"model": "qwen2.5-vl",
"max_dim": 640,
"max_tokens": 200,
"timeout": "45s",
"prompt": "Что тут?"
}
}`
var cfg Config
if err := json.Unmarshal([]byte(raw), &cfg); err != nil {
t.Fatalf("unmarshal: %v", err)
}
if cfg.Media.StoreDir() != "media" {
t.Errorf("media dir = %q", cfg.Media.StoreDir())
}
if time.Duration(cfg.Media.Retention) != 48*time.Hour {
t.Errorf("retention = %v", time.Duration(cfg.Media.Retention))
}
if cfg.Media.MaxBytes != 1<<20 {
t.Errorf("max_bytes = %d", cfg.Media.MaxBytes)
}
if !cfg.Vision.LooksAtImages() {
t.Fatal("vision did not parse as enabled")
}
if cfg.Vision.MaxDim != 640 || cfg.Vision.MaxTokens != 200 {
t.Errorf("vision limits = %+v", cfg.Vision)
}
if time.Duration(cfg.Vision.Timeout) != 45*time.Second {
t.Errorf("vision timeout = %v", time.Duration(cfg.Vision.Timeout))
}
if cfg.Vision.Prompt != "Что тут?" {
t.Errorf("prompt = %q", cfg.Vision.Prompt)
}
}
// A media dir set with no vision block is a valid state, and the useful one on a
// box with no vision model: images can be kept, they just cannot be described.
func TestMediaWithoutVisionIsValid(t *testing.T) {
var cfg Config
if err := json.Unmarshal([]byte(`{"media":{"dir":"/srv/media"}}`), &cfg); err != nil {
t.Fatal(err)
}
if cfg.Media.StoreDir() != "/srv/media" {
t.Errorf("dir = %q", cfg.Media.StoreDir())
}
if cfg.Vision.LooksAtImages() {
t.Error("vision came on by itself")
}
}
+66
View File
@@ -176,6 +176,51 @@ type IngestMailResp struct {
Skipped bool `json:"skipped,omitempty"`
}
// DescribeImageReq — one image handed to core to look at (Vikunja #252).
//
// Data is the raw image file as received (png / jpeg / gif). Core sniffs it and
// refuses anything else; a declared content type is not part of this request
// because the sender's claim about its own bytes is not evidence. Base64 on the
// wire via the usual JSON marshal of []byte.
//
// Question is what he asked about the picture ("что тут написано?"). Empty ⇒
// core uses its configured default prompt.
//
// Source is provenance recorded on the stored blob: "telegram", "web:upload".
//
// Exactly one of Data or ID is set. ID re-describes an image core already has —
// a different question, or the first attempt that succeeds after a vision model
// finally lands on disk.
//
// The method exists only when core has both a media store and an enabled vision
// block; otherwise it answers ErrUnknownMethod, which is what "off unless
// configured" looks like at the wire. A surface cannot make Maven look at
// pictures by merely sending one.
type DescribeImageReq struct {
Data []byte `json:"data,omitempty"`
ID string `json:"id,omitempty"`
Source string `json:"source,omitempty"`
Question string `json:"question,omitempty"`
// SaveNote — also write the description as a note (source
// "media:image:<id-prefix>") so it is recallable later. Default false: a
// glance at a screenshot is not automatically a memory.
SaveNote bool `json:"save_note,omitempty"`
}
// DescribeImageResp — what she saw. ID is the stored blob's content address, and
// it is set even when Description is empty because the description failed: the
// bytes are on disk and the same id can be retried. NoteID is non-zero only when
// SaveNote was set and the write succeeded.
//
// The image itself is never echoed back.
type DescribeImageResp struct {
ID string `json:"id"`
Description string `json:"description,omitempty"`
Width int `json:"width,omitempty"`
Height int `json:"height,omitempty"`
NoteID int64 `json:"note_id,omitempty"`
}
// SwapModelReq — load another resident model without restarting the daemon
// (Vikunja #250). ModelPath must be one of the paths in phraser.swap_models;
// anything else is ErrForbidden, and an unconfigured allowlist makes the whole
@@ -319,6 +364,19 @@ type Tool struct {
Updated time.Time `json:"updated"`
}
// MCPServerStatus — one configured MCP server, as the web surface sees it.
// Target is the command or url; Tools is how many tools discovery kept after
// allow_tools / max_tools, not how many the server offers.
type MCPServerStatus struct {
Name string `json:"name"`
Transport string `json:"transport"` // "stdio" (a local subprocess) or "http"
Target string `json:"target"`
Connected bool `json:"connected"`
Server string `json:"server,omitempty"` // the server's own name + version
Tools int `json:"tools"`
Err string `json:"err,omitempty"`
}
// chatReq / chatResp — text chat round-trip for the IPC Chat method.
type chatReq struct {
Text string `json:"text"`
@@ -457,6 +515,14 @@ type CoreAPI interface {
// TickTrace.
MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error)
// MCPServers reports the configured MCP servers and their health
// (Vikunja #251). Read-only introspection for /tools — there is no
// "call this tool" method on purpose: an MCP tool runs through the same
// allowlist, confirm turn and act path as any other tool, and a second
// mutation path would be a second thing to get wrong. Empty when the
// mcp config block is absent, which is the default.
MCPServers(ctx context.Context) ([]MCPServerStatus, error)
// DayPlan returns today's ordered plan — calendar events, pending
// reminders and any morning checklist still outstanding (see
// internal/morning.BuildPlan) — plus the spoken RU rendering of it.
+22
View File
@@ -72,6 +72,7 @@ var readOnlyMethods = map[Method]bool{
MethodListTasks: true,
MethodTickTrace: true,
MethodMorningStatus: true,
MethodMCPServers: true,
MethodDayPlan: true,
}
@@ -459,6 +460,19 @@ func (c *Client) IngestMail(ctx context.Context, req IngestMailReq) (IngestMailR
return r, nil
}
// DescribeImage hands one image to core to look at (Vikunja #252).
// ErrUnknownMethod means core has no media store or vision is off — the caller
// should stop asking, not retry. A response with an ID and an empty Description
// means the bytes were stored but nothing could describe them yet, which is the
// expected state on a box with no vision model on disk.
func (c *Client) DescribeImage(ctx context.Context, req DescribeImageReq) (DescribeImageResp, error) {
var r DescribeImageResp
if err := c.call(ctx, MethodDescribeImage, req, &r); err != nil {
return DescribeImageResp{}, err
}
return r, nil
}
// SwapModel asks core to load another resident model (Vikunja #250).
// ErrUnknownMethod means core has no phraser.swap_models allowlist configured;
// ErrForbidden means the path is not on it, or step-up was not asserted. A
@@ -505,6 +519,14 @@ func (c *Client) TickTrace(ctx context.Context) (TickTrace, error) {
return t, nil
}
func (c *Client) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
var s []MCPServerStatus
if err := c.call(ctx, MethodMCPServers, nil, &s); err != nil {
return nil, err
}
return s, nil
}
func (c *Client) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error) {
var s []MorningRoutineStatus
if err := c.call(ctx, MethodMorningStatus, nil, &s); err != nil {
+45 -4
View File
@@ -211,6 +211,10 @@ func (a *storeAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, e
return nil, errors.New("store: morning status not available via direct store API")
}
func (a *storeAPI) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
return nil, nil // no manager behind a bare store: nothing configured
}
func (a *storeAPI) DayPlan(ctx context.Context) (DayPlan, error) {
return DayPlan{}, errors.New("store: day plan not available via direct store API")
}
@@ -445,6 +449,16 @@ type Server struct {
SwapModelFn SwapModelFunc
ModelStatusFn ModelStatusFunc
// DescribeImageFn — looks at one image (Vikunja #252). Set by the daemon only
// when a media store is configured AND vision is enabled with a local
// endpoint; nil ⇒ MethodDescribeImage answers ErrUnknownMethod, so a surface
// cannot make Maven accept a photo by merely sending one.
//
// It bypasses CoreAPI for the same reason IngestMailFn does: it needs a blob
// store and a vision server, neither of which is a store operation, and no
// other CoreAPI implementation should have to carry it.
DescribeImageFn DescribeImageFunc
// UnlockFn — unwraps the store encryption key from the wrapped blob using
// the passkey credential public key, opens the encrypted store, and wires
// the rest of the daemon (voice, loop, delivery). Set by the daemon when
@@ -472,6 +486,9 @@ type ModelStatusFunc func(ctx context.Context) (ModelStatusResp, error)
// IngestMailFunc — core-side mail extraction. Returns what was captured.
type IngestMailFunc func(ctx context.Context, req IngestMailReq) (IngestMailResp, error)
// DescribeImageFunc — core-side image intake + description.
type DescribeImageFunc func(ctx context.Context, req DescribeImageReq) (DescribeImageResp, error)
// CheckFunc — the auth hook signature. Wired by the daemon (auth.Gate.Check
// satisfies this); dispatch calls it once per request after param-unmarshal
// independence (it gets the raw params, may unmarshal what it needs — ipc
@@ -640,10 +657,10 @@ func withoutParams[R any](fn func(ctx context.Context, api CoreAPI) (R, error))
// is still honored on the very next request with no extra plumbing here.
//
// MethodAssertStepUp, MethodStoreEncryptionKey, MethodUnlock,
// MethodIngestMail, MethodSwapModel and MethodModelStatus are NOT in this
// table: they bypass CoreAPI entirely
// (s.StepUp / s.WrapKeyFn / s.UnlockFn / s.IngestMailFn), so dispatch
// special-cases them before consulting the table.
// MethodIngestMail, MethodSwapModel, MethodModelStatus and
// MethodDescribeImage are NOT in this table: they bypass CoreAPI entirely
// (s.StepUp / s.WrapKeyFn / s.UnlockFn / s.IngestMailFn / s.DescribeImageFn),
// so dispatch special-cases them before consulting the table.
var methodTable = map[Method]handlerFunc{
MethodWriteFact: withParams(func(ctx context.Context, api CoreAPI, p WriteFactReq) (idResp, error) {
id, err := api.WriteFact(ctx, p)
@@ -832,6 +849,16 @@ var methodTable = map[Method]handlerFunc{
MethodMorningStatus: withoutParams(func(ctx context.Context, api CoreAPI) ([]MorningRoutineStatus, error) {
return api.MorningStatus(ctx)
}),
MethodMCPServers: withoutParams(func(ctx context.Context, api CoreAPI) ([]MCPServerStatus, error) {
out, err := api.MCPServers(ctx)
if err != nil {
return nil, err
}
if out == nil {
out = []MCPServerStatus{}
}
return out, nil
}),
}
// dispatch unmarshals params for req.Method and calls the matching CoreAPI
@@ -909,6 +936,20 @@ func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, er
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodDescribeImage:
if s.DescribeImageFn != nil {
var p DescribeImageReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.DescribeImageFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodModelStatus:
if s.ModelStatusFn != nil {
resp, err := s.ModelStatusFn(ctx)
+3
View File
@@ -122,6 +122,9 @@ func (UnimplementedCoreAPI) TickTrace(ctx context.Context) (TickTrace, error) {
func (UnimplementedCoreAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error) {
return nil, ErrNotImplemented
}
func (UnimplementedCoreAPI) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
return nil, ErrNotImplemented
}
func (UnimplementedCoreAPI) DayPlan(ctx context.Context) (DayPlan, error) {
return DayPlan{}, ErrNotImplemented
}
+2
View File
@@ -45,6 +45,7 @@ const (
MethodRevertFact Method = "revert_fact"
MethodTickTrace Method = "tick_trace"
MethodMorningStatus Method = "morning_status"
MethodMCPServers Method = "mcp_servers"
MethodDayPlan Method = "day_plan"
MethodChat Method = "chat"
MethodCaptureTask Method = "capture_task"
@@ -53,6 +54,7 @@ const (
MethodIngestMail Method = "ingest_mail"
MethodSwapModel Method = "swap_model"
MethodModelStatus Method = "model_status"
MethodDescribeImage Method = "describe_image"
)
// Request — one frame from module to core. Params is the JSON-encoded argument
+61
View File
@@ -0,0 +1,61 @@
package mcp
import (
"regexp"
"strings"
)
// CmdPrefix is the reserved first argv element that marks an allowlist row as
// an MCP call rather than a process. An MCP tool row looks like
//
// name: "vikunja_list_tasks" cmd: ["mcp", "vikunja", "list_tasks"]
//
// which is why there is no new column and no migration: the store, the /tools
// page, ProposeTool, EnableTool, DisableTool, the act matcher and the confirm
// turn all keep working unchanged. The executor is the only place that has to
// know the difference, and it is one branch on Cmd[0].
//
// The rest of the allowlist discipline is inherited whole: a row that is not
// status='enabled' does not run, and a row marked destructive does not run on
// first hearing. Nothing here can enable itself — discovery only proposes.
const CmdPrefix = "mcp"
// Cmd builds the argv encoding for a discovered tool.
func Cmd(server, tool string) []string { return []string{CmdPrefix, server, tool} }
// ParseCmd recognises an MCP allowlist row. ok=false for an ordinary process
// tool, which is what almost every row is.
func ParseCmd(cmd []string) (server, tool string, ok bool) {
if len(cmd) != 3 || cmd[0] != CmdPrefix {
return "", "", false
}
if cmd[1] == "" || cmd[2] == "" {
return "", "", false
}
return cmd[1], cmd[2], true
}
var notName = regexp.MustCompile(`[^a-z0-9_]+`)
// LocalName is the allowlist name for a discovered tool: the server handle, an
// underscore, the remote name, lowercased and stripped of anything that is not
// a word character. Namespacing by server is what keeps two servers that both
// offer "search" from colliding, and what makes the provenance of a row on the
// /tools page obvious without opening the diff.
func LocalName(server, tool string) string {
clean := func(s string) string {
return strings.Trim(notName.ReplaceAllString(strings.ToLower(strings.TrimSpace(s)), "_"), "_")
}
s, t := clean(server), clean(tool)
switch {
case s == "":
return t
case t == "":
return s
}
return s + "_" + t
}
// Scope is the store scope for a server's rows, so the /tools page can group
// them and a human can tell at a glance where a capability came from.
func Scope(server string) string { return "mcp:" + server }
+267
View File
@@ -0,0 +1,267 @@
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
}
+154
View File
@@ -0,0 +1,154 @@
package mcp
import (
"bufio"
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"sync"
)
// Poster is the HTTP seam: internal/webfetch.Fetcher satisfies it. The
// transport takes it as an interface so a test can serve a fake without a
// listener, and so that the ONLY implementation wired in production is the
// guarded fetcher — an MCP endpoint cannot get a bare http.Client this way.
type Poster interface {
Post(ctx context.Context, rawURL, contentType string, body []byte, hdr map[string]string) (*PostResponse, error)
}
// PostResponse is the shape webfetch returns, restated here so this package
// does not depend on it structurally.
type PostResponse struct {
Status int
ContentType string
Body []byte
Header map[string]string
}
// httpTransport speaks streamable HTTP: every request is a POST to one
// endpoint, and the reply is either a JSON object or a text/event-stream frame
// carrying one. Both are accepted — servers pick per response, and the two the
// LAN runs disagree about which.
type httpTransport struct {
poster Poster
url string
mu sync.Mutex
session string // Mcp-Session-Id, echoed back when the server issues one
}
func newHTTPTransport(post Poster, endpoint string) *httpTransport {
return &httpTransport{poster: post, url: endpoint}
}
func (t *httpTransport) Call(ctx context.Context, req *rpcRequest) (*rpcResponse, error) {
body, err := t.send(ctx, req)
if err != nil {
return nil, err
}
frame, err := decodeFrame(body)
if err != nil {
return nil, err
}
var resp rpcResponse
if err := json.Unmarshal(frame, &resp); err != nil {
return nil, fmt.Errorf("mcp: decode response: %w", err)
}
return &resp, nil
}
func (t *httpTransport) Notify(ctx context.Context, method string, params any) error {
_, err := t.send(ctx, &rpcRequest{JSONRPC: "2.0", Method: method, Params: params})
return err
}
func (t *httpTransport) send(ctx context.Context, req *rpcRequest) ([]byte, error) {
req.JSONRPC = "2.0"
raw, err := json.Marshal(req)
if err != nil {
return nil, err
}
hdr := map[string]string{"Accept": "application/json, text/event-stream"}
t.mu.Lock()
if t.session != "" {
hdr["Mcp-Session-Id"] = t.session
}
t.mu.Unlock()
resp, err := t.poster.Post(ctx, t.url, "application/json", raw, hdr)
if err != nil {
return nil, err
}
if sid := headerGet(resp.Header, "Mcp-Session-Id"); sid != "" {
t.mu.Lock()
t.session = sid
t.mu.Unlock()
}
return resp.Body, nil
}
func (t *httpTransport) Close() error {
t.mu.Lock()
t.session = ""
t.mu.Unlock()
return nil
}
func headerGet(h map[string]string, key string) string {
if h == nil {
return ""
}
if v, ok := h[key]; ok {
return v
}
lower := strings.ToLower(key)
for k, v := range h {
if strings.ToLower(k) == lower {
return v
}
}
return ""
}
// decodeFrame pulls the JSON object out of a body that is either raw JSON or
// SSE. For SSE we take the LAST data: payload that parses, which is the
// response — earlier frames on the stream are progress notifications.
func decodeFrame(body []byte) ([]byte, error) {
trimmed := bytes.TrimSpace(body)
if len(trimmed) == 0 {
return nil, errors.New("mcp: empty response body")
}
if trimmed[0] == '{' || trimmed[0] == '[' {
return trimmed, nil
}
var last []byte
sc := bufio.NewScanner(bytes.NewReader(trimmed))
sc.Buffer(make([]byte, 0, 64<<10), maxLine)
for sc.Scan() {
line := strings.TrimSpace(sc.Text())
if !strings.HasPrefix(line, "data:") {
continue
}
payload := strings.TrimSpace(strings.TrimPrefix(line, "data:"))
if payload == "" {
continue
}
var probe map[string]json.RawMessage
if json.Unmarshal([]byte(payload), &probe) != nil {
continue
}
if _, isResp := probe["id"]; isResp {
last = []byte(payload)
}
}
if err := sc.Err(); err != nil {
return nil, fmt.Errorf("mcp: read event stream: %w", err)
}
if last == nil {
return nil, errors.New("mcp: no JSON-RPC response in event stream")
}
return last, nil
}
+72
View File
@@ -0,0 +1,72 @@
// Package mcp is Maven's Model Context Protocol CLIENT. She is a host: she
// connects OUT to MCP servers, discovers the tools and resources they offer,
// and hands them to the parts of her that already exist for this — the tool
// allowlist in the store, the confirm turn for anything that mutates, the
// stage-3 gate that makes an uncertain act ask instead of run.
//
// She is not an MCP server. Nothing here exposes her own capabilities to an
// outside caller; docs/plans/06-mcp-support.md asks for the host direction only.
//
// Boundaries, in code rather than in prose:
//
// - OFF unless configured. No mcp_servers block ⇒ no manager, no goroutine,
// no socket.
// - A remote server is reached through internal/webfetch, so the SSRF guard,
// the size cap, the redirect cap and the per-host rate limit all apply to
// an MCP endpoint exactly as they do to a news feed. Reaching a loopback
// or LAN server means explicitly setting allow_private on THAT server —
// a different trust level, spelled out per server rather than globally.
// - Only the tool name and the arguments the router produced are sent. This
// package never sees his notes, facts, history or the persona block, and
// has no API through which a caller could pass them.
// - Discovery proposes, it does not enable. A discovered tool lands as a
// 'proposed' row; a human enables it on the authed surface.
package mcp
import (
"context"
"encoding/json"
"fmt"
)
// ProtocolVersion — the spec revision we ask for in the initialize handshake.
// A server that answers with a different one is accepted (the spec says the
// client may proceed if it can support what came back); we only refuse when it
// answers with no version at all, which means it is not an MCP server.
const ProtocolVersion = "2025-06-18"
// rpcRequest / rpcResponse — JSON-RPC 2.0. Deliberately hand-rolled: the wire
// format is four fields, and the repo vendors its dependencies, so pulling a
// library in for this would cost more than it saves.
type rpcRequest struct {
JSONRPC string `json:"jsonrpc"`
ID int64 `json:"id,omitempty"`
Method string `json:"method"`
Params any `json:"params,omitempty"`
}
type rpcResponse struct {
JSONRPC string `json:"jsonrpc"`
ID *int64 `json:"id"`
Result json.RawMessage `json:"result,omitempty"`
Error *rpcError `json:"error,omitempty"`
}
type rpcError struct {
Code int `json:"code"`
Message string `json:"message"`
}
func (e *rpcError) Error() string { return fmt.Sprintf("mcp: rpc error %d: %s", e.Code, e.Message) }
// transport carries one JSON-RPC conversation. Implementations: stdioTransport
// (a subprocess on this box) and httpTransport (streamable HTTP, guarded by
// webfetch). Both must be safe for concurrent use by the Client.
type transport interface {
// Call sends a request and returns the matching response.
Call(ctx context.Context, req *rpcRequest) (*rpcResponse, error)
// Notify sends a notification (no id, no reply expected).
Notify(ctx context.Context, method string, params any) error
// Close releases the transport (kills the subprocess, drops the session).
Close() error
}
+523
View File
@@ -0,0 +1,523 @@
package mcp
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"sort"
"strconv"
"strings"
"sync"
"time"
)
// Defaults for a server block. Small numbers on purpose — see MaxTools.
const (
// DefaultTimeout bounds one JSON-RPC call. A tool that takes longer than
// this is not usable in a spoken turn anyway.
DefaultTimeout = 15 * time.Second
// DefaultMaxTools caps how many tools ONE server may contribute. The
// resident model is a 1.7B with a 4096-token context: a catalogue of forty
// tool names does not fit in its head, and a name it half-remembers is a
// wrong act. Twelve per server is already generous.
DefaultMaxTools = 12
// DefaultReconnectEvery is how long the manager waits before re-dialing a
// server whose connection died.
DefaultReconnectEvery = 30 * time.Second
)
// ErrNoServer — the named server is not configured or not connected.
var ErrNoServer = errors.New("mcp: no such server")
// ServerConfig is one configured MCP server. Off unless present.
//
// Exactly one of Command (a subprocess on this box) or URL (a remote or
// loopback HTTP endpoint) must be set.
type ServerConfig struct {
// Name is the local handle. It prefixes every tool this server
// contributes, so it must be short and a valid identifier-ish word.
Name string `json:"name"`
// Command + Args + Env + Dir describe a stdio server: a child process of
// mavend, on this box, under this user. argv, never a shell string.
Command string `json:"command,omitempty"`
Args []string `json:"args,omitempty"`
Env []string `json:"env,omitempty"`
Dir string `json:"dir,omitempty"`
// URL is a streamable-HTTP endpoint. It goes through internal/webfetch, so
// it inherits the SSRF guard, the size cap and the per-host rate limit.
URL string `json:"url,omitempty"`
// AllowPrivate lets THIS server be a loopback or LAN address
// (http://localhost:9100/mcp is the Vikunja server on homesrv). It is a
// per-server hole in the private-address guard and it is not the same trust
// level as a public endpoint: whatever is behind it is inside the network,
// so an argument the router got wrong reaches something that matters. Set
// it only for a server you run.
AllowPrivate bool `json:"allow_private,omitempty"`
// AllowTools, when non-empty, is the ONLY set of remote tool names taken
// from this server. This is the knob for keeping the catalogue small and
// deliberate rather than "whatever the server grew this week".
AllowTools []string `json:"allow_tools,omitempty"`
// MaxTools caps the contribution (0 ⇒ DefaultMaxTools).
MaxTools int `json:"max_tools,omitempty"`
// Timeout bounds one call (0 ⇒ DefaultTimeout).
Timeout time.Duration `json:"-"`
// Enabled=false keeps a configured server described but dark.
Enabled bool `json:"enabled"`
}
// PosterFactory builds the HTTP door for one server. It is a factory rather
// than a single shared Poster because allow_private is per server: the fetcher
// that may reach http://localhost:9100/mcp must NOT be the same fetcher another
// server's public URL goes through, or one loopback exemption would quietly
// unlock the LAN for all of them.
type PosterFactory func(cfg ServerConfig) (Poster, error)
// Manager owns the connections. Nothing here starts unless at least one server
// is configured and enabled.
type Manager struct {
newPoster PosterFactory
mu sync.Mutex
conns map[string]*conn
order []string
}
type conn struct {
cfg ServerConfig
client *Client
tools []Tool
lastErr error
lastTry time.Time
dialedAt time.Time
}
// NewManager builds a manager for the enabled servers in cfgs. newPoster is
// the guarded HTTP door factory for url servers; pass nil only when no url
// server is configured (a nil factory with a url server is reported per server
// at dial time rather than fatally, so one bad block never stops the daemon).
//
// Dialing is lazy: NewManager validates and records, Connect dials.
func NewManager(newPoster PosterFactory, cfgs []ServerConfig) (*Manager, error) {
m := &Manager{newPoster: newPoster, conns: map[string]*conn{}}
for _, c := range cfgs {
if !c.Enabled {
continue
}
if err := validate(c); err != nil {
return nil, err
}
if _, dup := m.conns[c.Name]; dup {
return nil, fmt.Errorf("mcp: duplicate server name %q", c.Name)
}
if c.Timeout <= 0 {
c.Timeout = DefaultTimeout
}
if c.MaxTools <= 0 {
c.MaxTools = DefaultMaxTools
}
m.conns[c.Name] = &conn{cfg: c}
m.order = append(m.order, c.Name)
}
sort.Strings(m.order)
return m, nil
}
// Validate checks a set of server blocks without dialling anything, so a typo
// fails at startup rather than at the first turn that needed the tool.
func Validate(cfgs []ServerConfig) error {
seen := map[string]bool{}
for _, c := range cfgs {
if err := validate(c); err != nil {
return err
}
if seen[c.Name] {
return fmt.Errorf("mcp: duplicate server name %q", c.Name)
}
seen[c.Name] = true
}
return nil
}
func validate(c ServerConfig) error {
if strings.TrimSpace(c.Name) == "" {
return errors.New("mcp: server needs a name")
}
if strings.ContainsAny(c.Name, " \t/:") {
return fmt.Errorf("mcp: server name %q must be one word without spaces, slashes or colons", c.Name)
}
hasCmd, hasURL := c.Command != "", c.URL != ""
if hasCmd == hasURL {
return fmt.Errorf("mcp: server %q needs exactly one of command or url", c.Name)
}
if hasURL && !strings.HasPrefix(c.URL, "http://") && !strings.HasPrefix(c.URL, "https://") {
return fmt.Errorf("mcp: server %q url must be http or https", c.Name)
}
return nil
}
// Servers — the configured, enabled server names, sorted.
func (m *Manager) Servers() []string {
m.mu.Lock()
defer m.mu.Unlock()
return append([]string(nil), m.order...)
}
// Empty reports whether nothing is configured. The daemon uses it to skip
// wiring entirely.
func (m *Manager) Empty() bool {
m.mu.Lock()
defer m.mu.Unlock()
return len(m.conns) == 0
}
// Connect dials every configured server, handshakes, and discovers tools.
// A server that fails is recorded and retried later by Refresh — one bad
// server never blocks the others, and never blocks boot.
func (m *Manager) Connect(ctx context.Context) {
for _, name := range m.Servers() {
if err := m.dial(ctx, name); err != nil {
log.Printf("mcp: %s: %v", name, err)
}
}
}
func (m *Manager) dial(ctx context.Context, name string) error {
m.mu.Lock()
c, ok := m.conns[name]
if !ok {
m.mu.Unlock()
return ErrNoServer
}
cfg := c.cfg
c.lastTry = time.Now()
m.mu.Unlock()
var tr transport
var err error
if cfg.Command != "" {
tr, err = newStdioTransport(ctx, append([]string{cfg.Command}, cfg.Args...), cfg.Env, cfg.Dir)
} else if m.newPoster == nil {
err = fmt.Errorf("server %q has a url but no http door was wired", name)
} else {
var poster Poster
if poster, err = m.newPoster(cfg); err == nil {
tr = newHTTPTransport(poster, cfg.URL)
}
}
if err != nil {
m.fail(name, err)
return err
}
cl := newClient(name, tr)
ictx, cancel := context.WithTimeout(ctx, cfg.Timeout)
defer cancel()
if err := cl.Initialize(ictx); err != nil {
_ = cl.Close()
m.fail(name, err)
return err
}
tools, err := cl.ListTools(ictx)
if err != nil {
// A server with no tools capability is still a usable resource server.
log.Printf("mcp: %s: list tools: %v", name, err)
tools = nil
}
tools = filterTools(cfg, tools)
m.mu.Lock()
if old := m.conns[name].client; old != nil {
_ = old.Close()
}
m.conns[name].client = cl
m.conns[name].tools = tools
m.conns[name].lastErr = nil
m.conns[name].dialedAt = time.Now()
m.mu.Unlock()
log.Printf("mcp: %s connected (%s %s), %d tool(s)", name, cl.Info().Name, cl.Info().Version, len(tools))
return nil
}
func (m *Manager) fail(name string, err error) {
m.mu.Lock()
defer m.mu.Unlock()
if c := m.conns[name]; c != nil {
c.lastErr = err
c.client = nil
c.tools = nil
}
}
// filterTools applies AllowTools and MaxTools, and drops nameless entries.
// Sorted first, so the cap is deterministic rather than "whatever order the
// server felt like".
func filterTools(cfg ServerConfig, in []Tool) []Tool {
sort.Slice(in, func(i, j int) bool { return in[i].Name < in[j].Name })
out := make([]Tool, 0, len(in))
for _, t := range in {
if len(cfg.AllowTools) > 0 && !contains(cfg.AllowTools, t.Name) {
continue
}
out = append(out, t)
}
if cfg.MaxTools > 0 && len(out) > cfg.MaxTools {
log.Printf("mcp: %s offers %d tools, taking the first %d (raise max_tools or set allow_tools)",
cfg.Name, len(out), cfg.MaxTools)
out = out[:cfg.MaxTools]
}
return out
}
func contains(hay []string, needle string) bool {
for _, h := range hay {
if h == needle {
return true
}
}
return false
}
// Refresh re-dials any server that is down, if enough time has passed since the
// last attempt. Call it from the daemon's periodic tick — it is cheap when
// everything is up.
func (m *Manager) Refresh(ctx context.Context) {
now := time.Now()
var stale []string
m.mu.Lock()
for _, name := range m.order {
c := m.conns[name]
down := c.client == nil || !c.client.alive()
if down && now.Sub(c.lastTry) >= DefaultReconnectEvery {
stale = append(stale, name)
}
}
m.mu.Unlock()
for _, name := range stale {
if err := m.dial(ctx, name); err != nil {
log.Printf("mcp: %s: reconnect: %v", name, err)
}
}
}
// Tools — every discovered tool across connected servers, sorted by
// server then name.
func (m *Manager) Tools() []Tool {
m.mu.Lock()
defer m.mu.Unlock()
var out []Tool
for _, name := range m.order {
out = append(out, m.conns[name].tools...)
}
return out
}
// Status is one server's health, for the web surface.
type Status struct {
Name string
Transport string // "stdio" or "http"
Target string // command or url
Connected bool
Server string // the server's own name+version
Tools int
Err string
}
// Status reports every configured server.
func (m *Manager) Status() []Status {
m.mu.Lock()
defer m.mu.Unlock()
out := make([]Status, 0, len(m.order))
for _, name := range m.order {
c := m.conns[name]
s := Status{Name: name, Tools: len(c.tools)}
if c.cfg.Command != "" {
s.Transport, s.Target = "stdio", strings.Join(append([]string{c.cfg.Command}, c.cfg.Args...), " ")
} else {
s.Transport, s.Target = "http", c.cfg.URL
}
if c.client != nil {
s.Connected = true
s.Server = strings.TrimSpace(c.client.Info().Name + " " + c.client.Info().Version)
}
if c.lastErr != nil {
s.Err = c.lastErr.Error()
}
out = append(out, s)
}
return out
}
// Call runs server's tool with args. Args come from the router and nothing
// else; there is no path here through which a note or a fact could travel.
func (m *Manager) Call(ctx context.Context, server, tool string, args map[string]any) (string, error) {
m.mu.Lock()
c := m.conns[server]
m.mu.Unlock()
if c == nil {
return "", fmt.Errorf("%w: %s", ErrNoServer, server)
}
m.mu.Lock()
cl, timeout, known := c.client, c.cfg.Timeout, false
for _, t := range c.tools {
if t.Name == tool {
known = true
break
}
}
m.mu.Unlock()
if cl == nil {
return "", fmt.Errorf("mcp: %s is not connected", server)
}
// The discovered-and-filtered set is the second allowlist: even an enabled
// store row cannot reach a tool the server stopped offering, or one
// allow_tools excludes.
if !known {
return "", fmt.Errorf("mcp: %s offers no tool %q", server, tool)
}
cctx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
return cl.CallTool(cctx, tool, args)
}
// Resources lists resources across connected servers.
func (m *Manager) Resources(ctx context.Context) []Resource {
m.mu.Lock()
clients := make([]*Client, 0, len(m.order))
for _, name := range m.order {
if cl := m.conns[name].client; cl != nil {
clients = append(clients, cl)
}
}
m.mu.Unlock()
var out []Resource
for _, cl := range clients {
rs, err := cl.ListResources(ctx)
if err != nil {
continue // no resources capability; not an error worth logging per tick
}
out = append(out, rs...)
}
return out
}
// ReadResource reads one resource from one server.
func (m *Manager) ReadResource(ctx context.Context, server, uri string) (string, error) {
m.mu.Lock()
c := m.conns[server]
var cl *Client
var timeout time.Duration
if c != nil {
cl, timeout = c.client, c.cfg.Timeout
}
m.mu.Unlock()
if cl == nil {
return "", fmt.Errorf("%w: %s", ErrNoServer, server)
}
cctx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
return cl.ReadResource(cctx, uri)
}
// Close shuts every connection down.
func (m *Manager) Close() error {
m.mu.Lock()
defer m.mu.Unlock()
for _, name := range m.order {
if cl := m.conns[name].client; cl != nil {
_ = cl.Close()
m.conns[name].client = nil
}
}
return nil
}
// ErrNeedsArgs — the tool requires arguments that a voice verb cannot supply.
var ErrNeedsArgs = errors.New("mcp: tool needs named arguments")
// CallPositional is the voice path's way in. The router gives an act a verb and
// a tail of positional words; an MCP tool wants a named-argument object. There
// is no general mapping between those two, and inventing one is exactly the
// improvisation this codebase refuses, so the rule is deliberately narrow:
//
// - a tool with no required properties runs with no arguments (a spare tail
// is ignored — "покажи проекты пожалуйста" should still list projects);
// - a READ-ONLY tool with exactly one required property, of type string or
// integer/number, gets the tail bound to it;
// - anything else is refused with ErrNeedsArgs. Such a tool is still callable
// with explicit arguments from the authed surface, where a human types
// them.
//
// The refusal is the point, and the read-only condition on it was learned the
// hard way while testing against the Vikunja server: `update_task` requires
// only `task_id` and takes every other field as optional, so calling it with
// one guessed argument and no others BLANKED the fields it did not receive. A
// mutating tool therefore never gets a guessed argument — the one thing a
// partially-filled write can do is destroy what it did not mention. A mutating
// tool with nothing required is still fine: nothing was guessed, and it still
// goes through the confirm turn.
func (m *Manager) CallPositional(ctx context.Context, server, tool string, args []string) (string, error) {
m.mu.Lock()
c := m.conns[server]
var schema json.RawMessage
found, readOnly := false, false
if c != nil {
for _, t := range c.tools {
if t.Name == tool {
schema, readOnly, found = t.InputSchema, t.ReadOnly, true
break
}
}
}
m.mu.Unlock()
if !found {
return "", fmt.Errorf("mcp: %s offers no tool %q", server, tool)
}
named, err := bindPositional(schema, args, readOnly)
if err != nil {
return "", err
}
return m.Call(ctx, server, tool, named)
}
// bindPositional implements the rule documented on CallPositional.
func bindPositional(schema json.RawMessage, args []string, readOnly bool) (map[string]any, error) {
var s struct {
Required []string `json:"required"`
Properties map[string]struct {
Type string `json:"type"`
} `json:"properties"`
}
if len(schema) > 0 {
if err := json.Unmarshal(schema, &s); err != nil {
return nil, fmt.Errorf("mcp: unreadable input schema: %w", err)
}
}
switch len(s.Required) {
case 0:
return map[string]any{}, nil
case 1:
name := s.Required[0]
if !readOnly {
return nil, fmt.Errorf("%w: %q, and a tool that writes never gets a guessed one", ErrNeedsArgs, name)
}
tail := strings.TrimSpace(strings.Join(args, " "))
if tail == "" {
return nil, fmt.Errorf("%w: %q", ErrNeedsArgs, name)
}
switch s.Properties[name].Type {
case "string", "":
return map[string]any{name: tail}, nil
case "integer", "number":
n, err := strconv.ParseFloat(tail, 64)
if err != nil {
return nil, fmt.Errorf("%w: %q wants a number, got %q", ErrNeedsArgs, name, tail)
}
return map[string]any{name: n}, nil
default:
return nil, fmt.Errorf("%w: %q is a %s", ErrNeedsArgs, name, s.Properties[name].Type)
}
default:
return nil, fmt.Errorf("%w: %s", ErrNeedsArgs, strings.Join(s.Required, ", "))
}
}
+565
View File
@@ -0,0 +1,565 @@
package mcp
import (
"context"
"encoding/json"
"fmt"
"strings"
"sync"
"testing"
"time"
)
// fakePoster answers POSTs from a canned handler, in either JSON or SSE form.
type fakePoster struct {
mu sync.Mutex
handler func(method string, params json.RawMessage) (any, *rpcError)
sse bool
session string
seen []map[string]string // headers of each request, for the session test
calls []string
}
func (f *fakePoster) Post(_ context.Context, _, _ string, body []byte, hdr map[string]string) (*PostResponse, error) {
var req struct {
ID *int64 `json:"id"`
Method string `json:"method"`
Params json.RawMessage `json:"params"`
}
if err := json.Unmarshal(body, &req); err != nil {
return nil, err
}
f.mu.Lock()
f.seen = append(f.seen, hdr)
f.calls = append(f.calls, req.Method)
f.mu.Unlock()
if req.ID == nil { // notification
return &PostResponse{Status: 202, Body: []byte(`{}`)}, nil
}
result, rerr := f.handler(req.Method, req.Params)
resp := map[string]any{"jsonrpc": "2.0", "id": *req.ID}
if rerr != nil {
resp["error"] = map[string]any{"code": rerr.Code, "message": rerr.Message}
} else {
resp["result"] = result
}
raw, _ := json.Marshal(resp)
out := &PostResponse{Status: 200, Body: raw, ContentType: "application/json", Header: map[string]string{}}
if f.sse {
out.ContentType = "text/event-stream"
out.Body = []byte("event: message\ndata: {\"jsonrpc\":\"2.0\",\"method\":\"notifications/progress\"}\n\nevent: message\ndata: " + string(raw) + "\n\n")
}
if f.session != "" {
out.Header["Mcp-Session-Id"] = f.session
}
return out, nil
}
// echoServer is a handler with two tools, one read-only and one not.
func echoServer() func(string, json.RawMessage) (any, *rpcError) {
return func(method string, params json.RawMessage) (any, *rpcError) {
switch method {
case "initialize":
return map[string]any{
"protocolVersion": ProtocolVersion,
"serverInfo": map[string]any{"name": "fake", "version": "0.1"},
}, nil
case "tools/list":
return map[string]any{"tools": []any{
map[string]any{
"name": "read_thing", "description": "reads",
"inputSchema": map[string]any{"type": "object"},
"annotations": map[string]any{"readOnlyHint": true},
},
map[string]any{"name": "break_thing", "description": "mutates"},
}}, nil
case "tools/call":
var p struct {
Name string `json:"name"`
Args map[string]any `json:"arguments"`
}
_ = json.Unmarshal(params, &p)
if p.Name == "break_thing" {
return map[string]any{"isError": true, "content": []any{
map[string]any{"type": "text", "text": "не вышло"}}}, nil
}
return map[string]any{"content": []any{
map[string]any{"type": "text", "text": fmt.Sprintf("%s:%v", p.Name, p.Args["q"])},
map[string]any{"type": "image", "text": "ignored"},
}}, nil
case "resources/list":
return map[string]any{"resources": []any{
map[string]any{"uri": "note://one", "name": "one", "mimeType": "text/plain"},
map[string]any{"uri": "", "name": "nameless"},
}}, nil
case "resources/read":
return map[string]any{"contents": []any{map[string]any{"text": "тело ресурса"}}}, nil
}
return nil, &rpcError{Code: -32601, Message: "method not found"}
}
}
func dialFake(t *testing.T, p *fakePoster) *Client {
t.Helper()
c := newClient("fake", newHTTPTransport(p, "http://example.test/mcp"))
if err := c.Initialize(context.Background()); err != nil {
t.Fatalf("initialize: %v", err)
}
return c
}
func TestHandshakeAndDiscovery(t *testing.T) {
for _, sse := range []bool{false, true} {
name := "json"
if sse {
name = "sse"
}
t.Run(name, func(t *testing.T) {
p := &fakePoster{handler: echoServer(), sse: sse}
c := dialFake(t, p)
if got := c.Info().Name; got != "fake" {
t.Fatalf("server name = %q", got)
}
if got := c.Info().ProtocolVersion; got != ProtocolVersion {
t.Fatalf("protocol = %q", got)
}
tools, err := c.ListTools(context.Background())
if err != nil {
t.Fatalf("list tools: %v", err)
}
if len(tools) != 2 {
t.Fatalf("tools = %+v", tools)
}
byName := map[string]Tool{}
for _, tl := range tools {
byName[tl.Name] = tl
}
if !byName["read_thing"].ReadOnly {
t.Error("read_thing should be read-only (readOnlyHint true)")
}
// The important direction: no annotation ⇒ assume it mutates.
if byName["break_thing"].ReadOnly {
t.Error("break_thing has no readOnlyHint, must NOT be treated as read-only")
}
if byName["read_thing"].Server != "fake" {
t.Error("tool should carry its server handle")
}
})
}
}
func TestCallToolTextOnly(t *testing.T) {
c := dialFake(t, &fakePoster{handler: echoServer()})
out, err := c.CallTool(context.Background(), "read_thing", map[string]any{"q": "привет"})
if err != nil {
t.Fatalf("call: %v", err)
}
if out != "read_thing:привет" {
t.Fatalf("out = %q (non-text content must be dropped)", out)
}
}
func TestCallToolErrorResult(t *testing.T) {
c := dialFake(t, &fakePoster{handler: echoServer()})
out, err := c.CallTool(context.Background(), "break_thing", nil)
if err == nil {
t.Fatal("isError result must surface as an error")
}
if out != "не вышло" {
t.Fatalf("text should still come back, got %q", out)
}
}
func TestResources(t *testing.T) {
c := dialFake(t, &fakePoster{handler: echoServer()})
rs, err := c.ListResources(context.Background())
if err != nil {
t.Fatalf("list resources: %v", err)
}
if len(rs) != 1 || rs[0].URI != "note://one" {
t.Fatalf("resources = %+v (a uri-less entry must be dropped)", rs)
}
body, err := c.ReadResource(context.Background(), "note://one")
if err != nil {
t.Fatalf("read: %v", err)
}
if body != "тело ресурса" {
t.Fatalf("body = %q", body)
}
}
func TestCallBeforeInitializeRefused(t *testing.T) {
c := newClient("fake", newHTTPTransport(&fakePoster{handler: echoServer()}, "http://example.test/mcp"))
if _, err := c.CallTool(context.Background(), "read_thing", nil); err != ErrNotInitialized {
t.Fatalf("err = %v, want ErrNotInitialized", err)
}
}
func TestSessionIDEchoed(t *testing.T) {
p := &fakePoster{handler: echoServer(), session: "sess-1"}
c := dialFake(t, p)
if _, err := c.ListTools(context.Background()); err != nil {
t.Fatal(err)
}
p.mu.Lock()
defer p.mu.Unlock()
last := p.seen[len(p.seen)-1]
if last["Mcp-Session-Id"] != "sess-1" {
t.Fatalf("session header not echoed: %+v", last)
}
if !strings.Contains(last["Accept"], "text/event-stream") {
t.Fatalf("Accept must offer both forms: %q", last["Accept"])
}
}
func TestHandshakeWithoutProtocolVersionRefused(t *testing.T) {
p := &fakePoster{handler: func(m string, _ json.RawMessage) (any, *rpcError) {
return map[string]any{"serverInfo": map[string]any{"name": "not-mcp"}}, nil
}}
c := newClient("x", newHTTPTransport(p, "http://example.test/mcp"))
if err := c.Initialize(context.Background()); err == nil {
t.Fatal("a reply with no protocolVersion is not an MCP server")
}
}
func TestRPCErrorSurfaces(t *testing.T) {
c := dialFake(t, &fakePoster{handler: echoServer()})
if _, err := c.callRaw(context.Background(), "nope/nope"); err == nil {
t.Fatal("want an rpc error")
} else if !strings.Contains(err.Error(), "method not found") {
t.Fatalf("err = %v", err)
}
}
// callRaw is a test-only shim so the rpc-error path can be exercised without a
// typed wrapper for a method the server does not implement.
func (c *Client) callRaw(ctx context.Context, method string) (any, error) {
var out any
err := c.call(ctx, method, map[string]any{}, &out)
return out, err
}
func TestDecodeFrame(t *testing.T) {
cases := []struct {
name, in, want string
wantErr bool
}{
{name: "plain json", in: `{"id":1,"result":{}}`, want: `{"id":1,"result":{}}`},
{name: "sse single", in: "event: message\ndata: {\"id\":1,\"result\":1}\n\n", want: `{"id":1,"result":1}`},
{
name: "sse picks the response not the notification",
in: "data: {\"method\":\"notifications/progress\"}\n\ndata: {\"id\":2,\"result\":2}\n\n",
want: `{"id":2,"result":2}`,
},
{name: "empty", in: " ", wantErr: true},
{name: "sse with no response", in: "data: {\"method\":\"x\"}\n\n", wantErr: true},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
got, err := decodeFrame([]byte(tc.in))
if tc.wantErr {
if err == nil {
t.Fatalf("want error, got %q", got)
}
return
}
if err != nil {
t.Fatal(err)
}
if string(got) != tc.want {
t.Fatalf("got %q want %q", got, tc.want)
}
})
}
}
func TestValidate(t *testing.T) {
cases := []struct {
name string
cfg ServerConfig
wantErr bool
}{
{name: "stdio ok", cfg: ServerConfig{Name: "a", Command: "echo"}},
{name: "http ok", cfg: ServerConfig{Name: "a", URL: "http://x.test/mcp"}},
{name: "no name", cfg: ServerConfig{Command: "echo"}, wantErr: true},
{name: "spacey name", cfg: ServerConfig{Name: "a b", Command: "echo"}, wantErr: true},
{name: "neither", cfg: ServerConfig{Name: "a"}, wantErr: true},
{name: "both", cfg: ServerConfig{Name: "a", Command: "echo", URL: "http://x.test"}, wantErr: true},
{name: "bad scheme", cfg: ServerConfig{Name: "a", URL: "file:///etc/passwd"}, wantErr: true},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
err := Validate([]ServerConfig{tc.cfg})
if (err != nil) != tc.wantErr {
t.Fatalf("err = %v, wantErr = %v", err, tc.wantErr)
}
})
}
if err := Validate([]ServerConfig{{Name: "a", Command: "x"}, {Name: "a", Command: "y"}}); err == nil {
t.Error("duplicate names must be refused")
}
}
func TestManagerOffWhenNothingEnabled(t *testing.T) {
m, err := NewManager(nil, []ServerConfig{{Name: "a", Command: "echo"}}) // Enabled=false
if err != nil {
t.Fatal(err)
}
if !m.Empty() {
t.Fatal("a server that is not enabled must not be wired")
}
m.Connect(context.Background())
if got := m.Tools(); len(got) != 0 {
t.Fatalf("tools = %+v", got)
}
}
func TestManagerDiscoversAndCalls(t *testing.T) {
p := &fakePoster{handler: echoServer()}
m, err := NewManager(func(ServerConfig) (Poster, error) { return p, nil },
[]ServerConfig{{Name: "fake", URL: "http://example.test/mcp", Enabled: true}})
if err != nil {
t.Fatal(err)
}
m.Connect(context.Background())
defer m.Close()
tools := m.Tools()
if len(tools) != 2 {
t.Fatalf("tools = %+v", tools)
}
out, err := m.Call(context.Background(), "fake", "read_thing", map[string]any{"q": "да"})
if err != nil {
t.Fatalf("call: %v", err)
}
if out != "read_thing:да" {
t.Fatalf("out = %q", out)
}
// The discovered set is a second allowlist.
if _, err := m.Call(context.Background(), "fake", "not_offered", nil); err == nil {
t.Error("a tool the server does not offer must be refused")
}
if _, err := m.Call(context.Background(), "other", "read_thing", nil); err == nil {
t.Error("an unconfigured server must be refused")
}
st := m.Status()
if len(st) != 1 || !st[0].Connected || st[0].Transport != "http" || st[0].Tools != 2 {
t.Fatalf("status = %+v", st)
}
}
func TestManagerAllowToolsAndMaxTools(t *testing.T) {
p := &fakePoster{handler: echoServer()}
mk := func(cfg ServerConfig) *Manager {
cfg.Name, cfg.URL, cfg.Enabled = "fake", "http://example.test/mcp", true
m, err := NewManager(func(ServerConfig) (Poster, error) { return p, nil }, []ServerConfig{cfg})
if err != nil {
t.Fatal(err)
}
m.Connect(context.Background())
return m
}
m := mk(ServerConfig{AllowTools: []string{"read_thing"}})
defer m.Close()
if got := m.Tools(); len(got) != 1 || got[0].Name != "read_thing" {
t.Fatalf("allow_tools ignored: %+v", got)
}
if _, err := m.Call(context.Background(), "fake", "break_thing", nil); err == nil {
t.Error("a tool excluded by allow_tools must be unreachable")
}
m2 := mk(ServerConfig{MaxTools: 1})
defer m2.Close()
if got := m2.Tools(); len(got) != 1 || got[0].Name != "break_thing" {
t.Fatalf("max_tools should keep the first name-sorted tool: %+v", got)
}
}
func TestManagerURLServerWithoutHTTPDoor(t *testing.T) {
m, err := NewManager(nil, []ServerConfig{{Name: "fake", URL: "http://example.test/mcp", Enabled: true}})
if err != nil {
t.Fatal(err)
}
m.Connect(context.Background())
st := m.Status()
if len(st) != 1 || st[0].Connected || st[0].Err == "" {
t.Fatalf("a url server with no poster must be recorded as failed: %+v", st)
}
}
func TestManagerReconnectAfterFailure(t *testing.T) {
var mu sync.Mutex
fail := true
m, err := NewManager(func(ServerConfig) (Poster, error) {
mu.Lock()
defer mu.Unlock()
if fail {
return nil, fmt.Errorf("down")
}
return &fakePoster{handler: echoServer()}, nil
}, []ServerConfig{{Name: "fake", URL: "http://example.test/mcp", Enabled: true}})
if err != nil {
t.Fatal(err)
}
defer m.Close()
m.Connect(context.Background())
if m.Status()[0].Connected {
t.Fatal("should be down")
}
mu.Lock()
fail = false
mu.Unlock()
// Refresh honours the backoff, so pretend the last attempt was long ago.
m.mu.Lock()
m.conns["fake"].lastTry = time.Now().Add(-2 * DefaultReconnectEvery)
m.mu.Unlock()
m.Refresh(context.Background())
if !m.Status()[0].Connected {
t.Fatalf("should have reconnected: %+v", m.Status())
}
}
func TestLocalNameAndCmd(t *testing.T) {
cases := [][3]string{
{"vikunja", "list_tasks", "vikunja_list_tasks"},
{"Vikunja", "Get Task Details", "vikunja_get_task_details"},
{"fs", "read-file", "fs_read_file"},
{"", "search", "search"},
}
for _, c := range cases {
if got := LocalName(c[0], c[1]); got != c[2] {
t.Errorf("LocalName(%q,%q) = %q want %q", c[0], c[1], got, c[2])
}
}
server, tool, ok := ParseCmd(Cmd("vikunja", "list_tasks"))
if !ok || server != "vikunja" || tool != "list_tasks" {
t.Fatalf("ParseCmd round-trip: %q %q %v", server, tool, ok)
}
for _, bad := range [][]string{nil, {"systemctl", "restart", "nginx"}, {"mcp", "vikunja"}, {"mcp", "", "x"}} {
if _, _, ok := ParseCmd(bad); ok {
t.Errorf("ParseCmd(%v) must not claim an ordinary tool row", bad)
}
}
if Scope("vikunja") != "mcp:vikunja" {
t.Error("scope")
}
}
func TestBindPositional(t *testing.T) {
cases := []struct {
name, schema string
args []string
mutating bool
want map[string]any
wantErr bool
}{
{
name: "no required runs with nothing",
// A spare tail is fine: "покажи проекты пожалуйста" still lists them.
schema: `{"type":"object","properties":{},"required":[]}`,
args: []string{"пожалуйста"},
want: map[string]any{},
},
{
name: "empty schema",
schema: ``,
want: map[string]any{},
},
{
name: "one required string gets the tail",
schema: `{"properties":{"q":{"type":"string"}},"required":["q"]}`,
args: []string{"почему", "небо", "синее"},
want: map[string]any{"q": "почему небо синее"},
},
{
name: "one required string with no tail",
schema: `{"properties":{"q":{"type":"string"}},"required":["q"]}`,
wantErr: true,
},
{
name: "one required integer parses",
schema: `{"properties":{"task_id":{"type":"integer"}},"required":["task_id"]}`,
args: []string{"251"},
want: map[string]any{"task_id": float64(251)},
},
{
name: "one required integer with words",
schema: `{"properties":{"task_id":{"type":"integer"}},"required":["task_id"]}`,
args: []string{"двести", "пятьдесят", "один"},
wantErr: true,
},
{
name: "two required is refused rather than guessed",
schema: `{"properties":{"a":{"type":"string"},"b":{"type":"string"}},"required":["a","b"]}`,
args: []string{"что-то"},
wantErr: true,
},
{
name: "one required object is refused",
schema: `{"properties":{"payload":{"type":"object"}},"required":["payload"]}`,
args: []string{"что-то"},
wantErr: true,
},
{
// Learned from Vikunja's update_task: required ["task_id"], every
// other field optional, so one guessed argument blanks the rest.
name: "one required on a mutating tool is refused",
schema: `{"properties":{"task_id":{"type":"integer"}},"required":["task_id"]}`,
args: []string{"251"},
mutating: true,
wantErr: true,
},
{
// Nothing was guessed, so there is nothing to get wrong. It still
// goes through the confirm turn upstream.
name: "no required on a mutating tool still runs",
schema: `{"properties":{},"required":[]}`,
mutating: true,
want: map[string]any{},
},
{
name: "unreadable schema",
schema: `not json`,
wantErr: true,
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
got, err := bindPositional(json.RawMessage(tc.schema), tc.args, !tc.mutating)
if tc.wantErr {
if err == nil {
t.Fatalf("want an error, got %v", got)
}
return
}
if err != nil {
t.Fatal(err)
}
if fmt.Sprint(got) != fmt.Sprint(tc.want) {
t.Fatalf("got %v want %v", got, tc.want)
}
})
}
}
func TestCallPositionalThroughManager(t *testing.T) {
p := &fakePoster{handler: echoServer()}
m, err := NewManager(func(ServerConfig) (Poster, error) { return p, nil },
[]ServerConfig{{Name: "fake", URL: "http://example.test/mcp", Enabled: true}})
if err != nil {
t.Fatal(err)
}
m.Connect(context.Background())
defer m.Close()
// echoServer's tools declare no required properties.
out, err := m.CallPositional(context.Background(), "fake", "read_thing", []string{"хвост"})
if err != nil {
t.Fatalf("call: %v", err)
}
if out != "read_thing:<nil>" {
t.Fatalf("out = %q", out)
}
if _, err := m.CallPositional(context.Background(), "fake", "absent", nil); err == nil {
t.Error("an unknown tool must be refused")
}
}
+149
View File
@@ -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
}
+149
View File
@@ -0,0 +1,149 @@
package mcp
import (
"bufio"
"context"
"encoding/json"
"fmt"
"os"
"os/exec"
"strings"
"testing"
)
// The stdio transport is tested against a real subprocess — this test binary,
// re-executed with MAVEN_MCP_FAKE set, acting as a minimal MCP server. No
// python, no fixture file, no network.
func TestMain(m *testing.M) {
if os.Getenv("MAVEN_MCP_FAKE") != "" {
fakeStdioServer()
return
}
os.Exit(m.Run())
}
func fakeStdioServer() {
h := echoServer()
sc := bufio.NewScanner(os.Stdin)
out := bufio.NewWriter(os.Stdout)
defer out.Flush()
for sc.Scan() {
line := strings.TrimSpace(sc.Text())
if line == "" {
continue
}
var req struct {
ID *int64 `json:"id"`
Method string `json:"method"`
Params json.RawMessage `json:"params"`
}
if json.Unmarshal([]byte(line), &req) != nil {
continue
}
if req.ID == nil {
// A notification gets no reply, but we emit an unrelated
// notification so the client's frame-skipping is exercised.
_, _ = out.WriteString("{\"jsonrpc\":\"2.0\",\"method\":\"notifications/message\"}\n")
_ = out.Flush()
continue
}
result, rerr := h(req.Method, req.Params)
resp := map[string]any{"jsonrpc": "2.0", "id": *req.ID}
if rerr != nil {
resp["error"] = map[string]any{"code": rerr.Code, "message": rerr.Message}
} else {
resp["result"] = result
}
raw, _ := json.Marshal(resp)
_, _ = out.Write(append(raw, '\n'))
_ = out.Flush()
if os.Getenv("MAVEN_MCP_FAKE") == "die" && req.Method == "tools/list" {
return // hang up, so the reconnect path has something to see
}
}
}
func stdioManager(t *testing.T, mode string) *Manager {
t.Helper()
self, err := os.Executable()
if err != nil {
t.Skipf("no executable path: %v", err)
}
if _, err := exec.LookPath(self); err != nil && !strings.Contains(self, "/") {
t.Skip("test binary not executable")
}
m, err := NewManager(nil, []ServerConfig{{
Name: "fake",
Command: self,
Env: []string{"MAVEN_MCP_FAKE=" + mode},
Enabled: true,
}})
if err != nil {
t.Fatal(err)
}
m.Connect(context.Background())
return m
}
func TestStdioTransportEndToEnd(t *testing.T) {
m := stdioManager(t, "1")
defer m.Close()
st := m.Status()
if len(st) != 1 || !st[0].Connected {
t.Fatalf("status = %+v", st)
}
if st[0].Transport != "stdio" {
t.Fatalf("transport = %q", st[0].Transport)
}
if got := len(m.Tools()); got != 2 {
t.Fatalf("tools = %d", got)
}
out, err := m.Call(context.Background(), "fake", "read_thing", map[string]any{"q": "стдио"})
if err != nil {
t.Fatalf("call: %v", err)
}
if out != "read_thing:стдио" {
t.Fatalf("out = %q", out)
}
res := m.Resources(context.Background())
if len(res) != 1 || res[0].URI != "note://one" {
t.Fatalf("resources = %+v", res)
}
body, err := m.ReadResource(context.Background(), "fake", "note://one")
if err != nil {
t.Fatal(err)
}
if body != "тело ресурса" {
t.Fatalf("body = %q", body)
}
}
func TestStdioServerThatDiesIsNotUsable(t *testing.T) {
m := stdioManager(t, "die")
defer m.Close()
// The server hung up after tools/list; the next call must fail cleanly
// rather than hang or panic.
if _, err := m.Call(context.Background(), "fake", "read_thing", nil); err == nil {
t.Fatal("a call into a dead server must error")
}
}
func TestStdioMissingCommand(t *testing.T) {
m, err := NewManager(nil, []ServerConfig{{
Name: "nope", Command: "/nonexistent/mcp-server-that-is-not-there", Enabled: true,
}})
if err != nil {
t.Fatal(err)
}
m.Connect(context.Background())
st := m.Status()
if st[0].Connected || st[0].Err == "" {
t.Fatalf("a missing binary must be recorded, not fatal: %+v", st)
}
if got := len(m.Tools()); got != 0 {
t.Fatalf("tools = %d", got)
}
if !strings.Contains(fmt.Sprint(st[0].Err), "start") {
t.Logf("err = %q", st[0].Err)
}
}
+46
View File
@@ -0,0 +1,46 @@
package mcp
import (
"context"
"fmt"
"github.com/kami/maven/internal/webfetch"
)
// WebfetchDoor builds the PosterFactory used in production: one guarded
// webfetch.Fetcher per url server, with that server's allow_private and the
// shared host lists and limits.
//
// One fetcher PER server is the point. allow_private is a hole in the
// private-address guard, and a hole punched for the Vikunja server on loopback
// must not become a hole for some public endpoint that happens to redirect at
// the LAN. Rate limiting is per fetcher too, which is the right shape here:
// separate servers are separate hosts.
func WebfetchDoor(limits webfetch.Config) PosterFactory {
return func(cfg ServerConfig) (Poster, error) {
c := limits
c.AllowPrivate = cfg.AllowPrivate
if c.Timeout <= 0 && cfg.Timeout > 0 {
c.Timeout = cfg.Timeout
}
return fetcherPoster{webfetch.New(c)}, nil
}
}
// fetcherPoster adapts webfetch.Fetcher to Poster. It exists so this package
// does not have to know webfetch's Response type, and so a test can substitute
// a fake without a listener.
type fetcherPoster struct{ f *webfetch.Fetcher }
func (p fetcherPoster) Post(ctx context.Context, rawURL, contentType string, body []byte, hdr map[string]string) (*PostResponse, error) {
resp, err := p.f.Post(ctx, rawURL, contentType, body, hdr)
if err != nil {
return nil, fmt.Errorf("mcp: post %s: %w", rawURL, err)
}
return &PostResponse{
Status: resp.Status,
ContentType: resp.ContentType,
Body: resp.Body,
Header: resp.Header,
}, nil
}
+188
View File
@@ -0,0 +1,188 @@
package media
import (
"bytes"
"encoding/base64"
"errors"
"fmt"
"image"
"image/draw"
"image/gif"
"image/jpeg"
"image/png"
"strings"
)
// DefaultMaxDim — the longest edge an image is scaled down to before it goes to
// a vision model. 896 is the tile size the current crop of small
// vision-language models (Qwen2.5-VL, SmolVLM, moondream) work in; sending a
// 12-megapixel phone photo instead just costs the box minutes of prefill for
// tiles that get pooled away anyway.
const DefaultMaxDim = 896
// JPEGQuality for the re-encode. 85 is the usual "no visible artefacts" point,
// and the re-encode exists to shrink the payload, not to archive it — the
// original bytes stay in the blob store untouched.
const JPEGQuality = 85
// ErrUnsupportedImage — the bytes are not an image format this build can
// decode. Notably webp: the stdlib has no webp decoder and this repo takes no
// new dependencies, so a webp arriving from Telegram is refused here with a
// clear error rather than handed to a model as garbage.
var ErrUnsupportedImage = errors.New("media: unsupported image format")
// SniffImage identifies image bytes by magic number and returns the mime. It
// exists because a caller-declared content type is a claim, and the store's file
// extension (and the vision provider's data URI) should follow the bytes.
//
// Returns ErrUnsupportedImage for anything unrecognised, including webp — which
// is recognised well enough to name in the error, so the log says "webp is not
// supported" instead of "not an image".
func SniffImage(data []byte) (string, error) {
switch {
case len(data) >= 3 && data[0] == 0xFF && data[1] == 0xD8 && data[2] == 0xFF:
return "image/jpeg", nil
case len(data) >= 8 && string(data[:8]) == "\x89PNG\r\n\x1a\n":
return "image/png", nil
case len(data) >= 6 && (string(data[:6]) == "GIF87a" || string(data[:6]) == "GIF89a"):
return "image/gif", nil
case len(data) >= 12 && string(data[:4]) == "RIFF" && string(data[8:12]) == "WEBP":
return "", fmt.Errorf("%w: webp (no decoder in this build)", ErrUnsupportedImage)
}
return "", ErrUnsupportedImage
}
// Image — an image prepared for a vision model: JPEG bytes, downscaled, with
// the dimensions it ended up at. It is deliberately a separate type from Blob:
// a Blob is what he sent, an Image is what the model sees, and the two are not
// the same bytes.
type Image struct {
JPEG []byte
Width int
Height int
// Source names where the original came from ("telegram", "web:upload"),
// carried through only so a log line can say what was looked at.
Source string
}
// DataURI renders the image as a `data:image/jpeg;base64,...` URI, which is how
// every OpenAI-compatible multimodal endpoint takes an image. The string is
// large (roughly 4/3 of the JPEG); nothing caches it.
func (im Image) DataURI() string {
return "data:image/jpeg;base64," + base64.StdEncoding.EncodeToString(im.JPEG)
}
// PrepareImage decodes data, scales it so its longest edge is at most maxDim
// (never up — a small image is left alone), and re-encodes it as JPEG.
// maxDim ≤ 0 ⇒ DefaultMaxDim.
//
// An image with an alpha channel is composited onto white rather than having
// alpha dropped to black, because the common case is a screenshot or a
// transparent-background diagram, and text on black-on-black is unreadable to
// the model for no reason.
func PrepareImage(data []byte, source string, maxDim int) (Image, error) {
if len(data) == 0 {
return Image{}, ErrEmpty
}
if maxDim <= 0 {
maxDim = DefaultMaxDim
}
mime, err := SniffImage(data)
if err != nil {
return Image{}, err
}
src, err := decode(data, mime)
if err != nil {
return Image{}, fmt.Errorf("media: decode %s: %w", mime, err)
}
dst := flattenAndScale(src, maxDim)
var buf bytes.Buffer
if err := jpeg.Encode(&buf, dst, &jpeg.Options{Quality: JPEGQuality}); err != nil {
return Image{}, fmt.Errorf("media: encode jpeg: %w", err)
}
b := dst.Bounds()
return Image{JPEG: buf.Bytes(), Width: b.Dx(), Height: b.Dy(), Source: source}, nil
}
func decode(data []byte, mime string) (image.Image, error) {
r := bytes.NewReader(data)
switch strings.ToLower(mime) {
case "image/jpeg":
return jpeg.Decode(r)
case "image/png":
return png.Decode(r)
case "image/gif":
return gif.Decode(r)
}
return nil, ErrUnsupportedImage
}
// flattenAndScale composites onto white and box-scales down to maxDim. The
// scaler is a plain area average over the source pixels mapping to each
// destination pixel — nearest-neighbour would alias small text into noise,
// which defeats the point of reading a screenshot, and an area average is a
// dozen lines against pulling in golang.org/x/image on an offline box.
func flattenAndScale(src image.Image, maxDim int) *image.RGBA {
sb := src.Bounds()
sw, sh := sb.Dx(), sb.Dy()
dw, dh := fit(sw, sh, maxDim)
flat := image.NewRGBA(image.Rect(0, 0, sw, sh))
draw.Draw(flat, flat.Bounds(), image.NewUniform(image.White), image.Point{}, draw.Src)
draw.Draw(flat, flat.Bounds(), src, sb.Min, draw.Over)
if dw == sw && dh == sh {
return flat
}
dst := image.NewRGBA(image.Rect(0, 0, dw, dh))
for y := 0; y < dh; y++ {
y0, y1 := y*sh/dh, (y+1)*sh/dh
if y1 <= y0 {
y1 = y0 + 1
}
for x := 0; x < dw; x++ {
x0, x1 := x*sw/dw, (x+1)*sw/dw
if x1 <= x0 {
x1 = x0 + 1
}
var r, g, b, n uint32
for sy := y0; sy < y1; sy++ {
for sx := x0; sx < x1; sx++ {
i := flat.PixOffset(sx, sy)
r += uint32(flat.Pix[i])
g += uint32(flat.Pix[i+1])
b += uint32(flat.Pix[i+2])
n++
}
}
o := dst.PixOffset(x, y)
dst.Pix[o] = uint8(r / n)
dst.Pix[o+1] = uint8(g / n)
dst.Pix[o+2] = uint8(b / n)
dst.Pix[o+3] = 0xFF
}
}
return dst
}
// fit returns the largest w×h with the same aspect ratio whose longest edge is
// at most maxDim, never enlarging. Both edges are clamped to at least 1 so a
// 2000×1 strip does not scale to zero height.
func fit(w, h, maxDim int) (int, int) {
if w <= maxDim && h <= maxDim {
return w, h
}
if w >= h {
nh := h * maxDim / w
if nh < 1 {
nh = 1
}
return maxDim, nh
}
nw := w * maxDim / h
if nw < 1 {
nw = 1
}
return nw, maxDim
}
+191
View File
@@ -0,0 +1,191 @@
package media
import (
"bytes"
"errors"
"image"
"image/color"
"image/gif"
"image/jpeg"
"image/png"
"strings"
"testing"
)
// pngBytes builds a w×h test image: left half red, right half a light grey, so
// a downscale that averages produces a predictable mid value and a scaler that
// silently returns the wrong region is visible.
func pngBytes(t *testing.T, w, h int) []byte {
t.Helper()
img := image.NewRGBA(image.Rect(0, 0, w, h))
for y := 0; y < h; y++ {
for x := 0; x < w; x++ {
if x < w/2 {
img.Set(x, y, color.RGBA{255, 0, 0, 255})
} else {
img.Set(x, y, color.RGBA{200, 200, 200, 255})
}
}
}
var buf bytes.Buffer
if err := png.Encode(&buf, img); err != nil {
t.Fatal(err)
}
return buf.Bytes()
}
func TestSniffImage(t *testing.T) {
cases := []struct {
name string
data []byte
want string
}{
{"png", pngBytes(t, 4, 4), "image/png"},
{"jpeg", jpegBytes(t, 4, 4), "image/jpeg"},
{"gif", gifBytes(t, 4, 4), "image/gif"},
}
for _, c := range cases {
got, err := SniffImage(c.data)
if err != nil {
t.Errorf("%s: %v", c.name, err)
continue
}
if got != c.want {
t.Errorf("%s: got %q want %q", c.name, got, c.want)
}
}
}
// webp is common from Telegram and there is no stdlib decoder, so it must be
// refused by name rather than mis-sniffed or fed to a model as noise.
func TestSniffRefusesWebpByName(t *testing.T) {
webp := append([]byte("RIFF\x00\x00\x00\x00WEBP"), make([]byte, 8)...)
_, err := SniffImage(webp)
if !errors.Is(err, ErrUnsupportedImage) {
t.Fatalf("got %v, want ErrUnsupportedImage", err)
}
if !strings.Contains(err.Error(), "webp") {
t.Errorf("error does not name the format: %v", err)
}
}
func TestSniffRefusesGarbage(t *testing.T) {
for _, data := range [][]byte{nil, []byte("hello"), []byte("\x00\x01\x02\x03")} {
if _, err := SniffImage(data); !errors.Is(err, ErrUnsupportedImage) {
t.Errorf("SniffImage(%q) = %v", data, err)
}
}
}
func TestPrepareImageDownscalesLongestEdge(t *testing.T) {
im, err := PrepareImage(pngBytes(t, 2000, 1000), "web:upload", 500)
if err != nil {
t.Fatalf("prepare: %v", err)
}
if im.Width != 500 || im.Height != 250 {
t.Errorf("got %dx%d, want 500x250", im.Width, im.Height)
}
if _, err := jpeg.Decode(bytes.NewReader(im.JPEG)); err != nil {
t.Errorf("output is not decodable jpeg: %v", err)
}
if im.Source != "web:upload" {
t.Errorf("source lost: %q", im.Source)
}
}
// Tall images scale on the other axis; a scaler that only handles landscape is
// the classic version of this bug.
func TestPrepareImageHandlesPortrait(t *testing.T) {
im, err := PrepareImage(pngBytes(t, 400, 1600), "telegram", 800)
if err != nil {
t.Fatalf("prepare: %v", err)
}
if im.Height != 800 || im.Width != 200 {
t.Errorf("got %dx%d, want 200x800", im.Width, im.Height)
}
}
func TestPrepareImageNeverEnlarges(t *testing.T) {
im, err := PrepareImage(pngBytes(t, 64, 32), "telegram", 896)
if err != nil {
t.Fatalf("prepare: %v", err)
}
if im.Width != 64 || im.Height != 32 {
t.Errorf("got %dx%d, want the original 64x32", im.Width, im.Height)
}
}
// A degenerate strip must not scale to zero on the short axis — jpeg.Encode
// fails on a zero-height image, which would turn a weird screenshot into a
// hard error.
func TestPrepareImageClampsDegenerateAspect(t *testing.T) {
im, err := PrepareImage(pngBytes(t, 2000, 2), "web:upload", 100)
if err != nil {
t.Fatalf("prepare: %v", err)
}
if im.Height < 1 || im.Width != 100 {
t.Errorf("got %dx%d", im.Width, im.Height)
}
}
// Transparent pixels composite onto white, not black: the common case is a
// screenshot or a diagram, and dark-on-black is unreadable to the model.
func TestPrepareImageFlattensAlphaOntoWhite(t *testing.T) {
img := image.NewRGBA(image.Rect(0, 0, 8, 8)) // fully transparent
var buf bytes.Buffer
if err := png.Encode(&buf, img); err != nil {
t.Fatal(err)
}
im, err := PrepareImage(buf.Bytes(), "web:upload", 8)
if err != nil {
t.Fatalf("prepare: %v", err)
}
decoded, err := jpeg.Decode(bytes.NewReader(im.JPEG))
if err != nil {
t.Fatal(err)
}
r, g, b, _ := decoded.At(4, 4).RGBA()
if r>>8 < 240 || g>>8 < 240 || b>>8 < 240 {
t.Errorf("transparent pixel became rgb(%d,%d,%d), want near-white", r>>8, g>>8, b>>8)
}
}
func TestPrepareImageRejectsEmpty(t *testing.T) {
if _, err := PrepareImage(nil, "x", 0); !errors.Is(err, ErrEmpty) {
t.Errorf("got %v, want ErrEmpty", err)
}
}
func TestDataURIIsAJPEGDataURI(t *testing.T) {
im, err := PrepareImage(pngBytes(t, 16, 16), "x", 0)
if err != nil {
t.Fatal(err)
}
uri := im.DataURI()
if !strings.HasPrefix(uri, "data:image/jpeg;base64,") {
t.Fatalf("bad prefix: %.40s", uri)
}
if len(uri) <= len("data:image/jpeg;base64,") {
t.Error("data uri carries no payload")
}
}
func jpegBytes(t *testing.T, w, h int) []byte {
t.Helper()
img := image.NewRGBA(image.Rect(0, 0, w, h))
var buf bytes.Buffer
if err := jpeg.Encode(&buf, img, nil); err != nil {
t.Fatal(err)
}
return buf.Bytes()
}
func gifBytes(t *testing.T, w, h int) []byte {
t.Helper()
img := image.NewPaletted(image.Rect(0, 0, w, h), []color.Color{color.Black, color.White})
var buf bytes.Buffer
if err := gif.Encode(&buf, img, nil); err != nil {
t.Fatal(err)
}
return buf.Bytes()
}
+119
View File
@@ -0,0 +1,119 @@
// Package media is the intake for everything Maven sees or hears that is not
// text: a photo he sends her, a meeting she was asked to record, a voice sample
// used to enrol a speaker. All three senses (vision, hearing, speaker
// recognition) share one problem — a blob arrives, it has to be stored, and
// something has to describe it — so the storing half lives here once instead of
// three times.
//
// # What this package is
//
// A content-addressed blob store on the local filesystem. Put returns a Blob
// keyed by the sha256 of its bytes, so the same photo sent twice is one file.
// Each blob gets a sidecar `.json` with its kind, mime, size, source and
// creation time; the sidecar is the whole index, because at personal scale a
// directory walk is cheaper than another sqlite table and the store has to be
// readable with `ls` when something goes wrong.
//
// Blobs are NOT in the sqlite database. The database is small, encrypted, and
// read on every tick; a 40 MB meeting recording has no business in it. What
// goes in the database is the *text* a blob produced — a transcript, a
// description — written as an ordinary note, which is the durable artefact and
// the only part worth recalling later.
//
// # Invariants (these are the point of the package, not decoration)
//
// - Nothing is captured that was not asked for. This package never records;
// it stores what a caller hands it, and every caller is an explicit act
// with a start and a stop. There is no ambient path in, and none may be
// added: see the refusal recorded in docs/plans/08-hearing.md.
// - A blob never leaves the box. No provider in this repo may upload one, and
// the vision provider refuses a non-private endpoint for exactly that
// reason (internal/vision).
// - A blob is never search input and never embedded. His photos and the audio
// of his meetings are not corpus. Only text derived from them, once he can
// see it as a note, participates in recall.
// - Storage is bounded. Retention is a config knob with a default, Prune
// enforces it, and an unpruned store is a bug: audio of people accumulating
// forever on disk is the failure mode this capability has to avoid.
//
// # Layout
//
// <dir>/<kind>/<aa>/<sha256>.<ext> the bytes
// <dir>/<kind>/<aa>/<sha256>.json the sidecar metadata
//
// `aa` is the first two hex chars of the digest — one fan-out level, enough to
// keep a directory listing usable after a few thousand blobs.
package media
import (
"errors"
"fmt"
"time"
)
// Kind — what a blob is. Two values today; the kind is a directory name and a
// retention bucket, so adding a third is additive.
type Kind string
const (
// KindImage — a still image (png / jpeg / gif / webp bytes as received).
KindImage Kind = "image"
// KindAudio — raw PCM in the canonical internal/audio format, or a WAV
// container. Meeting captures and enrolment samples both land here.
KindAudio Kind = "audio"
)
// Valid reports whether k is a kind this package will store. An unknown kind is
// refused at Put rather than creating a stray directory.
func (k Kind) Valid() bool { return k == KindImage || k == KindAudio }
// Errors callers distinguish. ErrNotFound is the only one a caller usually
// handles; the rest mean the call was wrong.
var (
// ErrNotFound — no blob with that id in this store.
ErrNotFound = errors.New("media: not found")
// ErrEmpty — Put was handed zero bytes. Storing an empty capture would
// leave a sidecar claiming a recording exists when it does not.
ErrEmpty = errors.New("media: empty payload")
// ErrTooLarge — the payload is over the store's cap. The cap exists so a
// runaway capture cannot fill the disk that mavend's database lives on.
ErrTooLarge = errors.New("media: payload too large")
// ErrBadKind — unknown Kind.
ErrBadKind = errors.New("media: unknown kind")
// ErrBadID — the id is not a 64-char lowercase hex digest, so it cannot
// have come from this store and must not be turned into a path.
ErrBadID = errors.New("media: malformed id")
)
// Blob — one stored item. ID is the sha256 of the bytes in lowercase hex, which
// makes it both the primary key and the dedupe mechanism. Path is absolute and
// local; it is a debugging affordance and the argument a subprocess (whisper,
// llama-server) is pointed at, never something handed to a network client.
type Blob struct {
ID string `json:"id"`
Kind Kind `json:"kind"`
MIME string `json:"mime"`
Size int64 `json:"size"`
Source string `json:"source"` // provenance: "telegram", "web:upload", "capture:meeting", "enroll"
Created time.Time `json:"created"` // UTC
Path string `json:"-"` // filled by the store; not part of the sidecar
}
// Age is how long ago the blob was stored, measured against now. Prune uses it;
// it is exported because the /media surface will want to show it.
func (b Blob) Age(now time.Time) time.Duration { return now.Sub(b.Created) }
// String is a one-line summary for logs. Deliberately does not include Path:
// a log line is not the place to spell out where his meeting audio lives.
func (b Blob) String() string {
return fmt.Sprintf("%s %s %dB from %s", b.Kind, shortID(b.ID), b.Size, b.Source)
}
// shortID trims a digest to something readable in a log line. Twelve hex chars
// is unambiguous at personal scale and short enough to fit next to the rest.
func shortID(id string) string {
if len(id) <= 12 {
return id
}
return id[:12]
}
+363
View File
@@ -0,0 +1,363 @@
package media
import (
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io/fs"
"os"
"path/filepath"
"sort"
"strings"
"time"
)
// DefaultMaxBytes — the per-blob cap when a store is built without one. 64 MiB
// is about an hour of 16 kHz mono PCM, which is also the hearing capture's own
// ceiling; a single item bigger than that is a mistake, not a meeting.
const DefaultMaxBytes int64 = 64 << 20
// DefaultRetention — how long a blob is kept when no retention is configured.
// Seven days is long enough to re-run a transcription that came out wrong and
// short enough that "she has a month of my meetings on disk" is never true.
const DefaultRetention = 7 * 24 * time.Hour
// Store — a content-addressed blob directory. Zero value is not usable; build
// one with Open, which creates the directory 0700. The store holds no lock and
// no cache: every operation is a filesystem call, and two writers of the same
// bytes produce the same file, so concurrent Puts do not need coordinating.
type Store struct {
dir string
maxBytes int64
retention time.Duration
now func() time.Time
}
// Open prepares a blob store rooted at dir. maxBytes ≤ 0 ⇒ DefaultMaxBytes;
// retention ≤ 0 ⇒ DefaultRetention. The directory (and every kind subdirectory
// created later) is 0700: these are recordings of people, and the daemon's user
// is the only reader.
func Open(dir string, maxBytes int64, retention time.Duration) (*Store, error) {
if strings.TrimSpace(dir) == "" {
return nil, errors.New("media: empty dir")
}
abs, err := filepath.Abs(dir)
if err != nil {
return nil, fmt.Errorf("media: resolve dir: %w", err)
}
if err := os.MkdirAll(abs, 0o700); err != nil {
return nil, fmt.Errorf("media: create dir: %w", err)
}
if maxBytes <= 0 {
maxBytes = DefaultMaxBytes
}
if retention <= 0 {
retention = DefaultRetention
}
return &Store{dir: abs, maxBytes: maxBytes, retention: retention, now: time.Now}, nil
}
// Dir is the store root. Exported for logs and for pointing a subprocess at a
// path under it.
func (s *Store) Dir() string { return s.dir }
// Retention is the configured age limit Prune enforces.
func (s *Store) Retention() time.Duration { return s.retention }
// Put stores data and returns its Blob. The id is the sha256 of data, so
// storing the same bytes twice is idempotent: the second call rewrites the
// sidecar (keeping the ORIGINAL creation time, so a re-send cannot extend
// retention indefinitely) and returns the same id.
//
// mime is recorded as given and used only to pick a file extension; nothing
// dispatches on it. Callers that need the mime to be trustworthy sniff it
// first — see SniffImage.
func (s *Store) Put(kind Kind, mime, source string, data []byte) (Blob, error) {
if !kind.Valid() {
return Blob{}, ErrBadKind
}
if len(data) == 0 {
return Blob{}, ErrEmpty
}
if int64(len(data)) > s.maxBytes {
return Blob{}, fmt.Errorf("%w: %d > %d", ErrTooLarge, len(data), s.maxBytes)
}
sum := sha256.Sum256(data)
id := hex.EncodeToString(sum[:])
blobPath, metaPath, err := s.paths(kind, id, mime)
if err != nil {
return Blob{}, err
}
if err := os.MkdirAll(filepath.Dir(blobPath), 0o700); err != nil {
return Blob{}, fmt.Errorf("media: create bucket: %w", err)
}
b := Blob{ID: id, Kind: kind, MIME: mime, Size: int64(len(data)), Source: source,
Created: s.now().UTC(), Path: blobPath}
// A blob already here keeps its first-seen time. Re-sending the same photo
// every hour must not keep it alive past retention.
if prev, err := readMeta(metaPath); err == nil && !prev.Created.IsZero() {
b.Created = prev.Created
}
if err := writeFile(blobPath, data); err != nil {
return Blob{}, err
}
if err := writeMeta(metaPath, b); err != nil {
return Blob{}, err
}
return b, nil
}
// Get returns the blob's metadata without reading its bytes.
func (s *Store) Get(id string) (Blob, error) {
if !validID(id) {
return Blob{}, ErrBadID
}
for _, kind := range []Kind{KindImage, KindAudio} {
metaPath := filepath.Join(s.dir, string(kind), id[:2], id+".json")
b, err := readMeta(metaPath)
if err != nil {
continue
}
p, err := s.locate(kind, id)
if err != nil {
continue
}
b.Path = p
return b, nil
}
return Blob{}, ErrNotFound
}
// Read returns the blob's bytes together with its metadata. This is the only
// way out of the store, and it is a local read: nothing in this package can
// send bytes anywhere.
func (s *Store) Read(id string) (Blob, []byte, error) {
b, err := s.Get(id)
if err != nil {
return Blob{}, nil, err
}
data, err := os.ReadFile(b.Path)
if err != nil {
return Blob{}, nil, fmt.Errorf("media: read %s: %w", shortID(id), err)
}
return b, data, nil
}
// List returns every blob of the given kind, newest first. An empty kind lists
// both. It walks the directory; at personal volumes (tens to hundreds of items
// inside the retention window) that is cheap, and it means the sidecars are the
// single source of truth with no index to fall out of sync.
func (s *Store) List(kind Kind) ([]Blob, error) {
kinds := []Kind{KindImage, KindAudio}
if kind != "" {
if !kind.Valid() {
return nil, ErrBadKind
}
kinds = []Kind{kind}
}
var out []Blob
for _, k := range kinds {
root := filepath.Join(s.dir, string(k))
err := filepath.WalkDir(root, func(path string, d fs.DirEntry, err error) error {
if err != nil {
if errors.Is(err, fs.ErrNotExist) {
return nil // kind never used; not an error
}
return err
}
if d.IsDir() || !strings.HasSuffix(path, ".json") {
return nil
}
b, err := readMeta(path)
if err != nil {
return nil // a corrupt sidecar is skipped, not fatal
}
if p, err := s.locate(b.Kind, b.ID); err == nil {
b.Path = p
}
out = append(out, b)
return nil
})
if err != nil {
return nil, fmt.Errorf("media: list %s: %w", k, err)
}
}
sort.Slice(out, func(i, j int) bool {
if out[i].Created.Equal(out[j].Created) {
return out[i].ID < out[j].ID
}
return out[i].Created.After(out[j].Created)
})
return out, nil
}
// Delete removes a blob and its sidecar. Missing is not an error: the caller
// asked for it gone and it is gone.
func (s *Store) Delete(id string) error {
if !validID(id) {
return ErrBadID
}
for _, kind := range []Kind{KindImage, KindAudio} {
bucket := filepath.Join(s.dir, string(kind), id[:2])
entries, err := os.ReadDir(bucket)
if err != nil {
continue
}
for _, e := range entries {
if strings.HasPrefix(e.Name(), id) {
if err := os.Remove(filepath.Join(bucket, e.Name())); err != nil && !errors.Is(err, fs.ErrNotExist) {
return fmt.Errorf("media: delete %s: %w", shortID(id), err)
}
}
}
}
return nil
}
// Prune deletes every blob older than the store's retention and reports how
// many went. It is the enforcement half of the retention promise; a caller that
// never runs it has a store that grows without bound, which is why the daemon
// runs it on the digestion tick rather than leaving it to a cron the operator
// might not add.
func (s *Store) Prune() (int, error) {
blobs, err := s.List("")
if err != nil {
return 0, err
}
now := s.now()
deleted := 0
for _, b := range blobs {
if b.Age(now) <= s.retention {
continue
}
if err := s.Delete(b.ID); err != nil {
return deleted, err
}
deleted++
}
return deleted, nil
}
// paths returns the blob and sidecar paths for an id.
func (s *Store) paths(kind Kind, id, mime string) (blobPath, metaPath string, err error) {
if !validID(id) {
return "", "", ErrBadID
}
bucket := filepath.Join(s.dir, string(kind), id[:2])
return filepath.Join(bucket, id+extFor(mime, kind)), filepath.Join(bucket, id+".json"), nil
}
// locate finds the stored bytes for an id whose extension we do not know,
// because the extension came from the mime at Put time.
func (s *Store) locate(kind Kind, id string) (string, error) {
if !validID(id) {
return "", ErrBadID
}
bucket := filepath.Join(s.dir, string(kind), id[:2])
entries, err := os.ReadDir(bucket)
if err != nil {
return "", ErrNotFound
}
for _, e := range entries {
name := e.Name()
if strings.HasPrefix(name, id) && !strings.HasSuffix(name, ".json") {
return filepath.Join(bucket, name), nil
}
}
return "", ErrNotFound
}
// validID guards every path built from an id. Without it a caller-supplied id
// is a path traversal: Get("../../etc/passwd") would read outside the store.
func validID(id string) bool {
if len(id) != 64 {
return false
}
for i := 0; i < len(id); i++ {
c := id[i]
if (c < '0' || c > '9') && (c < 'a' || c > 'f') {
return false
}
}
return true
}
// extFor maps a mime to a file extension, defaulting per kind. The extension is
// cosmetic — the id is the key — but it is what makes the store browsable and
// lets a subprocess that sniffs by name (piper, some image tools) cope.
func extFor(mime string, kind Kind) string {
switch strings.ToLower(strings.TrimSpace(mime)) {
case "image/jpeg", "image/jpg":
return ".jpg"
case "image/png":
return ".png"
case "image/gif":
return ".gif"
case "image/webp":
return ".webp"
case "audio/wav", "audio/x-wav", "audio/wave":
return ".wav"
case "audio/l16", "audio/pcm":
return ".pcm"
}
if kind == KindImage {
return ".bin"
}
return ".pcm"
}
// writeFile writes data 0600 via a temp file in the same directory, so a
// crash mid-write cannot leave a truncated blob under a digest that claims
// to describe the whole thing.
func writeFile(path string, data []byte) error {
tmp, err := os.CreateTemp(filepath.Dir(path), ".tmp-*")
if err != nil {
return fmt.Errorf("media: temp: %w", err)
}
defer os.Remove(tmp.Name())
if err := tmp.Chmod(0o600); err != nil {
tmp.Close()
return fmt.Errorf("media: chmod: %w", err)
}
if _, err := tmp.Write(data); err != nil {
tmp.Close()
return fmt.Errorf("media: write: %w", err)
}
if err := tmp.Close(); err != nil {
return fmt.Errorf("media: close: %w", err)
}
if err := os.Rename(tmp.Name(), path); err != nil {
return fmt.Errorf("media: rename: %w", err)
}
return nil
}
func writeMeta(path string, b Blob) error {
data, err := json.Marshal(b)
if err != nil {
return fmt.Errorf("media: marshal meta: %w", err)
}
return writeFile(path, data)
}
func readMeta(path string) (Blob, error) {
data, err := os.ReadFile(path)
if err != nil {
return Blob{}, err
}
var b Blob
if err := json.Unmarshal(data, &b); err != nil {
return Blob{}, err
}
if !validID(b.ID) || !b.Kind.Valid() {
return Blob{}, errors.New("media: corrupt sidecar")
}
b.Created = b.Created.UTC()
return b, nil
}
+214
View File
@@ -0,0 +1,214 @@
package media
import (
"errors"
"os"
"path/filepath"
"strings"
"testing"
"time"
)
func testStore(t *testing.T) *Store {
t.Helper()
s, err := Open(t.TempDir(), 0, 0)
if err != nil {
t.Fatalf("open: %v", err)
}
return s
}
func TestPutAndRead(t *testing.T) {
s := testStore(t)
b, err := s.Put(KindImage, "image/png", "web:upload", []byte("pretend png"))
if err != nil {
t.Fatalf("put: %v", err)
}
if len(b.ID) != 64 {
t.Fatalf("id is not a sha256 hex digest: %q", b.ID)
}
if b.Size != int64(len("pretend png")) {
t.Errorf("size = %d", b.Size)
}
if !strings.HasSuffix(b.Path, ".png") {
t.Errorf("extension not taken from mime: %s", b.Path)
}
got, data, err := s.Read(b.ID)
if err != nil {
t.Fatalf("read: %v", err)
}
if string(data) != "pretend png" {
t.Errorf("data = %q", data)
}
if got.Source != "web:upload" || got.Kind != KindImage {
t.Errorf("metadata not round-tripped: %+v", got)
}
}
// The same bytes twice must be one file, and must NOT get a fresh creation
// time — otherwise re-sending a photo keeps it alive past retention forever.
func TestPutIsIdempotentAndKeepsFirstSeenTime(t *testing.T) {
s := testStore(t)
base := time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC)
s.now = func() time.Time { return base }
first, err := s.Put(KindAudio, "audio/wav", "capture:meeting", []byte("pcm"))
if err != nil {
t.Fatalf("put: %v", err)
}
s.now = func() time.Time { return base.Add(72 * time.Hour) }
second, err := s.Put(KindAudio, "audio/wav", "capture:meeting", []byte("pcm"))
if err != nil {
t.Fatalf("re-put: %v", err)
}
if first.ID != second.ID {
t.Fatalf("same bytes produced two ids")
}
if !second.Created.Equal(base) {
t.Errorf("re-put moved created time to %v, want %v", second.Created, base)
}
list, err := s.List(KindAudio)
if err != nil {
t.Fatalf("list: %v", err)
}
if len(list) != 1 {
t.Errorf("got %d blobs, want 1", len(list))
}
}
func TestPutRejects(t *testing.T) {
s, err := Open(t.TempDir(), 8, 0)
if err != nil {
t.Fatal(err)
}
if _, err := s.Put(KindImage, "image/png", "x", nil); !errors.Is(err, ErrEmpty) {
t.Errorf("empty payload: %v", err)
}
if _, err := s.Put("video", "video/mp4", "x", []byte("ab")); !errors.Is(err, ErrBadKind) {
t.Errorf("bad kind: %v", err)
}
if _, err := s.Put(KindImage, "image/png", "x", []byte("way too many bytes")); !errors.Is(err, ErrTooLarge) {
t.Errorf("over cap: %v", err)
}
}
// A caller-supplied id becomes a path, so a traversal attempt must be refused
// before it touches the filesystem rather than escaping the store root.
func TestMalformedIDIsRefused(t *testing.T) {
s := testStore(t)
for _, id := range []string{"", "../../etc/passwd", strings.Repeat("z", 64), strings.Repeat("a", 63)} {
if _, err := s.Get(id); !errors.Is(err, ErrBadID) && !errors.Is(err, ErrNotFound) {
t.Errorf("Get(%q) = %v, want a refusal", id, err)
}
if _, _, err := s.Read(id); err == nil {
t.Errorf("Read(%q) succeeded", id)
}
if err := s.Delete(id); err == nil && id != "" {
// Delete of a well-formed but absent id is fine; these are not
// well-formed.
t.Errorf("Delete(%q) succeeded", id)
}
}
}
func TestGetMissingIsNotFound(t *testing.T) {
s := testStore(t)
if _, err := s.Get(strings.Repeat("a", 64)); !errors.Is(err, ErrNotFound) {
t.Errorf("got %v, want ErrNotFound", err)
}
}
func TestPruneEnforcesRetention(t *testing.T) {
s, err := Open(t.TempDir(), 0, 48*time.Hour)
if err != nil {
t.Fatal(err)
}
now := time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC)
s.now = func() time.Time { return now.Add(-96 * time.Hour) }
old, _ := s.Put(KindAudio, "audio/wav", "capture:meeting", []byte("old meeting"))
s.now = func() time.Time { return now.Add(-1 * time.Hour) }
fresh, _ := s.Put(KindImage, "image/png", "telegram", []byte("recent photo"))
s.now = func() time.Time { return now }
n, err := s.Prune()
if err != nil {
t.Fatalf("prune: %v", err)
}
if n != 1 {
t.Errorf("pruned %d, want 1", n)
}
if _, err := s.Get(old.ID); !errors.Is(err, ErrNotFound) {
t.Errorf("stale blob survived prune: %v", err)
}
if _, err := s.Get(fresh.ID); err != nil {
t.Errorf("fresh blob was pruned: %v", err)
}
}
func TestListIsNewestFirstAcrossKinds(t *testing.T) {
s := testStore(t)
base := time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC)
s.now = func() time.Time { return base }
_, _ = s.Put(KindImage, "image/png", "telegram", []byte("one"))
s.now = func() time.Time { return base.Add(time.Hour) }
newest, _ := s.Put(KindAudio, "audio/wav", "capture:meeting", []byte("two"))
all, err := s.List("")
if err != nil {
t.Fatalf("list: %v", err)
}
if len(all) != 2 {
t.Fatalf("got %d, want 2", len(all))
}
if all[0].ID != newest.ID {
t.Errorf("list is not newest-first")
}
}
// Recordings of people are 0700/0600 and nothing else.
func TestPermissionsAreOwnerOnly(t *testing.T) {
dir := t.TempDir()
s, err := Open(filepath.Join(dir, "blobs"), 0, 0)
if err != nil {
t.Fatal(err)
}
b, err := s.Put(KindAudio, "audio/wav", "capture:meeting", []byte("pcm"))
if err != nil {
t.Fatal(err)
}
di, err := os.Stat(s.Dir())
if err != nil {
t.Fatal(err)
}
if di.Mode().Perm() != 0o700 {
t.Errorf("store dir mode = %o, want 700", di.Mode().Perm())
}
fi, err := os.Stat(b.Path)
if err != nil {
t.Fatal(err)
}
if fi.Mode().Perm() != 0o600 {
t.Errorf("blob mode = %o, want 600", fi.Mode().Perm())
}
}
func TestDeleteRemovesBytesAndSidecar(t *testing.T) {
s := testStore(t)
b, _ := s.Put(KindImage, "image/png", "telegram", []byte("bytes"))
if err := s.Delete(b.ID); err != nil {
t.Fatalf("delete: %v", err)
}
if _, err := os.Stat(b.Path); !os.IsNotExist(err) {
t.Errorf("bytes survived delete")
}
if _, err := s.Get(b.ID); !errors.Is(err, ErrNotFound) {
t.Errorf("sidecar survived delete: %v", err)
}
}
func TestOpenRejectsEmptyDir(t *testing.T) {
if _, err := Open(" ", 0, 0); err == nil {
t.Error("empty dir accepted")
}
}
+39
View File
@@ -56,6 +56,45 @@ func (s *Store) ProposeTool(ctx context.Context, name, utterance, scope string,
return n > 0, nil
}
// ProposeMCPTool is ProposeTool for a tool discovered on an MCP server
// (Vikunja #251): the proposal already knows what it would run, so cmd and
// destructive are written with it and Kami only has to press enable.
//
// It is still a PROPOSAL. Discovery cannot grant a capability — that is the
// whole reason a server can be configured without its tools becoming live.
// Like ProposeTool it never touches an existing row, so re-discovery on every
// restart is idempotent and cannot silently re-arm a tool that was disabled or
// change the cmd of one already enabled.
func (s *Store) ProposeMCPTool(ctx context.Context, name, scope string, cmd []string, destructive bool, utterance string, ts time.Time) (bool, error) {
if len(cmd) == 0 {
return false, ErrToolCmd
}
if scope == "" {
scope = "homelab"
}
raw, err := json.Marshal(cmd)
if err != nil {
return false, fmt.Errorf("propose mcp tool: %w", err)
}
d := 0
if destructive {
d = 1
}
res, err := s.db.ExecContext(ctx, `
INSERT INTO tools (name, scope, cmd, destructive, status, utterance, created_ts, updated_ts)
VALUES (?, ?, ?, ?, 'proposed', ?, ?, ?)
ON CONFLICT(name) DO NOTHING`,
name, scope, string(raw), d, utterance, ts.UnixMilli(), ts.UnixMilli())
if err != nil {
return false, fmt.Errorf("propose mcp tool: %w", err)
}
n, err := res.RowsAffected()
if err != nil {
return false, fmt.Errorf("propose mcp tool: rows affected: %w", err)
}
return n > 0, nil
}
// EnableTool fills cmd + destructive and flips status to 'enabled'. This is the
// human "enable" act (the authed surface calls it); it upserts so enabling a
// name that was never proposed still works. An empty cmd is refused — an
+62
View File
@@ -2,6 +2,7 @@ package store
import (
"context"
"errors"
"testing"
"time"
)
@@ -60,3 +61,64 @@ func TestToolLifecycle(t *testing.T) {
t.Fatalf("disable absent must be no-op: %v", err)
}
}
// A discovered MCP tool arrives as a proposal that already knows its cmd, so
// enabling it is one click rather than one retyped argv.
func TestProposeMCPTool(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
now := time.Now()
cmd := []string{"mcp", "vikunja", "list_tasks"}
fresh, err := s.ProposeMCPTool(ctx, "vikunja_list_tasks", "mcp:vikunja", cmd, false, "mcp vikunja/list_tasks: List tasks", now)
if err != nil {
t.Fatal(err)
}
if !fresh {
t.Fatal("first proposal should be new")
}
got, err := s.LookupTool(ctx, "vikunja_list_tasks")
if err != nil {
t.Fatal(err)
}
if got.Status != "proposed" {
t.Fatalf("status = %q — discovery must never enable", got.Status)
}
if len(got.Cmd) != 3 || got.Cmd[0] != "mcp" || got.Cmd[2] != "list_tasks" {
t.Fatalf("cmd = %v", got.Cmd)
}
if got.Scope != "mcp:vikunja" || got.Utterance == "" {
t.Fatalf("provenance lost: %+v", got)
}
// Re-discovery on the next boot is idempotent.
fresh, err = s.ProposeMCPTool(ctx, "vikunja_list_tasks", "mcp:vikunja", cmd, true, "changed", now)
if err != nil {
t.Fatal(err)
}
if fresh {
t.Error("re-proposing an existing row must report nothing new")
}
// And it must not re-arm or rewrite a row a human already acted on.
if err := s.EnableTool(ctx, "vikunja_list_tasks", cmd, false, "mcp:vikunja", now); err != nil {
t.Fatal(err)
}
if _, err := s.ProposeMCPTool(ctx, "vikunja_list_tasks", "mcp:vikunja", []string{"mcp", "vikunja", "delete_task"}, true, "x", now); err != nil {
t.Fatal(err)
}
got, err = s.LookupTool(ctx, "vikunja_list_tasks")
if err != nil {
t.Fatal(err)
}
if got.Status != "enabled" || got.Cmd[2] != "list_tasks" || got.Destructive {
t.Fatalf("an enabled row was modified by discovery: %+v", got)
}
}
func TestProposeMCPToolNeedsCmd(t *testing.T) {
s := newTestStore(t)
if _, err := s.ProposeMCPTool(context.Background(), "x", "mcp:y", nil, false, "", time.Now()); !errors.Is(err, ErrToolCmd) {
t.Fatalf("err = %v, want ErrToolCmd", err)
}
}
+33
View File
@@ -13,6 +13,11 @@
// - Args are passed as argv, NEVER through a shell. STT text lands as
// positional arguments to Cmd; there is no `sh -c`, so "restart nginx;
// rm -rf" can't inject — the tail is one argv element to the named binary.
// - An enabled row whose cmd is ["mcp", "<server>", "<tool>"] is a call to a
// configured MCP server instead of a process (Vikunja #251). It goes
// through every rule above unchanged — enabled, and confirmed if it
// mutates — because the store is still the allowlist; only the dispatch at
// the bottom of Exec differs.
// - Destructive tools don't run on first hearing: Exec returns ErrNeedsConfirm
// and the handler runs a confirm turn ("выполнить X? да/нет"); only a
// confirmed re-Exec runs them. A gate assumes a fully-formed action, which
@@ -31,6 +36,7 @@ import (
"time"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/mcp"
"github.com/kami/maven/internal/router"
)
@@ -50,12 +56,21 @@ var (
ErrNeedsConfirm = errors.New("destructive tool needs confirmation")
)
// MCPCaller is the seam for an act that is an MCP tool call rather than a
// process (Vikunja #251). internal/mcp.Manager satisfies it via CallPositional.
// nil ⇒ MCP is not configured, and an MCP row refuses to run rather than
// silently doing nothing.
type MCPCaller interface {
CallPositional(ctx context.Context, server, tool string, args []string) (string, error)
}
// Executor runs enabled tools. run is the exec seam (default: real process);
// tests swap it. timeout bounds each invocation.
type Executor struct {
api API
timeout time.Duration
run func(ctx context.Context, argv []string) (string, error)
mcp MCPCaller
}
// NewExecutor builds the executor. timeout<=0 defaults to 30s.
@@ -66,6 +81,13 @@ func NewExecutor(api API, timeout time.Duration) *Executor {
return &Executor{api: api, timeout: timeout, run: runProcess}
}
// WithMCP attaches the MCP caller. Called once at wiring time when the mcp
// config block is present; without it, a row whose cmd is ["mcp", …] refuses.
func (e *Executor) WithMCP(m MCPCaller) *Executor {
e.mcp = m
return e
}
// Exec looks up name in the store and runs Cmd+args as argv (no shell).
// confirmed=true is the second turn of a destructive act (the user said "да");
// it bypasses the ErrNeedsConfirm gate. Non-enabled ⇒ ErrNotEnabled; a
@@ -84,6 +106,17 @@ func (e *Executor) Exec(ctx context.Context, name string, args []string, confirm
if t.Destructive && !confirmed {
return "", ErrNeedsConfirm
}
// An MCP row is a call to a configured server, not a process. Everything
// above still applied: it had to be enabled, and a mutating one had to be
// confirmed. Only the dispatch differs.
if server, remote, ok := mcp.ParseCmd(t.Cmd); ok {
if e.mcp == nil {
return "", ErrNotEnabled
}
ctx, cancel := context.WithTimeout(ctx, e.timeout)
defer cancel()
return e.mcp.CallPositional(ctx, server, remote, args)
}
argv := append(append([]string(nil), t.Cmd...), args...)
if len(argv) == 0 {
return "", ErrNotEnabled
+100
View File
@@ -85,3 +85,103 @@ func TestExec(t *testing.T) {
t.Fatal("proposed tool must not match (not enabled)")
}
}
// fakeMCP records what the executor asked it to call.
type fakeMCP struct {
server, tool string
args []string
out string
err error
calls int
}
func (f *fakeMCP) CallPositional(_ context.Context, server, tool string, args []string) (string, error) {
f.calls++
f.server, f.tool, f.args = server, tool, args
return f.out, f.err
}
// An MCP row dispatches to the caller instead of a process, and the process
// seam is never touched.
func TestExecMCPRowDispatchesToMCP(t *testing.T) {
api := fakeAPI{tools: map[string]ipc.Tool{
"vikunja_list_tasks": {
Name: "vikunja_list_tasks", Status: "enabled", Scope: "mcp:vikunja",
Cmd: []string{"mcp", "vikunja", "list_tasks"},
},
}}
m := &fakeMCP{out: "две задачи"}
ran := false
e := NewExecutor(api, time.Second).WithMCP(m)
e.run = func(context.Context, []string) (string, error) { ran = true; return "", nil }
out, err := e.Exec(context.Background(), "vikunja_list_tasks", []string{"мавен"}, false)
if err != nil {
t.Fatalf("exec: %v", err)
}
if out != "две задачи" {
t.Fatalf("out = %q", out)
}
if ran {
t.Fatal("an MCP row must not be executed as a process")
}
if m.server != "vikunja" || m.tool != "list_tasks" || len(m.args) != 1 || m.args[0] != "мавен" {
t.Fatalf("dispatched wrong: %+v", m)
}
}
// The allowlist rules still apply to an MCP row: destructive means a confirm
// turn first, and nothing is called until the second turn.
func TestExecMCPRowStillNeedsConfirm(t *testing.T) {
api := fakeAPI{tools: map[string]ipc.Tool{
"vikunja_delete_task": {
Name: "vikunja_delete_task", Status: "enabled", Destructive: true,
Cmd: []string{"mcp", "vikunja", "delete_task"},
},
}}
m := &fakeMCP{out: "удалила"}
e := NewExecutor(api, time.Second).WithMCP(m)
if _, err := e.Exec(context.Background(), "vikunja_delete_task", nil, false); !errors.Is(err, ErrNeedsConfirm) {
t.Fatalf("err = %v, want ErrNeedsConfirm", err)
}
if m.calls != 0 {
t.Fatal("a destructive MCP tool must not reach the server before confirmation")
}
if _, err := e.Exec(context.Background(), "vikunja_delete_task", nil, true); err != nil {
t.Fatalf("confirmed exec: %v", err)
}
if m.calls != 1 {
t.Fatalf("calls = %d", m.calls)
}
}
// A proposed MCP row does not run, exactly like a proposed shell tool.
func TestExecMCPRowNotEnabled(t *testing.T) {
api := fakeAPI{tools: map[string]ipc.Tool{
"vikunja_list_tasks": {Name: "vikunja_list_tasks", Status: "proposed", Cmd: []string{"mcp", "vikunja", "list_tasks"}},
}}
m := &fakeMCP{}
e := NewExecutor(api, time.Second).WithMCP(m)
if _, err := e.Exec(context.Background(), "vikunja_list_tasks", nil, false); !errors.Is(err, ErrNotEnabled) {
t.Fatalf("err = %v", err)
}
if m.calls != 0 {
t.Fatal("a proposal must not call anything")
}
}
// With MCP unconfigured, an MCP row refuses rather than trying to exec "mcp".
func TestExecMCPRowWithoutCallerRefuses(t *testing.T) {
api := fakeAPI{tools: map[string]ipc.Tool{
"vikunja_list_tasks": {Name: "vikunja_list_tasks", Status: "enabled", Cmd: []string{"mcp", "vikunja", "list_tasks"}},
}}
ran := false
e := NewExecutor(api, time.Second)
e.run = func(context.Context, []string) (string, error) { ran = true; return "", nil }
if _, err := e.Exec(context.Background(), "vikunja_list_tasks", nil, false); !errors.Is(err, ErrNotEnabled) {
t.Fatalf("err = %v, want ErrNotEnabled", err)
}
if ran {
t.Fatal(`"mcp" must never be run as a binary`)
}
}
+111
View File
@@ -0,0 +1,111 @@
package vision
import (
"context"
"fmt"
"strings"
"github.com/kami/maven/internal/media"
)
// Intake is the whole path from "bytes arrived" to "here is what she saw",
// in one place, so that every surface that can receive an image — a Telegram
// photo, a mavweb upload, a file path he names — goes through the same steps in
// the same order:
//
// 1. sniff the bytes (the sender's declared content type is not trusted);
// 2. store them content-addressed, so the same photo twice is one file and the
// original is still on disk if the description came out wrong;
// 3. prepare a downscaled JPEG for the model;
// 4. describe it.
//
// Step 2 happens BEFORE step 4 deliberately. If the vision model is missing or
// broken — which is today's actual state on this box — the image is still safely
// stored and describable later, and the failure is "I can't look at it yet", not
// "it's gone".
//
// Writing the description as a note is NOT done here. That needs the store and
// the embedder and belongs to the daemon; Intake returns the text and lets the
// caller decide whether it becomes a note, a reply, or both.
type Intake struct {
store *media.Store
provider Provider
maxDim int
}
// NewIntake wires an intake. provider may be Disabled — storing still works,
// which is the point. maxDim ≤ 0 ⇒ media.DefaultMaxDim.
func NewIntake(store *media.Store, provider Provider, maxDim int) *Intake {
if provider == nil {
provider = Disabled{}
}
return &Intake{store: store, provider: provider, maxDim: maxDim}
}
// Result — what an intake produced. Blob is always set when Store succeeded, so
// a caller that got an error from the description still knows what was kept and
// can retry against the same id later.
type Result struct {
Blob media.Blob
Image media.Image
Description string
}
// Accept stores data and describes it. source is provenance recorded on the
// blob ("telegram", "web:upload"); question is what he asked about the image, or
// empty for the default "what is this".
//
// A description failure is returned alongside a populated Result: the caller
// gets the blob id for the log and the reply, and the error to explain why there
// are no words yet.
func (in *Intake) Accept(ctx context.Context, data []byte, source, question string) (Result, error) {
if in == nil || in.store == nil {
return Result{}, fmt.Errorf("vision: intake not wired")
}
mime, err := media.SniffImage(data)
if err != nil {
return Result{}, err
}
blob, err := in.store.Put(media.KindImage, mime, source, data)
if err != nil {
return Result{}, err
}
im, err := media.PrepareImage(data, source, in.maxDim)
if err != nil {
return Result{Blob: blob}, err
}
res := Result{Blob: blob, Image: im}
text, err := in.provider.Describe(ctx, im, question)
if err != nil {
return res, err
}
res.Description = strings.TrimSpace(text)
return res, nil
}
// Rerun describes an already-stored image again — a different question, or the
// first successful attempt after the model finally landed on disk. It is the
// reason step 2 comes before step 4.
func (in *Intake) Rerun(ctx context.Context, id, question string) (Result, error) {
if in == nil || in.store == nil {
return Result{}, fmt.Errorf("vision: intake not wired")
}
blob, data, err := in.store.Read(id)
if err != nil {
return Result{}, err
}
if blob.Kind != media.KindImage {
return Result{Blob: blob}, fmt.Errorf("vision: %s is %s, not an image", id[:12], blob.Kind)
}
im, err := media.PrepareImage(data, blob.Source, in.maxDim)
if err != nil {
return Result{Blob: blob}, err
}
res := Result{Blob: blob, Image: im}
text, err := in.provider.Describe(ctx, im, question)
if err != nil {
return res, err
}
res.Description = strings.TrimSpace(text)
return res, nil
}
+147
View File
@@ -0,0 +1,147 @@
package vision
import (
"bytes"
"context"
"errors"
"image"
"image/png"
"testing"
"github.com/kami/maven/internal/media"
)
type fakeProvider struct {
reply string
err error
seen int
lastQ string
lastDim int
}
func (f *fakeProvider) Describe(_ context.Context, im media.Image, prompt string) (string, error) {
f.seen++
f.lastQ = prompt
f.lastDim = im.Width
return f.reply, f.err
}
func pngPayload(t *testing.T, w, h int) []byte {
t.Helper()
var buf bytes.Buffer
if err := png.Encode(&buf, image.NewRGBA(image.Rect(0, 0, w, h))); err != nil {
t.Fatal(err)
}
return buf.Bytes()
}
func testIntake(t *testing.T, p Provider) (*Intake, *media.Store) {
t.Helper()
s, err := media.Open(t.TempDir(), 0, 0)
if err != nil {
t.Fatal(err)
}
return NewIntake(s, p, 64), s
}
func TestAcceptStoresThenDescribes(t *testing.T) {
fp := &fakeProvider{reply: "кот на подоконнике"}
in, store := testIntake(t, fp)
res, err := in.Accept(context.Background(), pngPayload(t, 200, 100), "telegram", "кто это?")
if err != nil {
t.Fatalf("accept: %v", err)
}
if res.Description != "кот на подоконнике" {
t.Errorf("description = %q", res.Description)
}
if fp.lastQ != "кто это?" {
t.Errorf("question not passed through: %q", fp.lastQ)
}
if fp.lastDim != 64 {
t.Errorf("image not downscaled to maxDim: width %d", fp.lastDim)
}
// The sniffed mime wins over anything a sender claimed.
got, _, err := store.Read(res.Blob.ID)
if err != nil {
t.Fatalf("blob not stored: %v", err)
}
if got.MIME != "image/png" || got.Source != "telegram" {
t.Errorf("blob metadata = %+v", got)
}
}
// The ordering promise: with no vision model on the box — today's real state —
// the image is still on disk and the id is still reported, so it can be
// described later instead of being lost.
func TestAcceptKeepsBlobWhenDescribeFails(t *testing.T) {
in, store := testIntake(t, Disabled{})
res, err := in.Accept(context.Background(), pngPayload(t, 32, 32), "web:upload", "")
if !errors.Is(err, ErrDisabled) {
t.Fatalf("got %v, want ErrDisabled", err)
}
if res.Blob.ID == "" {
t.Fatal("no blob id reported on a description failure")
}
if _, _, err := store.Read(res.Blob.ID); err != nil {
t.Errorf("blob was not kept: %v", err)
}
}
func TestRerunDescribesAStoredBlob(t *testing.T) {
fp := &fakeProvider{reply: "текст: ошибка E24"}
in, _ := testIntake(t, fp)
first, err := in.Accept(context.Background(), pngPayload(t, 40, 40), "telegram", "")
if err != nil {
t.Fatal(err)
}
res, err := in.Rerun(context.Background(), first.Blob.ID, "прочитай текст")
if err != nil {
t.Fatalf("rerun: %v", err)
}
if res.Description != "текст: ошибка E24" {
t.Errorf("description = %q", res.Description)
}
if fp.lastQ != "прочитай текст" {
t.Errorf("new question not used: %q", fp.lastQ)
}
if fp.seen != 2 {
t.Errorf("provider called %d times, want 2", fp.seen)
}
}
func TestRerunRefusesAudioBlob(t *testing.T) {
in, store := testIntake(t, &fakeProvider{reply: "x"})
b, err := store.Put(media.KindAudio, "audio/wav", "capture:meeting", []byte("pcm bytes"))
if err != nil {
t.Fatal(err)
}
if _, err := in.Rerun(context.Background(), b.ID, ""); err == nil {
t.Error("audio blob was accepted as an image")
}
}
func TestRerunUnknownID(t *testing.T) {
in, _ := testIntake(t, &fakeProvider{})
if _, err := in.Rerun(context.Background(), "nope", ""); err == nil {
t.Error("malformed id accepted")
}
}
func TestAcceptRefusesNonImage(t *testing.T) {
in, _ := testIntake(t, &fakeProvider{})
if _, err := in.Accept(context.Background(), []byte("this is a text file"), "web:upload", ""); !errors.Is(err, media.ErrUnsupportedImage) {
t.Errorf("got %v, want ErrUnsupportedImage", err)
}
}
func TestNilProviderDegradesToDisabled(t *testing.T) {
s, err := media.Open(t.TempDir(), 0, 0)
if err != nil {
t.Fatal(err)
}
in := NewIntake(s, nil, 0)
if _, err := in.Accept(context.Background(), pngPayload(t, 8, 8), "x", ""); !errors.Is(err, ErrDisabled) {
t.Errorf("got %v, want ErrDisabled", err)
}
}
+279
View File
@@ -0,0 +1,279 @@
// Package vision is Maven's image-understanding seam (Vikunja #252,
// docs/plans/07-vision.md).
//
// One interface, Provider, with one method: describe an image, in words, in
// Russian, with an optional question about it. Text extraction is not a second
// method — "прочитай текст с картинки" is a prompt, and a vision-language model
// does not have a separate OCR mode to select.
//
// # What is deliberately NOT here
//
// The plan document called for a `RemoteProvider` calling "an OpenAI-compatible
// vision API endpoint". That step is refused: CLAUDE.md's surviving hard
// constraint after "never phones home" was deprecated is *no cloud model,
// inference stays on the box*, and a photo of his flat is the single worst thing
// to make an exception for. Endpoint is therefore checked at construction and
// must be a loopback or private address — a public host is a config error, not a
// deployment option. That check is the reason this package does not simply reuse
// internal/llm.Client.
//
// # State on this box, honestly
//
// The resident model is Qwen3-1.7B, which is text-only, and as of 2026-08-01
// there is no vision-capable gguf and no mmproj file anywhere under
// /mnt/hdd1/llms. So LocalProvider is written, tested against a fake server, and
// currently has nothing real to talk to: the describing half is BLOCKED on a
// model download (see docs/plans/07-vision.md for the candidates and the
// recipe). What works today without any download is the intake — an image
// arrives, is stored, is prepared — and the config seam that turns the rest on.
//
// Provider is nil-safe through Disabled, and vision is OFF unless configured,
// like the weather and telegram.
package vision
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"net"
"net/http"
"net/url"
"strings"
"time"
"github.com/kami/maven/internal/media"
"github.com/kami/maven/internal/webfetch"
)
// DefaultTimeout — budget for one description. A small VLM doing prefill over
// an 896px image on a Vega iGPU is slow; 90s is generous because nobody is
// holding a conversation open on this path — the answer arrives as a reply or a
// note, and a too-tight timeout just means it never arrives at all.
const DefaultTimeout = 90 * time.Second
// DefaultMaxTokens — cap on the description. A paragraph is what a spoken
// answer can carry; a page is not.
const DefaultMaxTokens = 300
// DefaultPrompt — what she is asked when he did not ask anything specific,
// only sent a picture. Russian, because that is the channel language, and
// feminine self-reference is not needed here (the prompt is an instruction, the
// persona block is added by the caller that phrases the reply).
const DefaultPrompt = "Опиши, что на этом изображении. Коротко, 2-3 предложения. Если на нём есть текст, приведи его."
// Errors callers distinguish.
var (
// ErrDisabled — vision is not configured. Returned by Disabled, which is
// what the daemon wires when the config block is absent.
ErrDisabled = errors.New("vision: not configured")
// ErrNotPrivate — the configured endpoint is not on this box or its
// network. Refused at construction; see the package comment.
ErrNotPrivate = errors.New("vision: endpoint must be a local or private address")
// ErrEmptyReply — the model returned nothing usable.
ErrEmptyReply = errors.New("vision: empty description")
)
// Provider — the image-understanding contract. Describe takes an image already
// prepared by internal/media (decoded, downscaled, JPEG) and a prompt; an empty
// prompt means DefaultPrompt.
type Provider interface {
Describe(ctx context.Context, im media.Image, prompt string) (string, error)
}
// Disabled — the floor Provider. Every call fails with ErrDisabled, which the
// caller turns into "я не умею смотреть картинки — зрение не настроено". It
// exists so that no call site needs a nil check and switching vision off cannot
// crash a turn.
type Disabled struct{}
// Describe always fails. The signature matches Provider.
func (Disabled) Describe(context.Context, media.Image, string) (string, error) {
return "", ErrDisabled
}
// Config — how to reach the local vision server. Built from
// config.VisionConfig by the daemon; kept separate so this package does not
// import internal/config.
type Config struct {
// Endpoint — base URL of a llama-server started with a vision model and its
// mmproj (`llama-server -m model.gguf --mmproj mmproj.gguf`). Must be
// loopback or private. The path is appended by the provider; give it
// "http://127.0.0.1:8081".
Endpoint string
// Model — the model name to send. llama-server ignores it; it matters if the
// endpoint is something else OpenAI-shaped on the same box.
Model string
// Timeout — per-description budget. 0 ⇒ DefaultTimeout.
Timeout time.Duration
// MaxTokens — cap on the reply. 0 ⇒ DefaultMaxTokens.
MaxTokens int
// Prompt — the default question. Empty ⇒ DefaultPrompt.
Prompt string
}
// LocalProvider talks to a llama-server on this box over its
// /v1/chat/completions endpoint, sending the image as a data URI content part.
// It is the only real Provider, and it is a plain HTTP client: no subprocess
// spawning, because the daemon already owns llama-server lifecycle for the
// resident model and a second managed process is a bigger change than this task.
type LocalProvider struct {
endpoint string
model string
prompt string
maxTokens int
http *http.Client
}
// NewLocal builds a LocalProvider, refusing a non-private endpoint. A bad URL
// or a public host is an error at construction so the daemon logs it once at
// startup instead of failing every turn.
func NewLocal(cfg Config) (*LocalProvider, error) {
base := strings.TrimRight(strings.TrimSpace(cfg.Endpoint), "/")
if base == "" {
return nil, errors.New("vision: empty endpoint")
}
if err := checkPrivate(base); err != nil {
return nil, err
}
timeout := cfg.Timeout
if timeout <= 0 {
timeout = DefaultTimeout
}
maxTokens := cfg.MaxTokens
if maxTokens <= 0 {
maxTokens = DefaultMaxTokens
}
prompt := strings.TrimSpace(cfg.Prompt)
if prompt == "" {
prompt = DefaultPrompt
}
return &LocalProvider{
endpoint: base,
model: cfg.Model,
prompt: prompt,
maxTokens: maxTokens,
http: &http.Client{Timeout: timeout},
}, nil
}
// Endpoint is the server this provider talks to. For logs and /dash.
func (p *LocalProvider) Endpoint() string { return p.endpoint }
// checkPrivate refuses any endpoint that is not on this box or its LAN. A
// hostname that is not an IP literal is refused too: "vision.example.com" could
// resolve anywhere, and resolving it here would be trusting DNS with his photos.
// localhost is the one name allowed, because it is the common case.
func checkPrivate(raw string) error {
u, err := url.Parse(raw)
if err != nil {
return fmt.Errorf("vision: parse endpoint: %w", err)
}
if u.Scheme != "http" && u.Scheme != "https" {
return fmt.Errorf("vision: endpoint scheme %q not supported", u.Scheme)
}
host := u.Hostname()
if host == "" {
return errors.New("vision: endpoint has no host")
}
if strings.EqualFold(host, "localhost") {
return nil
}
ip := net.ParseIP(host)
if ip == nil {
return fmt.Errorf("%w: %q is a name, not an address", ErrNotPrivate, host)
}
if !webfetch.IsPrivateIP(ip) {
return fmt.Errorf("%w: %s", ErrNotPrivate, host)
}
return nil
}
// chat request shapes. Content is the OpenAI multimodal array form: a text part
// and an image_url part whose url is a data URI.
type textPart struct {
Type string `json:"type"`
Text string `json:"text"`
}
type imageURL struct {
URL string `json:"url"`
}
type imagePart struct {
Type string `json:"type"`
ImageURL imageURL `json:"image_url"`
}
type chatReq struct {
Model string `json:"model,omitempty"`
Messages []any `json:"messages"`
MaxTokens int `json:"max_tokens,omitempty"`
Temp float64 `json:"temperature"`
}
type userMsg struct {
Role string `json:"role"`
Content []any `json:"content"`
}
type chatResp struct {
Choices []struct {
Message struct {
Content string `json:"content"`
ReasoningContent string `json:"reasoning_content"`
} `json:"message"`
} `json:"choices"`
}
// Describe sends the image and prompt and returns the model's answer. An empty
// prompt uses the configured default. Errors are wrapped, never fatal: the
// caller says she could not make out the picture and the turn continues.
func (p *LocalProvider) Describe(ctx context.Context, im media.Image, prompt string) (string, error) {
if len(im.JPEG) == 0 {
return "", media.ErrEmpty
}
q := strings.TrimSpace(prompt)
if q == "" {
q = p.prompt
}
body, err := json.Marshal(chatReq{
Model: p.model,
MaxTokens: p.maxTokens,
Messages: []any{userMsg{Role: "user", Content: []any{
textPart{Type: "text", Text: q},
imagePart{Type: "image_url", ImageURL: imageURL{URL: im.DataURI()}},
}}},
})
if err != nil {
return "", fmt.Errorf("vision: marshal: %w", err)
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
p.endpoint+"/v1/chat/completions", bytes.NewReader(body))
if err != nil {
return "", fmt.Errorf("vision: request: %w", err)
}
req.Header.Set("Content-Type", "application/json")
resp, err := p.http.Do(req)
if err != nil {
return "", fmt.Errorf("vision: post: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return "", fmt.Errorf("vision: status %d", resp.StatusCode)
}
var out chatResp
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
return "", fmt.Errorf("vision: decode: %w", err)
}
if len(out.Choices) == 0 {
return "", ErrEmptyReply
}
text := strings.TrimSpace(out.Choices[0].Message.Content)
if text == "" {
// Same fallback as internal/llm: a Thinking model sometimes puts the
// whole answer in reasoning_content and leaves content empty.
text = strings.TrimSpace(out.Choices[0].Message.ReasoningContent)
}
if text == "" {
return "", ErrEmptyReply
}
return text, nil
}
+198
View File
@@ -0,0 +1,198 @@
package vision
import (
"bytes"
"context"
"encoding/json"
"errors"
"image"
"image/png"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/media"
)
func testImage(t *testing.T) media.Image {
t.Helper()
var buf bytes.Buffer
if err := png.Encode(&buf, image.NewRGBA(image.Rect(0, 0, 32, 32))); err != nil {
t.Fatal(err)
}
im, err := media.PrepareImage(buf.Bytes(), "test", 32)
if err != nil {
t.Fatal(err)
}
return im
}
func TestDisabledAlwaysRefuses(t *testing.T) {
_, err := Disabled{}.Describe(context.Background(), testImage(t), "что тут?")
if !errors.Is(err, ErrDisabled) {
t.Fatalf("got %v, want ErrDisabled", err)
}
}
// The whole reason this package has its own HTTP client instead of reusing
// internal/llm.Client: a vision endpoint that is not on this box is refused.
func TestNewLocalRefusesNonPrivateEndpoints(t *testing.T) {
bad := []string{
"https://api.openai.com",
"http://8.8.8.8:8080",
"https://vision.example.com", // a name could resolve anywhere
"ftp://127.0.0.1:8080", // wrong scheme
"", // nothing to talk to
}
for _, ep := range bad {
if _, err := NewLocal(Config{Endpoint: ep}); err == nil {
t.Errorf("NewLocal(%q) was accepted", ep)
}
}
}
func TestNewLocalAcceptsLocalEndpoints(t *testing.T) {
for _, ep := range []string{"http://127.0.0.1:8081", "http://localhost:8081/", "http://192.168.1.104:8081", "http://[::1]:8081"} {
p, err := NewLocal(Config{Endpoint: ep})
if err != nil {
t.Errorf("NewLocal(%q): %v", ep, err)
continue
}
if strings.HasSuffix(p.Endpoint(), "/") {
t.Errorf("trailing slash kept: %q", p.Endpoint())
}
}
}
func TestDescribeSendsImageAsDataURIAndReturnsText(t *testing.T) {
var gotBody map[string]any
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/v1/chat/completions" {
t.Errorf("path = %s", r.URL.Path)
}
raw, _ := io.ReadAll(r.Body)
if err := json.Unmarshal(raw, &gotBody); err != nil {
t.Errorf("unmarshal request: %v", err)
}
_, _ = w.Write([]byte(`{"choices":[{"message":{"content":" На картинке кот "}}]}`))
}))
defer srv.Close()
p, err := NewLocal(Config{Endpoint: srv.URL, Model: "qwen-vl"})
if err != nil {
t.Fatal(err)
}
text, err := p.Describe(context.Background(), testImage(t), "кто на фото?")
if err != nil {
t.Fatalf("describe: %v", err)
}
if text != "На картинке кот" {
t.Errorf("text = %q (should be trimmed)", text)
}
msgs, ok := gotBody["messages"].([]any)
if !ok || len(msgs) != 1 {
t.Fatalf("messages = %#v", gotBody["messages"])
}
parts, ok := msgs[0].(map[string]any)["content"].([]any)
if !ok || len(parts) != 2 {
t.Fatalf("content parts = %#v", msgs[0])
}
if got := parts[0].(map[string]any)["text"]; got != "кто на фото?" {
t.Errorf("prompt = %v", got)
}
url := parts[1].(map[string]any)["image_url"].(map[string]any)["url"].(string)
if !strings.HasPrefix(url, "data:image/jpeg;base64,") {
t.Errorf("image not sent as a jpeg data uri: %.40s", url)
}
}
func TestDescribeUsesDefaultPromptWhenNoQuestion(t *testing.T) {
var sentPrompt string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var body struct {
Messages []struct {
Content []struct {
Text string `json:"text"`
} `json:"content"`
} `json:"messages"`
}
_ = json.NewDecoder(r.Body).Decode(&body)
sentPrompt = body.Messages[0].Content[0].Text
_, _ = w.Write([]byte(`{"choices":[{"message":{"content":"ок"}}]}`))
}))
defer srv.Close()
p, err := NewLocal(Config{Endpoint: srv.URL, Prompt: "Опиши по-русски."})
if err != nil {
t.Fatal(err)
}
if _, err := p.Describe(context.Background(), testImage(t), " "); err != nil {
t.Fatal(err)
}
if sentPrompt != "Опиши по-русски." {
t.Errorf("prompt = %q", sentPrompt)
}
}
// A Thinking model sometimes leaves content empty and puts the answer in
// reasoning_content; internal/llm has the same fallback and vision needs it too.
func TestDescribeFallsBackToReasoningContent(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
_, _ = w.Write([]byte(`{"choices":[{"message":{"content":"","reasoning_content":"схема платы"}}]}`))
}))
defer srv.Close()
p, _ := NewLocal(Config{Endpoint: srv.URL})
text, err := p.Describe(context.Background(), testImage(t), "")
if err != nil {
t.Fatal(err)
}
if text != "схема платы" {
t.Errorf("text = %q", text)
}
}
func TestDescribeErrors(t *testing.T) {
t.Run("no choices", func(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
_, _ = w.Write([]byte(`{"choices":[]}`))
}))
defer srv.Close()
p, _ := NewLocal(Config{Endpoint: srv.URL})
if _, err := p.Describe(context.Background(), testImage(t), ""); !errors.Is(err, ErrEmptyReply) {
t.Errorf("got %v, want ErrEmptyReply", err)
}
})
t.Run("server error", func(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
}))
defer srv.Close()
p, _ := NewLocal(Config{Endpoint: srv.URL})
if _, err := p.Describe(context.Background(), testImage(t), ""); err == nil {
t.Error("500 was not an error")
}
})
t.Run("empty image", func(t *testing.T) {
p, _ := NewLocal(Config{Endpoint: "http://127.0.0.1:1"})
if _, err := p.Describe(context.Background(), media.Image{}, ""); !errors.Is(err, media.ErrEmpty) {
t.Errorf("got %v, want media.ErrEmpty", err)
}
})
t.Run("context cancelled", func(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
time.Sleep(200 * time.Millisecond)
_, _ = w.Write([]byte(`{"choices":[{"message":{"content":"поздно"}}]}`))
}))
defer srv.Close()
p, _ := NewLocal(Config{Endpoint: srv.URL})
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond)
defer cancel()
if _, err := p.Describe(ctx, testImage(t), ""); err == nil {
t.Error("cancelled context returned no error")
}
})
}
+53 -6
View File
@@ -28,6 +28,7 @@
package webfetch
import (
"bytes"
"context"
"errors"
"fmt"
@@ -92,6 +93,9 @@ type Response struct {
Status int
ContentType string
Body []byte
// Header — the response headers, one value each (the first). Populated for
// every request; MCP needs Mcp-Session-Id, nothing else reads it.
Header map[string]string
}
// Fetcher performs guarded GETs. Safe for concurrent use; the per-host rate
@@ -159,6 +163,37 @@ func New(cfg Config) *Fetcher {
// Get fetches rawURL. The body is capped: a larger response is an error, not a
// truncation, because half an XML document is worse than none.
func (f *Fetcher) Get(ctx context.Context, rawURL string) (*Response, error) {
return f.do(ctx, http.MethodGet, rawURL, nil, nil)
}
// Post sends body to rawURL and returns the reply, under exactly the same
// guards as Get: scheme rule, host lists, the dialer's private-address check on
// every hop, the size cap and the per-host rate limit.
//
// It exists for JSON-RPC over HTTP (internal/mcp), which cannot be expressed as
// a GET. That an outbound request now carries a body does not widen the
// address policy one bit — a POST to the LAN is refused for the same reason a
// GET is, unless AllowPrivate was set for that specific fetcher.
//
// hdr is merged over the defaults; a caller may not override User-Agent or
// Accept-Encoding, because identity encoding and an honest UA are part of the
// contract with whatever is on the other end.
func (f *Fetcher) Post(ctx context.Context, rawURL, contentType string, body []byte, hdr map[string]string) (*Response, error) {
if contentType == "" {
contentType = "application/json"
}
if hdr == nil {
hdr = map[string]string{}
}
merged := make(map[string]string, len(hdr)+1)
for k, v := range hdr {
merged[k] = v
}
merged["Content-Type"] = contentType
return f.do(ctx, http.MethodPost, rawURL, body, merged)
}
func (f *Fetcher) do(ctx context.Context, method, rawURL string, body []byte, hdr map[string]string) (*Response, error) {
u, err := url.Parse(strings.TrimSpace(rawURL))
if err != nil {
return nil, fmt.Errorf("webfetch: bad url %q: %w", rawURL, err)
@@ -170,10 +205,17 @@ func (f *Fetcher) Get(ctx context.Context, rawURL string) (*Response, error) {
return nil, err
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u.String(), nil)
var rdr io.Reader
if body != nil {
rdr = bytes.NewReader(body)
}
req, err := http.NewRequestWithContext(ctx, method, u.String(), rdr)
if err != nil {
return nil, err
}
for k, v := range hdr {
req.Header.Set(k, v)
}
req.Header.Set("User-Agent", f.cfg.UserAgent)
req.Header.Set("Accept-Encoding", "identity")
@@ -190,22 +232,27 @@ func (f *Fetcher) Get(ctx context.Context, rawURL string) (*Response, error) {
}
defer resp.Body.Close()
body, err := io.ReadAll(io.LimitReader(resp.Body, f.cfg.MaxBytes+1))
respBody, err := io.ReadAll(io.LimitReader(resp.Body, f.cfg.MaxBytes+1))
if err != nil {
return nil, err
}
if int64(len(body)) > f.cfg.MaxBytes {
if int64(len(respBody)) > f.cfg.MaxBytes {
return nil, fmt.Errorf("%w (%d bytes)", ErrTooLarge, f.cfg.MaxBytes)
}
if resp.StatusCode < 200 || resp.StatusCode > 299 {
return nil, fmt.Errorf("%w: %d", ErrStatus, resp.StatusCode)
}
return &Response{
out := &Response{
URL: resp.Request.URL.String(),
Status: resp.StatusCode,
ContentType: resp.Header.Get("Content-Type"),
Body: body,
}, nil
Body: respBody,
Header: map[string]string{},
}
for k := range resp.Header {
out.Header[k] = resp.Header.Get(k)
}
return out, nil
}
// checkURL applies the scheme rule and the host lists. The address rule is the
+72
View File
@@ -3,6 +3,7 @@ package webfetch
import (
"context"
"errors"
"io"
"net"
"net/http"
"net/http/httptest"
@@ -225,3 +226,74 @@ func TestUserAgentIsSent(t *testing.T) {
t.Fatalf("user-agent = %q", ua)
}
}
func TestPostSendsBodyAndHeaders(t *testing.T) {
type seen struct {
method, ctype, accept, ua, custom string
body []byte
}
ch := make(chan seen, 1)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
b, _ := io.ReadAll(r.Body)
ch <- seen{r.Method, r.Header.Get("Content-Type"), r.Header.Get("Accept"),
r.Header.Get("User-Agent"), r.Header.Get("X-Thing"), b}
w.Header().Set("Mcp-Session-Id", "sess-9")
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"ok":true}`))
}))
defer srv.Close()
f := testFetcher(t, Config{UserAgent: "Maven/test"})
resp, err := f.Post(context.Background(), srv.URL, "application/json",
[]byte(`{"jsonrpc":"2.0"}`), map[string]string{"Accept": "text/event-stream", "X-Thing": "1"})
if err != nil {
t.Fatal(err)
}
if string(resp.Body) != `{"ok":true}` {
t.Fatalf("body = %q", resp.Body)
}
if resp.Header["Mcp-Session-Id"] != "sess-9" {
t.Fatalf("response headers not surfaced: %+v", resp.Header)
}
s := <-ch
if s.method != http.MethodPost {
t.Fatalf("method = %s", s.method)
}
if string(s.body) != `{"jsonrpc":"2.0"}` {
t.Fatalf("request body = %q", s.body)
}
if s.ctype != "application/json" {
t.Fatalf("content-type = %q", s.ctype)
}
if s.accept != "text/event-stream" || s.custom != "1" {
t.Fatalf("caller headers dropped: %+v", s)
}
if s.ua != "Maven/test" {
t.Fatalf("user-agent = %q — a caller must not be able to override it", s.ua)
}
}
// The whole point of routing MCP through webfetch: a POST is guarded exactly
// like a GET. A body does not buy a caller a way onto the LAN.
func TestPostRefusesPrivateAddress(t *testing.T) {
f := New(Config{}) // no AllowPrivate
_, err := f.Post(context.Background(), "http://127.0.0.1:9100/mcp", "application/json", []byte(`{}`), nil)
if !errors.Is(err, ErrPrivate) {
t.Fatalf("error = %v, want ErrPrivate", err)
}
}
func TestPostRefusesNonHTTPScheme(t *testing.T) {
f := New(Config{})
if _, err := f.Post(context.Background(), "file:///etc/passwd", "application/json", nil, nil); !errors.Is(err, ErrScheme) {
t.Fatalf("error = %v, want ErrScheme", err)
}
}
func TestPostObeysDenylist(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}))
defer srv.Close()
f := testFetcher(t, Config{DenyHosts: []string{"127.0.0.1"}})
if _, err := f.Post(context.Background(), srv.URL, "application/json", []byte(`{}`), nil); !errors.Is(err, ErrBlocked) {
t.Fatalf("error = %v, want ErrBlocked", err)
}
}