Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d92349ca6e | |||
| 8d5e357b57 |
@@ -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 != "" {
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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) -----
|
||||
|
||||
@@ -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
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
@@ -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.
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -18,6 +18,7 @@ import (
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/delivery/ntfysink"
|
||||
@@ -193,6 +194,17 @@ type Config struct {
|
||||
// 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
|
||||
@@ -474,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
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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.
|
||||
|
||||
@@ -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
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -2,10 +2,12 @@ package mcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
@@ -429,3 +431,93 @@ func (m *Manager) Close() error {
|
||||
}
|
||||
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, ", "))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -444,3 +444,122 @@ func TestLocalNameAndCmd(t *testing.T) {
|
||||
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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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()
|
||||
}
|
||||
@@ -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]
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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`)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user