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 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 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.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() } 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() } 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 repairFactVectors(dataStore, emb) checkStoredEmbedder(dataStore, emb) // ----- 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)) // ----- 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(), 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 { 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 { return resident, nil } ws := cfg.Workstation pair := llm.NewPair( llm.New(ws.URL, time.Duration(ws.Timeout)), 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 } 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/.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) *router.Router { cls := router.NewClassifier(emb) seedClassifier(cls) grammars := router.DefaultGrammars(acts) grammars = append(grammars, router.SystemTimeDateGrammars()...) // After the time/date rules on purpose: "какой сегодня день" is a clock // question and must keep reaching replySystem, while "что у меня сегодня" // is an agenda question and must not. grammars = append(grammars, router.AgendaQueryGrammars()...) // Same reason as the agenda rules, for the feeds: "что нового в лентах?" // routed system and answered "пока не умею" (Vikunja #474). grammars = append(grammars, router.FeedQueryGrammar()) // The list side of the same exposure: a phrasing with no possessive in it // ("список дел") routed system and never reached queryTasks (Vikunja #467). grammars = append(grammars, router.TaskListGrammar()) grammars = append(grammars, router.ListGrammars()...) grammars = append(grammars, router.ReminderGrammar()) // Before the capture marker, because "отметь" is a capture verb and "отметь // второй пункт" is not a note. The Praxis rules are the narrower claim — a // lifecycle verb AND an item named — so they get first refusal (Vikunja #516). grammars = append(grammars, router.PraxisGrammars()...) // Last, and it matches any utterance shape — its Build is the filter. An // explicit capture marker beats the model, which called it an act and // rewrote the task text (Vikunja #467). After the rules above because a // marker never collides with a clock or agenda question. // After Praxis, whose bare "закрой" claim this rule cannot reach (it needs the // board noun), and before the capture marker, which would otherwise read // "убери из задач купить молоко" as a new task (Vikunja #512). grammars = append(grammars, router.TaskStatusGrammar()) // Before the capture markers, which all need an object. A capture verb // alone is a fact with no key, and the clarify path asks for it rather than // letting the model invent an answer (Vikunja #557). grammars = append(grammars, router.BareCaptureGrammar()...) grammars = append(grammars, router.TaskCaptureGrammar()) // After the capture marker, so "запиши" still wins over "расскажи", and // last overall because it matches on the first word alone: "расскажи про // X" is a world question the model called a fact (Vikunja #498). grammars = append(grammars, router.NarrativeQueryGrammars()...) return router.New(router.Config{ Grammars: grammars, Classifier: cls, Extractor: router.Extractor{ Time: router.NewPythonDateParser(), Acts: acts, Facts: router.DefaultFactParser{}, }, Threshold: threshold, LLM: llmR, }) } // seedDir is the directory containing intent seed files, relative to the repo // root. Each file is named .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) } }