3dbf67f8f9
The user only ever asks the time in his own zone, so answering other cities was code kept in step with the weather city list for no gain. Any named place now gets the honest "local time only" answer that was already there for unknown cities. Removes the 22-entry table, the lookup and the embedded tz database. Closes Vikunja #389 — there is only one city list again. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
1947 lines
73 KiB
Go
1947 lines
73 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"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"log"
|
||
"os"
|
||
"path/filepath"
|
||
"strconv"
|
||
"strings"
|
||
"sync"
|
||
"time"
|
||
|
||
hexisclient "github.com/kami/hexis/pkg/client"
|
||
"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)
|
||
}
|
||
|
||
// ----- 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
|
||
}
|
||
switch dec.Intent {
|
||
case router.IntentFact:
|
||
if !dec.Slots.HasKey {
|
||
return "не разобрала, что записать — попробуй иначе."
|
||
}
|
||
now := h.now()
|
||
req := ipc.WriteFactReq{
|
||
Ts: now,
|
||
Kind: "self",
|
||
Key: dec.Slots.Key,
|
||
Value: dec.Slots.Value,
|
||
Source: "tap:voice",
|
||
Confidence: 1.0,
|
||
// Subject: the key doubles as the entity-resolution candidate —
|
||
// a voice-tapped fact's key is usually the thing/person it's
|
||
// about ("espresso_machine", "kate"), so queueing it for Nexus
|
||
// resolution costs one async lookup and is a no-op (not_found)
|
||
// for the abstract self-state keys (mood, water) that aren't
|
||
// entities at all.
|
||
Subject: dec.Slots.Key,
|
||
}
|
||
factID, err := h.api.WriteFact(ctx, req)
|
||
if err != nil {
|
||
log.Printf("voice: write fact: %v", err)
|
||
return "не получилось сохранить факт."
|
||
}
|
||
// Index the fact utterance in long-term memory (best-effort, must not
|
||
// fail the fact write). Facts aren't in the notes table, so this is the
|
||
// only recall path for them — "когда я пил воду?" reads back from here.
|
||
if h.memStore != nil {
|
||
if vec, err := router.EmbedPassage(ctx, h.embedder, dec.Utterance); err != nil {
|
||
log.Printf("voice: embed fact for memory: %v", err)
|
||
} else if err := h.memStore.Insert(ctx, "fact:"+dec.Slots.Key+":"+strconv.FormatInt(now.Unix(), 10), vec, map[string]string{
|
||
"source": "voice",
|
||
"type": "fact",
|
||
"text": dec.Utterance,
|
||
"ts": strconv.FormatInt(now.Unix(), 10),
|
||
}); err != nil {
|
||
log.Printf("voice: memory insert fact: %v", err)
|
||
}
|
||
}
|
||
// Event extraction + pattern detection (best-effort, must not fail the
|
||
// fact write). If the fact describes a recognizable action, it becomes a
|
||
// normalized event; if ≥3 events for the same action+object show stable
|
||
// intervals, a proposed routine is created and parked for confirmation.
|
||
if h.dataStore != nil {
|
||
if phrase := h.detectPattern(ctx, factID, dec.Slots.Key, dec.Slots.Value, now); phrase != "" {
|
||
return phrase // "ты заправляешь ... напоминать?"
|
||
}
|
||
}
|
||
return "" // replier phrases the success reply
|
||
|
||
case router.IntentReminder:
|
||
if !dec.Slots.HasTime {
|
||
// Stage-0 (reminder-wakeword grammar) skips the extractor, so the
|
||
// time wasn't parsed. Run the parser as a fallback.
|
||
if dec.Stage == 0 && h.timeParser != nil {
|
||
t, ok, err := h.timeParser.Parse(ctx, dec.Utterance, h.now())
|
||
if err == nil && ok {
|
||
dec.Slots.Time = t
|
||
dec.Slots.HasTime = true
|
||
}
|
||
}
|
||
if !dec.Slots.HasTime {
|
||
return "не получилось разобрать время напоминания."
|
||
}
|
||
}
|
||
payload := `{"text":` + jsonString(dec.Utterance) + `}`
|
||
if _, err := h.api.CreateReminder(ctx, dec.Slots.Time, payload, ""); err != nil {
|
||
log.Printf("voice: create reminder: %v", err)
|
||
return "не получилось поставить напоминание."
|
||
}
|
||
return ""
|
||
|
||
case router.IntentAct:
|
||
// tool executor: run the matched fn against the enabled allowlist.
|
||
// HasFn=false ⇒ try the matcher (for LLM-routed acts where the verb
|
||
// didn't go through the stage-0 act grammar).
|
||
if !dec.Slots.HasFn && dec.Slots.Text != "" && h.matcher != nil {
|
||
if fn, args, ok := h.matcher.Match(dec.Slots.Text); ok {
|
||
dec.Slots.Fn, dec.Slots.Args, dec.Slots.HasFn = fn, args, true
|
||
}
|
||
}
|
||
|
||
// Praxis ecosystem tools: intercept before the system command executor.
|
||
if h.ecosystem != nil && h.ecosystem.praxis != nil && dec.Slots.HasFn {
|
||
if reply := h.handlePraxisAct(ctx, dec); reply != "" {
|
||
return reply
|
||
}
|
||
}
|
||
|
||
// Hexis ecosystem action: if ecosystem is configured and we have a verb
|
||
// + entity text, try to resolve the entity and execute via Hexis.
|
||
if h.ecosystem != nil && h.ecosystem.hexis != nil && dec.Slots.Text != "" {
|
||
if reply := h.handleHexisAct(ctx, dec); reply != "" {
|
||
return reply
|
||
}
|
||
}
|
||
|
||
// HasFn still false ⇒ no allowlist match: scaffold a 'proposed' tool
|
||
// the user can enable on the authed surface ("earn the right to ask").
|
||
if !dec.Slots.HasFn {
|
||
return h.proposeGap(ctx, dec)
|
||
}
|
||
out, err := h.tools.Exec(ctx, dec.Slots.Fn, dec.Slots.Args, false)
|
||
if err != nil {
|
||
switch {
|
||
case errors.Is(err, tool.ErrNeedsConfirm):
|
||
// destructive: park it and ask. The next utterance answers.
|
||
phrase := actPhrase(dec.Slots.Fn, dec.Slots.Args)
|
||
h.park(dec.Slots.Fn, dec.Slots.Args, phrase)
|
||
return "выполнить «" + phrase + "»? скажи «да» или «нет»."
|
||
case errors.Is(err, tool.ErrNotEnabled):
|
||
return h.proposeGap(ctx, dec)
|
||
}
|
||
log.Printf("voice: tool %s: %v", dec.Slots.Fn, err)
|
||
if out != "" {
|
||
return "не получилось выполнить команду: " + firstLine(out)
|
||
}
|
||
return "не получилось выполнить команду."
|
||
}
|
||
if out != "" {
|
||
return "готово: " + firstLine(out)
|
||
}
|
||
return "готово."
|
||
|
||
case router.IntentChat:
|
||
// Conversational: build history from dialogue session (prior user turns)
|
||
// and let the LLM respond from general knowledge + context.
|
||
history := h.chatHistory()
|
||
reply, err := h.phraser.PhraseChat(ctx, dec.Utterance, history)
|
||
if err != nil {
|
||
log.Printf("voice: chat: %v", err)
|
||
return "поговорили."
|
||
}
|
||
return reply
|
||
|
||
case router.IntentSystem:
|
||
return h.replySystem(ctx, dec)
|
||
|
||
case router.IntentNote:
|
||
// embed the note text with the same model the classifier uses, persist
|
||
// via CoreAPI (source=tap:voice). Semantic recall lives in `notes`, not
|
||
// facts — no predicate reads it (spec's two-memory split).
|
||
vec, err := router.EmbedPassage(ctx, h.embedder, dec.Utterance)
|
||
if err != nil {
|
||
log.Printf("voice: embed note: %v", err)
|
||
return "не получилось сохранить заметку."
|
||
}
|
||
noteTs := h.now()
|
||
noteID, err := h.api.WriteNote(ctx, noteTs, dec.Utterance, vec, "tap:voice")
|
||
if err != nil {
|
||
log.Printf("voice: write note: %v", err)
|
||
return "не получилось сохранить заметку."
|
||
}
|
||
// Insert into long-term memory (best-effort, must not fail the note write).
|
||
// text/ts in the meta make a Search hit self-describing (see bestRecall).
|
||
if h.memStore != nil {
|
||
if err := h.memStore.Insert(ctx, "note:"+strconv.FormatInt(noteID, 10), vec, map[string]string{
|
||
"source": "voice",
|
||
"type": "note",
|
||
"text": dec.Utterance,
|
||
"ts": strconv.FormatInt(noteTs.Unix(), 10),
|
||
}); err != nil {
|
||
log.Printf("voice: memory insert: %v", err)
|
||
}
|
||
}
|
||
return "" // replier phrases the "saved" reply
|
||
|
||
case router.IntentQuery:
|
||
// Fact-by-key lookup: when the dialogue layer resolved an anaphoric
|
||
// reference to a prior fact's key (e.g. "когда я это сделал?" after
|
||
// "запиши что я пил воду"), look up the fact's value directly.
|
||
if dec.Slots.HasKey && dec.Slots.Key != "" {
|
||
if f, err := h.api.LatestFact(ctx, dec.Slots.Key); err == nil {
|
||
if dec.Slots.HasTime {
|
||
// The query asks about timing — the fact's own timestamp
|
||
// is the answer it's looking for. Format as a natural reply.
|
||
reply := fmt.Sprintf("я записала это %s", formatTime(f.Ts))
|
||
return reply
|
||
}
|
||
// General fact reference: describe what we know.
|
||
if dec.Utterance == "" {
|
||
return fmt.Sprintf("вот что я знаю: %s — %s", dec.Slots.Key, f.Value)
|
||
}
|
||
// The utterance still carries the question; fall through to
|
||
// normal RAG with the resolved key in context.
|
||
}
|
||
}
|
||
|
||
// Calendar questions: "что у меня сегодня?", "планы на завтра?"
|
||
// h.now(), not time.Now(): the handler's clock is the injected one, so
|
||
// this arm can be tested at a fixed time like the rest.
|
||
if date, ok := router.ParseCalendarDate(dec.Utterance, h.now()); ok {
|
||
events, err := h.api.CalendarEvents(ctx, date, date.Add(24*time.Hour))
|
||
if err != nil {
|
||
log.Printf("voice: calendar events: %v", err)
|
||
return "не получилось проверить календарь."
|
||
}
|
||
values := make([]string, len(events))
|
||
for i, e := range events {
|
||
values[i] = e.Value
|
||
}
|
||
var f router.CalendarEventFormatter
|
||
return f.Format(values, date)
|
||
}
|
||
|
||
// Weather questions
|
||
if isWeatherQuery(dec.Utterance) {
|
||
loc := extractWeatherLocation(dec.Utterance, h.weatherLocation)
|
||
ctxWT, cancel := context.WithTimeout(ctx, 5*time.Second)
|
||
defer cancel()
|
||
w, err := h.weatherProvider.CurrentWeather(ctxWT, loc)
|
||
if errors.Is(err, weather.ErrNotConfigured) {
|
||
return "погода не настроена."
|
||
}
|
||
if err != nil {
|
||
log.Printf("voice: weather: %v", err)
|
||
return "не получилось узнать погоду."
|
||
}
|
||
return fmt.Sprintf("в %s сейчас %.0f градусов, %s.", w.Location, w.Temperature, w.Condition)
|
||
}
|
||
|
||
vec, err := router.EmbedQuery(ctx, h.embedder, dec.Utterance)
|
||
if err != nil {
|
||
log.Printf("voice: embed query: %v", err)
|
||
return "не получилось найти ответ."
|
||
}
|
||
// Long-term memory first: ONE search over everything Maven remembers
|
||
// (notes and facts share this index) and ONE confidence gate, so the
|
||
// memory that is clearly the best match answers — a note just as much
|
||
// as a fact.
|
||
//
|
||
// This used to run only after the notes-only gate below had already
|
||
// rejected the same note at the same score, which no note could ever
|
||
// survive a second time: the branch could only return a fact (#373).
|
||
// Order, not the gate, was the bug — the set of questions Maven answers
|
||
// is unchanged, only which memory gets to answer them.
|
||
if h.memStore != nil {
|
||
if hits, herr := h.memStore.Search(ctx, vec, 3); herr == nil {
|
||
if hit, ok := bestRecall(hits, h.queryMinScore, h.queryMinMargin); ok {
|
||
text := hit.Meta["text"]
|
||
// A note is phrased in Maven's voice; a fact is read back
|
||
// as it was stored.
|
||
if hit.Meta["type"] == "note" {
|
||
if reply, perr := h.phraser.PhraseQuery(ctx, dec.Utterance, []string{text}); perr == nil && reply != "" {
|
||
return reply
|
||
}
|
||
}
|
||
return text
|
||
}
|
||
} else {
|
||
log.Printf("voice: memory search: %v", herr)
|
||
}
|
||
}
|
||
|
||
notes, err := h.api.QueryNotes(ctx, vec, 5)
|
||
if err != nil {
|
||
log.Printf("voice: query notes: %v", err)
|
||
return "не получилось найти ответ."
|
||
}
|
||
// Notes-only pass, for notes the vector index above does not hold (an
|
||
// older note written before it existed). Same gate, notes-only
|
||
// candidates.
|
||
//
|
||
// Confidence gate: below it, say "I don't know" rather than read back
|
||
// the least-unrelated note — a confident wrong recall is worse than a
|
||
// gap (spec's "not a guesser-of-truth"). Same instinct as the loop's
|
||
// since(key)==null → don't fire. Two parts: an absolute cosine floor,
|
||
// and a margin over the runner-up, which is the part that works with
|
||
// the e5 embedder's narrow score band. See memory.Confident.
|
||
noteScores := make([]float64, len(notes))
|
||
for i, n := range notes {
|
||
noteScores[i] = n.Score
|
||
}
|
||
if !memory.ConfidentScores(noteScores, h.queryMinScore, h.queryMinMargin) {
|
||
// Try general knowledge from the phraser before giving up
|
||
reply, err := h.phraser.PhraseQuery(ctx, dec.Utterance, nil)
|
||
if err != nil || reply == "" {
|
||
return "не знаю."
|
||
}
|
||
return reply
|
||
}
|
||
texts := make([]string, len(notes))
|
||
for i, n := range notes {
|
||
texts[i] = n.Text
|
||
}
|
||
reply, err := h.phraser.PhraseQuery(ctx, dec.Utterance, texts)
|
||
if err != nil {
|
||
log.Printf("voice: phrase query: %v", err)
|
||
}
|
||
if reply == "" {
|
||
reply = "вот что я нашла: " + texts[0]
|
||
}
|
||
return reply
|
||
}
|
||
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 ""
|
||
}
|
||
events, err := h.dataStore.EventsFor(ctx, ev.Action, ev.Object)
|
||
if err != nil {
|
||
log.Printf("voice: events for %s/%s: %v", ev.Action, ev.Object, err)
|
||
return ""
|
||
}
|
||
// Convert store.Events to pattern.Events for the detector.
|
||
patEvents := make([]pattern.Event, len(events))
|
||
for i, e := range events {
|
||
patEvents[i] = pattern.Event{
|
||
FactID: e.FactID,
|
||
Action: e.Action,
|
||
Object: e.Object,
|
||
Ts: e.Ts,
|
||
}
|
||
}
|
||
r, err := pattern.Detect(patEvents)
|
||
if err != nil {
|
||
log.Printf("voice: pattern detect: %v", err)
|
||
return ""
|
||
}
|
||
if r == nil {
|
||
return "" // not enough data or intervals too irregular
|
||
}
|
||
// Check if already proposed/accepted/dismissed for this pair.
|
||
existing, err := h.dataStore.LookupProposedRoutine(ctx, r.Action, r.Object)
|
||
if err != nil {
|
||
log.Printf("voice: lookup proposed routine: %v", err)
|
||
return ""
|
||
}
|
||
if existing != nil {
|
||
return "" // already proposed, accepted, or dismissed
|
||
}
|
||
id, err := h.dataStore.CreateProposedRoutine(ctx, r.Action, r.Object, r.IntervalDays, ts)
|
||
if err != nil {
|
||
log.Printf("voice: create proposed routine: %v", err)
|
||
return ""
|
||
}
|
||
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
|
||
}
|
||
|
||
var ruWeekdays = []string{
|
||
"воскресенье", "понедельник", "вторник", "среда",
|
||
"четверг", "пятница", "суббота",
|
||
}
|
||
|
||
var ruMonths = []string{
|
||
"января", "февраля", "марта", "апреля", "мая", "июня",
|
||
"июля", "августа", "сентября", "октября", "ноября", "декабря",
|
||
}
|
||
|
||
// onlyLocalTimeReply — the honest answer when the user asks the time somewhere
|
||
// other than here. She only keeps one clock, and saying so is better than
|
||
// naming the wrong city's time.
|
||
//
|
||
// There used to be a city→time-zone table here. It was removed on purpose: the
|
||
// user only ever asks for local time, so the table was a second list of cities
|
||
// to keep in step with the weather one for no gain.
|
||
const onlyLocalTimeReply = "я знаю только местное время, про другие города пока не скажу."
|
||
|
||
// notPlaceAfterV — words that follow "в" without naming a place, so
|
||
// mentionsUnknownPlace does not mistake them for a city.
|
||
var notPlaceAfterV = map[string]bool{
|
||
"данный": true, "данную": true, "этот": true, "эту": true,
|
||
"котором": true, "какое": true, "какой": true, "который": true,
|
||
"общем": true, "точности": true, "курсе": true, "сутках": true,
|
||
"часах": true, "минутах": true, "секундах": true, "неделе": true,
|
||
}
|
||
|
||
// mentionsUnknownPlace reports whether the question has a "в <слово>" phrase
|
||
// that looks like a place we do not know ("который час в киеве"). Used only to
|
||
// pick the honest "local time only" reply instead of answering local time as
|
||
// if it were the city's.
|
||
func mentionsUnknownPlace(u string) bool {
|
||
toks := strings.Fields(u)
|
||
for i := 0; i+1 < len(toks); i++ {
|
||
if toks[i] != "в" && toks[i] != "во" {
|
||
continue
|
||
}
|
||
next := strings.Trim(toks[i+1], ".,?!")
|
||
if next == "" || notPlaceAfterV[next] {
|
||
continue
|
||
}
|
||
// A number after "в" is a clock ("в 5 часов"), not a place.
|
||
if _, err := strconv.Atoi(strings.SplitN(next, ":", 2)[0]); err == nil {
|
||
continue
|
||
}
|
||
return true
|
||
}
|
||
return false
|
||
}
|
||
|
||
// onlyNearDaysReply — she can work out today, tomorrow, the day after and
|
||
// yesterday, and nothing further. Said out loud instead of answering today's
|
||
// date for a day she did not understand.
|
||
const onlyNearDaysReply = "я считаю только сегодня, завтра, послезавтра и вчера — про другие дни пока не скажу."
|
||
|
||
// dayWords — day references the calendar parser cannot resolve. A weekday name
|
||
// or a "через …" phrase means he asked about a specific other day.
|
||
var dayWords = []string{
|
||
"понедельник", "вторник", "сред", "четверг", "пятниц", "суббот", "воскресен",
|
||
"через", "monday", "tuesday", "wednesday", "thursday", "friday", "saturday", "sunday",
|
||
}
|
||
|
||
// mentionsUnknownDay reports whether the question names a day the calendar
|
||
// parser could not resolve. Mirror of mentionsUnknownPlace: it exists only to
|
||
// pick an honest reply over a confidently wrong one.
|
||
//
|
||
// Only called after ParseCalendarDate has already failed, so "завтра" and the
|
||
// other words it does know never reach here.
|
||
func mentionsUnknownDay(u string) bool {
|
||
for _, w := range dayWords {
|
||
if strings.Contains(u, w) {
|
||
return true
|
||
}
|
||
}
|
||
return false
|
||
}
|
||
|
||
// ruClock renders the clock part of the time reply: "15 часов 4 минуты".
|
||
func ruClock(t time.Time) string {
|
||
h, m := t.Hour(), t.Minute()
|
||
hourWord := ruPlural(h, "час", "часа", "часов")
|
||
if m == 0 {
|
||
return fmt.Sprintf("%d %s ровно", h, hourWord)
|
||
}
|
||
return fmt.Sprintf("%d %s %d %s", h, hourWord, m, ruPlural(m, "минута", "минуты", "минут"))
|
||
}
|
||
|
||
// dayPrefix names the day relative to now ("завтра", "вчера", …) so the date
|
||
// reply opens the way a person would say it.
|
||
func dayPrefix(now, day time.Time) string {
|
||
base := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, now.Location())
|
||
switch int(day.Sub(base).Hours() / 24) {
|
||
case -1:
|
||
return "вчера"
|
||
case 0:
|
||
return "сегодня"
|
||
case 1:
|
||
return "завтра"
|
||
case 2:
|
||
return "послезавтра"
|
||
}
|
||
return "это"
|
||
}
|
||
|
||
func ruPlural(n int, one, two, many string) string {
|
||
n = n % 100
|
||
if n > 10 && n < 20 {
|
||
return many
|
||
}
|
||
n = n % 10
|
||
switch n {
|
||
case 1:
|
||
return one
|
||
case 2, 3, 4:
|
||
return two
|
||
default:
|
||
return many
|
||
}
|
||
}
|
||
|
||
// 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
|
||
}
|
||
|
||
// hasDurationWords checks whether u is asking about elapsed/remaining time
|
||
// rather than the current clock — guards replySystem from replying "сейчас
|
||
// X часов" to "сколько времени прошло". Mirrors the stage0.go build filter.
|
||
func hasDurationWords(u string) bool {
|
||
s := strings.ToLower(strings.TrimSpace(u))
|
||
// First-word duration markers (same keywords as timeQueryBuild in stage0).
|
||
first := strings.Fields(s)
|
||
if len(first) > 0 {
|
||
switch first[0] {
|
||
case "прошло", "осталось", "пройдет", "минуло", "проходит":
|
||
return true
|
||
}
|
||
}
|
||
// Broader duration keywords appearing anywhere in the utterance.
|
||
if strings.Contains(s, "прошло") || strings.Contains(s, "осталось") {
|
||
return true
|
||
}
|
||
if strings.Contains(s, " до ") {
|
||
return true
|
||
}
|
||
return false
|
||
}
|
||
|
||
// 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
|
||
}
|
||
|
||
// handlePraxisAct — dispatches ecosystem tool acts through the Praxis tools API.
|
||
// Returns "" when the act is not a Praxis verb (the caller falls through to the
|
||
// system command executor). Returns a reply string otherwise.
|
||
func (h *reactiveHandler) handlePraxisAct(ctx context.Context, dec router.Decision) string {
|
||
if h.ecosystem == nil || h.ecosystem.praxis == nil {
|
||
return ""
|
||
}
|
||
px := h.ecosystem.praxis
|
||
fn := dec.Slots.Fn
|
||
|
||
// Map verbs and Russian aliases to Praxis tool calls.
|
||
// Each case: if the verb matches, call the tool and return a user-facing reply.
|
||
switch fn {
|
||
case "list_attention", "attention", "внимание", "что требует внимания", "что нового":
|
||
items, err := px.ListAttention(ctx, 20)
|
||
if err != nil {
|
||
log.Printf("ecosystem: praxis attention: %v", err)
|
||
return "не могу сейчас узнать, что требует внимания."
|
||
}
|
||
if len(items) == 0 {
|
||
return "ничего не требует внимания."
|
||
}
|
||
h.recordPraxisTrace(ctx, "list_attention", map[string]any{"count": len(items)})
|
||
var parts []string
|
||
for _, item := range items {
|
||
title, _ := item["title"].(string)
|
||
// importance arrives as JSON number ⇒ float64 over the HTTP contract.
|
||
importance, _ := item["importance"].(float64)
|
||
rule, _ := item["rule"].(string)
|
||
s := title
|
||
if importance > 0 {
|
||
s += fmt.Sprintf(" (важность %d", int(importance))
|
||
if rule != "" {
|
||
s += ": " + rule
|
||
}
|
||
s += ")"
|
||
}
|
||
parts = append(parts, s)
|
||
|
||
// Speaking an item surfaces it, it does not acknowledge it
|
||
// (ECOSYSTEM-SPEC.md §2.3: surfaced != acknowledged). Best-effort:
|
||
// a failed surface call must not block delivering the digest.
|
||
if id, ok := item["id"].(string); ok && id != "" {
|
||
if _, err := px.Surface(ctx, id); err != nil {
|
||
log.Printf("ecosystem: praxis surface %s: %v", id, err)
|
||
}
|
||
}
|
||
}
|
||
return "требует внимания: " + strings.Join(parts, "; ")
|
||
|
||
case "acknowledge_item", "принято", "понял", "поняла":
|
||
id := dec.Slots.Value
|
||
if id == "" {
|
||
return "какой пункт отметить принятым?"
|
||
}
|
||
if _, err := px.Acknowledge(ctx, id); err != nil {
|
||
log.Printf("ecosystem: praxis acknowledge %s: %v", id, err)
|
||
return "не получилось отметить принятым."
|
||
}
|
||
h.recordPraxisTrace(ctx, "acknowledge", map[string]any{"item_id": id})
|
||
return "принято."
|
||
|
||
case "resolve_item", "сделано", "готово", "решено":
|
||
id := dec.Slots.Value
|
||
if id == "" {
|
||
return "какой пункт отметить сделанным?"
|
||
}
|
||
if _, err := px.Resolve(ctx, id); err != nil {
|
||
log.Printf("ecosystem: praxis resolve %s: %v", id, err)
|
||
return "не получилось отметить сделанным."
|
||
}
|
||
h.recordPraxisTrace(ctx, "resolve", map[string]any{"item_id": id})
|
||
return "отмечено как сделано."
|
||
|
||
case "ignore_item", "игнорировать", "неважно":
|
||
id := dec.Slots.Value
|
||
if id == "" {
|
||
return "какой пункт игнорировать?"
|
||
}
|
||
if _, err := px.Ignore(ctx, id); err != nil {
|
||
log.Printf("ecosystem: praxis ignore %s: %v", id, err)
|
||
return "не получилось проигнорировать."
|
||
}
|
||
h.recordPraxisTrace(ctx, "ignore", map[string]any{"item_id": id})
|
||
return "проигнорировано."
|
||
|
||
case "pin_item", "закрепить":
|
||
id := dec.Slots.Value
|
||
if id == "" {
|
||
return "какой пункт закрепить?"
|
||
}
|
||
if _, err := px.Pin(ctx, id, true); err != nil {
|
||
log.Printf("ecosystem: praxis pin %s: %v", id, err)
|
||
return "не получилось закрепить."
|
||
}
|
||
h.recordPraxisTrace(ctx, "pin", map[string]any{"item_id": id})
|
||
return "закреплено."
|
||
|
||
case "list_changes", "changes", "изменения", "что изменилось":
|
||
changes, err := px.ListChanges(ctx, 20)
|
||
if err != nil {
|
||
log.Printf("ecosystem: praxis changes: %v", err)
|
||
return "не могу сейчас узнать об изменениях."
|
||
}
|
||
if len(changes) == 0 {
|
||
return "нет изменений."
|
||
}
|
||
h.recordPraxisTrace(ctx, "list_changes", map[string]any{"count": len(changes)})
|
||
var parts []string
|
||
for _, c := range changes {
|
||
title, _ := c["title"].(string)
|
||
typ, _ := c["change_type"].(string)
|
||
parts = append(parts, fmt.Sprintf("%s (%s)", title, typ))
|
||
}
|
||
return "изменения: " + strings.Join(parts, "; ")
|
||
|
||
default:
|
||
// Not a Praxis verb — let the caller fall through.
|
||
return ""
|
||
}
|
||
}
|
||
|
||
// recordPraxisTrace — writes a fact recording a cross-service ecosystem call.
|
||
// The fact is stored with source "praxis:trace" so the proactive loop can
|
||
// reference it and the dashboard can display recent ecosystem activity.
|
||
func (h *reactiveHandler) recordPraxisTrace(ctx context.Context, operation string, details map[string]any) {
|
||
now := h.now()
|
||
value := operation
|
||
if len(details) > 0 {
|
||
if b, err := json.Marshal(details); err == nil {
|
||
value = operation + " " + string(b)
|
||
}
|
||
}
|
||
_, _ = h.api.WriteFact(ctx, ipc.WriteFactReq{
|
||
Ts: now,
|
||
Kind: "system",
|
||
Key: "praxis:" + operation,
|
||
Value: value,
|
||
Source: "praxis:trace",
|
||
Confidence: 1.0,
|
||
})
|
||
}
|
||
|
||
// handleHexisAct — resolves entity references through Nexus and executes
|
||
// matching capabilities through Hexis. Returns a reply string when handled,
|
||
// or "" to fall through to the system command executor.
|
||
func (h *reactiveHandler) handleHexisAct(ctx context.Context, dec router.Decision) string {
|
||
if h.ecosystem == nil {
|
||
return ""
|
||
}
|
||
|
||
// Resolve the utterance text as an entity reference through Nexus. An
|
||
// ambiguous match must stop and clarify — never guess a mutation target.
|
||
entityID, displayName, ambiguous, err := h.ecosystem.resolveEntityReference(ctx, dec.Slots.Text, nil)
|
||
if err != nil {
|
||
// A genuine Nexus dependency failure, not "no such entity" — stop here
|
||
// and report degradation rather than silently falling through to the
|
||
// local command executor (ECOSYSTEM-SPEC.md: services degrade
|
||
// independently, never a silent all-clear).
|
||
return "экосистема недоступна, попробуй ещё раз."
|
||
}
|
||
if len(ambiguous) > 0 {
|
||
return "уточни, что именно: " + strings.Join(ambiguous, ", ") + "?"
|
||
}
|
||
if entityID == "" {
|
||
return ""
|
||
}
|
||
|
||
// Discover Hexis capabilities for this entity. A resolved entity with a
|
||
// genuine Hexis failure must not be treated as "no capabilities" and
|
||
// fall through to unrelated local execution.
|
||
caps, err := h.ecosystem.discoverCapabilities(ctx, entityID)
|
||
if err != nil {
|
||
return "экосистема недоступна, попробуй ещё раз."
|
||
}
|
||
if len(caps) == 0 {
|
||
return ""
|
||
}
|
||
|
||
// Match the user's verb to a capability by name/description. Collect all
|
||
// matches: more than one is itself ambiguous, so we ask rather than pick
|
||
// the first (ecosystem invariant: no arbitrary target for mutation).
|
||
verb := dec.Slots.Fn
|
||
if verb == "" {
|
||
verb = dec.Slots.Text
|
||
}
|
||
verbLower := strings.ToLower(verb)
|
||
|
||
var matches []*hexisclient.Capability
|
||
for i, c := range caps {
|
||
if strings.Contains(strings.ToLower(c.Name), verbLower) ||
|
||
(c.Description != "" && strings.Contains(strings.ToLower(c.Description), verbLower)) {
|
||
matches = append(matches, &caps[i])
|
||
}
|
||
}
|
||
if len(matches) == 0 {
|
||
return ""
|
||
}
|
||
if len(matches) > 1 {
|
||
var names []string
|
||
for _, m := range matches {
|
||
names = append(names, m.Name)
|
||
}
|
||
return "какую команду для " + displayName + ": " + strings.Join(names, ", ") + "?"
|
||
}
|
||
matched := matches[0]
|
||
|
||
// Read-only capabilities run immediately; mutating ones are parked for an
|
||
// explicit spoken confirm bound to this capability + target.
|
||
if !matched.ReadOnly {
|
||
h.mu.Lock()
|
||
h.pendingHexis = &pendingHexisExec{
|
||
capabilityID: matched.ID,
|
||
capName: matched.Name,
|
||
entityID: entityID,
|
||
displayName: displayName,
|
||
expiry: h.now().Add(confirmTTL),
|
||
}
|
||
h.mu.Unlock()
|
||
return "выполнить «" + matched.Name + "» для " + displayName + "? скажи «да» или «нет»."
|
||
}
|
||
|
||
return h.execHexis(ctx, matched.ID, matched.Name, entityID, displayName)
|
||
}
|
||
|
||
// execHexis runs a resolved capability and records a cross-service trace with
|
||
// the correlation ID. It reports command success, never operational recovery
|
||
// (Praxis observes recovery independently).
|
||
func (h *reactiveHandler) execHexis(ctx context.Context, capID, capName, entityID, displayName string) string {
|
||
correlationID, err := h.ecosystem.executeCapability(ctx, capID, entityID, nil)
|
||
if err != nil {
|
||
log.Printf("ecosystem: hexis execute error (cor=%s): %v", correlationID, err)
|
||
return "не получилось выполнить команду для " + displayName + "."
|
||
}
|
||
h.recordPraxisTrace(ctx, "hexis:"+capName, map[string]any{
|
||
"entity_id": entityID,
|
||
"entity_name": displayName,
|
||
"capability": capName,
|
||
"correlation_id": correlationID,
|
||
})
|
||
return "команда выполнена для " + displayName + "."
|
||
}
|
||
|
||
// jsonString — a one-line JSON string encoder without dragging encoding/json
|
||
// into the top of this file. Used to wrap a reminder payload's text field;
|
||
// the router's reminder Slots are already absolute (DateTimeParser resolved
|
||
// relative→absolute), the payload shape is conventional {"text":...}.
|
||
func jsonString(s string) string {
|
||
return jsonStringImpl(s)
|
||
}
|
||
|
||
// 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, " ")
|
||
}
|
||
|
||
// stripWake removes a leading wake token (any script the STT phonetically
|
||
// transcribes "Maven" as) so the verb is the first word.
|
||
func stripWake(u string) string {
|
||
stripped, had := router.StripWakeToken(u)
|
||
if !had {
|
||
return strings.TrimSpace(u)
|
||
}
|
||
return stripped
|
||
}
|
||
|
||
// firstWord returns the first whitespace-delimited token (lowercased) — the
|
||
// proposed tool's name.
|
||
func firstWord(s string) string {
|
||
f := strings.Fields(s)
|
||
if len(f) == 0 {
|
||
return ""
|
||
}
|
||
return strings.ToLower(f[0])
|
||
}
|
||
|
||
// 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)
|
||
}
|
||
|
||
// firstLine — the first non-empty line of a tool's output, for a short spoken
|
||
// reply (the full output goes to the log, not the TTS). Trimmed to keep the
|
||
// utterance sane if a command dumps a wall of text.
|
||
func firstLine(s string) string {
|
||
for _, line := range strings.Split(s, "\n") {
|
||
line = strings.TrimSpace(line)
|
||
if line != "" {
|
||
if len(line) > 200 {
|
||
line = line[:200]
|
||
}
|
||
return line
|
||
}
|
||
}
|
||
return ""
|
||
}
|
||
|
||
// isWeatherQuery returns true if the utterance is about weather.
|
||
func isWeatherQuery(u string) bool {
|
||
lower := strings.ToLower(u)
|
||
return strings.Contains(lower, "погод") ||
|
||
strings.Contains(lower, "градус") ||
|
||
strings.Contains(lower, "температур") ||
|
||
strings.Contains(lower, "дожд") ||
|
||
strings.Contains(lower, "холод") ||
|
||
strings.Contains(lower, "тепл") ||
|
||
strings.Contains(lower, "weather") ||
|
||
strings.Contains(lower, "temperature")
|
||
}
|
||
|
||
// extractWeatherLocation parses a location from the utterance, or falls back
|
||
// to the configured default. Very basic: just checks for known city names.
|
||
func extractWeatherLocation(u, defaultLoc string) string {
|
||
lower := strings.ToLower(u)
|
||
cities := map[string]string{
|
||
"москв": "Moscow",
|
||
"moscow": "Moscow",
|
||
"питер": "Saint Petersburg",
|
||
"spb": "Saint Petersburg",
|
||
"петербур": "Saint Petersburg",
|
||
"лондон": "London",
|
||
"london": "London",
|
||
"париж": "Paris",
|
||
"paris": "Paris",
|
||
"берлин": "Berlin",
|
||
"berlin": "Berlin",
|
||
"нью-йорк": "New York",
|
||
"new york": "New York",
|
||
}
|
||
for substr, name := range cities {
|
||
if strings.Contains(lower, substr) {
|
||
return name
|
||
}
|
||
}
|
||
if defaultLoc != "" {
|
||
return defaultLoc
|
||
}
|
||
return "Moscow"
|
||
}
|
||
|
||
// formatTime returns a human-readable Russian time string for a fact timestamp.
|
||
// Used by the query handler when answering "когда я это сделал?"-style questions.
|
||
func formatTime(t time.Time) string {
|
||
now := time.Now()
|
||
if t.After(now.Add(-2*time.Minute)) && t.Before(now.Add(2*time.Minute)) {
|
||
return "только что"
|
||
}
|
||
diff := now.Sub(t)
|
||
switch {
|
||
case diff < 10*time.Minute:
|
||
return "несколько минут назад"
|
||
case diff < 60*time.Minute:
|
||
return fmt.Sprintf("%d минут назад", int(diff.Minutes()))
|
||
case diff < 2*time.Hour:
|
||
return "час назад"
|
||
case diff < 24*time.Hour:
|
||
return fmt.Sprintf("%d часа назад", int(diff.Hours()))
|
||
default:
|
||
return t.Format("2 января 15:04")
|
||
}
|
||
}
|
||
|
||
func jsonStringImpl(s string) string {
|
||
// minimal JSON string escape — quotes + backslash + control chars.
|
||
// adequate for the reminder payload's text field; not a general JSON
|
||
// encoder. The chroma / RAG modules (when they land) use a real json
|
||
// encoder for richer payloads. Keep it inline here so the import
|
||
// direction stays narrow.
|
||
var b []byte
|
||
b = append(b, '"')
|
||
for _, r := range s {
|
||
switch r {
|
||
case '"':
|
||
b = append(b, '\\', '"')
|
||
case '\\':
|
||
b = append(b, '\\', '\\')
|
||
case '\n':
|
||
b = append(b, '\\', 'n')
|
||
case '\r':
|
||
b = append(b, '\\', 'r')
|
||
case '\t':
|
||
b = append(b, '\\', 't')
|
||
default:
|
||
if r < 0x20 {
|
||
b = append(b, []byte(fmt.Sprintf("\\u%04x", r))...)
|
||
} else {
|
||
b = append(b, []byte(string(r))...)
|
||
}
|
||
}
|
||
}
|
||
b = append(b, '"')
|
||
return string(b)
|
||
}
|
||
|
||
// 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)
|
||
}
|
||
}
|