35c6ff5a71
Persist reminder presentations and retry state, atomically complete collapsed deliveries, fall back across away reaches, and block permanent failures visibly (V-715, V-678). Fail closed when enabled integrations lack credentials and keep remote arms explicitly dark (V-691). Give mavweb one sanitized, request-correlated error contract (V-689). Owner explicitly requested direct commits to master.
710 lines
28 KiB
Go
710 lines
28 KiB
Go
package main
|
|
|
|
import (
|
|
"bufio"
|
|
"context"
|
|
"fmt"
|
|
"log"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/kami/maven/internal/config"
|
|
"github.com/kami/maven/internal/decision"
|
|
"github.com/kami/maven/internal/delivery"
|
|
"github.com/kami/maven/internal/delivery/voicesink"
|
|
"github.com/kami/maven/internal/dialogue"
|
|
"github.com/kami/maven/internal/ipc"
|
|
"github.com/kami/maven/internal/llm"
|
|
"github.com/kami/maven/internal/memory"
|
|
"github.com/kami/maven/internal/phraser"
|
|
"github.com/kami/maven/internal/router"
|
|
"github.com/kami/maven/internal/store"
|
|
"github.com/kami/maven/internal/stt"
|
|
"github.com/kami/maven/internal/tool"
|
|
"github.com/kami/maven/internal/tts"
|
|
"github.com/kami/maven/internal/voice"
|
|
"github.com/kami/maven/internal/weather"
|
|
"github.com/kami/maven/internal/worker"
|
|
)
|
|
|
|
// voiceWiring — everything the daemon needs to run the audio path. Held by
|
|
// cmd/mavend/main.go alongside the other wirings; closed on shutdown.
|
|
type voiceWiring struct {
|
|
server *voice.Server
|
|
sessions *voice.Sessions
|
|
voiceSink delivery.Sink
|
|
embedder router.Embedder
|
|
// heads — the routing heads, nil unless embedder.heads_path is set.
|
|
heads *router.RouterHeads
|
|
handler *reactiveHandler // the reactive handler for IPC Chat
|
|
// worker clients (set when configured as Remote): closed on shutdown so
|
|
// mavsttd / mavttsd don't keep a stale conn into a restarting daemon.
|
|
sttClient *worker.Client
|
|
ttsClient *worker.Client
|
|
// transcriber — the STT in use, exposed so the meeting recorder
|
|
// (cmd/mavend/capture.go) can reuse it. Maven has exactly one STT and does
|
|
// not grow a second one for capture: this is the same whisper.cpp worker the
|
|
// voice path talks to.
|
|
transcriber stt.Transcriber
|
|
// 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.
|
|
// pair — the workstation model with the resident one as the floor, nil
|
|
// unless a `workstation` block names an address. Held here only so the
|
|
// prober is stopped on shutdown; callers were handed it at build time.
|
|
pair *llm.Pair
|
|
// sttPair — CrisperWhisper 2.0 on the workstation with mavsttd as the
|
|
// floor, nil unless the `workstation.stt` block names an address. Held for
|
|
// the same reason as pair: to stop its prober on shutdown.
|
|
sttPair *stt.Pair
|
|
mcp *mcpWiring
|
|
// home — the Home Assistant client, nil unless the `smarthome` block is
|
|
// enabled (Vikunja #256). Its devices land in the same allowlist as every
|
|
// other act, so nothing else here has to know about it.
|
|
home *homeWiring
|
|
// netscan — the LAN scanner, nil unless the `netscan` block is enabled
|
|
// (Vikunja #257).
|
|
netscan *netWiring
|
|
}
|
|
|
|
// close releases the listener + worker conns. Safe to call on nil (when
|
|
// voice is not wired — wireVoice returns nil,nil).
|
|
func (w *voiceWiring) close() {
|
|
if w == nil {
|
|
return
|
|
}
|
|
if w.embedder != nil {
|
|
_ = w.embedder.Close()
|
|
}
|
|
if w.heads != nil {
|
|
_ = w.heads.Close()
|
|
}
|
|
if w.server != nil {
|
|
_ = w.server.Close()
|
|
}
|
|
if w.sttClient != nil {
|
|
_ = w.sttClient.Close()
|
|
}
|
|
if w.ttsClient != nil {
|
|
_ = w.ttsClient.Close()
|
|
}
|
|
if w.pair != nil {
|
|
w.pair.Stop()
|
|
}
|
|
if w.sttPair != nil {
|
|
w.sttPair.Stop()
|
|
}
|
|
w.mcp.close()
|
|
}
|
|
|
|
// wireVoice builds the audio path from cfg + a CoreAPI + a router. Returns
|
|
// nil wiring + nil error when voice isn't enabled (the caller's voice sink
|
|
// stays nil; the dispatcher's ChannelVoice routing drops silently).
|
|
//
|
|
// When voice is enabled, MUST wire a voicesink into the dispatcher's Voice
|
|
// slot using w.sessions (the caller does that — see main.go).
|
|
func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, memStore memory.Store, dataStore *store.Store, eco *ecosystemWiring) (*voiceWiring, error) {
|
|
if cfg.Voice == nil || !cfg.Voice.Enabled {
|
|
return nil, nil
|
|
}
|
|
w := &voiceWiring{}
|
|
|
|
// ----- stt (Stub in-process OR Remote via worker socket) -----
|
|
var transcriber stt.Transcriber
|
|
if cfg.Voice.Stt != nil && cfg.Voice.Stt.Socket != "" {
|
|
c := worker.Dial(cfg.Voice.Stt.Socket)
|
|
w.sttClient = c
|
|
lang := cfg.Voice.Stt.Lang
|
|
if lang == "" {
|
|
lang = cfg.Voice.Lang
|
|
}
|
|
transcriber = stt.NewRemote(c, lang)
|
|
} else {
|
|
transcriber = stt.NewStub()
|
|
}
|
|
transcriber, w.sttPair = sttSeam(cfg, transcriber)
|
|
w.transcriber = transcriber
|
|
|
|
// ----- tts (Stub in-process OR Remote) -----
|
|
var synthesizer tts.Synthesizer
|
|
if cfg.Voice.Tts != nil && cfg.Voice.Tts.Socket != "" {
|
|
c := worker.Dial(cfg.Voice.Tts.Socket)
|
|
w.ttsClient = c
|
|
lang := cfg.Voice.Tts.Lang
|
|
if lang == "" {
|
|
lang = cfg.Voice.Lang
|
|
}
|
|
synthesizer = tts.NewRemote(c, lang, cfg.Voice.Tts.Voice)
|
|
} else {
|
|
synthesizer = tts.NewStub()
|
|
}
|
|
|
|
// ----- router: embedder (ONNX when configured, floor HashEmbedder otherwise) -----
|
|
var emb router.Embedder
|
|
if cfg.Voice.Embedder != nil {
|
|
onnx, err := router.NewONNXEmbedder(
|
|
cfg.Voice.Embedder.ModelPath,
|
|
cfg.Voice.Embedder.TokenizerPath,
|
|
cfg.Voice.Embedder.LibPath,
|
|
)
|
|
if err != nil {
|
|
w.close()
|
|
return nil, fmt.Errorf("embedder: %w", err)
|
|
}
|
|
log.Printf("voice: onnx embedder loaded (%d dim)", onnx.Dim())
|
|
emb = onnx
|
|
} else {
|
|
log.Printf("voice: embedder not configured, using HashEmbedder floor")
|
|
emb = router.NewHashEmbedder(1024)
|
|
}
|
|
w.embedder = emb
|
|
|
|
// ----- router: routing heads (only when configured, and never fatal) -----
|
|
// A missing or broken weights file logs and leaves w.heads nil, which is
|
|
// byte-for-byte the cascade that shipped before V-664. Refusing to start
|
|
// over a routing accelerator would trade a working box for a better one.
|
|
if cfg.Voice.Embedder != nil && cfg.Voice.Embedder.HeadsPath != "" {
|
|
h, err := router.NewRouterHeads(
|
|
cfg.Voice.Embedder.HeadsPath,
|
|
cfg.Voice.Embedder.TokenizerPath,
|
|
)
|
|
if err != nil {
|
|
log.Printf("voice: routing heads unavailable, cascade unchanged: %v", err)
|
|
} else {
|
|
log.Printf("voice: routing heads loaded from %s", cfg.Voice.Embedder.HeadsPath)
|
|
w.heads = h
|
|
}
|
|
}
|
|
|
|
repairFactVectors(dataStore, emb)
|
|
checkStoredEmbedder(dataStore, emb)
|
|
// Retention is enforced on write, which is not enough on its own: a box that
|
|
// goes quiet keeps every trace until the next sixty-fourth turn (V-629).
|
|
pruneTracesOnStart(dataStore, time.Now())
|
|
|
|
// ----- tool executor (the enabled act allowlist, store-backed) -----
|
|
// Config tools are the declarative bootstrap: seed them into the store as
|
|
// enabled (editing mavend.json IS the human enable act). Ad-hoc tools are
|
|
// enabled later through the authed mavweb surface. The executor + matcher
|
|
// both read the store live, so a newly-enabled tool is runnable without a
|
|
// 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())
|
|
}
|
|
// The house (Vikunja #256): same story as MCP. Discovery PROPOSES a row per
|
|
// controllable device, always destructive, and Kami enables the ones he
|
|
// wants on /tools. Off unless the `smarthome` block is enabled.
|
|
w.home = wireSmartHome(cfg, dataStore)
|
|
if w.home != nil {
|
|
exec = exec.WithHome(w.home.caller())
|
|
}
|
|
// The LAN scanner (Vikunja #257): a read, bounded to the configured
|
|
// subnets and rate-limited. Off unless the `netscan` block is enabled.
|
|
w.netscan = wireNetScan(cfg, coreAPI)
|
|
matcher := tool.NewMatcher(coreAPI).WithAliases(toolAliases(cfg.Voice.Tools))
|
|
|
|
// ----- weather provider (Open-Meteo when configured, Stub otherwise) -----
|
|
var weatherProvider weather.Provider
|
|
var weatherLocation string
|
|
if cfg.Voice.Weather != nil && cfg.Voice.Weather.Provider == "open-meteo" {
|
|
weatherProvider = weather.NewOpenMeteoProvider()
|
|
weatherLocation = cfg.Voice.Weather.DefaultLocation
|
|
log.Printf("voice: weather provider: open-meteo (default location: %s)", cfg.Voice.Weather.DefaultLocation)
|
|
} else {
|
|
weatherProvider = weather.NewStubProvider()
|
|
log.Printf("voice: weather provider: stub (not configured)")
|
|
}
|
|
|
|
// The replier uses the same llama-server as the phraser.
|
|
var llmClient *llm.Client
|
|
if lp, ok := phr.(*phraser.LLMPhraser); ok {
|
|
// llmClientFor, not llm.New: this client must follow the phraser onto
|
|
// the new llama-server when the resident model is swapped (Vikunja #250).
|
|
llmClient = llmClientFor(lp, 60*time.Second)
|
|
}
|
|
// The workstation model sits above that one when it is configured and its
|
|
// card is free. hot is what the router and the replier complete through:
|
|
// either the pair, or the resident client alone, or nothing at all.
|
|
hot, pair := modelSeam(cfg, llmClient)
|
|
w.pair = pair
|
|
// The phraser gets the same pair, which is what carries the workstation model
|
|
// into the paths that do not go through `hot`: world questions (the naming
|
|
// half), and the digestion worker's nudge and reminder phrasing (the silent
|
|
// half). Wiring, so it happens once and before the voice server listens.
|
|
if lp, ok := phr.(*phraser.LLMPhraser); ok && pair != nil {
|
|
lp.UseRemote(pair)
|
|
}
|
|
// ----- router (the cascade; floor examples seed the classifier) -----
|
|
// The act matcher's allowlist is exactly the enabled tool names — the
|
|
// router only matches acts the executor can run (one source of truth).
|
|
threshold := cfg.Voice.RouterThreshold
|
|
if threshold <= 0 {
|
|
threshold = config.DefaultRouterThreshold
|
|
}
|
|
// The resident model routes by default: 63.2% of held-out intents right
|
|
// against the classifier's 50.0%, at about 1s a turn instead of 30ms (see
|
|
// config.VoiceConfig.LLMRouter). The classifier always stays wired as the
|
|
// fallback, so a model error never breaks a turn.
|
|
rtr := buildRouter(emb, matcher, threshold,
|
|
pickLLMRouter(cfg.Voice.UseLLMRouter(), hot), w.heads)
|
|
|
|
// ----- sessions registry (shared with voicesink) -----
|
|
sessions := voice.NewSessions()
|
|
w.sessions = sessions
|
|
|
|
// ----- voice sink (proactive nudges: dispatcher → voicesink → tts → push to client) -----
|
|
w.voiceSink = voicesink.New(synthesizer, sessions)
|
|
|
|
// ----- memory (long-term vector storage) -----
|
|
// Persistent (store-backed, survives restarts) when the daemon passes one;
|
|
// falls back to the in-memory floor otherwise (tests / no-store paths).
|
|
if memStore == nil {
|
|
memStore = memory.NewInMemoryStore()
|
|
}
|
|
|
|
// ----- dialogue (multi-turn slot carry-over; dialogueSessionTTL follow-up window) -----
|
|
// Store-backed when the daemon passes a store, so a restart mid-conversation
|
|
// keeps the thread (Vikunja #363). Sessions past their TTL are dropped on
|
|
// load, never revived. Clarify's parked question stays in memory only, and
|
|
// that is a decision rather than an omission (Vikunja #385, docs/design.md):
|
|
// a restart expires it, so the thread comes back and the open question does
|
|
// not.
|
|
const dialogueSessionTTL = 2 * time.Minute
|
|
var dialogueSessions *dialogue.SessionStore
|
|
if dataStore != nil {
|
|
dialogueSessions = dialogue.NewPersistentSessionStore(dialogueSessionTTL, dataStore)
|
|
if err := dialogueSessions.Load(context.Background(), time.Now()); err != nil {
|
|
log.Printf("dialogue: load saved sessions: %v", err)
|
|
}
|
|
} else {
|
|
dialogueSessions = dialogue.NewSessionStore(dialogueSessionTTL)
|
|
}
|
|
clarifyStore := dialogue.NewClarifyStore(clarifyTTL)
|
|
timeParser := router.NewPythonDateParser()
|
|
|
|
// ----- replier (LLM-backed when the engine is on, Stub floor otherwise) -----
|
|
replier := voice.Replier(voice.NewStubReplier())
|
|
if hot != nil {
|
|
replier = newLLMReplier(hot, contextBlockFn(cfg, time.Now))
|
|
}
|
|
|
|
// ----- the handler (the reactive path; closes over stt / tts / router / coreAPI / memory) -----
|
|
h := &reactiveHandler{
|
|
stt: transcriber,
|
|
tts: synthesizer,
|
|
router: rtr,
|
|
api: coreAPI,
|
|
tools: exec,
|
|
matcher: matcher,
|
|
replier: replier,
|
|
phraser: phr,
|
|
now: time.Now,
|
|
feedsOn: cfg.Feeds != nil,
|
|
home: w.home,
|
|
netscan: w.netscan,
|
|
// nil unless `crawl.on_demand` is on: reading a page he names is a
|
|
// capability, and capabilities are off unless configured.
|
|
crawler: onDemandCrawler(cfg),
|
|
// nil unless a `search` block names a SearXNG instance. External search
|
|
// is off unless configured, and configuring it is the whole opt-in.
|
|
search: wireSearch(cfg),
|
|
// nil unless a `kiwix` block names a server. Same swap-aware client the
|
|
// router and replier use, so the rewriter follows a model swap.
|
|
kiwix: wireKiwix(cfg, llmClient),
|
|
weatherProvider: weatherProvider,
|
|
weatherLocation: weatherLocation,
|
|
recall: recallWiring{
|
|
embedder: emb,
|
|
memStore: memStore,
|
|
minScore: cfg.Voice.QueryMinScore,
|
|
minMargin: cfg.Voice.QueryMinMargin,
|
|
},
|
|
dataStore: dataStore,
|
|
dialogueSessions: dialogueSessions,
|
|
// Always on (V-564). The record is the instrument the rest of V-558 is
|
|
// measured with, and one that only runs when a flag is set is not there
|
|
// on the night the misroute happens.
|
|
decisions: decision.NewRing(),
|
|
// The second sink (V-629). Same records, persisted, because the routing
|
|
// heads cannot be fitted from a 25-turn ring. Nil store ⇒ ring only, and
|
|
// EmbedderID is the same string the vector marker uses, so a trace and a
|
|
// stored vector name their body the same way.
|
|
traces: traceSink(dataStore),
|
|
encoderID: router.EmbedderID(emb),
|
|
clarifyStore: clarifyStore,
|
|
// 0 here (unset config) ⇒ the dialogue default.
|
|
clarifyMaxAttempts: cfg.Voice.ClarifyMaxAttempts,
|
|
extractor: router.Extractor{Time: timeParser, Acts: matcher, Facts: router.DefaultFactParser{}},
|
|
timeParser: timeParser,
|
|
ecosystem: eco,
|
|
}
|
|
|
|
// ----- the server (TCP listener) -----
|
|
srv := voice.NewServer(cfg.Voice.Bind, h, sessions)
|
|
if err := srv.Listen(); err != nil {
|
|
w.close()
|
|
return nil, fmt.Errorf("voice listen: %w", err)
|
|
}
|
|
w.server = srv
|
|
w.handler = h
|
|
|
|
return w, nil
|
|
}
|
|
|
|
// pickLLMRouter returns the LLM router when the operator asked for it and there
|
|
// is a llama-server to talk to, and nil otherwise. nil is safe: the cascade then
|
|
// routes with the classifier, so an unusable setting costs accuracy, not turns.
|
|
// modelSeam builds the completion seam the hot paths use: routing and replies.
|
|
//
|
|
// With no `workstation` block it is the resident client and nothing probes
|
|
// anything, which is today's deploy exactly. With one, it is an llm.Pair that
|
|
// prefers the workstation and falls back to the resident model silently — the
|
|
// silent half of the degradation rule (docs/offload.md), because the big model
|
|
// is only better here and the 1.7B is today's shipping quality. He is never
|
|
// told which of the two phrased his reply.
|
|
//
|
|
// A nil resident client means the phraser is not an LLM phraser. There is then
|
|
// no floor, and a Pair with no floor is a configuration mistake rather than a
|
|
// degraded mode, so the seam is nil and the cascade routes with the classifier.
|
|
func modelSeam(cfg *config.Config, resident *llm.Client) (router.Completer, *llm.Pair) {
|
|
if resident == nil {
|
|
if cfg.Workstation != nil && !cfg.Workstation.ModelDisabled {
|
|
log.Printf("voice: a workstation is configured but there is no resident model to floor it with — ignoring the block")
|
|
}
|
|
return nil, nil
|
|
}
|
|
if cfg.Workstation == nil || cfg.Workstation.ModelDisabled {
|
|
return resident, nil
|
|
}
|
|
ws := cfg.Workstation
|
|
remote := llm.New(ws.URL, time.Duration(ws.Timeout))
|
|
remote.SetToken(ws.Token)
|
|
if ws.Token == "" {
|
|
log.Printf("voice: unauthenticated workstation model endpoint is loopback-only")
|
|
}
|
|
pair := llm.NewPair(
|
|
remote,
|
|
resident,
|
|
ws.Health,
|
|
time.Duration(ws.Probe),
|
|
)
|
|
pair.Start(context.Background())
|
|
log.Printf("voice: workstation model at %s, probed every %s, resident model as the floor",
|
|
ws.URL, time.Duration(ws.Probe))
|
|
return pair, pair
|
|
}
|
|
|
|
// sttSeam builds the transcription seam the voice path and the meeting
|
|
// recorder share. It is modelSeam for audio and follows the same rule.
|
|
//
|
|
// With no `workstation.stt` block it hands back the floor untouched, which is
|
|
// today's deploy exactly. With one, it is an stt.Pair preferring CrisperWhisper
|
|
// 2.0 on workpc, which scores 10.4% WER in Russian against the floor's 27.5%
|
|
// (docs/evals/2026-08-09-crisperwhisper2-russian-wer.md).
|
|
//
|
|
// Only the silent half of the degradation rule applies here. A worse transcript
|
|
// is still a turn, so there is nothing to name a gap about and the fallback is
|
|
// never spoken. That is why stt.Pair has no TranscribeRemote.
|
|
func sttSeam(cfg *config.Config, floor stt.Transcriber) (stt.Transcriber, *stt.Pair) {
|
|
if cfg.Workstation == nil || cfg.Workstation.Stt == nil {
|
|
return floor, nil
|
|
}
|
|
s := cfg.Workstation.Stt
|
|
lang := ""
|
|
if cfg.Voice != nil {
|
|
lang = cfg.Voice.Lang
|
|
if cfg.Voice.Stt != nil && cfg.Voice.Stt.Lang != "" {
|
|
lang = cfg.Voice.Stt.Lang
|
|
}
|
|
}
|
|
pair := stt.NewPair(
|
|
stt.NewHTTPTranscriber(s.URL, s.Token, lang, time.Duration(s.Timeout)),
|
|
floor,
|
|
s.Health,
|
|
time.Duration(s.Probe),
|
|
)
|
|
pair.Start(context.Background())
|
|
if s.Token == "" {
|
|
log.Print("voice: unauthenticated workstation transcriber endpoint is loopback-only")
|
|
}
|
|
log.Printf("voice: workstation transcriber at %s, probed every %s, mavsttd as the floor",
|
|
s.URL, time.Duration(s.Probe))
|
|
return pair, pair
|
|
}
|
|
|
|
func pickLLMRouter(enabled bool, c router.Completer) *router.LLMRouter {
|
|
if !enabled {
|
|
return nil
|
|
}
|
|
if c == nil {
|
|
log.Printf("voice: voice.llm_router is on but there is no llama-server to route with (the phraser is not an LLM phraser) — using the classifier instead")
|
|
return nil
|
|
}
|
|
log.Printf("voice: LLM router enabled")
|
|
return router.NewLLMRouter(c)
|
|
}
|
|
|
|
// buildRouter constructs the reactive-path router with the given embedder
|
|
// and confidence threshold.
|
|
// - stage-0 grammars from DefaultActMatcher whose fn allowlist is exactly
|
|
// the enabled tool names (actFns) — the router only matches acts the
|
|
// executor can run. Empty ⇒ every act refuses at the matcher.
|
|
// - The embedder is provided by wireVoice: HashEmbedder (floor) when no
|
|
// embedder config is present, or the ONNX multilingual model when
|
|
// configured — same interface, one constructor change.
|
|
// - The classifier is floored by seedClassifier, which loads one file per
|
|
// intent from seedDir (models/seeds/<intent>.txt) — see seedClassifier
|
|
// below for the current intent list and file names.
|
|
// - Threshold is from voice.router_threshold config (default 0.55).
|
|
func buildRouter(emb router.Embedder, acts router.ActMatcher, threshold float64,
|
|
llmR *router.LLMRouter, heads *router.RouterHeads) *router.Router {
|
|
cls := router.NewClassifier(emb)
|
|
seedClassifier(cls)
|
|
// The stage 0 set, in the router package, so the eval fixture runs the rules
|
|
// the daemon runs (V-693). Order and reasoning live with the list.
|
|
grammars := router.StageZeroGrammars(acts)
|
|
return router.New(router.Config{
|
|
Grammars: grammars,
|
|
Classifier: cls,
|
|
Extractor: router.Extractor{
|
|
Time: router.NewPythonDateParser(),
|
|
Acts: acts,
|
|
Facts: router.DefaultFactParser{},
|
|
},
|
|
Threshold: threshold,
|
|
LLM: llmR,
|
|
Heads: heads,
|
|
})
|
|
}
|
|
|
|
// seedDir is the directory containing intent seed files, relative to the repo
|
|
// root. Each file is named <intent>.txt and holds one training example per
|
|
// line (blank lines and lines starting with # are ignored).
|
|
const seedDir = "models/seeds"
|
|
|
|
// seedPath resolves seedDir against the working directory, walking up until it
|
|
// finds it. The daemon runs from the repo root and the first candidate hits.
|
|
//
|
|
// A test does not: `go test ./cmd/mavend/` runs with the working directory at
|
|
// cmd/mavend, so every open failed and the simulator scenarios replayed a whole
|
|
// scripted day against a classifier holding zero examples (Vikunja #465). They
|
|
// passed, which is the part that matters — a green simulator was not exercising
|
|
// the routing the deploy runs, and a regression in the seed set could not have
|
|
// shown up there.
|
|
//
|
|
// Bounded at five levels, so a daemon started somewhere without the seeds logs
|
|
// the same failure it always did rather than walking to the filesystem root.
|
|
func seedPath() string {
|
|
dir := seedDir
|
|
for i := 0; i < 5; i++ {
|
|
if st, err := os.Stat(dir); err == nil && st.IsDir() {
|
|
return dir
|
|
}
|
|
dir = filepath.Join("..", dir)
|
|
}
|
|
return seedDir
|
|
}
|
|
|
|
// seedClassifier floors the embedded examples so the cold-boot path
|
|
// doesn't return ErrNoIntents. Loads examples from seedDir — one file per
|
|
// intent (act.txt, reminder.txt, fact.txt, note.txt, query.txt, chat.txt,
|
|
// system.txt). When the classifier can't decide it falls through to
|
|
// Clarify — the last-resort path asks the user to rephrase rather than
|
|
// guessing wrong.
|
|
func seedClassifier(c *router.Classifier) {
|
|
intents := []router.Intent{
|
|
router.IntentAct,
|
|
router.IntentReminder,
|
|
router.IntentFact,
|
|
router.IntentNote,
|
|
router.IntentQuery,
|
|
router.IntentChat,
|
|
router.IntentSystem,
|
|
}
|
|
total := 0
|
|
for _, intent := range intents {
|
|
n, err := loadSeedFile(c, intent)
|
|
if err != nil {
|
|
log.Printf("voice: seed %s: %v", intent, err)
|
|
continue
|
|
}
|
|
total += n
|
|
}
|
|
log.Printf("voice: loaded %d seed examples from %s", total, seedPath())
|
|
}
|
|
|
|
func loadSeedFile(c *router.Classifier, intent router.Intent) (int, error) {
|
|
path := filepath.Join(seedPath(), string(intent)+".txt")
|
|
f, err := os.Open(path)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("open %s: %w", path, err)
|
|
}
|
|
defer f.Close()
|
|
|
|
var count int
|
|
sc := bufio.NewScanner(f)
|
|
for sc.Scan() {
|
|
line := strings.TrimSpace(sc.Text())
|
|
if line == "" || strings.HasPrefix(line, "#") {
|
|
continue
|
|
}
|
|
if err := c.AddExample(context.Background(), intent, line); err != nil {
|
|
log.Printf("voice: seed %s: skipping %q: %v", intent, line, err)
|
|
continue
|
|
}
|
|
count++
|
|
}
|
|
if err := sc.Err(); err != nil {
|
|
return count, fmt.Errorf("scan %s: %w", path, err)
|
|
}
|
|
return count, nil
|
|
}
|
|
|
|
// toolAliases collects the spoken phrases per tool name. Without them the act
|
|
// matcher only ever matched the English tool name, so no Russian utterance could
|
|
// reach a tool and every homelab act fell to proposeGap (V-633).
|
|
func toolAliases(tools []config.ToolConfig) map[string][]string {
|
|
out := make(map[string][]string, len(tools))
|
|
for _, tc := range tools {
|
|
if tc.Name == "" || len(tc.Aliases) == 0 {
|
|
continue
|
|
}
|
|
out[tc.Name] = tc.Aliases
|
|
}
|
|
return out
|
|
}
|
|
|
|
// seedTools upserts the config-declared tools into the store as enabled. Editing
|
|
// mavend.json is a human act, so a config tool is enabled by definition; this
|
|
// makes the declarative config the reproducible bootstrap while the store stays
|
|
// the single runtime source of truth (mavweb enables ad-hoc ones on top).
|
|
func seedTools(api ipc.CoreAPI, tools []config.ToolConfig) {
|
|
ctx := context.Background()
|
|
now := time.Now()
|
|
n := 0
|
|
for _, tc := range tools {
|
|
if tc.Name == "" || len(tc.Cmd) == 0 {
|
|
log.Printf("voice: skipping malformed tool config %+v", tc)
|
|
continue
|
|
}
|
|
if err := api.EnableTool(ctx, tc.Name, tc.Cmd, tc.Destructive, tc.Scope, now); err != nil {
|
|
log.Printf("voice: seed tool %q: %v", tc.Name, err)
|
|
continue
|
|
}
|
|
n++
|
|
}
|
|
log.Printf("voice: seeded %d act tools from config", n)
|
|
}
|
|
|
|
// repairFactVectors brings stored fact vectors in line with the facts they name
|
|
// (#493), once per box, before the embedder marker is even looked at.
|
|
//
|
|
// Automatic and not a flag, unlike -reembed: only voice-tapped facts are in
|
|
// this index, so the work is tens of embeddings rather than the thousands of
|
|
// notes that made the backfill a deliberate act. And the box that needs it is
|
|
// broken in a way nobody can see — recall answers with the wrong text and
|
|
// nothing logs an error — so waiting for an operator to know to run it is how
|
|
// the defect survived four restarts in the first place.
|
|
func repairFactVectors(dataStore *store.Store, emb router.Embedder) {
|
|
if dataStore == nil {
|
|
return
|
|
}
|
|
res, err := dataStore.RepairFactVectors(context.Background(),
|
|
// EmbedPassage, the stored side, same as every other writer of these
|
|
// vectors.
|
|
func(ctx context.Context, text string) ([]float32, error) {
|
|
return router.EmbedPassage(ctx, emb, text)
|
|
})
|
|
if err != nil {
|
|
log.Printf("voice: fact vector repair failed, no marker written and nothing half-done — retried next start: %v", err)
|
|
return
|
|
}
|
|
if res.Skipped || res.Rewritten+res.Dropped == 0 {
|
|
return
|
|
}
|
|
log.Printf("voice: fact vector repair — %d re-embedded from the fact they name, %d dropped as voided or superseded, %d already right, took %s (#493)",
|
|
res.Rewritten, res.Dropped, res.Kept, res.Took.Round(time.Millisecond))
|
|
}
|
|
|
|
// reembedOnStart is the -reembed flag (set in run()). Opt-in on purpose: see
|
|
// runReembed.
|
|
var reembedOnStart bool
|
|
|
|
// allowSeedOnStart is the -allow-seed flag (set in run()). Opt-in, and the
|
|
// default is the one that matters: a box nobody is testing has no live path to
|
|
// write a fact into the past. See seed.go and Vikunja #518.
|
|
var allowSeedOnStart bool
|
|
|
|
// checkStoredEmbedder compares the embedder we just loaded with the one that
|
|
// wrote the vectors already in the DB (Vikunja #378).
|
|
//
|
|
// The two models we have both make 384-dim vectors, so a size check catches
|
|
// nothing: after a swap, recall silently compares vectors from different
|
|
// spaces and the scores are noise. So we say it out loud. Recall itself is not
|
|
// changed here — the fix is `mavend -reembed`.
|
|
func checkStoredEmbedder(dataStore *store.Store, emb router.Embedder) {
|
|
if dataStore == nil {
|
|
return
|
|
}
|
|
current := router.EmbedderID(emb)
|
|
if reembedOnStart {
|
|
runReembed(dataStore, emb, current)
|
|
return
|
|
}
|
|
stored, mismatch, err := dataStore.CheckEmbedder(context.Background(), current)
|
|
if err != nil {
|
|
log.Printf("voice: embedder marker check failed: %v", err)
|
|
return
|
|
}
|
|
if mismatch {
|
|
log.Printf("voice: WARNING embedder MISMATCH — stored vectors were written by %q but the configured embedder is %q; recall scores are noise until the notes and facts are re-embedded — run `mavend -reembed` once (Vikunja #378)", stored, current)
|
|
return
|
|
}
|
|
log.Printf("voice: embedder marker ok (%s)", current)
|
|
}
|
|
|
|
// runReembed is the one-shot backfill behind -reembed.
|
|
//
|
|
// Why a flag and not automatic on mismatch: the embedder is ONNX on the
|
|
// laptop's CPU, so a few thousand notes is minutes of work. Doing that silently
|
|
// inside a normal start would look like the daemon hanging on boot. So the user
|
|
// runs it once, deliberately, after an embedder swap; the mismatch warning
|
|
// above tells them to. It re-embeds, logs what it did, and then the daemon
|
|
// carries on serving as usual — no separate binary, no second start needed.
|
|
func runReembed(dataStore *store.Store, emb router.Embedder, current string) {
|
|
log.Printf("voice: re-embedding stored notes and facts with %s — this can take a few minutes, do not interrupt", current)
|
|
res, err := dataStore.ReembedAll(context.Background(), current,
|
|
// EmbedPassage, not EmbedQuery: these are stored texts being searched
|
|
// FOR, which is the side they were written with.
|
|
func(ctx context.Context, text string) ([]float32, error) {
|
|
return router.EmbedPassage(ctx, emb, text)
|
|
})
|
|
if err != nil {
|
|
log.Printf("voice: re-embed FAILED, nothing was changed and no marker was written — safe to run again: %v", err)
|
|
return
|
|
}
|
|
if res.Skipped {
|
|
log.Printf("voice: re-embed skipped — the stored vectors were already written by %s", current)
|
|
return
|
|
}
|
|
log.Printf("voice: re-embed done — %d notes in the notes table, %d notes and %d facts in the memory index, took %s; stored vectors now belong to %s",
|
|
res.Notes, res.MemNotes, res.Facts, res.Took.Round(time.Second), current)
|
|
|
|
// A row with no text cannot be re-embedded, so its vector is still the old
|
|
// model's noise while the marker now says everything is current. Both write
|
|
// paths always store the text, so this should be zero — say it loudly
|
|
// rather than bury it in the line above if it ever isn't.
|
|
if res.NoText > 0 {
|
|
log.Printf("voice: WARNING %d stored rows had no text, so their vectors could not be re-embedded and are still noise; they will never match anything useful (Vikunja #378)", res.NoText)
|
|
}
|
|
}
|