5934110fa8
voice.go is still the biggest file in cmd/mavend and most of what is left has nothing to do with the audio path. The ecosystem integration is one such lump: it talks to Nexus, Praxis and Hexis over HTTP and only touches the handler for its store and clock. Lifting it into ecosystem_acts.go puts it next to ecosystem.go, where the clients it drives already live. Move-only: handlePraxisAct, recordPraxisTrace (called from nowhere else), handleHexisAct and execHexis verbatim, plus the two imports that became unused in voice.go.
1100 lines
43 KiB
Go
1100 lines
43 KiB
Go
// Package main is mavend's voice wiring + reactive handler.
|
|
//
|
|
// Two responsibilities for the audio path:
|
|
//
|
|
// 1. CONSTRUCTION: read cfg.Voice, build the stt/tts/transcribers (Stub
|
|
// in-process by default, Remote via worker socket when configured),
|
|
// the router (stage-0 grammar + HashEmbedder classifier seeded with
|
|
// floor examples — production swaps in the ONNX multilingual model
|
|
// later), the voice TCP listener, the sessions registry, the
|
|
// voicesink, and wire the voicesink into the dispatcher's Voice slot.
|
|
//
|
|
// 2. HANDLER: a concrete voice.Handler that processes PushToTalk
|
|
// requests: stt → router → action → replier → tts → reply. The
|
|
// handler is what makes the audio round-trip "live". It wires to the
|
|
// CoreAPI in-process (the daemon already has it as ipc.NewStoreAPI(st)
|
|
// for module-IPC — the reactive path uses the same CoreAPI off the
|
|
// same store; both are the "core = the only key-holder" path through
|
|
// the daemon-embedded adapter).
|
|
//
|
|
// The "actions" handled today (per spec order; some deferred):
|
|
//
|
|
// - IntentFact: WriteFact via CoreAPI. The router's Slots.Key/Value feed
|
|
// the write; Source = "tap:voice" (the voice path is a tap, value=1.0
|
|
// confidence — the user said it out loud, maven trusts the capture).
|
|
// - IntentReminder: CreateReminder via CoreAPI. The router already
|
|
// resolved relative→absolute at capture ("in 4h" → fire_ts); the
|
|
// CoreAPI stores it as-is.
|
|
// - IntentAct: the tool executor runs the matched fn against the store's
|
|
// ENABLED allowlist (internal/tool). A verb not on it is scaffolded as a
|
|
// 'proposed' tool a human enables on the authed mavweb surface (never
|
|
// voice). Destructive tools run only after a spoken confirm turn.
|
|
// - IntentNote: chroma/vector-store deferred. The handler replies
|
|
// "saved" without persisting — a stub on the way to chroma.
|
|
// - IntentQuery: RAG-over-chroma deferred. The handler replies "I'll
|
|
// look that up later" — same shape as the other deferred slots.
|
|
// - Clarify: the router's stage-3 confidence gate fired; reply "didn't
|
|
// catch that, can you rephrase?"
|
|
//
|
|
// The Replier (voice.StubReplier today) renders the reply TEXT across all
|
|
// these branches. The TTS synthesiser (tts.Stub today) renders that text
|
|
// to audio. The PushToTalkResp carries BOTH so the client can play (audio)
|
|
// AND log (text) for tests asserting the round-trip.
|
|
package main
|
|
|
|
import (
|
|
"bufio"
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/kami/maven/internal/audio"
|
|
"github.com/kami/maven/internal/config"
|
|
"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/pattern"
|
|
"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/ttsnorm"
|
|
"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
|
|
}
|
|
|
|
// 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()
|
|
}
|
|
}
|
|
|
|
// 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()
|
|
}
|
|
|
|
// ----- 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
|
|
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))
|
|
matcher := tool.NewMatcher(coreAPI)
|
|
|
|
// ----- 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 {
|
|
llmClient = llm.New(lp.BaseURL(), 60*time.Second)
|
|
}
|
|
// ----- 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(), llmClient))
|
|
|
|
// ----- 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; 2-min 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.
|
|
var dialogueSessions *dialogue.SessionStore
|
|
if dataStore != nil {
|
|
dialogueSessions = dialogue.NewPersistentSessionStore(2*time.Minute, dataStore)
|
|
if err := dialogueSessions.Load(context.Background(), time.Now()); err != nil {
|
|
log.Printf("dialogue: load saved sessions: %v", err)
|
|
}
|
|
} else {
|
|
dialogueSessions = dialogue.NewSessionStore(2 * time.Minute)
|
|
}
|
|
clarifyStore := dialogue.NewClarifyStore(clarifyTTL)
|
|
timeParser := router.NewPythonDateParser()
|
|
|
|
// ----- replier (LLM-backed when the engine is on, Stub floor otherwise) -----
|
|
replier := voice.Replier(voice.NewStubReplier())
|
|
if llmClient != nil {
|
|
replier = newLLMReplier(llmClient, contextBlockFn(cfg, time.Now))
|
|
}
|
|
|
|
// ----- the handler (the reactive path; closes over stt / tts / router / coreAPI / memory) -----
|
|
h := &reactiveHandler{
|
|
stt: transcriber,
|
|
tts: synthesizer,
|
|
router: rtr,
|
|
embedder: emb,
|
|
api: coreAPI,
|
|
tools: exec,
|
|
matcher: matcher,
|
|
replier: replier,
|
|
phraser: phr,
|
|
now: time.Now,
|
|
weatherProvider: weatherProvider,
|
|
weatherLocation: weatherLocation,
|
|
memStore: memStore,
|
|
dataStore: dataStore,
|
|
dialogueSessions: dialogueSessions,
|
|
clarifyStore: clarifyStore,
|
|
// 0 here (unset config) ⇒ the dialogue default.
|
|
clarifyMaxAttempts: cfg.Voice.ClarifyMaxAttempts,
|
|
extractor: router.Extractor{Time: timeParser, Acts: matcher, Facts: router.DefaultFactParser{}},
|
|
queryMinScore: cfg.Voice.QueryMinScore,
|
|
queryMinMargin: cfg.Voice.QueryMinMargin,
|
|
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
|
|
}
|
|
|
|
// reactiveHandler — voice.Handler implementation. One method: turn a
|
|
// PushToTalkReq into a reply (audio + text). The handler is concurrency-
|
|
// safe (the wired stt/tts/router/api all are); called from per-conn
|
|
// goroutines on the voice.Server.
|
|
type reactiveHandler struct {
|
|
stt stt.Transcriber
|
|
tts tts.Synthesizer
|
|
router *router.Router
|
|
embedder router.Embedder // reused for note write/query (same model as the classifier)
|
|
api ipc.CoreAPI
|
|
tools *tool.Executor
|
|
matcher *tool.Matcher
|
|
phraser phraser.Phraser
|
|
replier voice.Replier
|
|
now func() time.Time
|
|
|
|
weatherProvider weather.Provider
|
|
weatherLocation string // default location for weather queries
|
|
|
|
memStore memory.Store
|
|
dataStore *store.Store // direct store access for event extraction + pattern detection
|
|
|
|
// queryMinScore — the note-recall confidence gate. Top cosine below this ⇒
|
|
// "I don't know" instead of a guess. Tuned for the ONNX embedder; a knob, not
|
|
// load-bearing math (same posture as the presence thresholds). Set by
|
|
// wireVoice from VoiceConfig; default 0.55.
|
|
queryMinScore float64
|
|
// queryMinMargin — the second half of that gate: how far the top hit must
|
|
// beat the runner-up. 0 ⇒ margin off.
|
|
queryMinMargin float64
|
|
|
|
// timeParser — used as a fallback for stage-0 reminder grammar matches
|
|
// (where the extractor didn't run). Shared with the router's extractor.
|
|
// The production dateparser will replace StubDateTimeParser here too.
|
|
timeParser router.DateTimeParser
|
|
|
|
// dialogueSessions carries slots across turns for follow-ups (single-user
|
|
// box → one session slot, keyed voiceDialogueID). nil ⇒ no carry-over.
|
|
dialogueSessions *dialogue.SessionStore
|
|
|
|
// clarifyStore parks the request behind an open question she asked (see
|
|
// clarify.go). nil ⇒ she falls back to the canned "не поняла" reply.
|
|
clarifyStore *dialogue.ClarifyStore
|
|
|
|
// clarifyMaxAttempts — questions per request before she gives up out loud.
|
|
// 0 ⇒ dialogue.DefaultMaxAttempts (3). Set from VoiceConfig.
|
|
clarifyMaxAttempts int
|
|
|
|
// extractor parses the answer to an open question, with the same parsers
|
|
// the router's own stage-2 uses.
|
|
extractor router.Extractor
|
|
|
|
// pending destructive-act confirmation. A destructive act replies with a
|
|
// "выполнить X? да/нет" prompt and parks here; the NEXT utterance is read as
|
|
// the y/n answer. ponytail: single slot, single-user box — a second act
|
|
// while one waits overwrites it (last-asked wins); expires after confirmTTL.
|
|
mu sync.Mutex
|
|
pending *pendingAct
|
|
pendingRoutine *pendingRoutineConfirm // routine proposal awaiting y/n
|
|
pendingHexis *pendingHexisExec // mutating Hexis capability awaiting y/n
|
|
|
|
ecosystem *ecosystemWiring // nexus + hexis + praxis clients
|
|
}
|
|
|
|
// pendingHexisExec — a mutating Hexis capability parked awaiting a spoken
|
|
// confirm. The confirmation is bound to the resolved capability + canonical
|
|
// target entity so a later "да" can only execute exactly what was proposed
|
|
// (ecosystem invariant: protected actions require bound confirmation).
|
|
type pendingHexisExec struct {
|
|
capabilityID string
|
|
capName string
|
|
entityID string
|
|
displayName string
|
|
expiry time.Time
|
|
}
|
|
|
|
// pendingRoutineConfirm — a proposed routine awaiting a spoken y/n to become
|
|
// a recurring reminder. Set by detectPattern after creating a proposal.
|
|
type pendingRoutineConfirm struct {
|
|
routineID int64
|
|
action string
|
|
object string
|
|
interval float64
|
|
phrase string
|
|
expiry time.Time
|
|
}
|
|
|
|
// pendingAct — a destructive act awaiting a spoken confirm.
|
|
type pendingAct struct {
|
|
fn string
|
|
args []string
|
|
phrase string
|
|
expiry time.Time
|
|
}
|
|
|
|
// confirmTTL — how long a parked destructive confirm stays answerable. Short:
|
|
// a confirm is a same-breath gesture; a stale prompt shouldn't fire on an
|
|
// unrelated later "да".
|
|
const confirmTTL = 90 * time.Second
|
|
|
|
// HandlePushToTalk — the full reactive round-trip. Each step's failure
|
|
// surfaces as a short reply text + empty audio OR an error; the voice
|
|
// server translates an error into a wire RpcError. Today the handler
|
|
// prefers a canned error-reply over an error return (a user-facing "didn't
|
|
// catch that" is better than a wire error the client surfaces as
|
|
// "internal"); the only error returned is a synthesizer fault (no audio
|
|
// to ship back).
|
|
func (h *reactiveHandler) HandlePushToTalk(ctx context.Context, req voice.PushToTalkReq, _ uint64) (voice.PushToTalkResp, error) {
|
|
// 1. stt — transcribe the audio.
|
|
text, _, err := h.stt.Transcribe(ctx, req.Audio)
|
|
if err != nil {
|
|
log.Printf("voice: stt error: %v", err)
|
|
return h.reply(ctx, "не получилось разобрать речь — попробуй ещё раз.", nil)
|
|
}
|
|
if text == "" {
|
|
return h.reply(ctx, "ничего не услышала — попробуй ещё раз.", nil)
|
|
}
|
|
log.Printf("voice: stt → %q", text)
|
|
|
|
// 1b. confirm turn — if a destructive act is parked, this utterance is its
|
|
// y/n answer, not a fresh command. Handled before routing so "да" doesn't
|
|
// get classified as some other intent.
|
|
if reply, handled := h.resolveConfirm(ctx, text); handled {
|
|
return h.reply(ctx, reply, nil)
|
|
}
|
|
|
|
// 1b2. expired clarify — a question was parked but its TTL ran out, so the
|
|
// request behind it is gone. Say that out loud (see clarify.go) and carry
|
|
// on: these words are still routed as a fresh utterance below, with the
|
|
// notice glued in front of whatever the fresh routing answers. Checked
|
|
// BEFORE the answer path: reading a parked question drops an expired one.
|
|
expiredNotice := h.clarifyExpiredNotice()
|
|
|
|
// 1b3. clarify answer — if she asked a live question last turn, this
|
|
// utterance is its answer, not a fresh command. After the confirm check: a
|
|
// y/n gate is armed by her own prompt and is the narrower claim on the
|
|
// utterance.
|
|
if reply, handled := h.resolveClarifyAnswer(ctx, text); handled {
|
|
return h.reply(ctx, reply, nil)
|
|
}
|
|
|
|
// 1c. quiet-hours toggle — keyword match, not classifier-dependent.
|
|
// "тихий режим" / "quiet on" would route through the classifier
|
|
// unreliably (it's a command, not a free-form query), so we match it
|
|
// before routing. Same pattern as the confirm turn above.
|
|
if reply, handled := h.resolveQuietToggle(ctx, text); handled {
|
|
return h.reply(ctx, withNotice(expiredNotice, reply), nil)
|
|
}
|
|
|
|
// 2. router — classify the utterance.
|
|
dec, err := h.router.Route(ctx, text, h.now())
|
|
if err != nil {
|
|
// ErrNoIntents ⇒ classifier unseeded (cold boot). reply with a
|
|
// "still warming up" rather than a wire error.
|
|
if errors.Is(err, router.ErrNoIntents) {
|
|
return h.reply(ctx, withNotice(expiredNotice, "я ещё не понимаю свободную речь — скоро научусь."), nil)
|
|
}
|
|
log.Printf("voice: router error: %v", err)
|
|
return h.reply(ctx, withNotice(expiredNotice, "не получилось разобрать команду."), nil)
|
|
}
|
|
|
|
// 2b. dialogue — fill this turn's missing slots from a prior same-intent
|
|
// turn (follow-ups like «напомни завтра» → «…позвонить маме»), then remember
|
|
// this turn for the next follow-up. Only same-intent, non-expired, non-
|
|
// clarify turns carry (see followUpMerge). Best-effort: nil store ⇒ skipped.
|
|
if h.dialogueSessions != nil {
|
|
now := h.now()
|
|
prev := h.dialogueSessions.Get(voiceDialogueID, now)
|
|
dec = followUpMerge(prev, dec, now)
|
|
if !dec.Clarify {
|
|
h.rememberTurn(prev, dec, now)
|
|
}
|
|
}
|
|
|
|
// 2c. clarify — she is not sure. If one named thing is missing, ask about it
|
|
// and park the request (clarify.go); otherwise the replier's canned reply
|
|
// stands.
|
|
if dec.Clarify {
|
|
if question, asked := h.askClarify(dec); asked {
|
|
return h.reply(ctx, withNotice(expiredNotice, question), nil)
|
|
}
|
|
}
|
|
|
|
// 3. action — execute the decision's intent. errors here surface as
|
|
// short reply text (the user wants to know the action didn't land);
|
|
// the round-trip stays alive.
|
|
replyText := h.applyAction(ctx, dec)
|
|
|
|
// 4. replier — phrase the reply across the router decision.
|
|
if replyText == "" {
|
|
replyText = h.replier.Reply(dec)
|
|
}
|
|
|
|
// 5. tts — synthesise the reply text; return to the voice server which
|
|
// ships it back on the conn.
|
|
return h.reply(ctx, withNotice(expiredNotice, replyText), nil)
|
|
}
|
|
|
|
// handleText — the core reactive path without stt/tts: confirm check →
|
|
// route → dialogue → action → replier. Used by the IPC Chat endpoint
|
|
// (and eventually by telegram). Splits out the audio bookends from
|
|
// HandlePushToTalk so text channels share the same routing logic.
|
|
func (h *reactiveHandler) handleText(ctx context.Context, text string) string {
|
|
log.Printf("voice: handleText: %q", text)
|
|
// 1b. confirm turn — if a destructive act is parked, this utterance is its
|
|
// y/n answer. Same check as HandlePushToTalk.
|
|
if reply, handled := h.resolveConfirm(ctx, text); handled {
|
|
return reply
|
|
}
|
|
|
|
// 1b2/1b3. expired clarify then clarify answer — same order and reasons as
|
|
// HandlePushToTalk.
|
|
expiredNotice := h.clarifyExpiredNotice()
|
|
if reply, handled := h.resolveClarifyAnswer(ctx, text); handled {
|
|
return reply
|
|
}
|
|
|
|
// 2. router — classify the utterance.
|
|
dec, err := h.router.Route(ctx, text, h.now())
|
|
if err != nil {
|
|
if errors.Is(err, router.ErrNoIntents) {
|
|
return withNotice(expiredNotice, "я ещё не понимаю свободную речь — скоро научусь.")
|
|
}
|
|
log.Printf("voice: handleText router error: %v", err)
|
|
return withNotice(expiredNotice, "не получилось разобрать команду.")
|
|
}
|
|
log.Printf("voice: route result: intent=%s slots=%+v", dec.Intent, dec.Slots)
|
|
|
|
// 2b. dialogue — same as HandlePushToTalk.
|
|
if h.dialogueSessions != nil {
|
|
now := h.now()
|
|
prev := h.dialogueSessions.Get(voiceDialogueID, now)
|
|
dec = followUpMerge(prev, dec, now)
|
|
if !dec.Clarify {
|
|
h.rememberTurn(prev, dec, now)
|
|
}
|
|
}
|
|
|
|
// 2c. clarify — same as HandlePushToTalk: ask about the one missing thing.
|
|
if dec.Clarify {
|
|
if question, asked := h.askClarify(dec); asked {
|
|
return withNotice(expiredNotice, question)
|
|
}
|
|
}
|
|
|
|
// 3. action — execute the decision's intent.
|
|
replyText := h.applyAction(ctx, dec)
|
|
log.Printf("voice: applyAction returned: %q", replyText)
|
|
|
|
// 4. replier — phrase the reply when applyAction returned "".
|
|
if replyText == "" {
|
|
replyText = h.replier.Reply(dec)
|
|
}
|
|
return withNotice(expiredNotice, replyText)
|
|
}
|
|
|
|
// applyAction — executes the router's Decision. Intent-by-intent:
|
|
//
|
|
// - IntentFact: WriteFact via CoreAPI. Source = "tap:voice" (a voice
|
|
// capture is a tap; confidence 1.0).
|
|
// - IntentReminder: CreateReminder via CoreAPI.
|
|
// - IntentAct: tool-executor deferred (no-op today; the reply says so).
|
|
// - IntentNote / IntentQuery: chroma/RAG deferred (no-op; reply says so).
|
|
// - Clarify: the router's stage-3 fired; no action.
|
|
//
|
|
// Returns "" when the Replier should phrase the reply (the default path);
|
|
// returns a non-empty string when the action path wants to OVERRIDE the
|
|
// reply text (e.g. an action error the user should hear SPECIFICALLY, not
|
|
// a generic "ok"). Errors surface as a short reply text the user hears.
|
|
func (h *reactiveHandler) applyAction(ctx context.Context, dec router.Decision) string {
|
|
if dec.Clarify {
|
|
return "" // the Replier phrases clarify
|
|
}
|
|
if handler, ok := actionHandlers[dec.Intent]; ok {
|
|
return handler(h, ctx, dec)
|
|
}
|
|
return ""
|
|
}
|
|
|
|
// detectPattern extracts an event from the written fact and runs the pattern
|
|
// detector. If a stable recurring pattern is found and no proposed routine
|
|
// exists for this action+object yet, one is created and the user is prompted
|
|
// to confirm via the park() mechanism. Returns the suggestion phrase when a
|
|
// new proposal was created and parked; "" otherwise.
|
|
func (h *reactiveHandler) detectPattern(ctx context.Context, factID int64, key, value string, ts time.Time) string {
|
|
ev := pattern.Extract(factID, key, value, ts)
|
|
if ev == nil {
|
|
return "" // not an actionable event
|
|
}
|
|
if _, err := h.dataStore.CreateEvent(ctx, factID, ev.Action, ev.Object, ts); err != nil {
|
|
log.Printf("voice: create event: %v", err)
|
|
return ""
|
|
}
|
|
// Detect+propose (Vikunja #43) is shared with the digestion tick's
|
|
// proactive scan — see patterns.go. Event *extraction* above stays here,
|
|
// tied to this fact write; detection over the accumulated history does
|
|
// not need to happen right now for the voice path to have already done
|
|
// its job — it's dedupe-safe to also let the next tick find the same
|
|
// pattern independently.
|
|
r, id, err := detectAndPropose(ctx, h.dataStore, ev.Action, ev.Object, ts)
|
|
if err != nil {
|
|
log.Printf("voice: detect pattern %s/%s: %v", ev.Action, ev.Object, err)
|
|
return ""
|
|
}
|
|
if r == nil {
|
|
return "" // not enough data, too irregular, or already proposed/decided
|
|
}
|
|
log.Printf("voice: proposed routine: %s/%s every %.1f days", r.Action, r.Object, r.IntervalDays)
|
|
|
|
// Park the proposal for voice confirmation.
|
|
phrase := pattern.PhraseRoutine(r)
|
|
h.mu.Lock()
|
|
h.pendingRoutine = &pendingRoutineConfirm{
|
|
routineID: id,
|
|
action: r.Action,
|
|
object: r.Object,
|
|
interval: r.IntervalDays,
|
|
phrase: phrase,
|
|
expiry: ts.Add(confirmTTL),
|
|
}
|
|
h.mu.Unlock()
|
|
return phrase
|
|
}
|
|
|
|
// resolveQuietToggle — pre-route keyword check. Returns (reply, true) when
|
|
// the utterance is a quiet-on/off command; ("", false) otherwise. Called from
|
|
// HandlePushToTalk BEFORE the router so a classifier miscue can't drop it.
|
|
func (h *reactiveHandler) resolveQuietToggle(ctx context.Context, text string) (string, bool) {
|
|
u := strings.ToLower(strings.TrimSpace(text))
|
|
var on, off bool
|
|
// Match as whole-token phrases so "тихий" in "тихий режим включи" still
|
|
// catches, but "тихий" alone in "очень тихий сегодня день" doesn't fire.
|
|
// The confirm turn is handled above, so "да"/"нет" won't reach here.
|
|
for _, kw := range []string{"quiet on", "quiet mode", "тихий режим", "тихий", "не шуми", "не беспокоить", "тихо"} {
|
|
if strings.Contains(u, kw) {
|
|
on = true
|
|
break
|
|
}
|
|
}
|
|
if !on {
|
|
for _, kw := range []string{"quiet off", "quiet end", "громкий режим", "шумный режим", "отмени тихий", "выключи тихий", "не тихо"} {
|
|
if strings.Contains(u, kw) {
|
|
off = true
|
|
break
|
|
}
|
|
}
|
|
}
|
|
if !on && !off {
|
|
return "", false
|
|
}
|
|
val := "false"
|
|
reply := "тихий режим выключен."
|
|
if on {
|
|
val = "true"
|
|
reply = "тихий режим включён. буду реже напоминать."
|
|
}
|
|
if _, err := h.api.WriteFact(ctx, ipc.WriteFactReq{
|
|
Ts: h.now(),
|
|
Kind: "config",
|
|
Key: "quiet_hours",
|
|
Value: val,
|
|
Source: "tap:voice",
|
|
Confidence: 1.0,
|
|
}); err != nil {
|
|
log.Printf("voice: write quiet_hours: %v", err)
|
|
return "не получилось переключить тихий режим.", true
|
|
}
|
|
return reply, true
|
|
}
|
|
|
|
// replySystem answers system-observable queries using the handler's clock
|
|
// and (in future) system interfaces. The decision's utterance is parsed
|
|
// for keywords to determine what the user is asking about.
|
|
func (h *reactiveHandler) replySystem(ctx context.Context, dec router.Decision) string {
|
|
u := strings.ToLower(dec.Utterance)
|
|
now := h.now()
|
|
|
|
// stage-0 grammars catch the exact time/date patterns, but duration
|
|
// queries ("сколько времени прошло") bypass the grammar's build filter
|
|
// and can reach replySystem via the classifier path. Guard against them.
|
|
if hasDurationWords(u) {
|
|
return "пока не умею отвечать на этот вопрос."
|
|
}
|
|
|
|
switch {
|
|
case strings.Contains(u, "час") || strings.Contains(u, "врем"):
|
|
// "который час в киеве" — she keeps one clock, so any named place gets
|
|
// the honest answer. Never local time dressed up as the city's.
|
|
if mentionsUnknownPlace(u) {
|
|
return onlyLocalTimeReply
|
|
}
|
|
return "сейчас " + ruClock(now)
|
|
case strings.Contains(u, "день") || strings.Contains(u, "числ"):
|
|
// "какое число завтра" — answer for the day the user asked about,
|
|
// not today. Reuses the router's calendar day-word parser.
|
|
day := now
|
|
prefix := "сегодня"
|
|
if d, ok := router.ParseCalendarDate(u, now); ok {
|
|
day = d
|
|
prefix = dayPrefix(now, d)
|
|
} else if mentionsUnknownDay(u) {
|
|
// He named a day she cannot work out ("в пятницу", "через неделю").
|
|
// Answering today's date here would be the same silent wrong answer
|
|
// this arm was fixed for, so say what she can do instead.
|
|
return onlyNearDaysReply
|
|
}
|
|
dow := ruWeekdays[day.Weekday()]
|
|
month := ruMonths[day.Month()-1]
|
|
return fmt.Sprintf("%s %s, %d %s %d года", prefix, dow, day.Day(), month, day.Year())
|
|
case strings.Contains(u, "кто дома") || strings.Contains(u, "человек дома"):
|
|
return "присутствие пока не подключено к голосовому запросу."
|
|
case strings.Contains(u, "памят") || strings.Contains(u, "процессор") || strings.Contains(u, "загрузк") || strings.Contains(u, "статус") || strings.Contains(u, "работа") || strings.Contains(u, "сервис") || strings.Contains(u, "диск") || strings.Contains(u, "ip") || strings.Contains(u, "аптайм") || strings.Contains(u, "трафик") || strings.Contains(u, "интернет"):
|
|
return "системная статистика пока не подключена."
|
|
default:
|
|
return "пока не умею отвечать на этот вопрос."
|
|
}
|
|
}
|
|
|
|
// chatHistory collects dialogue turns from the session store for the current
|
|
// conversation. Returns prior user utterances (newest last) up to a depth of
|
|
// 4 turns. Returns nil when there's no session or no history.
|
|
func (h *reactiveHandler) chatHistory() []dialogue.Turn {
|
|
if h.dialogueSessions == nil {
|
|
return nil
|
|
}
|
|
now := h.now()
|
|
prev := h.dialogueSessions.Get(voiceDialogueID, now)
|
|
if prev == nil {
|
|
return nil
|
|
}
|
|
// History already includes the immediate prior turn (set by the dialogue
|
|
// merge at lines 373-395), plus up to 3 more from deeper history.
|
|
out := make([]dialogue.Turn, 0, 1+len(prev.History))
|
|
out = append(out, dialogue.Turn{
|
|
Intent: prev.Intent,
|
|
Slots: prev.Slots,
|
|
Text: prev.Slots.Text,
|
|
})
|
|
out = append(out, prev.History...)
|
|
return out
|
|
}
|
|
|
|
// reply wraps a text reply through TTS to produce a PushToTalkResp. If TTS
|
|
// fails, the response carries an empty audio + the text — the client can
|
|
// still display text if it can't play. The routedChannels field is
|
|
// reserved for a future "the dispatcher also forwarded to ntfy/telegram"
|
|
// reply (today the reactive path doesn't dispatch nudges; that's the loop
|
|
// tick's job).
|
|
func (h *reactiveHandler) reply(ctx context.Context, text string, _ []string) (voice.PushToTalkResp, error) {
|
|
spoken := ttsnorm.Speakable(text)
|
|
log.Printf("voice: reply → %q", text)
|
|
audioOut, err := h.tts.Synthesize(ctx, spoken)
|
|
if err != nil {
|
|
log.Printf("voice: tts error: %v", err)
|
|
return voice.PushToTalkResp{ReplyText: text, ReplyAudio: audio.Audio{}}, nil
|
|
}
|
|
return voice.PushToTalkResp{ReplyText: text, ReplyAudio: audioOut}, 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.
|
|
func pickLLMRouter(enabled bool, c *llm.Client) *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.
|
|
// - 6 bootstrap examples covering the 5 intents + one compound-capture
|
|
// placeholder. Spec calls for ~10 per intent at production; this is the
|
|
// bootstrapping floor swapped by tuning the seed set later.
|
|
// - 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()...)
|
|
grammars = append(grammars, router.ReminderGrammar())
|
|
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. Each file is named
|
|
// <intent>.txt and contains one training example per line (blank lines and
|
|
// lines starting with # are ignored). Relative to the working directory.
|
|
const seedDir = "models/seeds"
|
|
|
|
// 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). 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, seedDir)
|
|
}
|
|
|
|
func loadSeedFile(c *router.Classifier, intent router.Intent) (int, error) {
|
|
path := filepath.Join(seedDir, 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
|
|
}
|
|
|
|
// park stores a destructive act awaiting confirmation. Overwrites any prior
|
|
// pending (last-asked wins — single-user box).
|
|
func (h *reactiveHandler) park(fn string, args []string, phrase string) {
|
|
h.mu.Lock()
|
|
h.pending = &pendingAct{fn: fn, args: args, phrase: phrase, expiry: h.now().Add(confirmTTL)}
|
|
h.mu.Unlock()
|
|
}
|
|
|
|
// resolveConfirm interprets an utterance as the answer to a parked destructive
|
|
// act OR a parked routine proposal. Returns (reply, true) when it consumed the
|
|
// utterance as a y/n answer; ("", false) when there's nothing pending (or the
|
|
// parked act expired), so the caller routes the utterance normally. An
|
|
// unrecognised answer cancels the pending and routes normally — a confirm that
|
|
// can't be answered clearly is safer abandoned than left armed.
|
|
func (h *reactiveHandler) resolveConfirm(ctx context.Context, text string) (string, bool) {
|
|
h.mu.Lock()
|
|
defer h.mu.Unlock()
|
|
|
|
// Check routine proposal first (newer feature; checked before tool confirm
|
|
// so a routine confirm doesn't get eaten by a stale tool pending).
|
|
pr := h.pendingRoutine
|
|
if pr != nil && !h.now().After(pr.expiry) {
|
|
switch classifyConfirm(text) {
|
|
case confirmYes:
|
|
h.pendingRoutine = nil
|
|
// Only record the acceptance. The tick loop reads accepted
|
|
// routines and nudges on their own interval. Building a reminder
|
|
// here made a routine fire exactly once (Vikunja #366).
|
|
if err := h.dataStore.AcceptProposedRoutine(ctx, pr.routineID, h.now()); err != nil {
|
|
log.Printf("voice: accept proposed routine: %v", err)
|
|
return "не получилось запомнить рутину.", true
|
|
}
|
|
return "буду напоминать.", true
|
|
case confirmNo:
|
|
h.pendingRoutine = nil
|
|
if err := h.dataStore.DismissProposedRoutine(ctx, pr.routineID); err != nil {
|
|
log.Printf("voice: dismiss proposed routine: %v", err)
|
|
}
|
|
return "хорошо, не буду.", true
|
|
default:
|
|
// unclear: abandon the routine proposal, route normally.
|
|
h.pendingRoutine = nil
|
|
return "", false
|
|
}
|
|
}
|
|
// Clear expired routine if it existed.
|
|
if pr != nil {
|
|
h.pendingRoutine = nil
|
|
}
|
|
|
|
// Check pending Hexis execution confirm. Bound to the exact capability +
|
|
// target that was proposed; a stray "да" can only run that, nothing else.
|
|
if hx := h.pendingHexis; hx != nil {
|
|
if h.now().After(hx.expiry) {
|
|
h.pendingHexis = nil
|
|
} else {
|
|
switch classifyConfirm(text) {
|
|
case confirmYes:
|
|
h.pendingHexis = nil
|
|
return h.execHexis(ctx, hx.capabilityID, hx.capName, hx.entityID, hx.displayName), true
|
|
case confirmNo:
|
|
h.pendingHexis = nil
|
|
return "отменила.", true
|
|
default:
|
|
h.pendingHexis = nil
|
|
return "", false
|
|
}
|
|
}
|
|
}
|
|
|
|
// Check tool confirm (existing behavior).
|
|
p := h.pending
|
|
if p == nil {
|
|
return "", false
|
|
}
|
|
if h.now().After(p.expiry) {
|
|
h.pending = nil
|
|
return "", false
|
|
}
|
|
switch classifyConfirm(text) {
|
|
case confirmYes:
|
|
h.pending = nil
|
|
out, err := h.tools.Exec(ctx, p.fn, p.args, true) // confirmed
|
|
if err != nil {
|
|
log.Printf("voice: tool %s (confirmed): %v", p.fn, err)
|
|
if out != "" {
|
|
return "не получилось выполнить команду: " + firstLine(out), true
|
|
}
|
|
return "не получилось выполнить команду.", true
|
|
}
|
|
if out != "" {
|
|
return "готово: " + firstLine(out), true
|
|
}
|
|
return "готово.", true
|
|
case confirmNo:
|
|
h.pending = nil
|
|
return "отменила.", true
|
|
default:
|
|
// unclear answer: abandon the confirm, route this utterance normally.
|
|
h.pending = nil
|
|
return "", false
|
|
}
|
|
}
|
|
|
|
// proposeGap scaffolds a 'proposed' tool for an act whose verb isn't enabled.
|
|
// maven drafts the registration (name = the verb, provenance = the utterance);
|
|
// a human enables it on the authed surface. She suggests, never enables.
|
|
func (h *reactiveHandler) proposeGap(ctx context.Context, dec router.Decision) string {
|
|
name := firstWord(stripWake(dec.Utterance))
|
|
if name == "" {
|
|
return "не разобрала команду — попробуй иначе."
|
|
}
|
|
newly, err := h.api.ProposeTool(ctx, name, dec.Utterance, "", h.now())
|
|
if err != nil {
|
|
log.Printf("voice: propose tool %q: %v", name, err)
|
|
return "команды «" + name + "» нет в списке разрешённых."
|
|
}
|
|
if newly {
|
|
return "команды «" + name + "» нет в списке. Предложила её добавить — включи через клиент."
|
|
}
|
|
return "команды «" + name + "» пока нет в списке — она уже предложена, включи через клиент."
|
|
}
|
|
|
|
// confirmVerdict — the parse of a y/n confirm answer.
|
|
type confirmVerdict int
|
|
|
|
const (
|
|
confirmUnknown confirmVerdict = iota
|
|
confirmYes
|
|
confirmNo
|
|
)
|
|
|
|
// classifyConfirm reads a short ru/en yes-or-no answer. Substring match on the
|
|
// stems so inflections/fillers ("да, давай", "нет, отмени") still land.
|
|
func classifyConfirm(text string) confirmVerdict {
|
|
t := strings.ToLower(strings.TrimSpace(text))
|
|
// negatives first — "не надо" contains no "да", but check no-stems before
|
|
// yes so a leading "нет" isn't shadowed.
|
|
for _, no := range []string{"нет", "не надо", "отмен", "стоп", "no", "cancel", "stop", "don't"} {
|
|
if strings.Contains(t, no) {
|
|
return confirmNo
|
|
}
|
|
}
|
|
for _, yes := range []string{"да", "ага", "давай", "подтвер", "конечно", "yes", "yeah", "yep", "confirm", "ок", "okay", "ok"} {
|
|
if strings.Contains(t, yes) {
|
|
return confirmYes
|
|
}
|
|
}
|
|
return confirmUnknown
|
|
}
|
|
|
|
// actPhrase renders "fn arg1 arg2" for the confirm prompt.
|
|
func actPhrase(fn string, args []string) string {
|
|
if len(args) == 0 {
|
|
return fn
|
|
}
|
|
return fn + " " + strings.Join(args, " ")
|
|
}
|
|
|
|
// 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)
|
|
}
|
|
|
|
// reembedOnStart is the -reembed flag (set in run()). Opt-in on purpose: see
|
|
// runReembed.
|
|
var reembedOnStart 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)
|
|
}
|
|
}
|