Compare commits
39 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 766ca091a7 | |||
| c5317eb2b4 | |||
| 9190f897a3 | |||
| d29e7ba813 | |||
| f7e1187823 | |||
| ed48c59ba7 | |||
| b09967f9e6 | |||
| b4a3867479 | |||
| 88d07b5175 | |||
| c00e3003bf | |||
| c0f9834528 | |||
| ad5eb2d1cf | |||
| 5253123d99 | |||
| 2abf98dea6 | |||
| f5c71b87f3 | |||
| 5934110fa8 | |||
| 6e47a3d736 | |||
| 67a5eb3805 | |||
| 029449eefa | |||
| ac36216f5d | |||
| db17cfcc65 | |||
| 7d676eb941 | |||
| fe3a4e9514 | |||
| a906f2afad | |||
| a2031a31d1 | |||
| 9b8bdf73cc | |||
| 7ad3c9a408 | |||
| 84e1478823 | |||
| 9c8d0baffe | |||
| b0f5a16ec9 | |||
| 67563ed1f6 | |||
| f0f7ebc9b2 | |||
| d12de589a2 | |||
| 50cc17f33a | |||
| 73d13f1ea6 | |||
| 0ca5748699 | |||
| b43bb265b5 | |||
| b9a24334ea | |||
| c97aebf55a |
@@ -6,6 +6,7 @@
|
||||
/mavweb
|
||||
/mavpoll
|
||||
/mavcaldav
|
||||
/mavwaked
|
||||
|
||||
# Certs (private keys, don't commit)
|
||||
certs/
|
||||
|
||||
@@ -69,19 +69,41 @@ protocol; the config in `deploy/mavend.json` (with `${VAR}` env expansion from g
|
||||
|
||||
## Routing — read this before touching the router
|
||||
|
||||
`internal/router/` has TWO layered engines and the committed default is an **interim
|
||||
stopgap, not the intended design** (see memory `routing-architecture-target`):
|
||||
`internal/router/` has TWO layered engines. **The LLM router is now the default and it is
|
||||
on in deploy** — this section used to say it was wired `nil`, which stopped being true on
|
||||
2026-07-31.
|
||||
|
||||
- **Target (REARCH.md):** LLM-as-router. One resident Qwen3-1.7B (`llmrouter.go`) emits
|
||||
GBNF-constrained structured JSON, and the SAME model phrases replies. Embedder is demoted
|
||||
from a routing gate to a RAG hint.
|
||||
- **Current stopgap:** `llmrouter` is wired `nil` (around `voice.go`), so the
|
||||
`classifier.go` + `embedder.go` nearest-neighbour cascade actually runs. It routes by
|
||||
similarity to frozen seed phrases — the known cause of weak RU query handling.
|
||||
- **LLM router (the intended design, REARCH.md):** the resident Qwen3-1.7B (`llmrouter.go`)
|
||||
emits GBNF-constrained structured JSON, and the SAME model phrases replies. Embedder is
|
||||
demoted from a routing gate to a RAG hint. Wired at `voice.go:214` via
|
||||
`pickLLMRouter(cfg.Voice.UseLLMRouter(), llmClient)`; the flag is `voice.llm_router`
|
||||
(`config.go`), `DefaultLLMRouter` is **on**, and `deploy/mavend.json` sets it `true`.
|
||||
- **Classifier cascade (the failure floor, not dead code):** `classifier.go` +
|
||||
`embedder.go` nearest-neighbour over frozen seed phrases. It runs when the LLM router is
|
||||
off, when there is no llama-server to talk to (`pickLLMRouter` logs that and degrades),
|
||||
and on any per-turn LLM error. Do not delete it — routing by seed similarity is the known
|
||||
cause of weak RU query handling, but a turn must never break on the model.
|
||||
|
||||
Cascade order: `stage0.go` exact-match fast-path → LLM router (when non-nil) → classifier
|
||||
fallback. Any LLM error falls through to the classifier so a turn never breaks on the model.
|
||||
|
||||
Measured on the 77-case RU fixture (`MODEL-BAKEOFF-31-07-2026.md`): the classifier scores
|
||||
36.8% full accuracy at p50 31ms; Qwen3-1.7B scores 67.5% intent-only / 72.7% through the
|
||||
cascade at p50 ≈2.7s. Accuracy roughly doubled, latency is ~90× worse, and that trade was
|
||||
accepted deliberately. `Confidence: 1.0` used to be hardcoded in `llmrouter.go`, so the LLM
|
||||
path could never ask for clarification (6/6 refusal cases missed on the fixture) — Vikunja
|
||||
#359. Fixed 31-07-2026 with structural signal (single-token utterance, keyless fact, act with
|
||||
no allowlisted fn) feeding the same stage-3 gate the classifier path already had — see
|
||||
`gateLLMDecision` in `router.go`. Note the second half of that bug: the LLM branch never
|
||||
consulted `r.threshold` at all, so a correct low confidence would have been discarded anyway.
|
||||
|
||||
Re-measured on the fixture after the fix: **missed clarify 6/6 → 1**, at the cost of 3 false
|
||||
clarifies and 2.6pt of full accuracy (72.7% → 70.1%, intent-only 67.5% → 74.0%). Two of the
|
||||
three false clarifies are acts the model mis-routed and the gate caught — asking beats wrongly
|
||||
executing, so the fixture and the daemon disagree about what is correct there. The third,
|
||||
`"поужинал"`, is a real defect: **the single-token rule is an English intuition and does not
|
||||
transfer to Russian**, where one word is routinely a whole sentence. Narrow or drop it.
|
||||
|
||||
## LLM output contract
|
||||
|
||||
All phrasing paths emit `{"response":"...","mood":"..."}` (parsed in `replier_llm.go` and
|
||||
@@ -94,7 +116,11 @@ workspace enforces that the Go and relabelling prompts remain identical.
|
||||
## Non-goals (hard constraints)
|
||||
|
||||
Not a nag, not autonomous. Maven's persona is **feminine** — Russian
|
||||
self-reference must use feminine forms (the user is male; see memory `maven-persona-gender`).
|
||||
self-reference must use feminine forms — `рада`, not `рад`; `поняла`, not `понял`. The owner
|
||||
is male and is addressed informally: "ты", singular, never "вы"/"ваш" and never "он"/"его"
|
||||
(she talks TO him, not about him). Pet names ("милый", "дорогой") are forbidden; his name
|
||||
("Ками") is not. The eval enforces this: `CheckAddress`, `CheckFeminine` and `CheckCringe` in
|
||||
`internal/phraser/eval/checks.go`, scored by `make eval-phrasing`.
|
||||
|
||||
**"Never phones home" is DEPRECATED** (owner's call, 2026-07-31). It used to be a hard
|
||||
constraint and it is not one any more: a 0.8B — and a 1.7B — does not know enough to answer
|
||||
|
||||
@@ -19,7 +19,10 @@ Settles Vikunja **#278 / #250**.
|
||||
- Same fixture and scorer as `ROUTING-EVAL-31-07-2026.md`: `internal/router/eval/`
|
||||
(`ru_routing_v1.json`, 76 held-out cases).
|
||||
- Reproduce: `MAVEN_LLM_URL=http://127.0.0.1:<port> make eval-router`
|
||||
(`TestLLMRouterBaseline`). Note: there is no `make eval-models` target.
|
||||
(`TestLLMRouterBaseline`). (This line used to say there is no `make eval-models` target.
|
||||
There is one now — start a server with the gguf you want, then
|
||||
`make eval-models MAVEN_LLM_URL=http://127.0.0.1:<port>`. It runs only the LLM test, since
|
||||
the classifier baselines do not depend on the model.)
|
||||
- All three models served by the same `llama-server` flags — `-c 2048 -ngl 99 -t 6`, only
|
||||
`-m` and `--port` differ. One server at a time on an otherwise idle box, so latencies are
|
||||
real and not contention.
|
||||
@@ -207,10 +210,14 @@ swapped again when the CPT lands.
|
||||
- Routing is one run per model, not three. The gaps between families are far larger
|
||||
than the run-to-run spread seen on the talk fixture, but the 2B-vs-1.7B gap (62.3
|
||||
vs 67.5) is not safe to call on one run.
|
||||
- The routing numbers only reach production once the LLM router is wired on. It is
|
||||
still `nil`.
|
||||
- `/mnt/hdd1/llms/LFM2.5/Qwen3-1.7B-UD-Q4_K_XL.gguf` is a 293 MB truncated download
|
||||
in the wrong directory. The good 1.13 GB copy is in `qwen3/`. Delete the stray one.
|
||||
- ~~The routing numbers only reach production once the LLM router is wired on. It is
|
||||
still `nil`.~~ **Resolved the same evening:** the LLM router is wired at `voice.go:214`
|
||||
behind `voice.llm_router`, the default is on, and `deploy/mavend.json` sets it `true`.
|
||||
These numbers are the production path now, so the p50 ≈2.7s is a real per-turn cost and
|
||||
not a bench artifact.
|
||||
- ~~`/mnt/hdd1/llms/LFM2.5/Qwen3-1.7B-UD-Q4_K_XL.gguf` is a 293 MB truncated download
|
||||
in the wrong directory.~~ **Deleted 2026-07-31.** The good 1.13 GB copy in `qwen3/` is
|
||||
what `deploy/mavend.json` loads.
|
||||
- Harness: `scratchpad/bakeoff.sh`, one server at a time, health-checked before each
|
||||
run, `/v1/models` recorded per run. Never run two LLM consumers at once — see the
|
||||
contamination note in `TALK-EVAL-31-07-2026.md`.
|
||||
|
||||
@@ -12,7 +12,7 @@ import (
|
||||
)
|
||||
|
||||
type fakeCore struct {
|
||||
ipc.CoreAPI
|
||||
ipc.UnimplementedCoreAPI
|
||||
facts map[string]ipc.Fact // composite key "key|source" → Fact
|
||||
writeLog []ipc.WriteFactReq
|
||||
writeErr error
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
// actionTable dispatches applyAction's per-intent bodies. Each of the 7
|
||||
// intents (fact, reminder, note, query, act, chat, system) has one handler
|
||||
// here with the signature:
|
||||
//
|
||||
// func(h *reactiveHandler, ctx context.Context, dec router.Decision) string
|
||||
//
|
||||
// same contract as applyAction itself: "" means "let the Replier phrase the
|
||||
// reply", a non-empty string OVERRIDES it. This is a straight extraction of
|
||||
// applyAction's old switch cases (formerly ~300 lines in voice.go) — no
|
||||
// reordering of side effects, no new abstractions inside a handler.
|
||||
//
|
||||
// What does NOT belong in this table, because it is not per-intent:
|
||||
//
|
||||
// - the dec.Clarify short-circuit ("" when the router's stage-3 fired) —
|
||||
// stays in applyAction, before dispatch, since it applies to every
|
||||
// intent identically.
|
||||
// - the destructive-act confirm gate (park / resolveConfirm / confirmTTL)
|
||||
// and the enabled-tool allowlist. Both live entirely inside
|
||||
// actionAct/handleAct in actions_act.go, exactly where they lived in the old
|
||||
// switch's IntentAct case — they are act-specific (a fact or a note
|
||||
// can't be destructive), not shared across intents, so they do not need
|
||||
// to move to a separate layer. The important invariant, preserved
|
||||
// as-is: applyAction runs identically whether dec came from a fresh
|
||||
// route or from a completed clarify answer (see finishClarified in
|
||||
// clarify.go and its comment "filling in an argument never grants
|
||||
// authority") — a handler must never special-case a clarify-completed
|
||||
// decision to skip the confirm gate or the allowlist.
|
||||
// - detectPattern and dialogue-session bookkeeping (rememberTurn,
|
||||
// followUpMerge) run in the callers (runTurn,
|
||||
// finishClarified), not per-intent, and are untouched by this slice.
|
||||
//
|
||||
// Each handler lives in actions_<intent>.go; the small ones (chat, system)
|
||||
// and the table itself stay here.
|
||||
//
|
||||
// Adding an intent: write its handler in its own file, add one line to
|
||||
// actionHandlers. Do not grow applyAction's switch back.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
|
||||
// actionHandlers is the per-intent dispatch table used by applyAction.
|
||||
var actionHandlers = map[router.Intent]func(*reactiveHandler, context.Context, router.Decision) string{
|
||||
router.IntentFact: (*reactiveHandler).actionFact,
|
||||
router.IntentReminder: (*reactiveHandler).actionReminder,
|
||||
router.IntentAct: (*reactiveHandler).actionAct,
|
||||
router.IntentChat: (*reactiveHandler).actionChat,
|
||||
router.IntentSystem: (*reactiveHandler).actionSystem,
|
||||
router.IntentNote: (*reactiveHandler).actionNote,
|
||||
router.IntentQuery: (*reactiveHandler).actionQuery,
|
||||
}
|
||||
|
||||
func (h *reactiveHandler) actionChat(ctx context.Context, dec router.Decision) string {
|
||||
// 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
|
||||
}
|
||||
|
||||
func (h *reactiveHandler) actionSystem(ctx context.Context, dec router.Decision) string {
|
||||
return h.replySystem(ctx, dec)
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
|
||||
"github.com/kami/maven/internal/router"
|
||||
"github.com/kami/maven/internal/tool"
|
||||
)
|
||||
|
||||
// actionAct handles router.IntentAct: match a verb to an enabled tool, offer
|
||||
// it to the ecosystems first, and run it behind the confirm gate and the
|
||||
// allowlist. proposeGap and the confirm gate itself live in confirm.go.
|
||||
func (h *reactiveHandler) actionAct(ctx context.Context, dec router.Decision) string {
|
||||
// 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 "готово."
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"strconv"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
|
||||
// actionFact handles router.IntentFact: persist a tapped self-fact, index
|
||||
// it for recall, and let pattern detection propose a routine.
|
||||
func (h *reactiveHandler) actionFact(ctx context.Context, dec router.Decision) string {
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"strconv"
|
||||
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
|
||||
// actionNote handles router.IntentNote: embed the note, persist it, and
|
||||
// index it for recall.
|
||||
func (h *reactiveHandler) actionNote(ctx context.Context, dec router.Decision) string {
|
||||
// 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
|
||||
}
|
||||
@@ -0,0 +1,226 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/memory"
|
||||
"github.com/kami/maven/internal/router"
|
||||
"github.com/kami/maven/internal/weather"
|
||||
)
|
||||
|
||||
// queryTurn is the per-turn scratch a chain of query sources shares: the
|
||||
// decision being answered plus the work an earlier source already paid for
|
||||
// (the query embedding, the notes it pulled). Sources read and fill it in
|
||||
// order, so a later source never re-embeds.
|
||||
type queryTurn struct {
|
||||
dec router.Decision
|
||||
vec []float32
|
||||
notes []ipc.Note
|
||||
}
|
||||
|
||||
// querySource — one answer source in the chain actionQuery walks. answer
|
||||
// returns (reply, true) when this source claims the question, ("", false)
|
||||
// when it passes to the next one. name is for reading the table, not logged.
|
||||
//
|
||||
// A struct of one func rather than an interface: every source is a plain
|
||||
// method on *reactiveHandler with no state of its own (what state a turn has
|
||||
// lives in queryTurn), so an interface would mean one empty type per source
|
||||
// to satisfy it — ceremony for nothing. Same reasoning as confirmResolver in
|
||||
// confirm.go, and the table then reads like actionHandlers: a flat list of
|
||||
// method expressions you extend with one line.
|
||||
type querySource struct {
|
||||
name string
|
||||
answer func(*reactiveHandler, context.Context, *queryTurn) (string, bool)
|
||||
}
|
||||
|
||||
// querySources is the ordered chain actionQuery walks; first source to claim
|
||||
// answers the turn. THE ORDER IS LOAD-BEARING — see the memory-before-notes
|
||||
// comment on queryMemory: running the notes-only pass first was #373, and the
|
||||
// gate was never the bug. Adding a source (Kiwix, RSS, crawler, email) is one
|
||||
// line here plus its method; where you put the line is the whole decision.
|
||||
var querySources = []querySource{
|
||||
{"fact-by-key", (*reactiveHandler).queryFactByKey},
|
||||
{"calendar", (*reactiveHandler).queryCalendar},
|
||||
{"weather", (*reactiveHandler).queryWeather},
|
||||
{"embed", (*reactiveHandler).queryEmbed},
|
||||
{"memory", (*reactiveHandler).queryMemory},
|
||||
{"notes", (*reactiveHandler).queryNotes},
|
||||
{"general-knowledge", (*reactiveHandler).queryGeneral},
|
||||
}
|
||||
|
||||
func (h *reactiveHandler) actionQuery(ctx context.Context, dec router.Decision) string {
|
||||
t := &queryTurn{dec: dec}
|
||||
for _, src := range querySources {
|
||||
if reply, ok := src.answer(h, ctx, t); ok {
|
||||
return reply
|
||||
}
|
||||
}
|
||||
return "не знаю."
|
||||
}
|
||||
|
||||
// queryFactByKey — when the dialogue layer resolved an anaphoric reference to
|
||||
// a prior fact's key (e.g. "когда я это сделал?" after "запиши что я пил
|
||||
// воду"), look up the fact's value directly.
|
||||
func (h *reactiveHandler) queryFactByKey(ctx context.Context, t *queryTurn) (string, bool) {
|
||||
dec := t.dec
|
||||
if !dec.Slots.HasKey || dec.Slots.Key == "" {
|
||||
return "", false
|
||||
}
|
||||
f, err := h.api.LatestFact(ctx, dec.Slots.Key)
|
||||
if err != nil {
|
||||
return "", false
|
||||
}
|
||||
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.
|
||||
return fmt.Sprintf("я записала это %s", formatTime(f.Ts)), true
|
||||
}
|
||||
// General fact reference: describe what we know.
|
||||
if dec.Utterance == "" {
|
||||
return fmt.Sprintf("вот что я знаю: %s — %s", dec.Slots.Key, f.Value), true
|
||||
}
|
||||
// The utterance still carries the question; fall through to normal RAG
|
||||
// with the resolved key in context.
|
||||
return "", false
|
||||
}
|
||||
|
||||
// queryCalendar — "что у меня сегодня?", "планы на завтра?"
|
||||
// h.now(), not time.Now(): the handler's clock is the injected one, so this
|
||||
// source can be tested at a fixed time like the rest.
|
||||
func (h *reactiveHandler) queryCalendar(ctx context.Context, t *queryTurn) (string, bool) {
|
||||
date, ok := router.ParseCalendarDate(t.dec.Utterance, h.now())
|
||||
if !ok {
|
||||
return "", false
|
||||
}
|
||||
events, err := h.api.CalendarEvents(ctx, date, date.Add(24*time.Hour))
|
||||
if err != nil {
|
||||
log.Printf("voice: calendar events: %v", err)
|
||||
return "не получилось проверить календарь.", true
|
||||
}
|
||||
values := make([]string, len(events))
|
||||
for i, e := range events {
|
||||
values[i] = e.Value
|
||||
}
|
||||
var f router.CalendarEventFormatter
|
||||
return f.Format(values, date), true
|
||||
}
|
||||
|
||||
func (h *reactiveHandler) queryWeather(ctx context.Context, t *queryTurn) (string, bool) {
|
||||
if !isWeatherQuery(t.dec.Utterance) {
|
||||
return "", false
|
||||
}
|
||||
loc := extractWeatherLocation(t.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 "погода не настроена.", true
|
||||
}
|
||||
if err != nil {
|
||||
log.Printf("voice: weather: %v", err)
|
||||
return "не получилось узнать погоду.", true
|
||||
}
|
||||
return fmt.Sprintf("в %s сейчас %.0f градусов, %s.", w.Location, w.Temperature, w.Condition), true
|
||||
}
|
||||
|
||||
// queryEmbed isn't an answer source — it's the shared cost the two recall
|
||||
// sources below both need, run once, in the position it always ran in. It
|
||||
// only claims the turn when the embedder fails.
|
||||
func (h *reactiveHandler) queryEmbed(ctx context.Context, t *queryTurn) (string, bool) {
|
||||
vec, err := router.EmbedQuery(ctx, h.embedder, t.dec.Utterance)
|
||||
if err != nil {
|
||||
log.Printf("voice: embed query: %v", err)
|
||||
return "не получилось найти ответ.", true
|
||||
}
|
||||
t.vec = vec
|
||||
return "", false
|
||||
}
|
||||
|
||||
// queryMemory — 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 source 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.
|
||||
func (h *reactiveHandler) queryMemory(ctx context.Context, t *queryTurn) (string, bool) {
|
||||
if h.memStore == nil {
|
||||
return "", false
|
||||
}
|
||||
hits, herr := h.memStore.Search(ctx, t.vec, 3)
|
||||
if herr != nil {
|
||||
log.Printf("voice: memory search: %v", herr)
|
||||
return "", false
|
||||
}
|
||||
hit, ok := bestRecall(hits, h.queryMinScore, h.queryMinMargin)
|
||||
if !ok {
|
||||
return "", false
|
||||
}
|
||||
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, t.dec.Utterance, []string{text}); perr == nil && reply != "" {
|
||||
return reply, true
|
||||
}
|
||||
}
|
||||
return text, true
|
||||
}
|
||||
|
||||
// queryNotes — 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. Failing the gate passes the turn on to general
|
||||
// knowledge, which is what "don't read back the runner-up" means here.
|
||||
func (h *reactiveHandler) queryNotes(ctx context.Context, t *queryTurn) (string, bool) {
|
||||
notes, err := h.api.QueryNotes(ctx, t.vec, 5)
|
||||
if err != nil {
|
||||
log.Printf("voice: query notes: %v", err)
|
||||
return "не получилось найти ответ.", true
|
||||
}
|
||||
t.notes = notes
|
||||
noteScores := make([]float64, len(notes))
|
||||
for i, n := range notes {
|
||||
noteScores[i] = n.Score
|
||||
}
|
||||
if !memory.ConfidentScores(noteScores, h.queryMinScore, h.queryMinMargin) {
|
||||
return "", false
|
||||
}
|
||||
texts := make([]string, len(notes))
|
||||
for i, n := range notes {
|
||||
texts[i] = n.Text
|
||||
}
|
||||
reply, err := h.phraser.PhraseQuery(ctx, t.dec.Utterance, texts)
|
||||
if err != nil {
|
||||
log.Printf("voice: phrase query: %v", err)
|
||||
}
|
||||
if reply == "" {
|
||||
reply = "вот что я нашла: " + texts[0]
|
||||
}
|
||||
return reply, true
|
||||
}
|
||||
|
||||
// queryGeneral — general knowledge from the phraser, the last source before
|
||||
// giving up. It always claims: either the model answers or Maven says she
|
||||
// doesn't know.
|
||||
func (h *reactiveHandler) queryGeneral(ctx context.Context, t *queryTurn) (string, bool) {
|
||||
reply, err := h.phraser.PhraseQuery(ctx, t.dec.Utterance, nil)
|
||||
if err != nil || reply == "" {
|
||||
return "не знаю.", true
|
||||
}
|
||||
return reply, true
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
|
||||
// actionReminder handles router.IntentReminder: parse the time when stage-0
|
||||
// skipped the extractor, then create the reminder.
|
||||
func (h *reactiveHandler) actionReminder(ctx context.Context, dec router.Decision) string {
|
||||
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 ""
|
||||
}
|
||||
+57
-6
@@ -3,6 +3,8 @@ package main
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"math/rand"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/dialogue"
|
||||
@@ -21,8 +23,12 @@ const clarifyTTL = 90 * time.Second
|
||||
// raw utterance, chat and system have nothing to fill in. For those a clarify
|
||||
// decision keeps the canned "не поняла" reply — inventing a question for noise
|
||||
// is worse than admitting she missed it.
|
||||
// A reminder wants BOTH what to remind about and when. Subject first: "напомни
|
||||
// в 11" has a time and nothing to say at 11, and a reminder with no subject is
|
||||
// not worth setting. Order here is the order she asks in — she still only asks
|
||||
// about the first one missing.
|
||||
var wantedSlots = map[router.Intent][]dialogue.Slot{
|
||||
router.IntentReminder: {dialogue.SlotTime},
|
||||
router.IntentReminder: {dialogue.SlotText, dialogue.SlotTime},
|
||||
router.IntentFact: {dialogue.SlotKey},
|
||||
router.IntentAct: {dialogue.SlotFn},
|
||||
}
|
||||
@@ -35,7 +41,8 @@ var wantedSlots = map[router.Intent][]dialogue.Slot{
|
||||
// questions, so there is no gender agreement to get wrong; the feminine
|
||||
// self-reference lives in the reply she gives when she drops the request.
|
||||
var clarifyQuestions = map[dialogue.Slot]string{
|
||||
dialogue.SlotTime: "На когда напомнить?",
|
||||
dialogue.SlotTime: "Когда?",
|
||||
dialogue.SlotText: "О чём напомнить?",
|
||||
dialogue.SlotKey: "Что записать?",
|
||||
dialogue.SlotFn: "Что сделать?",
|
||||
}
|
||||
@@ -45,11 +52,55 @@ var clarifyQuestions = map[dialogue.Slot]string{
|
||||
// landed. Feminine self-reference ("поняла"), as everywhere.
|
||||
const clarifyGaveUp = "Прости, я не поняла. Скажи, пожалуйста, по-другому."
|
||||
|
||||
// clarifyExpired — his answer came after the TTL, so the parked request is
|
||||
// already gone. Same tone as clarifyGaveUp, different reason: too much time
|
||||
// clarifyExpiredVariants — his answer came after the TTL, so the parked request
|
||||
// is already gone. Same tone as clarifyGaveUp, different reason: too much time
|
||||
// passed, not "I did not understand". Feminine self-reference ("ждала",
|
||||
// "отпустила"); he is addressed with a plain imperative.
|
||||
const clarifyExpired = "Прости, я слишком долго ждала ответа и отпустила прошлую просьбу. Если она ещё нужна, скажи заново."
|
||||
//
|
||||
// Five phrasings, not one. This is the line he hears whenever he walks off
|
||||
// mid-request, so it is the line that repeats most — and the same sentence every
|
||||
// time is what makes a house assistant sound like a kiosk. They all carry the
|
||||
// same two facts (the old request is gone; say it again if it still matters),
|
||||
// because the wording may vary and the meaning may not.
|
||||
//
|
||||
// Fixed templates rather than model output, for the same reason as
|
||||
// clarifyQuestions: this text has to be right every time, and it is not worth a
|
||||
// generation to say something this small.
|
||||
var clarifyExpiredVariants = []string{
|
||||
"Прости, я слишком долго ждала ответа и отпустила прошлую просьбу. Если она ещё нужна, скажи заново.",
|
||||
"Кажется, прошлая просьба уже не важна — я её отпустила. Если я ошибаюсь, повтори.",
|
||||
"Ты как-то резко замолчал, и я не стала ждать дальше. Если та просьба ещё нужна, скажи заново.",
|
||||
"Я не дождалась ответа и убрала прошлую просьбу. Повтори, если она всё ещё нужна.",
|
||||
"Столько времени прошло, что я отпустила прошлую просьбу. Скажи заново, если она в силе.",
|
||||
}
|
||||
|
||||
// clarifyExpiredLine picks one of them at random.
|
||||
func clarifyExpiredLine() string {
|
||||
return clarifyExpiredVariants[rand.Intn(len(clarifyExpiredVariants))]
|
||||
}
|
||||
|
||||
// isClarifyExpired reports whether s opens with any of the expiry lines. The
|
||||
// notice is glued in front of this turn's reply (see withNotice), so a caller
|
||||
// checking for it has to match a prefix, not the whole string.
|
||||
func isClarifyExpired(s string) bool {
|
||||
for _, v := range clarifyExpiredVariants {
|
||||
if strings.HasPrefix(s, v) {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// trimClarifyExpired strips a leading expiry notice, leaving this turn's actual
|
||||
// reply. "" ⇒ the notice was the whole thing.
|
||||
func trimClarifyExpired(s string) string {
|
||||
for _, v := range clarifyExpiredVariants {
|
||||
if strings.HasPrefix(s, v) {
|
||||
return strings.TrimSpace(strings.TrimPrefix(s, v))
|
||||
}
|
||||
}
|
||||
return strings.TrimSpace(s)
|
||||
}
|
||||
|
||||
// clarifyExpiredNotice returns that line when a parked question had just timed
|
||||
// out, and "" when nothing was parked. Call it right after
|
||||
@@ -63,7 +114,7 @@ func (h *reactiveHandler) clarifyExpiredNotice() string {
|
||||
return ""
|
||||
}
|
||||
log.Printf("voice: clarify — parked question expired, telling him and routing the words fresh")
|
||||
return clarifyExpired
|
||||
return clarifyExpiredLine()
|
||||
}
|
||||
|
||||
// withNotice glues the expiry notice in front of this turn's reply. One turn
|
||||
|
||||
@@ -56,10 +56,13 @@ func TestClarifyQuestionForMissingSlot(t *testing.T) {
|
||||
want string
|
||||
asked bool
|
||||
}{
|
||||
{"reminder without a time", clarifyDec(router.IntentReminder, router.Slots{Text: "напомни позвонить маме"}, "напомни позвонить маме"), "На когда напомнить?", true},
|
||||
{"reminder without a time", clarifyDec(router.IntentReminder, router.Slots{Text: "напомни позвонить маме"}, "напомни позвонить маме"), "Когда?", true},
|
||||
{"fact without a key", clarifyDec(router.IntentFact, router.Slots{Text: "запиши"}, "запиши"), "Что записать?", true},
|
||||
{"act without a fn", clarifyDec(router.IntentAct, router.Slots{Text: "сделай это"}, "сделай это"), "Что сделать?", true},
|
||||
{"reminder that already has a time", clarifyDec(router.IntentReminder, router.Slots{HasTime: true}, "напомни в 11"), "", false},
|
||||
// A time with nothing to say at that time is still half a reminder, so
|
||||
// the subject is what she asks about — not silence.
|
||||
{"reminder that has a time but no subject", clarifyDec(router.IntentReminder, router.Slots{HasTime: true}, "напомни в 11"), "О чём напомнить?", true},
|
||||
{"reminder that has both", clarifyDec(router.IntentReminder, router.Slots{Text: "позвонить маме", HasTime: true}, "напомни в 11 позвонить маме"), "", false},
|
||||
{"chat is never worth a question", clarifyDec(router.IntentChat, router.Slots{Text: "мгм"}, "мгм"), "", false},
|
||||
{"query is never worth a question", clarifyDec(router.IntentQuery, router.Slots{Text: "а"}, "а"), "", false},
|
||||
}
|
||||
@@ -78,7 +81,7 @@ func TestClarifyReminderCompletesOnAnswer(t *testing.T) {
|
||||
h, st, _ := newClarifyHandler(t)
|
||||
|
||||
question, asked := h.askClarify(clarifyDec(router.IntentReminder, router.Slots{Text: "напомни позвонить маме"}, "напомни позвонить маме"))
|
||||
if !asked || question != "На когда напомнить?" {
|
||||
if !asked || question != "Когда?" {
|
||||
t.Fatalf("expected the time question, got %q asked=%v", question, asked)
|
||||
}
|
||||
|
||||
@@ -152,7 +155,7 @@ func TestClarifyAsksThreeTimesThenSaysSo(t *testing.T) {
|
||||
if !handled {
|
||||
t.Fatalf("answer %d must be consumed as an answer", i)
|
||||
}
|
||||
if reply != "На когда напомнить?" {
|
||||
if reply != "Когда?" {
|
||||
t.Fatalf("attempt %d should ask again, got %q", i, reply)
|
||||
}
|
||||
if h.clarifyStore.Get(voiceDialogueID, h.now()) == nil {
|
||||
@@ -304,17 +307,17 @@ func TestClarifyExpiryIsAnnouncedAndWordsStillRoute(t *testing.T) {
|
||||
*now = now.Add(clarifyTTL + time.Second)
|
||||
|
||||
reply := h.handleText(ctx, "как дела")
|
||||
if !strings.HasPrefix(reply, clarifyExpired) {
|
||||
if !isClarifyExpired(reply) {
|
||||
t.Fatalf("expired question must be announced first, got %q", reply)
|
||||
}
|
||||
if strings.TrimSpace(strings.TrimPrefix(reply, clarifyExpired)) == "" {
|
||||
if trimClarifyExpired(reply) == "" {
|
||||
t.Fatalf("the new words must still be answered, got only the notice: %q", reply)
|
||||
}
|
||||
if h.clarifyStore.Get(voiceDialogueID, h.now()) != nil {
|
||||
t.Fatal("the expired question must be gone")
|
||||
}
|
||||
// The notice is said once, not on every later utterance.
|
||||
if reply := h.handleText(ctx, "как дела"); strings.Contains(reply, clarifyExpired) {
|
||||
if reply := h.handleText(ctx, "как дела"); isClarifyExpired(reply) {
|
||||
t.Fatalf("notice repeated on a later turn: %q", reply)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,216 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
|
||||
// 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
|
||||
|
||||
// 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()
|
||||
|
||||
for _, r := range h.confirmResolvers(ctx) {
|
||||
if !r.claim() {
|
||||
continue
|
||||
}
|
||||
// The slot is already cleared by claim(): every branch below drops the
|
||||
// pending, including the unclear one — a confirm that can't be
|
||||
// answered clearly is safer abandoned than left armed.
|
||||
switch classifyConfirm(text) {
|
||||
case confirmYes:
|
||||
return r.yes(), true
|
||||
case confirmNo:
|
||||
return r.no(), true
|
||||
default:
|
||||
return "", false
|
||||
}
|
||||
}
|
||||
return "", false
|
||||
}
|
||||
|
||||
// confirmResolver — one parked-confirm slot in the chain. claim() reports
|
||||
// whether this slot holds a live pending, taking it (and dropping an expired
|
||||
// one) as it goes; yes/no then run the answer. Only ever called with h.mu held.
|
||||
type confirmResolver struct {
|
||||
claim func() bool
|
||||
yes func() string
|
||||
no func() string
|
||||
}
|
||||
|
||||
// confirmResolvers builds the ordered chain resolveConfirm walks. Order is
|
||||
// deliberate: the routine proposal is checked before the tool confirm so a
|
||||
// routine confirm doesn't get eaten by a stale tool pending.
|
||||
func (h *reactiveHandler) confirmResolvers(ctx context.Context) []confirmResolver {
|
||||
var pr *pendingRoutineConfirm
|
||||
var hx *pendingHexisExec
|
||||
var p *pendingAct
|
||||
|
||||
return []confirmResolver{
|
||||
// Routine proposal.
|
||||
{
|
||||
claim: func() bool {
|
||||
pr, h.pendingRoutine = h.pendingRoutine, nil
|
||||
return pr != nil && !h.now().After(pr.expiry)
|
||||
},
|
||||
yes: func() string {
|
||||
// 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 "не получилось запомнить рутину."
|
||||
}
|
||||
return "буду напоминать."
|
||||
},
|
||||
no: func() string {
|
||||
if err := h.dataStore.DismissProposedRoutine(ctx, pr.routineID); err != nil {
|
||||
log.Printf("voice: dismiss proposed routine: %v", err)
|
||||
}
|
||||
return "хорошо, не буду."
|
||||
},
|
||||
},
|
||||
// Hexis execution confirm. Bound to the exact capability + target that
|
||||
// was proposed; a stray "да" can only run that, nothing else.
|
||||
{
|
||||
claim: func() bool {
|
||||
hx, h.pendingHexis = h.pendingHexis, nil
|
||||
return hx != nil && !h.now().After(hx.expiry)
|
||||
},
|
||||
yes: func() string {
|
||||
return h.execHexis(ctx, hx.capabilityID, hx.capName, hx.entityID, hx.displayName)
|
||||
},
|
||||
no: func() string { return "отменила." },
|
||||
},
|
||||
// Tool confirm.
|
||||
{
|
||||
claim: func() bool {
|
||||
p, h.pending = h.pending, nil
|
||||
return p != nil && !h.now().After(p.expiry)
|
||||
},
|
||||
yes: func() string {
|
||||
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)
|
||||
}
|
||||
return "не получилось выполнить команду."
|
||||
}
|
||||
if out != "" {
|
||||
return "готово: " + firstLine(out)
|
||||
}
|
||||
return "готово."
|
||||
},
|
||||
no: func() string { return "отменила." },
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// 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, " ")
|
||||
}
|
||||
@@ -0,0 +1,158 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/loop"
|
||||
"github.com/kami/maven/internal/store"
|
||||
)
|
||||
|
||||
// Vikunja #281 — the fourth delivery outcome: a care candidate the restraint
|
||||
// gate suppresses (quiet hours / away / calendar-busy) is not necessarily
|
||||
// lost. If it's worth resurfacing (loop.DigestEligible), it's durably held
|
||||
// (internal/store's digest_entries) and spoken as one bundle once speaking
|
||||
// is appropriate again — never while the suppression reason still holds.
|
||||
|
||||
func breakTrace(blockedBy string) *loop.TickTrace {
|
||||
return &loop.TickTrace{
|
||||
RuleTraces: []loop.RuleTrace{{
|
||||
RuleName: "break",
|
||||
Severity: loop.Sev2,
|
||||
PredicateResult: true,
|
||||
GateResult: false,
|
||||
GateBlockedBy: blockedBy,
|
||||
}},
|
||||
}
|
||||
}
|
||||
|
||||
// TestSuppressedCareDigestsAcrossQuietHours — a Sev2 care candidate blocked
|
||||
// by quiet hours is enqueued into the durable digest, and is spoken as a
|
||||
// "digest" nudge only once quiet hours actually end — never while still
|
||||
// suppressed (that would just be a second way to nag through quiet hours).
|
||||
func TestSuppressedCareDigestsAcrossQuietHours(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
sink := &fakeSink{}
|
||||
tl := newTestTickLoop(t, st, sink, nil)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
|
||||
quiet := loop.State{Now: now, QuietHours: true, Presence: store.Present}
|
||||
tl.enqueueSuppressedDigest(ctx, breakTrace("quiet_hours"), quiet, now)
|
||||
|
||||
entries, err := st.PendingDigestEntries(ctx, now)
|
||||
if err != nil {
|
||||
t.Fatalf("pending: %v", err)
|
||||
}
|
||||
if len(entries) != 1 || entries[0].Rule != "break" {
|
||||
t.Fatalf("want 1 pending digest entry for break, got %+v", entries)
|
||||
}
|
||||
|
||||
// still quiet hours: draining now must not speak — the same restraint
|
||||
// that suppressed the live nudge must suppress the bundle too.
|
||||
tl.maybeDrainDigest(ctx, quiet, now)
|
||||
if len(sink.sends) != 0 {
|
||||
t.Fatalf("digest must not drain while quiet hours holds, got %+v", sink.sends)
|
||||
}
|
||||
|
||||
// quiet hours end: this is the moment speaking is appropriate again.
|
||||
after := now.Add(time.Hour)
|
||||
clear := loop.State{Now: after, QuietHours: false, Presence: store.Present}
|
||||
tl.maybeDrainDigest(ctx, clear, after)
|
||||
|
||||
if len(sink.sends) != 1 {
|
||||
t.Fatalf("want exactly 1 dispatched digest bundle, got %d: %+v", len(sink.sends), sink.sends)
|
||||
}
|
||||
if sink.sends[0].RuleName != "digest" {
|
||||
t.Fatalf("want RuleName digest, got %q", sink.sends[0].RuleName)
|
||||
}
|
||||
|
||||
remaining, err := st.PendingDigestEntries(ctx, after)
|
||||
if err != nil {
|
||||
t.Fatalf("pending after drain: %v", err)
|
||||
}
|
||||
if len(remaining) != 0 {
|
||||
t.Fatalf("drained entry must no longer be pending, got %+v", remaining)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSuppressedCareDigestDedupesAcrossTicks — quiet hours holding for
|
||||
// several ticks must not enqueue several copies of the same suppressed
|
||||
// nudge; he hears it once when the bundle finally drains.
|
||||
func TestSuppressedCareDigestDedupesAcrossTicks(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
sink := &fakeSink{}
|
||||
tl := newTestTickLoop(t, st, sink, nil)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
|
||||
quiet := loop.State{Now: now, QuietHours: true, Presence: store.Present}
|
||||
for i := 0; i < 3; i++ {
|
||||
tl.enqueueSuppressedDigest(ctx, breakTrace("quiet_hours"), quiet, now.Add(time.Duration(i)*time.Minute))
|
||||
}
|
||||
|
||||
entries, err := st.PendingDigestEntries(ctx, now)
|
||||
if err != nil {
|
||||
t.Fatalf("pending: %v", err)
|
||||
}
|
||||
if len(entries) != 1 {
|
||||
t.Fatalf("3 suppressions of the same nudge must collapse to 1 pending entry, got %d", len(entries))
|
||||
}
|
||||
}
|
||||
|
||||
// TestSuppressedCareDigestExpiresRatherThanDeliveringLate — an entry that
|
||||
// aged out before the suppression cleared is dropped, not spoken late: a
|
||||
// two-day-old "you skipped a break" is noise, not news.
|
||||
func TestSuppressedCareDigestExpiresRatherThanDeliveringLate(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
sink := &fakeSink{}
|
||||
tl := newTestTickLoop(t, st, sink, nil)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
|
||||
quiet := loop.State{Now: now, QuietHours: true, Presence: store.Present}
|
||||
tl.enqueueSuppressedDigest(ctx, breakTrace("quiet_hours"), quiet, now)
|
||||
|
||||
// well past digestExpiry (24h) before the suppression ever clears.
|
||||
stale := now.Add(48 * time.Hour)
|
||||
tl.expireStaleDigest(ctx, stale)
|
||||
|
||||
clear := loop.State{Now: stale, QuietHours: false, Presence: store.Present}
|
||||
tl.maybeDrainDigest(ctx, clear, stale)
|
||||
|
||||
if len(sink.sends) != 0 {
|
||||
t.Fatalf("a stale digest entry must be dropped, not delivered late; got %+v", sink.sends)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSuppressedCareDigestIgnoresHighSeverity — defense in depth at the
|
||||
// wiring layer: even if a RuleTrace somehow showed a high-severity rule
|
||||
// blocked by a care-only gate reason, the tick driver must not durably
|
||||
// digest it. Alarms bypass the gate and deliver now, unchanged; they must
|
||||
// never be silently delayed into a bundle.
|
||||
func TestSuppressedCareDigestIgnoresHighSeverity(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
sink := &fakeSink{}
|
||||
tl := newTestTickLoop(t, st, sink, nil)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
|
||||
trace := &loop.TickTrace{RuleTraces: []loop.RuleTrace{{
|
||||
RuleName: "service_down",
|
||||
Severity: loop.Sev4,
|
||||
PredicateResult: true,
|
||||
GateResult: false,
|
||||
GateBlockedBy: "quiet_hours",
|
||||
}}}
|
||||
quiet := loop.State{Now: now, QuietHours: true, Presence: store.Present}
|
||||
tl.enqueueSuppressedDigest(ctx, trace, quiet, now)
|
||||
|
||||
entries, err := st.PendingDigestEntries(ctx, now)
|
||||
if err != nil {
|
||||
t.Fatalf("pending: %v", err)
|
||||
}
|
||||
if len(entries) != 0 {
|
||||
t.Fatalf("high severity must never be digested, got %+v", entries)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,312 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
|
||||
hexisclient "github.com/kami/hexis/pkg/client"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
|
||||
// praxisCapability is one arm of the Praxis act dispatch. This is an interface
|
||||
// rather than a map[string]func because each arm carries its own state: the
|
||||
// verb aliases it answers to, the trace name it records, and its own reply
|
||||
// formatting. The dispatch grows an arm per Praxis capability, so a new one is
|
||||
// added to praxisCapabilities below and nothing else changes.
|
||||
type praxisCapability interface {
|
||||
// aliases are the verbs (router fn slots, EN and RU) this capability answers to.
|
||||
aliases() []string
|
||||
// handle runs the capability and returns the user-facing reply.
|
||||
handle(ctx context.Context, h *reactiveHandler, px *praxisClient, dec router.Decision) string
|
||||
}
|
||||
|
||||
// praxisCapabilities is the registry handlePraxisAct consults, in order.
|
||||
var praxisCapabilities = []praxisCapability{
|
||||
listAttentionCapability{},
|
||||
praxisItemAction{
|
||||
verbs: []string{"acknowledge_item", "принято", "понял", "поняла"},
|
||||
ask: "какой пункт отметить принятым?",
|
||||
op: "acknowledge",
|
||||
failure: "не получилось отметить принятым.",
|
||||
success: "принято.",
|
||||
call: func(ctx context.Context, px *praxisClient, id string) error {
|
||||
_, err := px.Acknowledge(ctx, id)
|
||||
return err
|
||||
},
|
||||
},
|
||||
praxisItemAction{
|
||||
verbs: []string{"resolve_item", "сделано", "готово", "решено"},
|
||||
ask: "какой пункт отметить сделанным?",
|
||||
op: "resolve",
|
||||
failure: "не получилось отметить сделанным.",
|
||||
success: "отмечено как сделано.",
|
||||
call: func(ctx context.Context, px *praxisClient, id string) error {
|
||||
_, err := px.Resolve(ctx, id)
|
||||
return err
|
||||
},
|
||||
},
|
||||
praxisItemAction{
|
||||
verbs: []string{"ignore_item", "игнорировать", "неважно"},
|
||||
ask: "какой пункт игнорировать?",
|
||||
op: "ignore",
|
||||
failure: "не получилось проигнорировать.",
|
||||
success: "проигнорировано.",
|
||||
call: func(ctx context.Context, px *praxisClient, id string) error {
|
||||
_, err := px.Ignore(ctx, id)
|
||||
return err
|
||||
},
|
||||
},
|
||||
praxisItemAction{
|
||||
verbs: []string{"pin_item", "закрепить"},
|
||||
ask: "какой пункт закрепить?",
|
||||
op: "pin",
|
||||
failure: "не получилось закрепить.",
|
||||
success: "закреплено.",
|
||||
call: func(ctx context.Context, px *praxisClient, id string) error {
|
||||
_, err := px.Pin(ctx, id, true)
|
||||
return err
|
||||
},
|
||||
},
|
||||
listChangesCapability{},
|
||||
}
|
||||
|
||||
// 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
|
||||
for _, capability := range praxisCapabilities {
|
||||
for _, alias := range capability.aliases() {
|
||||
if alias == dec.Slots.Fn {
|
||||
return capability.handle(ctx, h, px, dec)
|
||||
}
|
||||
}
|
||||
}
|
||||
// Not a Praxis verb — let the caller fall through.
|
||||
return ""
|
||||
}
|
||||
|
||||
// praxisItemAction is the shared shape of the item-lifecycle capabilities: take
|
||||
// an item id from the value slot, call one Praxis endpoint, trace the result.
|
||||
type praxisItemAction struct {
|
||||
verbs []string
|
||||
ask string // reply when no item id was given
|
||||
op string // trace + log name of the operation
|
||||
failure string // reply when the Praxis call errors
|
||||
success string
|
||||
call func(ctx context.Context, px *praxisClient, id string) error
|
||||
}
|
||||
|
||||
func (a praxisItemAction) aliases() []string { return a.verbs }
|
||||
|
||||
func (a praxisItemAction) handle(ctx context.Context, h *reactiveHandler, px *praxisClient, dec router.Decision) string {
|
||||
id := dec.Slots.Value
|
||||
if id == "" {
|
||||
return a.ask
|
||||
}
|
||||
if err := a.call(ctx, px, id); err != nil {
|
||||
log.Printf("ecosystem: praxis %s %s: %v", a.op, id, err)
|
||||
return a.failure
|
||||
}
|
||||
h.recordPraxisTrace(ctx, a.op, map[string]any{"item_id": id})
|
||||
return a.success
|
||||
}
|
||||
|
||||
// listAttentionCapability reads the attention digest and surfaces every item it speaks.
|
||||
type listAttentionCapability struct{}
|
||||
|
||||
func (listAttentionCapability) aliases() []string {
|
||||
return []string{"list_attention", "attention", "внимание", "что требует внимания", "что нового"}
|
||||
}
|
||||
|
||||
func (listAttentionCapability) handle(ctx context.Context, h *reactiveHandler, px *praxisClient, _ router.Decision) string {
|
||||
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, "; ")
|
||||
}
|
||||
|
||||
// listChangesCapability reads the recent-changes feed.
|
||||
type listChangesCapability struct{}
|
||||
|
||||
func (listChangesCapability) aliases() []string {
|
||||
return []string{"list_changes", "changes", "изменения", "что изменилось"}
|
||||
}
|
||||
|
||||
func (listChangesCapability) handle(ctx context.Context, h *reactiveHandler, px *praxisClient, _ router.Decision) string {
|
||||
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, "; ")
|
||||
}
|
||||
|
||||
// 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 + "."
|
||||
}
|
||||
+21
-96
@@ -97,96 +97,6 @@ func main() {
|
||||
}
|
||||
}
|
||||
|
||||
// lockedAPI is a dummy CoreAPI used while the daemon is locked. Every method
|
||||
// returns errLocked. The wire protocol's StoreAPI methods all go through the
|
||||
// Server dispatch on CoreAPI, so returning errLocked from each is correct.
|
||||
type lockedAPI struct{}
|
||||
|
||||
var _ ipc.CoreAPI = (*lockedAPI)(nil)
|
||||
|
||||
func (l *lockedAPI) WriteFact(ctx context.Context, req ipc.WriteFactReq) (int64, error) {
|
||||
return 0, errLocked
|
||||
}
|
||||
func (l *lockedAPI) LatestFact(ctx context.Context, key string) (ipc.Fact, error) {
|
||||
return ipc.Fact{}, errLocked
|
||||
}
|
||||
func (l *lockedAPI) LatestFactBySource(ctx context.Context, key, source string) (ipc.Fact, error) {
|
||||
return ipc.Fact{}, errLocked
|
||||
}
|
||||
func (l *lockedAPI) Since(ctx context.Context, key string, now time.Time) (time.Duration, error) {
|
||||
return 0, errLocked
|
||||
}
|
||||
func (l *lockedAPI) Presence(ctx context.Context) (ipc.Presence, error) {
|
||||
return ipc.Presence{}, errLocked
|
||||
}
|
||||
func (l *lockedAPI) CreateReminder(ctx context.Context, fire time.Time, payload, cron string) (int64, error) {
|
||||
return 0, errLocked
|
||||
}
|
||||
func (l *lockedAPI) MarkReminder(ctx context.Context, id int64, status string) error {
|
||||
return errLocked
|
||||
}
|
||||
func (l *lockedAPI) ListReminders(ctx context.Context, n int) ([]ipc.Reminder, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
func (l *lockedAPI) RecordNudge(ctx context.Context, rule, channel, message string, ts time.Time) (int64, error) {
|
||||
return 0, errLocked
|
||||
}
|
||||
func (l *lockedAPI) ResolveNudge(ctx context.Context, id int64, outcome string, ts time.Time) error {
|
||||
return errLocked
|
||||
}
|
||||
func (l *lockedAPI) RecentOutcomes(ctx context.Context, rule string, n int) ([]string, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
func (l *lockedAPI) RecentFacts(ctx context.Context, n int) ([]ipc.Fact, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
func (l *lockedAPI) CalendarEvents(ctx context.Context, from, to time.Time) ([]ipc.Fact, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
func (l *lockedAPI) RecentNudges(ctx context.Context, n int) ([]ipc.Nudge, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
func (l *lockedAPI) WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error) {
|
||||
return 0, errLocked
|
||||
}
|
||||
func (l *lockedAPI) QueryNotes(ctx context.Context, embedding []float32, k int) ([]ipc.Note, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
func (l *lockedAPI) RecentNotes(ctx context.Context, n int) ([]ipc.Note, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
func (l *lockedAPI) ProposeTool(ctx context.Context, name, utterance, scope string, ts time.Time) (bool, error) {
|
||||
return false, errLocked
|
||||
}
|
||||
func (l *lockedAPI) EnableTool(ctx context.Context, name string, cmd []string, destructive bool, scope string, ts time.Time) error {
|
||||
return errLocked
|
||||
}
|
||||
func (l *lockedAPI) DisableTool(ctx context.Context, name string) error { return errLocked }
|
||||
func (l *lockedAPI) DeleteTool(ctx context.Context, name string) error { return errLocked }
|
||||
func (l *lockedAPI) ListProposedRoutines(ctx context.Context) ([]ipc.ProposedRoutine, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
func (l *lockedAPI) DismissProposedRoutine(ctx context.Context, id int64) error { return errLocked }
|
||||
func (l *lockedAPI) AcceptProposedRoutine(ctx context.Context, id int64) error {
|
||||
return errLocked
|
||||
}
|
||||
func (l *lockedAPI) LookupTool(ctx context.Context, name string) (ipc.Tool, error) {
|
||||
return ipc.Tool{}, errLocked
|
||||
}
|
||||
func (l *lockedAPI) ListTools(ctx context.Context, status string) ([]ipc.Tool, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
func (l *lockedAPI) RevertFact(ctx context.Context, key string) (int64, error) { return 0, errLocked }
|
||||
func (l *lockedAPI) Chat(ctx context.Context, text string) (string, error) {
|
||||
return "", errLocked
|
||||
}
|
||||
func (l *lockedAPI) TickTrace(ctx context.Context) (ipc.TickTrace, error) {
|
||||
return ipc.TickTrace{}, errLocked
|
||||
}
|
||||
func (l *lockedAPI) MorningStatus(ctx context.Context) ([]ipc.MorningRoutineStatus, error) {
|
||||
return nil, errLocked
|
||||
}
|
||||
|
||||
func run(args []string) error {
|
||||
cfgPath := flag.String("config", defaultConfigPath(), "path to mavend JSON config")
|
||||
wrappedKeyPath := flag.String("wrapped-key-file", "", "path to wrapped encryption key blob (enables cold-start unlock)")
|
||||
@@ -351,7 +261,7 @@ func run(args []string) error {
|
||||
tickInterval := time.Duration(cfg.TickInterval)
|
||||
repeatInterval := time.Duration(cfg.RepeatInterval)
|
||||
autotuneInterval := time.Duration(cfg.AutotuneInterval)
|
||||
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines))
|
||||
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines), cfg.PatternProposals)
|
||||
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
|
||||
|
||||
coreAPI = &daemonAPI{
|
||||
@@ -364,8 +274,13 @@ func run(args []string) error {
|
||||
api.chatFn = voiceW.handler.handleText
|
||||
}
|
||||
} else {
|
||||
// locked mode: dummy CoreAPI that returns errLocked for everything
|
||||
coreAPI = &lockedAPI{}
|
||||
// locked mode: no real store yet, so there's no meaningful CoreAPI to
|
||||
// serve. srv.Check below is the actual guard — every CoreAPI call is
|
||||
// refused before it reaches this value. This is just a safe non-nil
|
||||
// placeholder: if the guard is ever bypassed by a bug, calls land
|
||||
// here and fail loudly with ipc.ErrNotImplemented instead of a nil
|
||||
// dereference or, worse, silently succeeding.
|
||||
coreAPI = ipc.UnimplementedCoreAPI{}
|
||||
}
|
||||
|
||||
// ----- IPC boundary (core ↔ modules) -----
|
||||
@@ -376,7 +291,17 @@ func run(args []string) error {
|
||||
|
||||
passkeySess := webauthn.NewPasskeySession(5 * time.Minute)
|
||||
|
||||
// Set Server.Check — in locked mode, block everything except unlock-path methods.
|
||||
// Set Server.Check — the single authorization guard, run once by
|
||||
// Server.dispatch before any CoreAPI method is called (see
|
||||
// internal/ipc/server.go). In locked mode this is the ONLY thing
|
||||
// standing between an unauthenticated caller and the store: it must
|
||||
// default-deny, with an explicit allowlist for the two methods the
|
||||
// unlock flow itself needs (MethodAssertStepUp, MethodUnlock — neither
|
||||
// of which touches CoreAPI; dispatch handles them directly via
|
||||
// srv.StepUp/srv.UnlockFn). Forgetting to allowlist a new unlock-path
|
||||
// method fails safe (denied); forgetting to guard a new CoreAPI method
|
||||
// is impossible because there is nothing left to forget — every method
|
||||
// not in the allowlist is refused by construction.
|
||||
if locked {
|
||||
srv.Check = func(ctx context.Context, m ipc.Method, _ json.RawMessage) error {
|
||||
switch m {
|
||||
@@ -513,10 +438,10 @@ func run(args []string) error {
|
||||
tickInterval := time.Duration(cfg.TickInterval)
|
||||
repeatInterval := time.Duration(cfg.RepeatInterval)
|
||||
autotuneInterval := time.Duration(cfg.AutotuneInterval)
|
||||
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines))
|
||||
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines), cfg.PatternProposals)
|
||||
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
|
||||
|
||||
// Swap the CoreAPI from lockedAPI to the real store adapter.
|
||||
// Swap the CoreAPI from the locked placeholder to the real store adapter.
|
||||
newAPI := &daemonAPI{
|
||||
CoreAPI: ipc.NewStoreAPI(st),
|
||||
getTrace: tl.trace,
|
||||
|
||||
@@ -0,0 +1,120 @@
|
||||
// mavend/patterns.go — the shared detect+propose step of pattern inference
|
||||
// (Vikunja #43). Event *extraction* (fact -> action/object) happens at fact-
|
||||
// write time in detectPattern below, tied to whichever channel wrote the
|
||||
// fact. Detection — turning a run of events into a proposed routine — is
|
||||
// channel-agnostic: it only needs what's already in the events table, so it
|
||||
// runs both right after a voice fact-write (for the immediate "напоминать?"
|
||||
// confirmation) and, proactively, from the digestion tick (tick.go's
|
||||
// detectPatterns) over every action+object pair on record, not just the one
|
||||
// that was just talked about.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/pattern"
|
||||
"github.com/kami/maven/internal/store"
|
||||
)
|
||||
|
||||
// detectAndPropose runs the pattern detector over every recorded event for
|
||||
// action+object and, if a stable pattern is found and nothing has been
|
||||
// proposed/accepted/dismissed for this pair yet, creates a proposed_routines
|
||||
// row. Returns (nil, 0, nil) — not an error — whenever there is nothing new
|
||||
// to report: too few events, irregular intervals, or a pair that already has
|
||||
// a row in any status. That last case is the one that matters most: it is
|
||||
// how a routine the owner already DISMISSED stays dismissed forever, because
|
||||
// the row survives dismissal (status flips in place, see
|
||||
// store.DismissProposedRoutine) and both the Lookup check here and the
|
||||
// table's UNIQUE(action, object) constraint refuse to create a second one.
|
||||
func detectAndPropose(ctx context.Context, ds *store.Store, action, object string, ts time.Time) (*pattern.ProposedRoutine, int64, error) {
|
||||
events, err := ds.EventsFor(ctx, action, object)
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("events for %s/%s: %w", action, object, err)
|
||||
}
|
||||
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 {
|
||||
return nil, 0, fmt.Errorf("detect %s/%s: %w", action, object, err)
|
||||
}
|
||||
if r == nil {
|
||||
return nil, 0, nil // not enough data or intervals too irregular
|
||||
}
|
||||
|
||||
// Belt: check first so the common "nothing new" case never even attempts
|
||||
// an insert. Suspenders: CreateProposedRoutine's ON CONFLICT DO NOTHING
|
||||
// (backed by the UNIQUE(action,object) constraint) is the actual
|
||||
// guarantee — this Lookup is an optimization, not the source of truth.
|
||||
existing, err := ds.LookupProposedRoutine(ctx, r.Action, r.Object)
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("lookup proposed routine %s/%s: %w", action, object, err)
|
||||
}
|
||||
if existing != nil {
|
||||
return nil, 0, nil // already proposed, accepted, or dismissed — say nothing
|
||||
}
|
||||
|
||||
id, err := ds.CreateProposedRoutine(ctx, r.Action, r.Object, r.IntervalDays, ts)
|
||||
if err != nil {
|
||||
if errors.Is(err, store.ErrProposedRoutineExists) {
|
||||
return nil, 0, nil // lost a race with another caller — not an error
|
||||
}
|
||||
return nil, 0, fmt.Errorf("create proposed routine %s/%s: %w", action, object, err)
|
||||
}
|
||||
return r, id, nil
|
||||
}
|
||||
|
||||
// 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 detectAndPropose above. Event *extraction* 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
|
||||
}
|
||||
@@ -0,0 +1,285 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/delivery"
|
||||
"github.com/kami/maven/internal/loop"
|
||||
"github.com/kami/maven/internal/pattern"
|
||||
"github.com/kami/maven/internal/store"
|
||||
)
|
||||
|
||||
// seedRefillEvents writes N weekly "refill/cat_water" events straight to the
|
||||
// events table — this is what the tick reads, independent of any utterance.
|
||||
func seedRefillEvents(t *testing.T, st *store.Store, ctx context.Context, base time.Time, n int) {
|
||||
t.Helper()
|
||||
for i := 0; i < n; i++ {
|
||||
factID, err := st.WriteFact(ctx, base.Add(time.Duration(i)*7*24*time.Hour), store.KindSelf,
|
||||
"cat_water", "refill", "test", 1.0, sql.NullInt64{})
|
||||
if err != nil {
|
||||
t.Fatalf("write fact %d: %v", i, err)
|
||||
}
|
||||
if _, err := st.CreateEvent(ctx, factID, "refill", "cat_water", base.Add(time.Duration(i)*7*24*time.Hour)); err != nil {
|
||||
t.Fatalf("create event %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestTickDetectsPatternFromStoredEvents proves the tick notices a pattern on
|
||||
// its own, reading straight from the store — not as a side effect of a live
|
||||
// utterance (Vikunja #43). MinEvents weekly events with no voice turn in
|
||||
// sight must produce exactly one proposed routine.
|
||||
func TestTickDetectsPatternFromStoredEvents(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
seedRefillEvents(t, st, ctx, now, pattern.MinEvents)
|
||||
|
||||
tl := newTestTickLoop(t, st, &fakeSink{}, nil)
|
||||
tl.detectPatterns(ctx, now, loop.State{})
|
||||
|
||||
rows, err := st.ListProposedRoutines(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("list proposed routines: %v", err)
|
||||
}
|
||||
if len(rows) != 1 {
|
||||
t.Fatalf("proposed routines = %d, want 1: %+v", len(rows), rows)
|
||||
}
|
||||
if rows[0].Action != "refill" || rows[0].Object != "cat_water" {
|
||||
t.Errorf("proposed routine = %s/%s, want refill/cat_water", rows[0].Action, rows[0].Object)
|
||||
}
|
||||
}
|
||||
|
||||
// TestTickPatternDetectionIsIdempotent proves running the tick's pattern scan
|
||||
// twice does not spam a second proposal for the same pair, and that the store
|
||||
// itself is what stops the duplicate (not tick-local state) — the whole point
|
||||
// of the guard, since the tick has no memory of what it proposed last time.
|
||||
func TestTickPatternDetectionIsIdempotent(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
seedRefillEvents(t, st, ctx, now, pattern.MinEvents)
|
||||
|
||||
tl := newTestTickLoop(t, st, &fakeSink{}, nil)
|
||||
tl.detectPatterns(ctx, now, loop.State{})
|
||||
tl.detectPatterns(ctx, now.Add(time.Hour), loop.State{})
|
||||
|
||||
rows, err := st.ListProposedRoutines(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("list proposed routines: %v", err)
|
||||
}
|
||||
if len(rows) != 1 {
|
||||
t.Fatalf("proposed routines after two ticks = %d, want 1 (no duplicate): %+v", len(rows), rows)
|
||||
}
|
||||
}
|
||||
|
||||
// TestTickPatternDetectionRespectsDismissal proves the single worst failure
|
||||
// mode here — a proposal the owner already said no to coming back on the next
|
||||
// tick — cannot happen. Dismissal flips the row's status in place; it must
|
||||
// still be there to block re-proposal.
|
||||
func TestTickPatternDetectionRespectsDismissal(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
seedRefillEvents(t, st, ctx, now, pattern.MinEvents)
|
||||
|
||||
tl := newTestTickLoop(t, st, &fakeSink{}, nil)
|
||||
tl.detectPatterns(ctx, now, loop.State{})
|
||||
|
||||
rows, err := st.ListProposedRoutines(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("list proposed routines: %v", err)
|
||||
}
|
||||
if len(rows) != 1 {
|
||||
t.Fatalf("setup: proposed routines = %d, want 1", len(rows))
|
||||
}
|
||||
if err := st.DismissProposedRoutine(ctx, rows[0].ID); err != nil {
|
||||
t.Fatalf("dismiss: %v", err)
|
||||
}
|
||||
|
||||
// More events for the same pair arrive, and the tick runs again — a
|
||||
// dismissed pattern must not resurface.
|
||||
seedRefillEvents(t, st, ctx, now.Add(30*24*time.Hour), pattern.MinEvents)
|
||||
tl.detectPatterns(ctx, now.Add(60*24*time.Hour), loop.State{})
|
||||
|
||||
proposed, err := st.ListProposedRoutinesByStatus(ctx, store.RoutineProposed)
|
||||
if err != nil {
|
||||
t.Fatalf("list proposed: %v", err)
|
||||
}
|
||||
if len(proposed) != 0 {
|
||||
t.Fatalf("a dismissed pattern came back: %+v", proposed)
|
||||
}
|
||||
all, err := st.ListProposedRoutinesByStatus(ctx, "")
|
||||
if err != nil {
|
||||
t.Fatalf("list all: %v", err)
|
||||
}
|
||||
if len(all) != 1 {
|
||||
t.Fatalf("total rows for the pair = %d, want 1 (still dismissed, not duplicated): %+v", len(all), all)
|
||||
}
|
||||
if all[0].Status != store.RoutineDismissed {
|
||||
t.Errorf("status = %s, want dismissed", all[0].Status)
|
||||
}
|
||||
}
|
||||
|
||||
// proposalRule — the rule name announceProposal uses for the seeded pair.
|
||||
const proposalRule = "proposal:refill cat_water"
|
||||
|
||||
// TestTickProposalSilentByDefault — detection is always on, announcing is not.
|
||||
// With no pattern_proposals block the tick still records the proposal, and says
|
||||
// nothing about it: Maven is not autonomous, so a behaviour that speaks without
|
||||
// being asked stays off until it is configured.
|
||||
func TestTickProposalSilentByDefault(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
seedRefillEvents(t, st, ctx, now, pattern.MinEvents)
|
||||
markPresent(t, st, ctx, now)
|
||||
|
||||
sink := &fakeSink{}
|
||||
tl := newTestTickLoop(t, st, sink, nil)
|
||||
tl.tick(ctx, now)
|
||||
|
||||
if n := countSends(sink, proposalRule); n != 0 {
|
||||
t.Fatalf("announced %d proposals with no config, want 0", n)
|
||||
}
|
||||
rows, err := st.ListProposedRoutinesByStatus(ctx, store.RoutineProposed)
|
||||
if err != nil {
|
||||
t.Fatalf("list proposed: %v", err)
|
||||
}
|
||||
if len(rows) != 1 {
|
||||
t.Fatalf("proposed routines = %d, want 1 (silent, but recorded)", len(rows))
|
||||
}
|
||||
}
|
||||
|
||||
// TestTickAnnouncesProposalWhenConfigured — with notify on, the proposal goes
|
||||
// out once through the ordinary delivery path, worded by the detector itself.
|
||||
// Later ticks stay quiet because the pair is already proposed: one pattern is
|
||||
// one announcement, ever.
|
||||
func TestTickAnnouncesProposalWhenConfigured(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
seedRefillEvents(t, st, ctx, now, pattern.MinEvents)
|
||||
markPresent(t, st, ctx, now)
|
||||
|
||||
sink := &fakeSink{}
|
||||
tl := newTestTickLoop(t, st, sink, nil)
|
||||
tl.proposalCfg = &config.PatternProposalConfig{Notify: true}
|
||||
tl.tick(ctx, now)
|
||||
|
||||
var got *delivery.Sendable
|
||||
for i := range sink.sends {
|
||||
if sink.sends[i].RuleName == proposalRule {
|
||||
got = &sink.sends[i]
|
||||
}
|
||||
}
|
||||
if got == nil {
|
||||
t.Fatalf("proposal was not announced; sends=%+v", sink.sends)
|
||||
}
|
||||
if !strings.Contains(got.Body, "напоминать?") {
|
||||
t.Errorf("body = %q, want the detector's own question", got.Body)
|
||||
}
|
||||
if got.Channel != delivery.ChannelVoice {
|
||||
t.Errorf("channel = %v, want voice (sev1, present)", got.Channel)
|
||||
}
|
||||
|
||||
// A month of further ticks: the pair already has a row, so there is
|
||||
// nothing new to detect and nothing more to say.
|
||||
sink.sends = nil
|
||||
later := now.Add(40 * 24 * time.Hour)
|
||||
markPresent(t, st, ctx, later)
|
||||
tl.tick(ctx, later)
|
||||
if n := countSends(sink, proposalRule); n != 0 {
|
||||
t.Fatalf("re-announced an existing proposal %d times, want 0", n)
|
||||
}
|
||||
}
|
||||
|
||||
// TestTickProposalRespectsGate — a proposal is the least urgent thing Maven can
|
||||
// say, so it is sev1 and the restraint gate suppresses it. Away presence means
|
||||
// it is not announced at all: it is not held, not retried, it just lives on
|
||||
// /routines. The proposal row is still written — noticing is never gated.
|
||||
func TestTickProposalRespectsGate(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
seedRefillEvents(t, st, ctx, now, pattern.MinEvents)
|
||||
// no presence probes ⇒ away ⇒ care-class gate blocks.
|
||||
|
||||
sink := &fakeSink{}
|
||||
tl := newTestTickLoop(t, st, sink, nil)
|
||||
tl.proposalCfg = &config.PatternProposalConfig{Notify: true}
|
||||
tl.tick(ctx, now)
|
||||
|
||||
if n := countSends(sink, proposalRule); n != 0 {
|
||||
t.Fatalf("away: announced %d proposals, want 0", n)
|
||||
}
|
||||
if !tl.lastProposalAt.IsZero() {
|
||||
t.Error("cooldown clock advanced on a suppressed announcement")
|
||||
}
|
||||
rows, err := st.ListProposedRoutinesByStatus(ctx, store.RoutineProposed)
|
||||
if err != nil {
|
||||
t.Fatalf("list proposed: %v", err)
|
||||
}
|
||||
if len(rows) != 1 {
|
||||
t.Fatalf("proposed routines = %d, want 1 (detection is never gated)", len(rows))
|
||||
}
|
||||
}
|
||||
|
||||
// TestTickProposalCooldownSpacesAnnouncements — two patterns detected on the
|
||||
// same tick must not become two interruptions. The second one waits for the
|
||||
// cooldown, and is on /routines meanwhile.
|
||||
func TestTickProposalCooldownSpacesAnnouncements(t *testing.T) {
|
||||
st := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := refNow()
|
||||
seedRefillEvents(t, st, ctx, now, pattern.MinEvents)
|
||||
for i := 0; i < pattern.MinEvents; i++ {
|
||||
ts := now.Add(time.Duration(i) * 3 * 24 * time.Hour)
|
||||
factID, err := st.WriteFact(ctx, ts, store.KindSelf, "litter_box", "clean", "test", 1.0, sql.NullInt64{})
|
||||
if err != nil {
|
||||
t.Fatalf("write fact: %v", err)
|
||||
}
|
||||
if _, err := st.CreateEvent(ctx, factID, "clean", "litter_box", ts); err != nil {
|
||||
t.Fatalf("create event: %v", err)
|
||||
}
|
||||
}
|
||||
markPresent(t, st, ctx, now)
|
||||
|
||||
sink := &fakeSink{}
|
||||
tl := newTestTickLoop(t, st, sink, nil)
|
||||
tl.proposalCfg = &config.PatternProposalConfig{Notify: true, Cooldown: config.Duration(24 * time.Hour)}
|
||||
tl.tick(ctx, now)
|
||||
|
||||
announced := 0
|
||||
for _, s := range sink.sends {
|
||||
if strings.HasPrefix(s.RuleName, "proposal:") {
|
||||
announced++
|
||||
}
|
||||
}
|
||||
if announced != 1 {
|
||||
t.Fatalf("announced %d proposals on one tick, want exactly 1", announced)
|
||||
}
|
||||
rows, err := st.ListProposedRoutinesByStatus(ctx, store.RoutineProposed)
|
||||
if err != nil {
|
||||
t.Fatalf("list proposed: %v", err)
|
||||
}
|
||||
if len(rows) != 2 {
|
||||
t.Fatalf("proposed routines = %d, want 2 (both recorded, one announced)", len(rows))
|
||||
}
|
||||
|
||||
// Still inside the cooldown: silence, even though a proposal is pending.
|
||||
sink.sends = nil
|
||||
soon := now.Add(time.Hour)
|
||||
markPresent(t, st, ctx, soon)
|
||||
tl.tick(ctx, soon)
|
||||
for _, s := range sink.sends {
|
||||
if strings.HasPrefix(s.RuleName, "proposal:") {
|
||||
t.Fatalf("announced %q inside the cooldown", s.RuleName)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,144 @@
|
||||
// Quiet-mode toggle recognition — the pre-route keyword check that lets
|
||||
// "тихий режим" flip the daemon-wide quiet_hours config without going through
|
||||
// the router. Moved out of voice.go unchanged (Vikunja #321); the tests live in
|
||||
// quiet_toggle_test.go.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"strings"
|
||||
"unicode"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
|
||||
// resolveQuietToggle — pre-route keyword check. Returns (reply, true) when
|
||||
// the utterance is a quiet-on/off command; ("", false) otherwise. Called from
|
||||
// runTurn BEFORE the router so a classifier miscue can't drop it — which means
|
||||
// both the voice path and the text path (mavweb /api/chat, telegram) reach it,
|
||||
// so a false positive here is a network-reachable way to flip a daemon-wide
|
||||
// setting. See classifyQuietToggle for the matching rule.
|
||||
func (h *reactiveHandler) resolveQuietToggle(ctx context.Context, text string) (string, bool) {
|
||||
on, off := classifyQuietToggle(text)
|
||||
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
|
||||
}
|
||||
|
||||
// quietInflections — the inflectional endings a stem may carry and still be
|
||||
// the same word. Adjective/adverb/noun/verb endings, all ≤3 letters. This is
|
||||
// what separates "тихий"/"тихом"/"тихо" (stem "тих" + a real ending) from
|
||||
// "тихонько"/"потихоньку", which are different words: "онько" is not an
|
||||
// ending, and "потихоньку" doesn't start with the stem at all.
|
||||
var quietInflections = []string{
|
||||
"", "а", "е", "и", "й", "о", "у", "ы", "ю", "я",
|
||||
"ая", "ее", "ей", "ем", "ие", "ий", "им", "их", "ия", "ию", "ое", "ой", "ом", "ую", "ые", "ый", "ым", "ых", "ья",
|
||||
"ами", "ого", "ому", "ыми", "ать", "ить", "ять",
|
||||
}
|
||||
|
||||
// quietStem reports whether tok is the given stem carrying at most one
|
||||
// inflectional ending. Word boundaries come from tokenisation (see
|
||||
// quietTokens), not from a regexp — Go's \b is ASCII-oriented and treats every
|
||||
// Cyrillic letter as a non-word character, so `\bтих\b` would happily match
|
||||
// inside "тихонько". Comparing whole tokens sidesteps that entirely.
|
||||
func quietStem(tok, stem string) bool {
|
||||
if !strings.HasPrefix(tok, stem) {
|
||||
return false
|
||||
}
|
||||
suffix := tok[len(stem):]
|
||||
for _, e := range quietInflections {
|
||||
if suffix == e {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// quietTokens splits an utterance into lowercase word tokens, dropping
|
||||
// punctuation and spacing. Unicode-aware, so Cyrillic words tokenise the same
|
||||
// way ASCII ones do.
|
||||
func quietTokens(text string) []string {
|
||||
return strings.FieldsFunc(strings.ToLower(strings.TrimSpace(text)), func(r rune) bool {
|
||||
return !unicode.IsLetter(r) && !unicode.IsDigit(r)
|
||||
})
|
||||
}
|
||||
|
||||
// quietPhrase matches a pattern (a sequence of stems) against the token list.
|
||||
// Multi-word patterns match any contiguous run of tokens — "включи тихий
|
||||
// режим" carries "тихий режим". Single-word patterns match ONLY when they are
|
||||
// the whole utterance: bare "тихо" is a command, but "в комнате тихо" is a
|
||||
// remark about the room and must not flip a daemon-wide setting.
|
||||
func quietPhrase(tokens, pattern []string) bool {
|
||||
if len(pattern) == 0 || len(tokens) < len(pattern) {
|
||||
return false
|
||||
}
|
||||
if len(pattern) == 1 {
|
||||
return len(tokens) == 1 && quietStem(tokens[0], pattern[0])
|
||||
}
|
||||
for i := 0; i+len(pattern) <= len(tokens); i++ {
|
||||
hit := true
|
||||
for j, stem := range pattern {
|
||||
if !quietStem(tokens[i+j], stem) {
|
||||
hit = false
|
||||
break
|
||||
}
|
||||
}
|
||||
if hit {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// quietOffPhrases / quietOnPhrases — the toggle vocabulary, as stem sequences.
|
||||
var (
|
||||
quietOffPhrases = [][]string{
|
||||
{"quiet", "off"}, {"quiet", "end"},
|
||||
{"громк", "режим"}, {"шумн", "режим"},
|
||||
{"отмен", "тих"}, {"выключ", "тих"}, {"не", "тих"},
|
||||
}
|
||||
quietOnPhrases = [][]string{
|
||||
{"quiet", "on"}, {"quiet", "mode"},
|
||||
{"тих", "режим"}, {"не", "шум"}, {"не", "беспоко"},
|
||||
{"тих"},
|
||||
}
|
||||
)
|
||||
|
||||
// classifyQuietToggle reads an utterance as a quiet-mode command. OFF is
|
||||
// resolved before ON for the same reason classifyConfirm checks negatives
|
||||
// first: the OFF phrases are built out of the ON words ("выключи тихий"
|
||||
// contains "тихий"), so scanning ON first would shadow them and "выключи
|
||||
// тихий режим" would turn quiet mode on. Negation wins.
|
||||
func classifyQuietToggle(text string) (on, off bool) {
|
||||
tokens := quietTokens(text)
|
||||
for _, p := range quietOffPhrases {
|
||||
if quietPhrase(tokens, p) {
|
||||
return false, true
|
||||
}
|
||||
}
|
||||
for _, p := range quietOnPhrases {
|
||||
if quietPhrase(tokens, p) {
|
||||
return true, false
|
||||
}
|
||||
}
|
||||
return false, false
|
||||
}
|
||||
@@ -0,0 +1,114 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
|
||||
// quietFakeAPI records the WriteFact the toggle performs.
|
||||
type quietFakeAPI struct {
|
||||
ipc.UnimplementedCoreAPI
|
||||
got ipc.WriteFactReq
|
||||
call int
|
||||
}
|
||||
|
||||
func (a *quietFakeAPI) WriteFact(_ context.Context, req ipc.WriteFactReq) (int64, error) {
|
||||
a.got, a.call = req, a.call+1
|
||||
return 1, nil
|
||||
}
|
||||
|
||||
// quietVerdict — what a phrase should do to the setting.
|
||||
type quietVerdict int
|
||||
|
||||
const (
|
||||
quietNone quietVerdict = iota
|
||||
quietOn
|
||||
quietOff
|
||||
)
|
||||
|
||||
func TestResolveQuietToggle(t *testing.T) {
|
||||
cases := []struct {
|
||||
text string
|
||||
want quietVerdict
|
||||
}{
|
||||
// ON vocabulary.
|
||||
{"quiet on", quietOn},
|
||||
{"quiet mode", quietOn},
|
||||
{"тихий режим", quietOn},
|
||||
{"тихий", quietOn},
|
||||
{"не шуми", quietOn},
|
||||
{"не беспокоить", quietOn},
|
||||
{"тихо", quietOn},
|
||||
// ON, inflected / embedded in a sentence.
|
||||
{"включи тихий режим", quietOn},
|
||||
{"побудь в тихом режиме", quietOn},
|
||||
{"Тихий Режим!", quietOn},
|
||||
{"тихая", quietOn},
|
||||
|
||||
// OFF vocabulary — all seven, incl. the three that used to say ON.
|
||||
{"quiet off", quietOff},
|
||||
{"quiet end", quietOff},
|
||||
{"громкий режим", quietOff},
|
||||
{"шумный режим", quietOff},
|
||||
{"отмени тихий", quietOff},
|
||||
{"выключи тихий", quietOff},
|
||||
{"не тихо", quietOff},
|
||||
// OFF wins over the ON words it contains.
|
||||
{"выключи тихий режим", quietOff},
|
||||
{"отмени тихий режим пожалуйста", quietOff},
|
||||
{"верни громкий режим", quietOff},
|
||||
|
||||
// False positives: "тихо"/"тихий" as ordinary Russian.
|
||||
{"очень тихий сегодня день", quietNone},
|
||||
{"в комнате тихо", quietNone},
|
||||
{"тихонько напомни", quietNone},
|
||||
{"потихоньку", quietNone},
|
||||
{"тихонько", quietNone},
|
||||
{"он говорил тихим голосом весь вечер", quietNone},
|
||||
|
||||
// Unrelated.
|
||||
{"напомни завтра позвонить маме", quietNone},
|
||||
{"какая погода", quietNone},
|
||||
{"", quietNone},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.text, func(t *testing.T) {
|
||||
api := &quietFakeAPI{}
|
||||
h := &reactiveHandler{api: api, now: func() time.Time { return time.Unix(0, 0).UTC() }}
|
||||
reply, handled := h.resolveQuietToggle(context.Background(), tc.text)
|
||||
|
||||
if tc.want == quietNone {
|
||||
if handled || reply != "" {
|
||||
t.Fatalf("%q: got (%q, %v), want no match", tc.text, reply, handled)
|
||||
}
|
||||
if api.call != 0 {
|
||||
t.Fatalf("%q: wrote a fact on a non-match", tc.text)
|
||||
}
|
||||
return
|
||||
}
|
||||
if !handled {
|
||||
t.Fatalf("%q: not handled, want %v", tc.text, tc.want)
|
||||
}
|
||||
wantReply, wantVal := "тихий режим выключен.", "false"
|
||||
if tc.want == quietOn {
|
||||
wantReply, wantVal = "тихий режим включён. буду реже напоминать.", "true"
|
||||
}
|
||||
if reply != wantReply {
|
||||
t.Errorf("%q: reply = %q, want %q", tc.text, reply, wantReply)
|
||||
}
|
||||
if api.call != 1 {
|
||||
t.Fatalf("%q: WriteFact called %d times, want 1", tc.text, api.call)
|
||||
}
|
||||
if api.got.Kind != "config" || api.got.Key != "quiet_hours" || api.got.Source != "tap:voice" || api.got.Confidence != 1.0 {
|
||||
t.Errorf("%q: request shape = %+v", tc.text, api.got)
|
||||
}
|
||||
if api.got.Value != wantVal {
|
||||
t.Errorf("%q: value = %q, want %q", tc.text, api.got.Value, wantVal)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -34,7 +34,7 @@ func newLLMReplier(c completer, block func() string) *llmReplier {
|
||||
return &llmReplier{c: c, stub: voice.NewStubReplier(), block: block}
|
||||
}
|
||||
|
||||
const replySystem = `Ты — Maven, домашняя ассистентка (о себе — в женском роде). Владелец — мужчина, говоришь с ним на "ты", в единственном числе; никогда не "вы"/"ваш" и не "он"/"его". Подтверди действие РОВНО ОДНИМ коротким предложением (≤120 символов), тепло и по-русски. Не задавай вопросов, не повторяй слова, не добавляй ничего после точки. Отвечай ТОЛЬКО одним объектом JSON с полями "response" (текст) и "mood" (ровно одно из: neutral, happy, thinking, tired, confused).
|
||||
const replySystem = `Ты — Maven, домашняя ассистентка (о себе — в женском роде). Владелец — мужчина, говоришь с ним на "ты", в единственном числе; никогда не "вы"/"ваш" и не "он"/"его". Подтверди действие РОВНО ОДНИМ коротким предложением (≤120 символов), по-русски, спокойно и без официальных формулировок. Не задавай вопросов, не повторяй слова, не добавляй ничего после точки. Отвечай ТОЛЬКО одним объектом JSON с полями "response" (текст) и "mood" (ровно одно из: neutral, happy, thinking, tired, confused).
|
||||
Пример: {"response": "Записала, что ты выпил стакан воды.", "mood": "neutral"}
|
||||
Никогда не пиши "..." в поле response.`
|
||||
|
||||
|
||||
@@ -0,0 +1,180 @@
|
||||
// Package main — ruwords.go holds Russian language + calendar/time formatting
|
||||
// helpers used by the voice reply paths (replySystem, the reminder/routine
|
||||
// phrasing, etc). Pure functions, no receivers: weekday/month name tables,
|
||||
// plural agreement, clock/date rendering, and the "do I actually know this
|
||||
// place/day" guards that pick an honest reply over a confidently wrong one.
|
||||
// Extend this file rather than voice.go for anything in that shape.
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
// 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")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,85 @@
|
||||
// Package main — strutil.go holds small, receiver-free string utilities used
|
||||
// across the voice reply paths: trimming a wake token, pulling out the first
|
||||
// word or first line, and a minimal JSON string encoder for the one payload
|
||||
// shape that needs it. Extend this file rather than voice.go for anything in
|
||||
// that shape.
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
|
||||
// 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])
|
||||
}
|
||||
|
||||
// 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 ""
|
||||
}
|
||||
|
||||
// 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 {
|
||||
// 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)
|
||||
}
|
||||
@@ -24,6 +24,7 @@ import (
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/loop"
|
||||
"github.com/kami/maven/internal/morning"
|
||||
"github.com/kami/maven/internal/pattern"
|
||||
"github.com/kami/maven/internal/phraser"
|
||||
"github.com/kami/maven/internal/routine"
|
||||
"github.com/kami/maven/internal/store"
|
||||
@@ -67,6 +68,14 @@ type tickLoop struct {
|
||||
morningRoutines []morning.Routine
|
||||
morningLast map[string]time.Time
|
||||
|
||||
// proposalCfg — announcement policy for routines the tick inferred itself.
|
||||
// nil ⇒ detect silently, never announce (the default). lastProposalAt is
|
||||
// the cooldown clock, in-memory on purpose: a restart is allowed to permit
|
||||
// one more announcement, and a restart-per-day loop is a bigger problem
|
||||
// than a duplicate proposal notice.
|
||||
proposalCfg *config.PatternProposalConfig
|
||||
lastProposalAt time.Time
|
||||
|
||||
// digestQ — in-memory queue of eligible nudges waiting for batch flush.
|
||||
// populated when digestCfg != nil && digestCfg.Enabled.
|
||||
digestQ []QueuedNudge
|
||||
@@ -92,6 +101,7 @@ func newTickLoop(
|
||||
digestCfg *config.DigestConfig,
|
||||
routines []routine.Routine,
|
||||
morningRoutines []morning.Routine,
|
||||
proposalCfg *config.PatternProposalConfig,
|
||||
) *tickLoop {
|
||||
return &tickLoop{
|
||||
store: st,
|
||||
@@ -108,6 +118,7 @@ func newTickLoop(
|
||||
routineLast: make(map[string]time.Time),
|
||||
morningRoutines: morningRoutines,
|
||||
morningLast: make(map[string]time.Time),
|
||||
proposalCfg: proposalCfg,
|
||||
lastPhrase: make(map[string]delivery.PhrasedNudge),
|
||||
}
|
||||
}
|
||||
@@ -180,6 +191,17 @@ func (t *tickLoop) tick(ctx context.Context, now time.Time) {
|
||||
// (with dedup) avoids re-queueing the same rule after a flush.
|
||||
t.maybeFlush(ctx, now, state)
|
||||
|
||||
// gate-suppressed digest (Vikunja #281): rules the restraint gate held
|
||||
// back this tick (quiet hours / away / calendar-busy), not because they
|
||||
// weren't due, but because it wasn't the moment. Some of those are worth
|
||||
// resurfacing later instead of just being lost — loop.DigestEligible
|
||||
// draws that line. This is a SEPARATE mechanism from the in-memory
|
||||
// digestQ above: that one batches candidates the gate already ALLOWED to
|
||||
// fire; this one durably holds candidates the gate BLOCKED.
|
||||
t.enqueueSuppressedDigest(ctx, trace, state, now)
|
||||
t.expireStaleDigest(ctx, now)
|
||||
t.maybeDrainDigest(ctx, state, now)
|
||||
|
||||
// routines: operator-declared scheduled behaviors. fire the ones whose cron
|
||||
// crossed since last fire, delivered through the normal routing (voice when
|
||||
// present, away channels otherwise). bodies are literal operator text — not
|
||||
@@ -195,6 +217,14 @@ func (t *tickLoop) tick(ctx context.Context, now time.Time) {
|
||||
// nudge time. See internal/morning for the "why not four timers" rationale.
|
||||
t.fireMorningRoutines(ctx, now, state)
|
||||
|
||||
// pattern detection: scan every action+object pair with recorded events
|
||||
// and propose a routine for any stable one not already decided (Vikunja
|
||||
// #43). This used to only run as a side effect of the voice fact-write
|
||||
// path, so a pattern already sitting in history went unnoticed until he
|
||||
// happened to mention it again by voice. See patterns.go and
|
||||
// detectPatterns below for how idempotence and dismissal are respected.
|
||||
t.detectPatterns(ctx, now, state)
|
||||
|
||||
// reminders: gate-bypassing class. fired once, marked after a successful
|
||||
// delivery. a failed send leaves the reminder pending — the next tick
|
||||
// re-gathers and re-attempts.
|
||||
@@ -342,6 +372,235 @@ func (t *tickLoop) flushDigest(ctx context.Context, now time.Time, state loop.St
|
||||
t.digestQ = nil
|
||||
}
|
||||
|
||||
// detectPatterns runs the pattern detector proactively over every
|
||||
// action+object pair that has ever produced an event, independent of
|
||||
// whichever fact write (or channel) last touched it (Vikunja #43). This is
|
||||
// what makes pattern inference actually proactive: it fires on the daemon's
|
||||
// own schedule reading accumulated history, not only as a side effect of a
|
||||
// live voice turn.
|
||||
//
|
||||
// Idempotence and noise are handled by the store, not here — this function
|
||||
// is safe to call every tick:
|
||||
// - Same pattern, tick after tick: detectAndPropose's LookupProposedRoutine
|
||||
// check plus proposed_routines' UNIQUE(action, object) constraint (with
|
||||
// CreateProposedRoutine's ON CONFLICT DO NOTHING) mean a pair that
|
||||
// already has a row — in ANY status — produces no second row and no log
|
||||
// spam beyond the one line at genuine creation.
|
||||
// - A DISMISSED proposal must never come back. DismissProposedRoutine flips
|
||||
// status in place; the row is never deleted. So the same Lookup check
|
||||
// that stops a duplicate "proposed" also stops a "dismissed" one from
|
||||
// resurrecting — there is nothing tick-specific to get right here beyond
|
||||
// calling the same shared path the voice route already used.
|
||||
//
|
||||
// By default this only creates a row for the /routines page to show: it does
|
||||
// not notify, ring, or speak. Detection is not the same act as disturbing him
|
||||
// about it, and Maven is "not a nag, not autonomous" (CLAUDE.md). Announcing
|
||||
// is opt-in through the pattern_proposals config block — see announceProposal
|
||||
// for the restraints that apply even then. A proposal only starts producing
|
||||
// recurring nudges once he accepts it (fireAcceptedRoutines).
|
||||
func (t *tickLoop) detectPatterns(ctx context.Context, now time.Time, state loop.State) {
|
||||
pairs, err := t.store.DistinctEventPairs(ctx)
|
||||
if err != nil {
|
||||
log.Printf("tick: distinct event pairs: %v", err)
|
||||
return
|
||||
}
|
||||
announced := false
|
||||
for _, p := range pairs {
|
||||
r, _, err := detectAndPropose(ctx, t.store, p.Action, p.Object, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: detect pattern %s/%s: %v", p.Action, p.Object, err)
|
||||
continue
|
||||
}
|
||||
if r == nil {
|
||||
continue // no stable pattern, or already proposed/accepted/dismissed
|
||||
}
|
||||
log.Printf("tick: proposed routine: %s/%s every %.1f days", r.Action, r.Object, r.IntervalDays)
|
||||
// One announcement per tick at most, whatever the scan turned up. The
|
||||
// rest are on /routines; they are not lost, they are just not shouted.
|
||||
if announced {
|
||||
continue
|
||||
}
|
||||
announced = t.announceProposal(ctx, r, now, state)
|
||||
}
|
||||
}
|
||||
|
||||
// announceProposal offers a freshly inferred routine through the ordinary
|
||||
// care-delivery path, if announcing is switched on at all. Returns true when
|
||||
// something was actually sent.
|
||||
//
|
||||
// Everything here is restraint. The feature is off unless configured; when on
|
||||
// it is sev1 (the lowest severity, so quiet hours, away presence and snooze
|
||||
// all suppress it via loop.Gate exactly like a care nudge); it is spaced by
|
||||
// proposalCfg.Cooldown across every pair, not per pair; and a suppressed or
|
||||
// dropped announcement is NOT retried — the cooldown clock advances only on a
|
||||
// real send, but the proposal row already exists, so the next tick will not
|
||||
// re-detect it and nothing queues up behind it. A missed announcement means
|
||||
// he reads it on /routines instead, which is the whole point of the page.
|
||||
//
|
||||
// The body is the detector's own literal Russian phrasing (pattern.PhraseRoutine
|
||||
// — "ты заправляешь поилку раз в 7 дней — напоминать?"), not LLM-generated, so
|
||||
// an inferred routine cannot arrive worded as something Maven never observed.
|
||||
func (t *tickLoop) announceProposal(ctx context.Context, r *pattern.ProposedRoutine, now time.Time, state loop.State) bool {
|
||||
if !t.proposalCfg.AnnounceProposals() {
|
||||
return false
|
||||
}
|
||||
cooldown := time.Duration(t.proposalCfg.Cooldown)
|
||||
if cooldown <= 0 {
|
||||
cooldown = config.DefaultProposalCooldown
|
||||
}
|
||||
if !t.lastProposalAt.IsZero() && now.Sub(t.lastProposalAt) < cooldown {
|
||||
return false
|
||||
}
|
||||
|
||||
rule := loop.Rule{Name: "proposal:" + r.Action + " " + r.Object, Severity: loop.Sev1}
|
||||
if !loop.Gate(state, rule) {
|
||||
return false
|
||||
}
|
||||
body := pattern.PhraseRoutine(r)
|
||||
pn := delivery.PhrasedNudge{
|
||||
Candidate: loop.Candidate{Rule: rule, Severity: rule.Severity, State: state},
|
||||
Body: body,
|
||||
Summary: body,
|
||||
}
|
||||
sent, err := t.dispatcher.DispatchNudge(ctx, pn, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: announce proposal %s/%s: %v", r.Action, r.Object, err)
|
||||
return false
|
||||
}
|
||||
if len(sent) == 0 {
|
||||
return false // routing dropped it — /routines still has it.
|
||||
}
|
||||
t.lastProposalAt = now
|
||||
return true
|
||||
}
|
||||
|
||||
// digestExpiry — how long a gate-suppressed care nudge stays worth
|
||||
// resurfacing. 24h: these are daily-cadence rules (water/meal/break run on
|
||||
// hour-scale cooldowns and re-derive from facts that reset every day), so a
|
||||
// digest entry that outlives one full day is describing a day that's already
|
||||
// over — "you skipped a break yesterday" said tomorrow evening is noise, not
|
||||
// news. Bounding at one day also means a digest can never silently span a
|
||||
// weekend of quiet hours into an unbounded backlog.
|
||||
const digestExpiry = 24 * time.Hour
|
||||
|
||||
// maxDigestSpokenItems — the bundle read-out is capped so "batched, not
|
||||
// dropped" cannot regress into "she dumps twelve things on me the moment I
|
||||
// walk in" — a digest that nags in bulk is worse than the drops it replaced.
|
||||
// Anything beyond the cap is still marked drained (it did get its moment;
|
||||
// the cap limits WORDS, not whether it counted) and folded into a trailing
|
||||
// count instead of being spoken in full.
|
||||
const maxDigestSpokenItems = 3
|
||||
|
||||
// enqueueSuppressedDigest scans this tick's trace for care candidates the
|
||||
// gate blocked for a genuine restraint reason and durably records the
|
||||
// digest-eligible ones (loop.DigestEligible). Phrasing happens once, here,
|
||||
// at enqueue time — not re-derived at drain time — the same way queueNudge
|
||||
// phrases once and caches, so a rule suppressed for hours isn't re-prompting
|
||||
// the LLM every tick it stays blocked (EnqueueDigestEntry's rule+body dedupe
|
||||
// makes repeat calls here harmless, but skipping the phrase call entirely
|
||||
// when a pending entry already exists avoids the LLM round-trip too).
|
||||
func (t *tickLoop) enqueueSuppressedDigest(ctx context.Context, trace *loop.TickTrace, state loop.State, now time.Time) {
|
||||
if trace == nil {
|
||||
return
|
||||
}
|
||||
for _, tr := range trace.RuleTraces {
|
||||
if !tr.PredicateResult || tr.GateResult {
|
||||
continue // didn't want to fire, or wasn't suppressed
|
||||
}
|
||||
if !loop.DigestEligible(tr.Severity, tr.GateBlockedBy) {
|
||||
continue
|
||||
}
|
||||
rule := loop.Rule{Name: tr.RuleName, Severity: tr.Severity}
|
||||
cand := loop.Candidate{Rule: rule, Severity: tr.Severity, State: state}
|
||||
pn, err := t.phraser.PhraseNudge(ctx, cand)
|
||||
if err != nil {
|
||||
log.Printf("tick: phrase digest candidate %s: %v", tr.RuleName, err)
|
||||
continue
|
||||
}
|
||||
expires := now.Add(digestExpiry)
|
||||
if _, deduped, err := t.store.EnqueueDigestEntry(ctx, tr.RuleName, int(tr.Severity), pn.Body, now, expires); err != nil {
|
||||
log.Printf("tick: enqueue digest entry %s: %v", tr.RuleName, err)
|
||||
} else if deduped {
|
||||
// same suppressed nudge already pending — nothing new to say.
|
||||
continue
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// expireStaleDigest sweeps entries past their expiry once per tick — cheap
|
||||
// bookkeeping, mirrors ReconcileStaleDeliveryAttempts's shape.
|
||||
func (t *tickLoop) expireStaleDigest(ctx context.Context, now time.Time) {
|
||||
n, err := t.store.ExpireStaleDigestEntries(ctx, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: expire stale digest entries: %v", err)
|
||||
return
|
||||
}
|
||||
if n > 0 {
|
||||
log.Printf("tick: expired %d stale digest entr(y/ies) unspoken", n)
|
||||
}
|
||||
}
|
||||
|
||||
// maybeDrainDigest speaks the pending digest bundle once the gate's
|
||||
// suppression reasons have actually cleared — quiet hours over, back from
|
||||
// away, out of the meeting. Draining while still suppressed would just be a
|
||||
// second way to nag through quiet hours; the bundle waits for the same "is
|
||||
// it allowed right now" condition a live nudge already waits for.
|
||||
func (t *tickLoop) maybeDrainDigest(ctx context.Context, state loop.State, now time.Time) {
|
||||
if state.QuietHours || state.CalendarBusy || state.Presence == store.Away {
|
||||
return
|
||||
}
|
||||
entries, err := t.store.PendingDigestEntries(ctx, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: pending digest entries: %v", err)
|
||||
return
|
||||
}
|
||||
if len(entries) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
spoken := entries
|
||||
extra := 0
|
||||
if len(spoken) > maxDigestSpokenItems {
|
||||
spoken = entries[:maxDigestSpokenItems]
|
||||
extra = len(entries) - maxDigestSpokenItems
|
||||
}
|
||||
var b strings.Builder
|
||||
maxSev := 0
|
||||
for i, e := range spoken {
|
||||
if i > 0 {
|
||||
b.WriteString(" · ")
|
||||
}
|
||||
b.WriteString(e.Body)
|
||||
if e.Severity > maxSev {
|
||||
maxSev = e.Severity
|
||||
}
|
||||
}
|
||||
if extra > 0 {
|
||||
fmt.Fprintf(&b, " · и ещё %d", extra)
|
||||
}
|
||||
body := b.String()
|
||||
summary := fmt.Sprintf("%d отложенных уведомлений", len(entries))
|
||||
|
||||
cand := loop.Candidate{
|
||||
Rule: loop.Rule{Name: "digest", Severity: loop.Severity(maxSev)},
|
||||
Severity: loop.Severity(maxSev),
|
||||
State: state,
|
||||
}
|
||||
pn := delivery.PhrasedNudge{Candidate: cand, Body: body, Summary: summary}
|
||||
t.cachePhrase(pn)
|
||||
if _, err := t.dispatcher.DispatchNudge(ctx, pn, now); err != nil {
|
||||
log.Printf("tick: dispatch digest bundle: %v", err)
|
||||
return // leave entries pending; retried next tick
|
||||
}
|
||||
ids := make([]int64, len(entries))
|
||||
for i, e := range entries {
|
||||
ids[i] = e.ID
|
||||
}
|
||||
if err := t.store.DrainDigestEntries(ctx, ids, now); err != nil {
|
||||
log.Printf("tick: drain digest entries: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// routinesFromConfig maps the config's routine blocks to the engine type.
|
||||
// Validation (cron parses, name/body present, severity defaulted) already ran
|
||||
// in config.Load, so this is a pure field copy.
|
||||
|
||||
@@ -46,7 +46,7 @@ func newTestTickLoop(t *testing.T, st *store.Store, sink delivery.Sink, digestCf
|
||||
Nudges: st,
|
||||
Reminders: st,
|
||||
})
|
||||
return newTickLoop(st, g, d, phraser.NewStub(), rules, time.Second, 5*time.Minute, 0, digestCfg, nil, nil)
|
||||
return newTickLoop(st, g, d, phraser.NewStub(), rules, time.Second, 5*time.Minute, 0, digestCfg, nil, nil, nil)
|
||||
}
|
||||
|
||||
func TestTickFiresRoutineWhenScheduleCrosses(t *testing.T) {
|
||||
@@ -63,7 +63,7 @@ func TestTickFiresRoutineWhenScheduleCrosses(t *testing.T) {
|
||||
sink := &fakeSink{}
|
||||
d := delivery.NewDispatcher(delivery.Config{Voice: sink, Ntfy: sink, Telegram: sink, Nudges: st, Reminders: st})
|
||||
rs := []routine.Routine{{Name: "morning", Cron: "0 12 * * *", Body: "полдень, время воды", Severity: 1}}
|
||||
tl := newTickLoop(st, g, d, phraser.NewStub(), rules, time.Second, 5*time.Minute, 0, nil, rs, nil)
|
||||
tl := newTickLoop(st, g, d, phraser.NewStub(), rules, time.Second, 5*time.Minute, 0, nil, rs, nil, nil)
|
||||
|
||||
// first tick: seeds, does not fire the routine.
|
||||
tl.tick(ctx, now)
|
||||
|
||||
+44
-1625
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,434 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"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/phraser"
|
||||
"github.com/kami/maven/internal/router"
|
||||
"github.com/kami/maven/internal/store"
|
||||
"github.com/kami/maven/internal/stt"
|
||||
"github.com/kami/maven/internal/tool"
|
||||
"github.com/kami/maven/internal/tts"
|
||||
"github.com/kami/maven/internal/voice"
|
||||
"github.com/kami/maven/internal/weather"
|
||||
"github.com/kami/maven/internal/worker"
|
||||
)
|
||||
|
||||
// voiceWiring — everything the daemon needs to run the audio path. Held by
|
||||
// cmd/mavend/main.go alongside the other wirings; closed on shutdown.
|
||||
type voiceWiring struct {
|
||||
server *voice.Server
|
||||
sessions *voice.Sessions
|
||||
voiceSink delivery.Sink
|
||||
embedder router.Embedder
|
||||
handler *reactiveHandler // the reactive handler for IPC Chat
|
||||
// worker clients (set when configured as Remote): closed on shutdown so
|
||||
// mavsttd / mavttsd don't keep a stale conn into a restarting daemon.
|
||||
sttClient *worker.Client
|
||||
ttsClient *worker.Client
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
// Package main — weatherq.go holds the weather-query keyword helpers: does
|
||||
// this utterance ask about weather at all, and which city (if any) did it
|
||||
// name. Both are plain substring/lookup matching, not NLU — extend this file
|
||||
// rather than voice.go for anything in that shape.
|
||||
package main
|
||||
|
||||
import "strings"
|
||||
|
||||
// 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"
|
||||
}
|
||||
@@ -17,11 +17,12 @@ import (
|
||||
)
|
||||
|
||||
// fakeCore records the mutating calls handleTools makes and returns canned
|
||||
// tool lists / errors. Embedding ipc.CoreAPI (nil) satisfies the large
|
||||
// tool lists / errors. Embedding ipc.UnimplementedCoreAPI satisfies the large
|
||||
// interface — only the methods the handlers touch are overridden; any other
|
||||
// call would nil-panic, which is fine since the handlers never make them.
|
||||
// call returns ipc.ErrNotImplemented instead of nil-panicking, so a test that
|
||||
// accidentally exercises an undeclared method fails loudly.
|
||||
type fakeCore struct {
|
||||
ipc.CoreAPI
|
||||
ipc.UnimplementedCoreAPI
|
||||
|
||||
proposed, enabled []ipc.Tool
|
||||
listErr error
|
||||
@@ -62,6 +63,18 @@ type fakeCore struct {
|
||||
// for handleTrace tests
|
||||
tickTrace ipc.TickTrace
|
||||
traceErr error
|
||||
|
||||
// for handleChatAPI tests
|
||||
chatText string
|
||||
chatErr error
|
||||
}
|
||||
|
||||
func (f *fakeCore) Chat(_ context.Context, text string) (string, error) {
|
||||
f.chatText = text
|
||||
if f.chatErr != nil {
|
||||
return "", f.chatErr
|
||||
}
|
||||
return "поняла", nil
|
||||
}
|
||||
|
||||
func (f *fakeCore) EnableTool(_ context.Context, name string, cmd []string, destructive bool, scope string, _ time.Time) error {
|
||||
@@ -1044,3 +1057,64 @@ func TestHandleRoutines_NilCore_503(t *testing.T) {
|
||||
t.Fatalf("status = %d, want 503", rr.Code)
|
||||
}
|
||||
}
|
||||
|
||||
// --- handleChatAPI step-up gate (Vikunja #317) ---
|
||||
//
|
||||
// POST /api/chat reaches the router, the LLM and the act path, so it carries
|
||||
// the same gate as POST /tools and POST /api/revert.
|
||||
|
||||
func postChat(text string) *http.Request {
|
||||
req := httptest.NewRequest(http.MethodPost, "/api/chat", strings.NewReader("text="+url.QueryEscape(text)))
|
||||
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
|
||||
return req
|
||||
}
|
||||
|
||||
func TestHandleChatAPI_RequireStepUp_FailsClosed(t *testing.T) {
|
||||
core := &fakeCore{}
|
||||
rr := httptest.NewRecorder()
|
||||
handleChatAPI(rr, postChat("выключи свет"), core, nil, true)
|
||||
if rr.Code != http.StatusForbidden {
|
||||
t.Fatalf("status = %d, want 403; body=%s", rr.Code, rr.Body.String())
|
||||
}
|
||||
if core.chatText != "" {
|
||||
t.Errorf("core.Chat called with %q, but -require-stepup should deny", core.chatText)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHandleChatAPI_UnassertedSession_Denied(t *testing.T) {
|
||||
core := &fakeCore{}
|
||||
rr := httptest.NewRecorder()
|
||||
handleChatAPI(rr, postChat("выключи свет"), core, webauthn.NewPasskeySession(5*time.Minute), false)
|
||||
if rr.Code != http.StatusForbidden {
|
||||
t.Fatalf("status = %d, want 403", rr.Code)
|
||||
}
|
||||
if core.chatText != "" {
|
||||
t.Errorf("core.Chat called with %q despite an unasserted session", core.chatText)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHandleChatAPI_AssertedSession_PassesGate(t *testing.T) {
|
||||
core := &fakeCore{}
|
||||
rr := httptest.NewRecorder()
|
||||
handleChatAPI(rr, postChat("привет"), core, stepUpSession(), true)
|
||||
if rr.Code != http.StatusSeeOther {
|
||||
t.Fatalf("status = %d, want 303; body=%s", rr.Code, rr.Body.String())
|
||||
}
|
||||
if core.chatText != "привет" {
|
||||
t.Errorf("core.Chat text = %q, want %q", core.chatText, "привет")
|
||||
}
|
||||
}
|
||||
|
||||
// Default deploy: WebAuthn unconfigured and -require-stepup off ⇒ chat keeps
|
||||
// working, resting on the transport-level auth in front of mavweb.
|
||||
func TestHandleChatAPI_FailOpenByDefault(t *testing.T) {
|
||||
core := &fakeCore{}
|
||||
rr := httptest.NewRecorder()
|
||||
handleChatAPI(rr, postChat("привет"), core, nil, false)
|
||||
if rr.Code != http.StatusSeeOther {
|
||||
t.Fatalf("status = %d, want 303", rr.Code)
|
||||
}
|
||||
if core.chatText != "привет" {
|
||||
t.Errorf("core.Chat text = %q, want %q", core.chatText, "привет")
|
||||
}
|
||||
}
|
||||
|
||||
+32
-9
@@ -314,7 +314,7 @@ func noCache(h http.Handler) http.Handler {
|
||||
}
|
||||
|
||||
func main() {
|
||||
addr := flag.String("addr", ":9200", "HTTP listen address")
|
||||
addr := flag.String("addr", "127.0.0.1:9200", "HTTP listen address (loopback by default; pass e.g. \":9200\" or a LAN IP deliberately for wider exposure — POST /chat and /routines are state-changing)")
|
||||
voiceAddr := flag.String("voice", "127.0.0.1:9100", "voice server TCP addr (host:port)")
|
||||
// ntfyWS: the ntfy WebSocket subscribe URL the PWA connects to for in-app
|
||||
// nudge delivery, e.g. wss://ntfy.kvmx.ru/maven/ws?auth=<base64-token>. The
|
||||
@@ -329,7 +329,7 @@ func main() {
|
||||
coreSock := flag.String("core", "", "mavend IPC socket path for presence-signal ingest (empty = disabled)")
|
||||
pkOrigin := flag.String("webauthn-origin", "", "WebAuthn origin URL (e.g. https://maven.kvmx.ru)")
|
||||
pkRPID := flag.String("webauthn-rpid", "", "WebAuthn RP ID (e.g. maven.kvmx.ru)")
|
||||
requireStepUp := flag.Bool("require-stepup", false, "fail closed on step-up-gated actions (/tools POST, /api/revert) when WebAuthn step-up cannot be asserted; default false preserves the historical fail-open behaviour")
|
||||
requireStepUp := flag.Bool("require-stepup", false, "fail closed on step-up-gated actions (POST /tools, /routines, /api/revert, /api/chat) when WebAuthn step-up cannot be asserted; default false preserves the historical fail-open behaviour")
|
||||
pkFile := flag.String("passkey-file", "./passkeys.json", "path to WebAuthn credential store (JSON)")
|
||||
nexusURL := flag.String("nexus", "", "Nexus base URL for the /ecosystem panel (empty = not configured)")
|
||||
praxisURL := flag.String("praxis", "", "Praxis base URL for the /ecosystem panel (empty = not configured)")
|
||||
@@ -434,9 +434,9 @@ func main() {
|
||||
}
|
||||
if stepUpSession == nil {
|
||||
if *requireStepUp {
|
||||
log.Printf("SECURITY: step-up verification is DISABLED (-webauthn-origin/-webauthn-rpid unset) and -require-stepup is set: POST /tools (tool enable/disable/dismiss — defines and executes arbitrary argv) and POST /api/revert will be DENIED (403). Set -webauthn-origin and -webauthn-rpid to enable passkey step-up.")
|
||||
log.Printf("SECURITY: step-up verification is DISABLED (-webauthn-origin/-webauthn-rpid unset) and -require-stepup is set: POST /tools (tool enable/disable/dismiss — defines and executes arbitrary argv), POST /routines (accepting schedules recurring firing), POST /api/revert and POST /api/chat (reaches the router, the LLM and the act path) will be DENIED (403). Set -webauthn-origin and -webauthn-rpid to enable passkey step-up.")
|
||||
} else {
|
||||
log.Printf("SECURITY WARNING: step-up verification is DISABLED because -webauthn-origin/-webauthn-rpid are unset. UNGUARDED SURFACES: POST /tools (defines arbitrary argv via name+cmd, which internal/tool then EXECUTES) and POST /api/revert (voids the latest fact for a key). These are protected only by whatever transport-level auth sits in front of mavweb (wg+nginx+auth) — do NOT expose -addr on a public interface. Set -webauthn-origin and -webauthn-rpid to require passkey step-up, or pass -require-stepup to fail closed instead.")
|
||||
log.Printf("SECURITY WARNING: step-up verification is DISABLED because -webauthn-origin/-webauthn-rpid are unset. UNGUARDED SURFACES: POST /tools (defines arbitrary argv via name+cmd, which internal/tool then EXECUTES), POST /routines (accepting schedules recurring firing), POST /api/revert (voids the latest fact for a key) and POST /api/chat (reaches the router, the LLM and, through applyAction, the act path). These are protected only by whatever transport-level auth sits in front of mavweb (wg+nginx+auth) — do NOT expose -addr on a public interface. Set -webauthn-origin and -webauthn-rpid to require passkey step-up, or pass -require-stepup to fail closed instead.")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -454,14 +454,26 @@ func main() {
|
||||
handleRoutines(w, r, core, stepUpSession, *requireStepUp)
|
||||
})
|
||||
|
||||
// /api/revert voids the latest fact for a key — a store mutation, so it
|
||||
// sits behind the same passkey step-up as tool enable (nil session ⇒
|
||||
// WebAuthn unconfigured ⇒ transport-level auth only, same as /tools).
|
||||
// State-changing routes on this server, and their gate (Vikunja #317):
|
||||
//
|
||||
// POST /tools step-up — defines argv that internal/tool executes
|
||||
// POST /routines step-up — accepting schedules recurring firing
|
||||
// POST /api/revert step-up — voids the latest fact for a key
|
||||
// POST /api/chat step-up — reaches the router, LLM and the act path
|
||||
// POST /api/signal none — appends a presence fact, no argv, no act
|
||||
// POST /api/ptt, /ws none — proxy audio to mavend's voice port, which
|
||||
// is itself only reachable inside the deploy
|
||||
//
|
||||
// "step-up" means stepUpOK: asserted passkey when WebAuthn is configured,
|
||||
// otherwise fail-open unless -require-stepup, which denies.
|
||||
//
|
||||
// GET /chat only renders the page and echoes back the q/r query params the
|
||||
// POST redirect set — nothing to gate.
|
||||
mux.HandleFunc("/chat", func(w http.ResponseWriter, r *http.Request) {
|
||||
handleChatPage(w, r, core)
|
||||
})
|
||||
mux.HandleFunc("/api/chat", func(w http.ResponseWriter, r *http.Request) {
|
||||
handleChatAPI(w, r, core)
|
||||
handleChatAPI(w, r, core, stepUpSession, *requireStepUp)
|
||||
})
|
||||
mux.HandleFunc("/api/revert", func(w http.ResponseWriter, r *http.Request) {
|
||||
handleRevert(w, r, core, stepUpSession, *requireStepUp)
|
||||
@@ -1258,7 +1270,14 @@ func handleChatPage(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI) {
|
||||
}
|
||||
|
||||
// handleChatAPI processes a chat message POST and redirects back to /chat.
|
||||
func handleChatAPI(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI) {
|
||||
//
|
||||
// State-changing, and the widest surface on this server: the text reaches the
|
||||
// router, the LLM, and through mavend's applyAction the whole action path
|
||||
// including `act` — so it is gated on the same step-up as POST /tools and
|
||||
// POST /api/revert (Vikunja #317). With WebAuthn unconfigured the gate is
|
||||
// fail-open exactly like the others (see stepUpOK); with -require-stepup it
|
||||
// denies, which is the point of that flag.
|
||||
func handleChatAPI(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI, session *webauthn.PasskeySession, requireStepUp bool) {
|
||||
if r.Method != http.MethodPost {
|
||||
http.Error(w, "POST only", http.StatusMethodNotAllowed)
|
||||
return
|
||||
@@ -1267,6 +1286,10 @@ func handleChatAPI(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI) {
|
||||
http.Error(w, "chat disabled (no -core)", http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
if !stepUpOK(session, requireStepUp) {
|
||||
http.Error(w, "step-up required: assert a passkey first", http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
text := strings.TrimSpace(r.FormValue("text"))
|
||||
if text == "" {
|
||||
http.Redirect(w, r, "/chat", http.StatusSeeOther)
|
||||
|
||||
@@ -8,6 +8,18 @@
|
||||
#
|
||||
# Maven's own compose joins this same network (add `ecosystem` as an external
|
||||
# network there) to reach nexus:9740 / praxis:8989 / hexis:9741 directly.
|
||||
#
|
||||
# NO RELEASE PINNING (Vikunja #354): each `build:` below points at a sibling
|
||||
# WORKING TREE, so `up --build` ships whatever is checked out there, including
|
||||
# uncommitted edits. Before bringing this up, check what you are about to
|
||||
# deploy:
|
||||
#
|
||||
# for r in nexus praxis hexis; do git -C ../../../$r status --short; \
|
||||
# git -C ../../../$r log -1 --oneline; done
|
||||
#
|
||||
# The host nginx that fronts these is deploy/ecosystem/nginx.conf — it binds
|
||||
# the wg and LAN addresses only, with allow/deny. Keep it that way: none of
|
||||
# these containers has auth of its own.
|
||||
name: ecosystem
|
||||
|
||||
services:
|
||||
|
||||
@@ -1,13 +1,71 @@
|
||||
# Reverse-proxy the three sibling admin UIs. Drop into your nginx sites (or the
|
||||
# nginx-panel app) and reload. Assumes the compose publishes each service on
|
||||
# 127.0.0.1:<port>. Add TLS (certbot / your existing cert block) per server.
|
||||
# Reverse-proxy Maven's own web UI plus the three sibling admin UIs. Drop into
|
||||
# your nginx sites (or the nginx-panel app) and reload. Assumes the compose
|
||||
# publishes each service on 127.0.0.1:<port>. Add TLS (certbot / your existing
|
||||
# cert block) per server.
|
||||
#
|
||||
# NOTE: hexis.<domain> previously pointed at the MCP tool — repoint that
|
||||
# elsewhere first (the app now owns hexis.*).
|
||||
#
|
||||
# 10.42.0.1 and 192.168.1.104 below are THIS BOX's WireGuard and LAN
|
||||
# addresses (homesrv) — these admin UIs have no auth of their own, so the
|
||||
# explicit bind + allow/deny below is what keeps them off the open internet.
|
||||
# On a different box, replace both addresses with that box's wg and LAN IPs.
|
||||
# Do NOT "fix" a failed bind by reverting to `listen 80` (all interfaces) —
|
||||
# that removes the only access control these containers have.
|
||||
|
||||
# maven.<domain> → mavweb (docker-compose.yml publishes it on 127.0.0.1:9201).
|
||||
# Same bind + ACL as the siblings, and for a stronger reason: mavweb serves
|
||||
# POST /tools, which defines argv that internal/tool EXECUTES, plus POST
|
||||
# /routines, /api/revert and /api/chat (Vikunja #317). Without
|
||||
# -webauthn-origin/-webauthn-rpid mavweb has no auth of its own, so this block
|
||||
# is the auth. If you add TLS and a basic-auth/oauth2-proxy layer, keep the
|
||||
# allow/deny anyway — belt and braces on an RCE surface.
|
||||
#
|
||||
# WebSocket upgrade matters here: /ws carries push-to-talk audio, so the
|
||||
# Upgrade/Connection headers below are required, not decoration. The map keeps
|
||||
# `Connection: upgrade` off plain requests; it sits in the http context, which
|
||||
# is where sites-available files are included — if your nginx already defines
|
||||
# $connection_upgrade, drop this block.
|
||||
map $http_upgrade $connection_upgrade {
|
||||
default upgrade;
|
||||
'' close;
|
||||
}
|
||||
|
||||
server {
|
||||
listen 80;
|
||||
listen 10.42.0.1:80;
|
||||
listen 192.168.1.104:80;
|
||||
server_name maven.kvmx.ru;
|
||||
|
||||
allow 10.42.0.0/24;
|
||||
allow 192.168.1.0/24;
|
||||
deny all;
|
||||
|
||||
# push-to-talk uploads raw PCM; the default 1m is enough for a short
|
||||
# utterance but not for a long one.
|
||||
client_max_body_size 32m;
|
||||
|
||||
location / {
|
||||
proxy_pass http://127.0.0.1:9201;
|
||||
proxy_http_version 1.1;
|
||||
proxy_set_header Upgrade $http_upgrade;
|
||||
proxy_set_header Connection $connection_upgrade;
|
||||
proxy_set_header Host $host;
|
||||
proxy_set_header X-Real-IP $remote_addr;
|
||||
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
|
||||
proxy_set_header X-Forwarded-Proto $scheme;
|
||||
proxy_read_timeout 300s; # an LLM turn can take minutes on the iGPU
|
||||
}
|
||||
}
|
||||
|
||||
server {
|
||||
listen 10.42.0.1:80;
|
||||
listen 192.168.1.104:80;
|
||||
server_name nexus.kvmx.ru;
|
||||
|
||||
allow 10.42.0.0/24;
|
||||
allow 192.168.1.0/24;
|
||||
deny all;
|
||||
|
||||
location / {
|
||||
proxy_pass http://127.0.0.1:9740;
|
||||
proxy_set_header Host $host;
|
||||
@@ -18,8 +76,14 @@ server {
|
||||
}
|
||||
|
||||
server {
|
||||
listen 80;
|
||||
listen 10.42.0.1:80;
|
||||
listen 192.168.1.104:80;
|
||||
server_name praxis.kvmx.ru;
|
||||
|
||||
allow 10.42.0.0/24;
|
||||
allow 192.168.1.0/24;
|
||||
deny all;
|
||||
|
||||
location / {
|
||||
proxy_pass http://127.0.0.1:8989;
|
||||
proxy_set_header Host $host;
|
||||
@@ -30,8 +94,14 @@ server {
|
||||
}
|
||||
|
||||
server {
|
||||
listen 80;
|
||||
listen 10.42.0.1:80;
|
||||
listen 192.168.1.104:80;
|
||||
server_name hexis.kvmx.ru;
|
||||
|
||||
allow 10.42.0.0/24;
|
||||
allow 192.168.1.0/24;
|
||||
deny all;
|
||||
|
||||
location / {
|
||||
proxy_pass http://127.0.0.1:9741;
|
||||
proxy_set_header Host $host;
|
||||
|
||||
@@ -26,6 +26,11 @@
|
||||
"severity_ceiling": 2
|
||||
},
|
||||
|
||||
"pattern_proposals": {
|
||||
"notify": false,
|
||||
"cooldown": "24h"
|
||||
},
|
||||
|
||||
"nexus": { "url": "http://nexus:9740" },
|
||||
"praxis": { "url": "http://praxis:8989" },
|
||||
"hexis": { "url": "http://hexis:9741" },
|
||||
|
||||
@@ -7,7 +7,6 @@ import (
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
@@ -359,8 +358,12 @@ func TestGate_IpcServer_ChatAllowedForEnrolledCaller(t *testing.T) {
|
||||
|
||||
// recordingAPI — a no-op CoreAPI that counts WriteFact invocations; the auth
|
||||
// check must reject before reaching it, otherwise the refusal leaks into the
|
||||
// fake's counts and we fail.
|
||||
// fake's counts and we fail. Embeds ipc.UnimplementedCoreAPI so every method
|
||||
// this test doesn't exercise returns ipc.ErrNotImplemented loudly instead of
|
||||
// being hand-stubbed to a canned value nobody checks.
|
||||
type recordingAPI struct {
|
||||
ipc.UnimplementedCoreAPI
|
||||
|
||||
writes int
|
||||
chats int
|
||||
}
|
||||
@@ -369,89 +372,7 @@ func (r *recordingAPI) WriteFact(_ context.Context, _ ipc.WriteFactReq) (int64,
|
||||
r.writes++
|
||||
return int64(r.writes), nil
|
||||
}
|
||||
func (r *recordingAPI) LatestFact(_ context.Context, _ string) (ipc.Fact, error) {
|
||||
return ipc.Fact{}, ipc.ErrNoFact
|
||||
}
|
||||
func (r *recordingAPI) LatestFactBySource(_ context.Context, _, _ string) (ipc.Fact, error) {
|
||||
return ipc.Fact{}, ipc.ErrNoFact
|
||||
}
|
||||
func (r *recordingAPI) Since(_ context.Context, _ string, _ time.Time) (time.Duration, error) {
|
||||
return 0, ipc.ErrNoFact
|
||||
}
|
||||
func (r *recordingAPI) Presence(_ context.Context) (ipc.Presence, error) {
|
||||
return ipc.Presence{}, nil
|
||||
}
|
||||
func (r *recordingAPI) CreateReminder(_ context.Context, _ time.Time, _, _ string) (int64, error) {
|
||||
return 1, nil
|
||||
}
|
||||
func (r *recordingAPI) MarkReminder(_ context.Context, _ int64, _ string) error { return nil }
|
||||
func (r *recordingAPI) ListReminders(_ context.Context, _ int) ([]ipc.Reminder, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (r *recordingAPI) TickTrace(_ context.Context) (ipc.TickTrace, error) {
|
||||
return ipc.TickTrace{}, nil
|
||||
}
|
||||
func (r *recordingAPI) MorningStatus(_ context.Context) ([]ipc.MorningRoutineStatus, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (r *recordingAPI) RecordNudge(_ context.Context, _, _, _ string, _ time.Time) (int64, error) {
|
||||
return 1, nil
|
||||
}
|
||||
func (r *recordingAPI) ResolveNudge(_ context.Context, _ int64, _ string, _ time.Time) error {
|
||||
return nil
|
||||
}
|
||||
func (r *recordingAPI) RecentOutcomes(_ context.Context, _ string, _ int) ([]string, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (r *recordingAPI) RecentFacts(_ context.Context, _ int) ([]ipc.Fact, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (r *recordingAPI) CalendarEvents(_ context.Context, _, _ time.Time) ([]ipc.Fact, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (r *recordingAPI) RecentNudges(_ context.Context, _ int) ([]ipc.Nudge, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (r *recordingAPI) WriteNote(_ context.Context, _ time.Time, _ string, _ []float32, _ string) (int64, error) {
|
||||
return 1, nil
|
||||
}
|
||||
func (r *recordingAPI) QueryNotes(_ context.Context, _ []float32, _ int) ([]ipc.Note, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (r *recordingAPI) RecentNotes(_ context.Context, _ int) ([]ipc.Note, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (r *recordingAPI) ProposeTool(_ context.Context, _, _, _ string, _ time.Time) (bool, error) {
|
||||
return false, nil
|
||||
}
|
||||
func (r *recordingAPI) EnableTool(_ context.Context, _ string, _ []string, _ bool, _ string, _ time.Time) error {
|
||||
return nil
|
||||
}
|
||||
func (r *recordingAPI) DisableTool(_ context.Context, _ string) error {
|
||||
return nil
|
||||
}
|
||||
func (r *recordingAPI) DeleteTool(_ context.Context, _ string) error {
|
||||
return nil
|
||||
}
|
||||
func (r *recordingAPI) LookupTool(_ context.Context, _ string) (ipc.Tool, error) {
|
||||
return ipc.Tool{}, ipc.ErrToolNotFound
|
||||
}
|
||||
func (r *recordingAPI) ListTools(_ context.Context, _ string) ([]ipc.Tool, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (r *recordingAPI) RevertFact(_ context.Context, _ string) (int64, error) {
|
||||
return 0, nil
|
||||
}
|
||||
func (r *recordingAPI) ListProposedRoutines(_ context.Context) ([]ipc.ProposedRoutine, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (r *recordingAPI) AcceptProposedRoutine(_ context.Context, _ int64) error {
|
||||
return nil
|
||||
}
|
||||
func (r *recordingAPI) DismissProposedRoutine(_ context.Context, _ int64) error {
|
||||
return nil
|
||||
}
|
||||
func (r *recordingAPI) Chat(_ context.Context, text string) (string, error) {
|
||||
r.chats++
|
||||
return "echo: " + text, nil
|
||||
|
||||
@@ -140,6 +140,12 @@ type Config struct {
|
||||
// item. See internal/morning for the evaluation engine. Empty ⇒ disabled.
|
||||
MorningRoutines []MorningRoutineConfig `json:"morning_routines,omitempty"`
|
||||
|
||||
// PatternProposals — whether a routine the digestion tick inferred on its
|
||||
// own may be announced, and how often. nil / absent ⇒ silent detection
|
||||
// only: proposals are written for /routines and never announced. See
|
||||
// PatternProposalConfig.
|
||||
PatternProposals *PatternProposalConfig `json:"pattern_proposals,omitempty"`
|
||||
|
||||
// Praxis — the ecosystem attention-state service. When configured, maven
|
||||
// calls the Praxis HTTP tools API for attention listing and item lifecycle.
|
||||
// Maven never touches Praxis's database directly (ecosystem invariant: no
|
||||
@@ -352,6 +358,40 @@ type DigestConfig struct {
|
||||
SeverityCeiling int `json:"severity_ceiling,omitempty"` // max sev batched
|
||||
}
|
||||
|
||||
// PatternProposalConfig — announcement policy for routines the digestion tick
|
||||
// inferred by itself (Vikunja #247, #43).
|
||||
//
|
||||
// Detection is always on and always silent by default: the tick writes a
|
||||
// proposed_routines row and the /routines page shows it. Notify is what turns
|
||||
// "she noticed" into "she said something", and it is OFF unless configured —
|
||||
// Maven is not a nag and not autonomous, so a behaviour that speaks without
|
||||
// being asked has to be switched on deliberately, like weather and telegram.
|
||||
//
|
||||
// When Notify is on, the announcement is still heavily restrained:
|
||||
// - at most one proposal per tick, however many were detected;
|
||||
// - at most one per Cooldown across all pairs (not per pair), so a batch of
|
||||
// freshly-detected patterns cannot turn into a queue of interruptions;
|
||||
// - through the ordinary care-class gate (quiet hours / away / snooze), at
|
||||
// sev1 — the lowest severity there is. A proposal is the least urgent
|
||||
// thing Maven can say.
|
||||
//
|
||||
// A pair is only ever announced once, because it is only ever proposed once:
|
||||
// proposed_routines is UNIQUE(action, object) and the row survives dismissal.
|
||||
type PatternProposalConfig struct {
|
||||
// Notify — announce newly inferred routines. Default false.
|
||||
Notify bool `json:"notify,omitempty"`
|
||||
|
||||
// Cooldown — minimum spacing between two proposal announcements. 0 ⇒
|
||||
// DefaultProposalCooldown (24h).
|
||||
Cooldown Duration `json:"cooldown,omitempty"`
|
||||
}
|
||||
|
||||
// AnnounceProposals reports whether inferred routines may be announced. Safe
|
||||
// on a nil receiver — an absent config block means silent detection.
|
||||
func (p *PatternProposalConfig) AnnounceProposals() bool {
|
||||
return p != nil && p.Notify
|
||||
}
|
||||
|
||||
// PhraserConfig — the LLM-backed phraser seam. The daemon spawns llama-server
|
||||
// as a managed subprocess and sends chat-completion requests to phrase nudge
|
||||
// and reminder messages. nil ⇒ the template-based Stub is used instead.
|
||||
@@ -448,6 +488,11 @@ const (
|
||||
DefaultLLMRouter = true
|
||||
|
||||
DefaultFactEnrichmentInterval = 30 * time.Second
|
||||
|
||||
// DefaultProposalCooldown — one inferred-routine announcement per day at
|
||||
// most. A proposal is never urgent; if two patterns surface in the same
|
||||
// hour, the second one waits, and the /routines page has it either way.
|
||||
DefaultProposalCooldown = 24 * time.Hour
|
||||
)
|
||||
|
||||
// Load reads the JSON config at path and applies defaults. A missing file is
|
||||
@@ -520,6 +565,12 @@ func (c *Config) applyDefaults() {
|
||||
c.Digest.SeverityCeiling = 2
|
||||
}
|
||||
|
||||
// Absent block stays nil (⇒ silent detection). Present-but-partial gets the
|
||||
// cooldown default, so `{"notify": true}` is enough to switch it on.
|
||||
if c.PatternProposals != nil && c.PatternProposals.Cooldown <= 0 {
|
||||
c.PatternProposals.Cooldown = Duration(DefaultProposalCooldown)
|
||||
}
|
||||
|
||||
if c.Voice != nil {
|
||||
if c.Voice.RouterThreshold <= 0 {
|
||||
c.Voice.RouterThreshold = DefaultRouterThreshold
|
||||
|
||||
@@ -410,95 +410,12 @@ func TestChatViaClient(t *testing.T) {
|
||||
}
|
||||
|
||||
// chatTestAPI — a minimal CoreAPI that only implements Chat for testing.
|
||||
type chatTestAPI struct{}
|
||||
// Embeds UnimplementedCoreAPI so every other method fails loudly with
|
||||
// ErrNotImplemented instead of needing 27 hand-written no-op stubs.
|
||||
type chatTestAPI struct {
|
||||
UnimplementedCoreAPI
|
||||
}
|
||||
|
||||
func (a *chatTestAPI) WriteFact(ctx context.Context, req WriteFactReq) (int64, error) {
|
||||
return 0, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) LatestFact(ctx context.Context, key string) (Fact, error) {
|
||||
return Fact{}, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) LatestFactBySource(ctx context.Context, key, source string) (Fact, error) {
|
||||
return Fact{}, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) Since(ctx context.Context, key string, now time.Time) (time.Duration, error) {
|
||||
return 0, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) Presence(ctx context.Context) (Presence, error) {
|
||||
return Presence{}, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) CreateReminder(ctx context.Context, fire time.Time, payload, cron string) (int64, error) {
|
||||
return 0, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) MarkReminder(ctx context.Context, id int64, status string) error {
|
||||
return ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) ListReminders(ctx context.Context, n int) ([]Reminder, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) RecordNudge(ctx context.Context, rule, channel, message string, ts time.Time) (int64, error) {
|
||||
return 0, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) ResolveNudge(ctx context.Context, id int64, outcome string, ts time.Time) error {
|
||||
return ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) RecentOutcomes(ctx context.Context, rule string, n int) ([]string, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) RecentFacts(ctx context.Context, n int) ([]Fact, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) CalendarEvents(ctx context.Context, from, to time.Time) ([]Fact, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) RecentNudges(ctx context.Context, n int) ([]Nudge, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error) {
|
||||
return 0, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) QueryNotes(ctx context.Context, embedding []float32, k int) ([]Note, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) RecentNotes(ctx context.Context, n int) ([]Note, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) ProposeTool(ctx context.Context, name, utterance, scope string, ts time.Time) (bool, error) {
|
||||
return false, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) EnableTool(ctx context.Context, name string, cmd []string, destructive bool, scope string, ts time.Time) error {
|
||||
return ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) DisableTool(ctx context.Context, name string) error {
|
||||
return ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) DeleteTool(ctx context.Context, name string) error {
|
||||
return ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) ListProposedRoutines(ctx context.Context) ([]ProposedRoutine, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) AcceptProposedRoutine(ctx context.Context, id int64) error {
|
||||
return nil
|
||||
}
|
||||
func (a *chatTestAPI) DismissProposedRoutine(ctx context.Context, id int64) error {
|
||||
return ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) LookupTool(ctx context.Context, name string) (Tool, error) {
|
||||
return Tool{}, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) ListTools(ctx context.Context, status string) ([]Tool, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) RevertFact(ctx context.Context, key string) (int64, error) {
|
||||
return 0, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) TickTrace(ctx context.Context) (TickTrace, error) {
|
||||
return TickTrace{}, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error) {
|
||||
return nil, ErrUnknownMethod
|
||||
}
|
||||
func (a *chatTestAPI) Chat(ctx context.Context, text string) (string, error) {
|
||||
if text == "привет" {
|
||||
return "и тебе привет!", nil
|
||||
|
||||
+241
-306
@@ -501,6 +501,239 @@ func (s *Server) safeDispatch(ctx context.Context, req Request) (result json.Raw
|
||||
return s.dispatch(ctx, req)
|
||||
}
|
||||
|
||||
// handlerFunc — one table entry's shape: unmarshal req.Params (if it wants
|
||||
// any), call the matching CoreAPI method against the api passed in, marshal
|
||||
// the result. api is a parameter, not a closed-over field, precisely so a
|
||||
// table built once at package init never pins a stale CoreAPI — see the note
|
||||
// on methodTable below about SetAPI.
|
||||
type handlerFunc func(ctx context.Context, api CoreAPI, raw json.RawMessage) (json.RawMessage, error)
|
||||
|
||||
// withParams adapts a (typed params, typed result) CoreAPI call into a
|
||||
// handlerFunc: unmarshal into P, call fn, marshal R. On error the result is
|
||||
// dropped (marshalResult's output is never read when err != nil — see
|
||||
// serveConn) so every entry can uniformly return early on error without
|
||||
// re-deriving what the pre-table per-arm code used to return in that case.
|
||||
func withParams[P any, R any](fn func(ctx context.Context, api CoreAPI, p P) (R, error)) handlerFunc {
|
||||
return func(ctx context.Context, api CoreAPI, raw json.RawMessage) (json.RawMessage, error) {
|
||||
var p P
|
||||
if err := unmarshalParams(raw, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
r, err := fn(ctx, api, p)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(r), nil
|
||||
}
|
||||
}
|
||||
|
||||
// withParamsVoid is withParams for the error-only methods (mark/resolve/
|
||||
// enable/disable/...): params in, no result out, wire reply is always null.
|
||||
func withParamsVoid[P any](fn func(ctx context.Context, api CoreAPI, p P) error) handlerFunc {
|
||||
return func(ctx context.Context, api CoreAPI, raw json.RawMessage) (json.RawMessage, error) {
|
||||
var p P
|
||||
if err := unmarshalParams(raw, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(nil), fn(ctx, api, p)
|
||||
}
|
||||
}
|
||||
|
||||
// withoutParams is withParams for the handful of methods that take no
|
||||
// params at all (Presence, TickTrace, MorningStatus, ListProposedRoutines).
|
||||
// It does NOT call unmarshalParams — matching the pre-table arms, which
|
||||
// never touched req.Params for these four methods.
|
||||
func withoutParams[R any](fn func(ctx context.Context, api CoreAPI) (R, error)) handlerFunc {
|
||||
return func(ctx context.Context, api CoreAPI, _ json.RawMessage) (json.RawMessage, error) {
|
||||
r, err := fn(ctx, api)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(r), nil
|
||||
}
|
||||
}
|
||||
|
||||
// methodTable — one entry per CoreAPI-backed method. Built once at package
|
||||
// init, not per-Server and not per-dispatch: entries close over nothing but
|
||||
// the CoreAPI method being called, and dispatch passes in the *current*
|
||||
// api (loaded fresh via s.api.Load() every call, same as before the table
|
||||
// existed) as an argument — so SetAPI's runtime swap (the unlock transition)
|
||||
// is still honored on the very next request with no extra plumbing here.
|
||||
//
|
||||
// MethodAssertStepUp, MethodStoreEncryptionKey and MethodUnlock are NOT in
|
||||
// this table: they bypass CoreAPI entirely (s.StepUp / s.WrapKeyFn /
|
||||
// s.UnlockFn), so dispatch special-cases them before consulting the table.
|
||||
var methodTable = map[Method]handlerFunc{
|
||||
MethodWriteFact: withParams(func(ctx context.Context, api CoreAPI, p WriteFactReq) (idResp, error) {
|
||||
id, err := api.WriteFact(ctx, p)
|
||||
return idResp{ID: id}, err
|
||||
}),
|
||||
MethodLatestFact: withParams(func(ctx context.Context, api CoreAPI, p keyReq) (Fact, error) {
|
||||
return api.LatestFact(ctx, p.Key)
|
||||
}),
|
||||
MethodLatestFactBySource: withParams(func(ctx context.Context, api CoreAPI, p keySourceReq) (Fact, error) {
|
||||
return api.LatestFactBySource(ctx, p.Key, p.Source)
|
||||
}),
|
||||
MethodSince: withParams(func(ctx context.Context, api CoreAPI, p sinceReq) (sinceResp, error) {
|
||||
d, err := api.Since(ctx, p.Key, p.Now)
|
||||
return sinceResp{Dur: d}, err
|
||||
}),
|
||||
MethodPresence: withoutParams(func(ctx context.Context, api CoreAPI) (Presence, error) {
|
||||
return api.Presence(ctx)
|
||||
}),
|
||||
MethodCreateReminder: withParams(func(ctx context.Context, api CoreAPI, p createReminderReq) (idResp, error) {
|
||||
id, err := api.CreateReminder(ctx, p.Fire, p.Payload, p.Cron)
|
||||
return idResp{ID: id}, err
|
||||
}),
|
||||
MethodMarkReminder: withParamsVoid(func(ctx context.Context, api CoreAPI, p markReminderReq) error {
|
||||
return api.MarkReminder(ctx, p.ID, p.Status)
|
||||
}),
|
||||
MethodListReminders: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]Reminder, error) {
|
||||
out, err := api.ListReminders(ctx, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Reminder{}
|
||||
}
|
||||
return out, nil
|
||||
}),
|
||||
MethodRecordNudge: withParams(func(ctx context.Context, api CoreAPI, p recordNudgeReq) (idResp, error) {
|
||||
id, err := api.RecordNudge(ctx, p.Rule, p.Channel, p.Message, p.Ts)
|
||||
return idResp{ID: id}, err
|
||||
}),
|
||||
MethodResolveNudge: withParamsVoid(func(ctx context.Context, api CoreAPI, p resolveNudgeReq) error {
|
||||
return api.ResolveNudge(ctx, p.ID, p.Outcome, p.Ts)
|
||||
}),
|
||||
MethodRecentOutcomes: withParams(func(ctx context.Context, api CoreAPI, p outcomesReq) ([]string, error) {
|
||||
out, err := api.RecentOutcomes(ctx, p.Rule, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []string{} // stable non-null on the wire
|
||||
}
|
||||
return out, nil
|
||||
}),
|
||||
MethodRecentFacts: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]Fact, error) {
|
||||
out, err := api.RecentFacts(ctx, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Fact{}
|
||||
}
|
||||
return out, nil
|
||||
}),
|
||||
MethodCalendarEvents: withParams(func(ctx context.Context, api CoreAPI, p calendarEventsReq) ([]Fact, error) {
|
||||
out, err := api.CalendarEvents(ctx, p.From, p.To)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Fact{}
|
||||
}
|
||||
return out, nil
|
||||
}),
|
||||
MethodRecentNudges: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]Nudge, error) {
|
||||
out, err := api.RecentNudges(ctx, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Nudge{}
|
||||
}
|
||||
return out, nil
|
||||
}),
|
||||
MethodWriteNote: withParams(func(ctx context.Context, api CoreAPI, p writeNoteReq) (idResp, error) {
|
||||
id, err := api.WriteNote(ctx, p.Ts, p.Text, p.Embedding, p.Source)
|
||||
return idResp{ID: id}, err
|
||||
}),
|
||||
MethodQueryNotes: withParams(func(ctx context.Context, api CoreAPI, p queryNotesReq) ([]Note, error) {
|
||||
out, err := api.QueryNotes(ctx, p.Embedding, p.K)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Note{}
|
||||
}
|
||||
return out, nil
|
||||
}),
|
||||
MethodRecentNotes: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]Note, error) {
|
||||
out, err := api.RecentNotes(ctx, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Note{}
|
||||
}
|
||||
return out, nil
|
||||
}),
|
||||
MethodProposeTool: withParams(func(ctx context.Context, api CoreAPI, p proposeToolReq) (proposeToolResp, error) {
|
||||
ok, err := api.ProposeTool(ctx, p.Name, p.Utterance, p.Scope, p.Ts)
|
||||
return proposeToolResp{Proposed: ok}, err
|
||||
}),
|
||||
MethodEnableTool: withParamsVoid(func(ctx context.Context, api CoreAPI, p enableToolReq) error {
|
||||
return api.EnableTool(ctx, p.Name, p.Cmd, p.Destructive, p.Scope, p.Ts)
|
||||
}),
|
||||
MethodDisableTool: withParamsVoid(func(ctx context.Context, api CoreAPI, p disableToolReq) error {
|
||||
return api.DisableTool(ctx, p.Name)
|
||||
}),
|
||||
MethodLookupTool: withParams(func(ctx context.Context, api CoreAPI, p lookupToolReq) (Tool, error) {
|
||||
return api.LookupTool(ctx, p.Name)
|
||||
}),
|
||||
MethodListTools: withParams(func(ctx context.Context, api CoreAPI, p listToolsReq) (listToolsResp, error) {
|
||||
out, err := api.ListTools(ctx, p.Status)
|
||||
if err != nil {
|
||||
return listToolsResp{}, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Tool{}
|
||||
}
|
||||
return listToolsResp{Tools: out}, nil
|
||||
}),
|
||||
// MethodDeleteTool shares disableToolReq — both take just a tool name.
|
||||
MethodDeleteTool: withParamsVoid(func(ctx context.Context, api CoreAPI, p disableToolReq) error {
|
||||
return api.DeleteTool(ctx, p.Name)
|
||||
}),
|
||||
MethodListProposedRoutines: withoutParams(func(ctx context.Context, api CoreAPI) (listProposedRoutinesResp, error) {
|
||||
out, err := api.ListProposedRoutines(ctx)
|
||||
if err != nil {
|
||||
return listProposedRoutinesResp{}, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []ProposedRoutine{}
|
||||
}
|
||||
return listProposedRoutinesResp{Routines: out}, nil
|
||||
}),
|
||||
MethodDismissProposedRoutine: withParamsVoid(func(ctx context.Context, api CoreAPI, p dismissProposedRoutineReq) error {
|
||||
return api.DismissProposedRoutine(ctx, p.ID)
|
||||
}),
|
||||
MethodAcceptProposedRoutine: withParamsVoid(func(ctx context.Context, api CoreAPI, p acceptProposedRoutineReq) error {
|
||||
return api.AcceptProposedRoutine(ctx, p.ID)
|
||||
}),
|
||||
MethodRevertFact: withParams(func(ctx context.Context, api CoreAPI, p revertReq) (map[string]int64, error) {
|
||||
newID, err := api.RevertFact(ctx, p.Key)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return map[string]int64{"new_id": newID}, nil
|
||||
}),
|
||||
MethodChat: withParams(func(ctx context.Context, api CoreAPI, p chatReq) (chatResp, error) {
|
||||
reply, err := api.Chat(ctx, p.Text)
|
||||
return chatResp{Reply: reply}, err
|
||||
}),
|
||||
MethodTickTrace: withoutParams(func(ctx context.Context, api CoreAPI) (TickTrace, error) {
|
||||
return api.TickTrace(ctx)
|
||||
}),
|
||||
// MorningStatus intentionally has no nil→[]T{} normalization here — the
|
||||
// pre-table arm marshaled api.MorningStatus's result as-is (a nil slice
|
||||
// serializes as JSON null), and this preserves that exact wire shape.
|
||||
MethodMorningStatus: withoutParams(func(ctx context.Context, api CoreAPI) ([]MorningRoutineStatus, error) {
|
||||
return api.MorningStatus(ctx)
|
||||
}),
|
||||
}
|
||||
|
||||
// dispatch unmarshals params for req.Method and calls the matching CoreAPI
|
||||
// method. Unknown method ⇒ ErrUnknownMethod; a malformed params payload ⇒
|
||||
// ErrBadParams with the underlying text (local, server-side, not shipped to
|
||||
@@ -517,312 +750,11 @@ func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, er
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
// These three bypass CoreAPI entirely — they drive Server fields set
|
||||
// directly by the daemon (StepUp / WrapKeyFn / UnlockFn), not store
|
||||
// state, so they can never be table entries keyed on a CoreAPI method.
|
||||
switch req.Method {
|
||||
case MethodWriteFact:
|
||||
var p WriteFactReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
id, err := api.WriteFact(ctx, p)
|
||||
return marshalResult(idResp{ID: id}), err
|
||||
|
||||
case MethodLatestFact:
|
||||
var p keyReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
f, err := api.LatestFact(ctx, p.Key)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(f), nil
|
||||
|
||||
case MethodLatestFactBySource:
|
||||
var p keySourceReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
f, err := api.LatestFactBySource(ctx, p.Key, p.Source)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(f), nil
|
||||
|
||||
case MethodSince:
|
||||
var p sinceReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
d, err := api.Since(ctx, p.Key, p.Now)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(sinceResp{Dur: d}), nil
|
||||
|
||||
case MethodPresence:
|
||||
pres, err := api.Presence(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(pres), nil
|
||||
|
||||
case MethodCreateReminder:
|
||||
var p createReminderReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
id, err := api.CreateReminder(ctx, p.Fire, p.Payload, p.Cron)
|
||||
return marshalResult(idResp{ID: id}), err
|
||||
|
||||
case MethodMarkReminder:
|
||||
var p markReminderReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
err := api.MarkReminder(ctx, p.ID, p.Status)
|
||||
return marshalResult(nil), err
|
||||
|
||||
case MethodListReminders:
|
||||
var p nReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out, err := api.ListReminders(ctx, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Reminder{}
|
||||
}
|
||||
return marshalResult(out), nil
|
||||
|
||||
case MethodRecordNudge:
|
||||
var p recordNudgeReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
id, err := api.RecordNudge(ctx, p.Rule, p.Channel, p.Message, p.Ts)
|
||||
return marshalResult(idResp{ID: id}), err
|
||||
|
||||
case MethodResolveNudge:
|
||||
var p resolveNudgeReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
err := api.ResolveNudge(ctx, p.ID, p.Outcome, p.Ts)
|
||||
return marshalResult(nil), err
|
||||
|
||||
case MethodRecentOutcomes:
|
||||
var p outcomesReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out, err := api.RecentOutcomes(ctx, p.Rule, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []string{} // stable non-null on the wire
|
||||
}
|
||||
return marshalResult(out), nil
|
||||
|
||||
case MethodRecentFacts:
|
||||
var p nReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out, err := api.RecentFacts(ctx, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Fact{}
|
||||
}
|
||||
return marshalResult(out), nil
|
||||
|
||||
case MethodCalendarEvents:
|
||||
var p calendarEventsReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out, err := api.CalendarEvents(ctx, p.From, p.To)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Fact{}
|
||||
}
|
||||
return marshalResult(out), nil
|
||||
|
||||
case MethodRecentNudges:
|
||||
var p nReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out, err := api.RecentNudges(ctx, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Nudge{}
|
||||
}
|
||||
return marshalResult(out), nil
|
||||
|
||||
case MethodWriteNote:
|
||||
var p writeNoteReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
id, err := api.WriteNote(ctx, p.Ts, p.Text, p.Embedding, p.Source)
|
||||
return marshalResult(idResp{ID: id}), err
|
||||
|
||||
case MethodQueryNotes:
|
||||
var p queryNotesReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out, err := api.QueryNotes(ctx, p.Embedding, p.K)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Note{}
|
||||
}
|
||||
return marshalResult(out), nil
|
||||
|
||||
case MethodRecentNotes:
|
||||
var p nReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out, err := api.RecentNotes(ctx, p.N)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Note{}
|
||||
}
|
||||
return marshalResult(out), nil
|
||||
|
||||
case MethodProposeTool:
|
||||
var p proposeToolReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ok, err := api.ProposeTool(ctx, p.Name, p.Utterance, p.Scope, p.Ts)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(proposeToolResp{Proposed: ok}), nil
|
||||
|
||||
case MethodEnableTool:
|
||||
var p enableToolReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(nil), api.EnableTool(ctx, p.Name, p.Cmd, p.Destructive, p.Scope, p.Ts)
|
||||
|
||||
case MethodDisableTool:
|
||||
var p disableToolReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(nil), api.DisableTool(ctx, p.Name)
|
||||
|
||||
case MethodLookupTool:
|
||||
var p lookupToolReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
t, err := api.LookupTool(ctx, p.Name)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(t), nil
|
||||
|
||||
case MethodListTools:
|
||||
var p listToolsReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out, err := api.ListTools(ctx, p.Status)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []Tool{}
|
||||
}
|
||||
return marshalResult(listToolsResp{Tools: out}), nil
|
||||
|
||||
case MethodDeleteTool:
|
||||
var p disableToolReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(nil), api.DeleteTool(ctx, p.Name)
|
||||
|
||||
case MethodListProposedRoutines:
|
||||
out, err := api.ListProposedRoutines(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if out == nil {
|
||||
out = []ProposedRoutine{}
|
||||
}
|
||||
return marshalResult(listProposedRoutinesResp{Routines: out}), nil
|
||||
|
||||
case MethodDismissProposedRoutine:
|
||||
var p dismissProposedRoutineReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(nil), api.DismissProposedRoutine(ctx, p.ID)
|
||||
|
||||
case MethodAcceptProposedRoutine:
|
||||
var p acceptProposedRoutineReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(nil), api.AcceptProposedRoutine(ctx, p.ID)
|
||||
|
||||
case MethodRevertFact:
|
||||
var p struct {
|
||||
Key string `json:"key"`
|
||||
}
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
newID, err := api.RevertFact(ctx, p.Key)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(map[string]int64{"new_id": newID}), nil
|
||||
|
||||
case MethodChat:
|
||||
var p chatReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
reply, err := api.Chat(ctx, p.Text)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(chatResp{Reply: reply}), nil
|
||||
|
||||
case MethodTickTrace:
|
||||
t, err := api.TickTrace(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(t), nil
|
||||
|
||||
case MethodMorningStatus:
|
||||
s, err := api.MorningStatus(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(s), nil
|
||||
|
||||
case MethodAssertStepUp:
|
||||
if s.StepUp != nil {
|
||||
return marshalResult(nil), s.StepUp(ctx)
|
||||
@@ -848,10 +780,13 @@ func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, er
|
||||
return marshalResult(nil), s.UnlockFn(ctx, p.PublicKey)
|
||||
}
|
||||
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
|
||||
}
|
||||
|
||||
default:
|
||||
h, ok := methodTable[req.Method]
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
|
||||
}
|
||||
return h(ctx, api, req.Params)
|
||||
}
|
||||
|
||||
func unmarshalParams(raw json.RawMessage, v any) error {
|
||||
|
||||
@@ -0,0 +1,118 @@
|
||||
package ipc
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"time"
|
||||
)
|
||||
|
||||
// ErrNotImplemented is returned by every UnimplementedCoreAPI method. It is
|
||||
// deliberately distinct from ErrUnknownMethod (a wire-level "no such
|
||||
// method exists" verdict) and from any daemon-level "locked" error: this one
|
||||
// means "this method exists on CoreAPI, but the fake/adapter embedding
|
||||
// UnimplementedCoreAPI never got a real implementation for it." A test that
|
||||
// exercises an undeclared method fails loudly on this text instead of
|
||||
// silently nil-panicking or being mistaken for a legitimate failure.
|
||||
var ErrNotImplemented = errors.New("ipc: not implemented (unimplemented CoreAPI stub)")
|
||||
|
||||
// UnimplementedCoreAPI is the gRPC Unimplemented*Server pattern applied to
|
||||
// CoreAPI: embed it in a test double or adapter and override only the
|
||||
// methods you actually exercise. Every method returns ErrNotImplemented, so
|
||||
// a call that reaches an undeclared method fails loudly and specifically,
|
||||
// rather than compiling to a silent no-op or nil-pointer panic. This
|
||||
// replaces the old pattern of hand-writing all 30 no-op stubs per double —
|
||||
// those were compiler-satisfying padding, not tests of anything.
|
||||
type UnimplementedCoreAPI struct{}
|
||||
|
||||
var _ CoreAPI = UnimplementedCoreAPI{}
|
||||
|
||||
func (UnimplementedCoreAPI) WriteFact(ctx context.Context, req WriteFactReq) (int64, error) {
|
||||
return 0, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) LatestFact(ctx context.Context, key string) (Fact, error) {
|
||||
return Fact{}, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) LatestFactBySource(ctx context.Context, key, source string) (Fact, error) {
|
||||
return Fact{}, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) Since(ctx context.Context, key string, now time.Time) (time.Duration, error) {
|
||||
return 0, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) Presence(ctx context.Context) (Presence, error) {
|
||||
return Presence{}, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) CreateReminder(ctx context.Context, fire time.Time, payload, cron string) (int64, error) {
|
||||
return 0, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) MarkReminder(ctx context.Context, id int64, status string) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) ListReminders(ctx context.Context, n int) ([]Reminder, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) RecordNudge(ctx context.Context, rule, channel, message string, ts time.Time) (int64, error) {
|
||||
return 0, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) ResolveNudge(ctx context.Context, id int64, outcome string, ts time.Time) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) RecentOutcomes(ctx context.Context, rule string, n int) ([]string, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) RecentFacts(ctx context.Context, n int) ([]Fact, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) CalendarEvents(ctx context.Context, from, to time.Time) ([]Fact, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) RecentNudges(ctx context.Context, n int) ([]Nudge, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error) {
|
||||
return 0, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) QueryNotes(ctx context.Context, embedding []float32, k int) ([]Note, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) RecentNotes(ctx context.Context, n int) ([]Note, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) ProposeTool(ctx context.Context, name, utterance, scope string, ts time.Time) (bool, error) {
|
||||
return false, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) EnableTool(ctx context.Context, name string, cmd []string, destructive bool, scope string, ts time.Time) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) DisableTool(ctx context.Context, name string) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) DeleteTool(ctx context.Context, name string) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) ListProposedRoutines(ctx context.Context) ([]ProposedRoutine, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) DismissProposedRoutine(ctx context.Context, id int64) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) AcceptProposedRoutine(ctx context.Context, id int64) error {
|
||||
return ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) LookupTool(ctx context.Context, name string) (Tool, error) {
|
||||
return Tool{}, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) ListTools(ctx context.Context, status string) ([]Tool, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) RevertFact(ctx context.Context, key string) (int64, error) {
|
||||
return 0, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) TickTrace(ctx context.Context) (TickTrace, error) {
|
||||
return TickTrace{}, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error) {
|
||||
return nil, ErrNotImplemented
|
||||
}
|
||||
func (UnimplementedCoreAPI) Chat(ctx context.Context, text string) (string, error) {
|
||||
return "", ErrNotImplemented
|
||||
}
|
||||
@@ -286,6 +286,54 @@ func TestRemindersStillHonourSnooze(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------- digest eligibility ------------------------------
|
||||
|
||||
// A Sev2 care candidate (break) suppressed for a genuine restraint reason is
|
||||
// worth resurfacing later.
|
||||
func TestDigestEligibleSev2SuppressedByRestraint(t *testing.T) {
|
||||
for _, reason := range []string{"quiet_hours", "calendar_busy", "presence"} {
|
||||
if !DigestEligible(Sev2, reason) {
|
||||
t.Errorf("sev2 blocked by %q: want digest-eligible", reason)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A Sev1 care candidate (water/meal) never digests — a biological timer
|
||||
// nudge is stale by the time anyone could resurface it, so it just drops.
|
||||
func TestDigestEligibleSev1NeverDigests(t *testing.T) {
|
||||
for _, reason := range []string{"quiet_hours", "calendar_busy", "presence"} {
|
||||
if DigestEligible(Sev1, reason) {
|
||||
t.Errorf("sev1 blocked by %q: want drop, got digest-eligible", reason)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Ops severities are never blocked by these reasons in practice (Gate only
|
||||
// applies quiet_hours/calendar_busy/presence to care severities), but the
|
||||
// boundary itself must refuse to digest a high severity even if asked —
|
||||
// alarms bypass the gate and deliver now, unchanged, never delayed.
|
||||
func TestDigestEligibleNeverDigestsHighSeverity(t *testing.T) {
|
||||
for _, sev := range []Severity{Sev3, Sev4} {
|
||||
for _, reason := range []string{"quiet_hours", "calendar_busy", "presence"} {
|
||||
if DigestEligible(sev, reason) {
|
||||
t.Errorf("sev%d blocked by %q: high severity must never digest", sev, reason)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// cooldown and snooze are not "suppression" in the digest sense — cooldown
|
||||
// means it was already said recently, snooze means the user asked to not
|
||||
// hear about it. Neither should resurface later just because the severity
|
||||
// matches.
|
||||
func TestDigestEligibleExcludesCooldownAndSnooze(t *testing.T) {
|
||||
for _, reason := range []string{"cooldown", "snooze", "inert_no_data", "predicate", ""} {
|
||||
if DigestEligible(Sev2, reason) {
|
||||
t.Errorf("sev2 blocked by %q: should not be digest-eligible", reason)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// GAP — the gate reads State.SnoozeUntil, but the Gatherer hard-codes it to nil
|
||||
// (internal/loop/gather.go:153), so snooze is dead in the running daemon: the
|
||||
// unit tests above pass while nothing can ever populate the map. This asserts
|
||||
|
||||
@@ -105,6 +105,36 @@ func Tick(s State, rules []Rule) *Candidate {
|
||||
return fire
|
||||
}
|
||||
|
||||
// DigestEligible decides digest-vs-drop for a care candidate the gate
|
||||
// suppressed this tick (see ExplainGate's blockedBy). Pure — no I/O, no
|
||||
// state, just the two facts that matter: why it was suppressed, and how
|
||||
// insistent it was.
|
||||
//
|
||||
// Only genuine RESTRAINT blocks are eligible at all — quiet_hours,
|
||||
// calendar_busy, presence(away). cooldown and snooze are not suppression in
|
||||
// this sense: cooldown means "you already heard this recently" (resurfacing
|
||||
// it later would be an actual repeat, not a rescue) and snooze is the user
|
||||
// explicitly saying "not this" (digesting it anyway would defeat the ask).
|
||||
// Ops severities (Sev3/4) never reach here — the gate never blocks them for
|
||||
// these reasons in the first place (see Gate), and even if a future rule
|
||||
// dropped Sev3+ into "care", digest still refuses them: alarms bypass the
|
||||
// gate on purpose and must never be silently delayed into a bundle.
|
||||
//
|
||||
// Within care (Sev1–2), the boundary is severity itself: Sev1 (water, meal —
|
||||
// biological timers with no "still relevant later" property; a water nudge
|
||||
// from 3 hours into quiet hours is just wrong by morning) drops. Sev2
|
||||
// (break — "you worked through a long stretch without a break while I
|
||||
// couldn't reach you") is information that stays true and useful after the
|
||||
// fact, so it digests.
|
||||
func DigestEligible(sev Severity, blockedBy string) bool {
|
||||
switch blockedBy {
|
||||
case "quiet_hours", "calendar_busy", "presence":
|
||||
default:
|
||||
return false
|
||||
}
|
||||
return sev == Sev2
|
||||
}
|
||||
|
||||
// ReminderDecision — a due reminder the daemon should deliver now.
|
||||
// NOT gated by the universal Gate (per spec: "wake me 7" fires in quiet hours;
|
||||
// that's the point). Snooze is the one part of restraint that still applies.
|
||||
|
||||
@@ -20,9 +20,19 @@ type ProposedRoutine struct {
|
||||
const MaxIntervalRatio = 1.5
|
||||
|
||||
// MinEvents is the minimum number of events needed to detect a pattern.
|
||||
// With N events, there are N-1 intervals; we need at least 2 intervals
|
||||
// before proposing anything.
|
||||
const MinEvents = 3
|
||||
// With N events there are N-1 intervals, so 4 events means 3 intervals.
|
||||
//
|
||||
// This used to be 3 (two intervals), which is not a pattern — it is a
|
||||
// coincidence with a mean. Two gaps of similar length happen constantly:
|
||||
// water the plants on a Sunday, again the next Sunday, once more the Sunday
|
||||
// after, and a detector with a ±50% band calls that a weekly routine. The
|
||||
// cost of being wrong is asymmetric now that the digestion tick scans all of
|
||||
// history on its own schedule and can announce what it finds: a false
|
||||
// positive is something the owner has to read and dismiss, and a dismissal
|
||||
// is permanent, so one bad guess burns that action+object pair forever.
|
||||
// Three intervals is the cheapest bar that makes a run distinguishable from
|
||||
// a repeat. False negatives cost one more observation and nothing else.
|
||||
const MinEvents = 4
|
||||
|
||||
// Detect checks whether a sequence of events for the same action+object
|
||||
// forms a stable recurring pattern. Returns a ProposedRoutine when:
|
||||
|
||||
@@ -6,12 +6,13 @@ import (
|
||||
)
|
||||
|
||||
func TestDetectEnoughEvents(t *testing.T) {
|
||||
// 3 events with 7-day intervals → stable pattern
|
||||
// MinEvents events with 7-day intervals → stable pattern
|
||||
base := time.Date(2026, 7, 1, 12, 0, 0, 0, time.UTC)
|
||||
events := []Event{
|
||||
{Action: "refill", Object: "cat_water", Ts: base},
|
||||
{Action: "refill", Object: "cat_water", Ts: base.Add(7 * 24 * time.Hour)},
|
||||
{Action: "refill", Object: "cat_water", Ts: base.Add(14 * 24 * time.Hour)},
|
||||
{Action: "refill", Object: "cat_water", Ts: base.Add(21 * 24 * time.Hour)},
|
||||
}
|
||||
|
||||
r, err := Detect(events)
|
||||
@@ -24,8 +25,8 @@ func TestDetectEnoughEvents(t *testing.T) {
|
||||
if r.Action != "refill" || r.Object != "cat_water" {
|
||||
t.Fatalf("action/object: want refill/cat_water, got %s/%s", r.Action, r.Object)
|
||||
}
|
||||
if r.N != 3 {
|
||||
t.Fatalf("want N=3, got %d", r.N)
|
||||
if r.N != 4 {
|
||||
t.Fatalf("want N=4, got %d", r.N)
|
||||
}
|
||||
// ~7 days
|
||||
if r.IntervalDays < 6.9 || r.IntervalDays > 7.1 {
|
||||
@@ -33,19 +34,28 @@ func TestDetectEnoughEvents(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestDetectNotEnoughEvents — two intervals are a coincidence, not a routine
|
||||
// (Vikunja #43). Three same-day-of-week events used to be enough to propose a
|
||||
// weekly reminder; MinEvents is 4 now so a repeat has to happen a third time
|
||||
// before Maven calls it a pattern.
|
||||
func TestDetectNotEnoughEvents(t *testing.T) {
|
||||
base := time.Date(2026, 7, 1, 12, 0, 0, 0, time.UTC)
|
||||
events := []Event{
|
||||
{Action: "refill", Object: "cat_water", Ts: base},
|
||||
{Action: "refill", Object: "cat_water", Ts: base.Add(7 * 24 * time.Hour)},
|
||||
}
|
||||
|
||||
r, err := Detect(events)
|
||||
if err != nil {
|
||||
t.Fatalf("Detect: %v", err)
|
||||
}
|
||||
if r != nil {
|
||||
t.Fatal("want nil for <3 events")
|
||||
for _, n := range []int{1, 2, MinEvents - 1} {
|
||||
events := make([]Event, n)
|
||||
for i := range events {
|
||||
events[i] = Event{
|
||||
Action: "refill",
|
||||
Object: "cat_water",
|
||||
Ts: base.Add(time.Duration(i) * 7 * 24 * time.Hour),
|
||||
}
|
||||
}
|
||||
r, err := Detect(events)
|
||||
if err != nil {
|
||||
t.Fatalf("Detect(%d events): %v", n, err)
|
||||
}
|
||||
if r != nil {
|
||||
t.Fatalf("Detect(%d events) proposed %+v, want nil below MinEvents=%d", n, r, MinEvents)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -68,12 +78,13 @@ func TestDetectEmpty(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestDetectIrregularRejects(t *testing.T) {
|
||||
// 3 events but wildly irregular: 1 day, then 14 days → ratio 14 > 1.5
|
||||
// wildly irregular: 1 day, then 14 days → ratio 14 > 1.5
|
||||
base := time.Date(2026, 7, 1, 12, 0, 0, 0, time.UTC)
|
||||
events := []Event{
|
||||
{Action: "refill", Object: "cat_water", Ts: base},
|
||||
{Action: "refill", Object: "cat_water", Ts: base.Add(1 * 24 * time.Hour)},
|
||||
{Action: "refill", Object: "cat_water", Ts: base.Add(15 * 24 * time.Hour)},
|
||||
{Action: "refill", Object: "cat_water", Ts: base.Add(16 * 24 * time.Hour)},
|
||||
}
|
||||
|
||||
r, err := Detect(events)
|
||||
@@ -117,6 +128,7 @@ func TestDetectSameTimestamp(t *testing.T) {
|
||||
{Action: "refill", Object: "cat_water", Ts: base},
|
||||
{Action: "refill", Object: "cat_water", Ts: base},
|
||||
{Action: "refill", Object: "cat_water", Ts: base.Add(7 * 24 * time.Hour)},
|
||||
{Action: "refill", Object: "cat_water", Ts: base.Add(14 * 24 * time.Hour)},
|
||||
}
|
||||
|
||||
r, err := Detect(events)
|
||||
|
||||
@@ -551,12 +551,21 @@ func (p *LLMPhraser) chatWithSystem(ctx context.Context, system, user string, ma
|
||||
// Russian only, feminine self-reference, second person masculine (the owner is
|
||||
// a man). She talks TO him, informally, singular — never "вы", never "он".
|
||||
// One short sentence — the nudge is spoken aloud.
|
||||
//
|
||||
// What the ban on обращения forbids is pet names ("дорогой", "милый"), not his
|
||||
// name: "Ками, ноутбук на трёх процентах" is exactly how she talks, and the
|
||||
// unqualified word read as forbidding that too. Hence "ласковые обращения".
|
||||
//
|
||||
// The examples also never claim a physical act. She has no hands and no smart
|
||||
// plug — she can tell him the battery is at three percent, she cannot put the
|
||||
// laptop on charge. An example that says she did teaches the model to invent
|
||||
// actions Maven never took, which is worse than a missing nudge.
|
||||
const nudgeSystem = `Ты — Maven, домашняя ассистентка. О себе говоришь в женском роде ("я проверила", "я записала"). Владелец — мужчина, обращайся к нему в мужском роде ("ты пил", "ты забыл").
|
||||
Говоришь с ним на "ты", в единственном числе ("выпей", "встань"). Никогда не "вы"/"вас"/"ваш" и никогда "он"/"его" — ты говоришь ему, а не о нём.
|
||||
|
||||
Пиши ОДНО короткое напоминание по-русски: не больше 120 символов и не больше 16 слов. Только по делу.
|
||||
|
||||
Запрещено: обращения ("дорогой", "милый"), эмодзи, извинения ("прости", "извини"), вопросы о самочувствии, похвала, больше одного восклицательного знака, английские слова кроме имён сервисов.
|
||||
Запрещено: ласковые обращения ("дорогой", "милый"), эмодзи, извинения ("прости", "извини"), вопросы о самочувствии, похвала, больше одного восклицательного знака, английские слова кроме имён сервисов.
|
||||
|
||||
Отвечай ТОЛЬКО одним объектом JSON с полями "response" и "mood".
|
||||
"response" — сам текст напоминания.
|
||||
@@ -564,7 +573,7 @@ const nudgeSystem = `Ты — Maven, домашняя ассистентка. О
|
||||
|
||||
Так выглядит правильный ответ по форме. Темы здесь посторонние — их в запросе не будет:
|
||||
{"response": "Стиральная машина закончила. Развесь бельё.", "mood": "neutral"}
|
||||
{"response": "Ноутбук на трёх процентах. Я поставила его на зарядку.", "mood": "confused"}
|
||||
{"response": "Ками, ноутбук на трёх процентах. Поставь его на зарядку.", "mood": "confused"}
|
||||
|
||||
Это примеры ФОРМЫ, а не темы. Пиши только про ту ситуацию, которую тебе дали в запросе. Не копируй примеры и никогда не пиши "..." в поле response.`
|
||||
|
||||
|
||||
@@ -118,6 +118,39 @@ const routeRepeatPenalty = 1.15
|
||||
// Route return ok=false so the caller drops to the classifier cascade.
|
||||
const routeIntentUnknown = "unknown"
|
||||
|
||||
// llmFullConfidence / llmThinConfidence — Vikunja #359. Confidence used to be
|
||||
// hardcoded to 1.0 for every LLM decision, so the stage-3 gate in router.go
|
||||
// never had anything to bite on and the LLM path could never produce a
|
||||
// Clarify: on the 77-case RU fixture, 6/6 want_clarify cases were missed by
|
||||
// EVERY model in the 31-07-2026 bake-off (0.8B through 2B) — proof this was a
|
||||
// code bug, not a capability ceiling.
|
||||
//
|
||||
// The fix does not touch the prompt (routeSystem is under
|
||||
// llm/check_prompt_parity.py in the training workspace; changing its text
|
||||
// creates a parity break that has to be fixed there too — see Vikunja #362).
|
||||
// Instead it reads structural signal that is already free:
|
||||
// - a single-token utterance is thin evidence for anything a grammar
|
||||
// didn't already catch at stage 0 — "вода" and "бэкап" alone don't say
|
||||
// fact-vs-query or act-vs-report;
|
||||
// - a fact with no key, or an act that never resolves to an allowlisted fn
|
||||
// (checked in router.go, after slot-fill has had its say), is a decision
|
||||
// with a hole in the one slot that makes it actionable.
|
||||
//
|
||||
// A model self-reporting confidence in the JSON was considered and rejected:
|
||||
// a sub-2B is not calibrated (nothing stops it saying "confident" on exactly
|
||||
// the cases it gets wrong today), and true logprobs would need a response
|
||||
// field internal/llm.Client's Complete does not currently return — see
|
||||
// internal/llm/client.go.
|
||||
//
|
||||
// llmThinConfidence sits below config.DefaultRouterThreshold (0.55) so the
|
||||
// existing stage-3 gate in Router.Route treats it exactly like a low-scoring
|
||||
// classifier result — same lane, same daemon-side clarify machinery
|
||||
// (cmd/mavend/clarify.go), no new consumer to build.
|
||||
const (
|
||||
llmFullConfidence = 1.0
|
||||
llmThinConfidence = 0.3
|
||||
)
|
||||
|
||||
type routeAction struct {
|
||||
Intent string `json:"intent"`
|
||||
Key string `json:"key"`
|
||||
@@ -160,7 +193,15 @@ func (lr *LLMRouter) Route(ctx context.Context, utterance string, now time.Time)
|
||||
if a.Intent == routeIntentUnknown {
|
||||
return Decision{}, false, nil
|
||||
}
|
||||
d := Decision{Utterance: utterance, Stage: 1, Confidence: 1.0}
|
||||
d := Decision{Utterance: utterance, Stage: 1, Confidence: llmFullConfidence}
|
||||
// A single-token utterance is thin evidence: the model had nothing to
|
||||
// disambiguate on ("вода" is a fact-or-query coin flip, "бэкап" an
|
||||
// act-or-report one) and stage 0 would already have won on anything
|
||||
// that pattern-matches cleanly. Flag it now; router.go's stage-3 gate
|
||||
// (Router.Route) decides whether that trips Clarify.
|
||||
if len(strings.Fields(utterance)) <= 1 {
|
||||
d.Confidence = llmThinConfidence
|
||||
}
|
||||
switch Intent(a.Intent) {
|
||||
case IntentFact:
|
||||
d.Intent = IntentFact
|
||||
|
||||
@@ -258,4 +258,101 @@ func TestLLMFactGetsKeyFromParser(t *testing.T) {
|
||||
if !d.Slots.HasKey || d.Slots.Key != "water" {
|
||||
t.Fatalf("want key=water, got %+v", d.Slots)
|
||||
}
|
||||
if d.Clarify {
|
||||
t.Fatalf("the parser resolved the key, this must not clarify: %+v", d)
|
||||
}
|
||||
}
|
||||
|
||||
// --- confidence / stage-3 gate on the LLM path (Vikunja #359) -----------------
|
||||
|
||||
// A single-token utterance is thin evidence on its own — "вода" alone is a
|
||||
// fact/query coin flip. The gate must ask rather than guess confidently.
|
||||
func TestLLMRouterSingleTokenTripsClarify(t *testing.T) {
|
||||
r := newLLMTestRouter(t, `{"intent":"query","text":"вода"}`)
|
||||
d, err := r.Route(context.Background(), "вода", refNow())
|
||||
if err != nil {
|
||||
t.Fatalf("route: %v", err)
|
||||
}
|
||||
if !d.Clarify {
|
||||
t.Fatalf("a bare single-token decision must clarify, got %+v", d)
|
||||
}
|
||||
}
|
||||
|
||||
// A multi-word utterance with a clean answer must not be punished — the
|
||||
// whole point is not trading the confident cases away for clarify coverage.
|
||||
func TestLLMRouterMultiWordStaysConfident(t *testing.T) {
|
||||
r := newLLMTestRouter(t, `{"intent":"reminder","text":"позвонить маме"}`)
|
||||
d, err := r.Route(context.Background(), "напомни позвонить маме", refNow())
|
||||
if err != nil {
|
||||
t.Fatalf("route: %v", err)
|
||||
}
|
||||
if d.Clarify {
|
||||
t.Fatalf("a clean multi-word decision must not clarify: %+v", d)
|
||||
}
|
||||
if d.Confidence != llmFullConfidence {
|
||||
t.Fatalf("want full confidence, got %v", d.Confidence)
|
||||
}
|
||||
}
|
||||
|
||||
// "бэкап" alone: the model guesses act, but nothing on the allowlist matches
|
||||
// "бэкап" as a verb — that must not fire a tool blind.
|
||||
func TestLLMRouterActWithoutFnTripsClarify(t *testing.T) {
|
||||
r := newLLMTestRouter(t, `{"intent":"act","verb":"бэкап"}`)
|
||||
d, err := r.Route(context.Background(), "бэкап сделай пожалуйста расписание", refNow())
|
||||
if err != nil {
|
||||
t.Fatalf("route: %v", err)
|
||||
}
|
||||
if d.Slots.HasFn {
|
||||
t.Fatalf("test setup drifted: %q now resolves to an fn", d.Slots.Fn)
|
||||
}
|
||||
if !d.Clarify {
|
||||
t.Fatalf("an unresolved act must clarify rather than guess: %+v", d)
|
||||
}
|
||||
}
|
||||
|
||||
// An act that DOES resolve to an allowlisted fn must stay confident even
|
||||
// though its own verb is single-word-ish in spirit — guard against the fn
|
||||
// check firing on the happy path.
|
||||
func TestLLMRouterActWithFnStaysConfident(t *testing.T) {
|
||||
r := newLLMTestRouter(t, `{"intent":"act","verb":"restart nginx"}`)
|
||||
d, err := r.Route(context.Background(), "слушай, restart nginx пожалуйста", refNow())
|
||||
if err != nil {
|
||||
t.Fatalf("route: %v", err)
|
||||
}
|
||||
if !d.Slots.HasFn {
|
||||
t.Fatalf("test setup drifted, want fn resolved: %+v", d.Slots)
|
||||
}
|
||||
if d.Clarify {
|
||||
t.Fatalf("a resolved act must not clarify: %+v", d)
|
||||
}
|
||||
}
|
||||
|
||||
// A fact where NEITHER the model NOR the deterministic parser can name a key
|
||||
// must clarify instead of silently writing under an empty/guessed key.
|
||||
func TestLLMRouterFactWithoutKeyTripsClarify(t *testing.T) {
|
||||
r := newLLMTestRouter(t, `{"intent":"fact","value":"что-то"}`)
|
||||
d, err := r.Route(context.Background(), "у меня какая-то фигня случилась вот прямо только что", refNow())
|
||||
if err != nil {
|
||||
t.Fatalf("route: %v", err)
|
||||
}
|
||||
if d.Slots.HasKey {
|
||||
t.Fatalf("test setup drifted: parser now resolves a key for this utterance")
|
||||
}
|
||||
if !d.Clarify {
|
||||
t.Fatalf("a keyless fact must clarify rather than guess: %+v", d)
|
||||
}
|
||||
}
|
||||
|
||||
// The whole point of #359: the classifier cascade cannot be traded away for
|
||||
// clarify coverage. A multi-word fact the parser CAN key must stay confident
|
||||
// through the full Router.Route path, not just the raw LLMRouter.
|
||||
func TestRouterLLMFactWithResolvedKeyStaysConfident(t *testing.T) {
|
||||
r := newLLMTestRouter(t, `{"intent":"fact","text":"я выпил воду"}`)
|
||||
d, err := r.Route(context.Background(), "я выпил воду", refNow())
|
||||
if err != nil {
|
||||
t.Fatalf("route: %v", err)
|
||||
}
|
||||
if d.Clarify {
|
||||
t.Fatalf("a fact the parser could key must not clarify: %+v", d)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -89,6 +89,7 @@ func (r *Router) Route(ctx context.Context, utterance string, now time.Time) (De
|
||||
if d, ok, err := r.llm.Route(ctx, utterance, now); err == nil && ok {
|
||||
d.Utterance = utterance
|
||||
r.fillSlots(ctx, &d, now)
|
||||
r.gateLLMDecision(&d)
|
||||
return d, nil
|
||||
} else if err != nil {
|
||||
log.Printf("router: llm route fell back to classifier: %v", err)
|
||||
@@ -152,6 +153,35 @@ func (r *Router) fillSlots(ctx context.Context, d *Decision, now time.Time) {
|
||||
// Stage stays 1: it says who decided the route, and that was the LLM.
|
||||
}
|
||||
|
||||
// gateLLMDecision — stage 3 for the LLM path (Vikunja #359). This used to be
|
||||
// the classifier's job alone (see the threshold check at the bottom of
|
||||
// Route): the LLM branch returned straight from fillSlots and never touched
|
||||
// r.threshold at all, so a hardcoded Confidence: 1.0 in llmrouter.go could
|
||||
// never gate. Two more structural holes are checked here, after fillSlots
|
||||
// has had a chance to fill them from the deterministic parsers — checking
|
||||
// before fillSlots would flag e.g. every keyless fact the fact parser goes
|
||||
// on to resolve (TestLLMFactGetsKeyFromParser):
|
||||
// - a fact with no key even after the parser tried — nothing to write, or
|
||||
// worse, a confident write under the wrong key;
|
||||
// - an act that never resolved to an allowlisted fn — a confident guess
|
||||
// here means either silently doing nothing or, if the daemon is lax,
|
||||
// running something never on the allowlist. Don't guess; ask.
|
||||
//
|
||||
// Anything below threshold gets the exact same Clarify=true treatment the
|
||||
// classifier path already produces — same field, same daemon-side consumer
|
||||
// (cmd/mavend/clarify.go), nothing new to wire.
|
||||
func (r *Router) gateLLMDecision(d *Decision) {
|
||||
if d.Intent == IntentFact && !d.Slots.HasKey && d.Confidence > llmThinConfidence {
|
||||
d.Confidence = llmThinConfidence
|
||||
}
|
||||
if d.Intent == IntentAct && !d.Slots.HasFn && d.Confidence > llmThinConfidence {
|
||||
d.Confidence = llmThinConfidence
|
||||
}
|
||||
if d.Confidence < r.threshold {
|
||||
d.Clarify = true
|
||||
}
|
||||
}
|
||||
|
||||
// CorrectMisroute — the user corrected a bad classification. Appends a new
|
||||
// example for the corrected intent (append-only — grows the classifier, no
|
||||
// retrain). Same shape as nudges.outcome tuning cooldowns: more reliable over
|
||||
|
||||
@@ -0,0 +1,142 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Digest entry statuses. pending = enqueued, waiting for a drain. drained =
|
||||
// spoken as part of a bundle. expired = the tick loop's expiry sweep found it
|
||||
// past its expires_ts before a drain happened — dropped, not delivered late.
|
||||
const (
|
||||
DigestPending = "pending"
|
||||
DigestDrained = "drained"
|
||||
DigestExpired = "expired"
|
||||
)
|
||||
|
||||
// DigestEntry — one gate-suppressed care candidate durably held for later
|
||||
// bundled delivery.
|
||||
type DigestEntry struct {
|
||||
ID int64
|
||||
Rule string
|
||||
Severity int
|
||||
Body string
|
||||
CreatedTs time.Time
|
||||
ExpiresTs time.Time
|
||||
}
|
||||
|
||||
// DigestBodyHash is the dedupe key for a digest entry: same rule, same
|
||||
// wording ⇒ the same suppressed nudge repeating across ticks, and he should
|
||||
// hear it once, not once per tick it kept getting suppressed.
|
||||
func DigestBodyHash(rule, body string) string {
|
||||
sum := sha256.Sum256([]byte(rule + "\x00" + body))
|
||||
return hex.EncodeToString(sum[:8])
|
||||
}
|
||||
|
||||
// EnqueueDigestEntry durably records a suppressed care candidate worth
|
||||
// resurfacing later. If a pending entry with the same rule+body already
|
||||
// exists, this is a no-op that returns the existing id and deduped=true —
|
||||
// the same suppressed nudge repeating across ticks must not pile up into
|
||||
// several copies of itself in the eventual bundle.
|
||||
func (s *Store) EnqueueDigestEntry(ctx context.Context, rule string, severity int, body string, now, expiresAt time.Time) (id int64, deduped bool, err error) {
|
||||
hash := DigestBodyHash(rule, body)
|
||||
var existing int64
|
||||
err = s.db.QueryRowContext(ctx,
|
||||
`SELECT id FROM digest_entries WHERE status = ? AND rule = ? AND body_hash = ? LIMIT 1`,
|
||||
DigestPending, rule, hash).Scan(&existing)
|
||||
if err == nil {
|
||||
return existing, true, nil
|
||||
}
|
||||
|
||||
res, err := s.db.ExecContext(ctx,
|
||||
`INSERT INTO digest_entries (rule, severity, body, body_hash, status, created_ts, expires_ts)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)`,
|
||||
rule, severity, body, hash, DigestPending, now.UnixMilli(), expiresAt.UnixMilli())
|
||||
if err != nil {
|
||||
return 0, false, fmt.Errorf("enqueue digest entry: %w", err)
|
||||
}
|
||||
id, err = res.LastInsertId()
|
||||
if err != nil {
|
||||
return 0, false, fmt.Errorf("enqueue digest entry: last insert id: %w", err)
|
||||
}
|
||||
return id, false, nil
|
||||
}
|
||||
|
||||
// PendingDigestEntries returns the live (not yet expired) pending entries,
|
||||
// oldest first — the order they were suppressed in, which is also the order
|
||||
// a bundled readout should mention them.
|
||||
func (s *Store) PendingDigestEntries(ctx context.Context, now time.Time) ([]DigestEntry, error) {
|
||||
rows, err := s.db.QueryContext(ctx,
|
||||
`SELECT id, rule, severity, body, created_ts, expires_ts
|
||||
FROM digest_entries WHERE status = ? AND expires_ts > ? ORDER BY created_ts ASC`,
|
||||
DigestPending, now.UnixMilli())
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("pending digest entries: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var out []DigestEntry
|
||||
for rows.Next() {
|
||||
var e DigestEntry
|
||||
var created, expires int64
|
||||
if err := rows.Scan(&e.ID, &e.Rule, &e.Severity, &e.Body, &created, &expires); err != nil {
|
||||
return nil, fmt.Errorf("pending digest entries: scan: %w", err)
|
||||
}
|
||||
e.CreatedTs = time.UnixMilli(created)
|
||||
e.ExpiresTs = time.UnixMilli(expires)
|
||||
out = append(out, e)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// ExpireStaleDigestEntries marks pending entries whose expires_ts has passed
|
||||
// as expired — stale information (yesterday's battery warning) is noise, not
|
||||
// news, so it is dropped rather than delivered late. Called once per tick,
|
||||
// mirroring ReconcileStaleDeliveryAttempts's "sweep, don't guess" shape.
|
||||
// Returns the count expired, for logging.
|
||||
func (s *Store) ExpireStaleDigestEntries(ctx context.Context, now time.Time) (int, error) {
|
||||
res, err := s.db.ExecContext(ctx,
|
||||
`UPDATE digest_entries SET status = ? WHERE status = ? AND expires_ts <= ?`,
|
||||
DigestExpired, DigestPending, now.UnixMilli())
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("expire stale digest entries: %w", err)
|
||||
}
|
||||
n, err := res.RowsAffected()
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("expire stale digest entries: rows affected: %w", err)
|
||||
}
|
||||
return int(n), nil
|
||||
}
|
||||
|
||||
// DrainDigestEntries marks the given entries drained — they were folded into
|
||||
// a bundle that was successfully dispatched. Called only after a successful
|
||||
// send, same rule as the delivery outbox: a failed dispatch must not mark
|
||||
// entries drained, or the bundle is lost along with the failed send.
|
||||
func (s *Store) DrainDigestEntries(ctx context.Context, ids []int64, now time.Time) error {
|
||||
if len(ids) == 0 {
|
||||
return nil
|
||||
}
|
||||
tx, err := s.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("drain digest entries: begin: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
stmt, err := tx.PrepareContext(ctx,
|
||||
`UPDATE digest_entries SET status = ? WHERE id = ? AND status = ?`)
|
||||
if err != nil {
|
||||
return fmt.Errorf("drain digest entries: prepare: %w", err)
|
||||
}
|
||||
defer stmt.Close()
|
||||
for _, id := range ids {
|
||||
if _, err := stmt.ExecContext(ctx, DigestDrained, id, DigestPending); err != nil {
|
||||
return fmt.Errorf("drain digest entry %d: %w", id, err)
|
||||
}
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return fmt.Errorf("drain digest entries: commit: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,202 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// TestDigestEntryRoundTrips — a suppressed care candidate lands durably and
|
||||
// comes back out of PendingDigestEntries with its severity and body intact.
|
||||
func TestDigestEntryRoundTrips(t *testing.T) {
|
||||
s := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := time.Now()
|
||||
|
||||
id, deduped, err := s.EnqueueDigestEntry(ctx, "break", 2, "ты долго не отдыхала", now, now.Add(24*time.Hour))
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
if deduped {
|
||||
t.Fatal("first enqueue must not report deduped")
|
||||
}
|
||||
if id == 0 {
|
||||
t.Fatal("want a nonzero id")
|
||||
}
|
||||
|
||||
entries, err := s.PendingDigestEntries(ctx, now)
|
||||
if err != nil {
|
||||
t.Fatalf("pending: %v", err)
|
||||
}
|
||||
if len(entries) != 1 || entries[0].ID != id {
|
||||
t.Fatalf("want 1 pending entry with id %d, got %+v", id, entries)
|
||||
}
|
||||
if entries[0].Rule != "break" || entries[0].Severity != 2 || entries[0].Body != "ты долго не отдыхала" {
|
||||
t.Fatalf("entry contents wrong: %+v", entries[0])
|
||||
}
|
||||
}
|
||||
|
||||
// TestDigestEntrySurvivesRestart — durability is the whole point: a fresh
|
||||
// Store handle on the same file must see the same pending entry, exactly
|
||||
// like the delivery outbox's crash-recovery promise.
|
||||
func TestDigestEntrySurvivesRestart(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
ctx := context.Background()
|
||||
now := time.Now()
|
||||
|
||||
s1, err := Open(ctx, dir+"/m.db")
|
||||
if err != nil {
|
||||
t.Fatalf("open: %v", err)
|
||||
}
|
||||
id, _, err := s1.EnqueueDigestEntry(ctx, "break", 2, "перерыв", now, now.Add(24*time.Hour))
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
if err := s1.Close(); err != nil {
|
||||
t.Fatalf("close: %v", err)
|
||||
}
|
||||
|
||||
// simulated restart: a brand new Store handle on the same file.
|
||||
s2, err := Open(ctx, dir+"/m.db")
|
||||
if err != nil {
|
||||
t.Fatalf("reopen: %v", err)
|
||||
}
|
||||
defer func() { _ = s2.Close() }()
|
||||
|
||||
entries, err := s2.PendingDigestEntries(ctx, now)
|
||||
if err != nil {
|
||||
t.Fatalf("pending after restart: %v", err)
|
||||
}
|
||||
if len(entries) != 1 || entries[0].ID != id {
|
||||
t.Fatalf("digest entry did not survive restart: %+v", entries)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDigestEntryDedupesSameRuleAndBody — the same suppressed nudge
|
||||
// repeating across ticks (quiet hours holding for hours) must not pile up
|
||||
// into several copies of itself; he hears it once.
|
||||
func TestDigestEntryDedupesSameRuleAndBody(t *testing.T) {
|
||||
s := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := time.Now()
|
||||
|
||||
id1, deduped1, err := s.EnqueueDigestEntry(ctx, "break", 2, "перерыв нужен", now, now.Add(24*time.Hour))
|
||||
if err != nil {
|
||||
t.Fatalf("first enqueue: %v", err)
|
||||
}
|
||||
if deduped1 {
|
||||
t.Fatal("first enqueue should not be deduped")
|
||||
}
|
||||
|
||||
for i := 0; i < 2; i++ {
|
||||
id2, deduped2, err := s.EnqueueDigestEntry(ctx, "break", 2, "перерыв нужен", now.Add(time.Minute), now.Add(25*time.Hour))
|
||||
if err != nil {
|
||||
t.Fatalf("repeat enqueue: %v", err)
|
||||
}
|
||||
if !deduped2 {
|
||||
t.Fatal("repeat enqueue of the same rule+body should report deduped")
|
||||
}
|
||||
if id2 != id1 {
|
||||
t.Fatalf("deduped enqueue should return the original id: want %d got %d", id1, id2)
|
||||
}
|
||||
}
|
||||
|
||||
entries, err := s.PendingDigestEntries(ctx, now)
|
||||
if err != nil {
|
||||
t.Fatalf("pending: %v", err)
|
||||
}
|
||||
if len(entries) != 1 {
|
||||
t.Fatalf("want exactly 1 pending entry after 3 enqueues of the same nudge, got %d", len(entries))
|
||||
}
|
||||
}
|
||||
|
||||
// TestDigestEntryExpiresRatherThanDeliversLate — a stale entry (past its
|
||||
// expires_ts) must not surface in PendingDigestEntries, and the sweep should
|
||||
// mark it expired instead of leaving it around to be delivered late.
|
||||
func TestDigestEntryExpiresRatherThanDeliversLate(t *testing.T) {
|
||||
s := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
created := time.Now()
|
||||
expiresAt := created.Add(time.Hour)
|
||||
|
||||
id, _, err := s.EnqueueDigestEntry(ctx, "water", 1, "стакан воды", created, expiresAt)
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
|
||||
afterExpiry := expiresAt.Add(time.Minute)
|
||||
|
||||
// even before the sweep runs, a stale entry must not be handed back as
|
||||
// pending — "not yet swept" must not mean "still deliverable".
|
||||
entries, err := s.PendingDigestEntries(ctx, afterExpiry)
|
||||
if err != nil {
|
||||
t.Fatalf("pending: %v", err)
|
||||
}
|
||||
if len(entries) != 0 {
|
||||
t.Fatalf("stale entry must not be returned as pending, got %+v", entries)
|
||||
}
|
||||
|
||||
n, err := s.ExpireStaleDigestEntries(ctx, afterExpiry)
|
||||
if err != nil {
|
||||
t.Fatalf("expire sweep: %v", err)
|
||||
}
|
||||
if n != 1 {
|
||||
t.Fatalf("want 1 entry expired, got %d", n)
|
||||
}
|
||||
|
||||
var status string
|
||||
if err := s.db.QueryRowContext(ctx, `SELECT status FROM digest_entries WHERE id = ?`, id).Scan(&status); err != nil {
|
||||
t.Fatalf("read back: %v", err)
|
||||
}
|
||||
if status != DigestExpired {
|
||||
t.Fatalf("status: want %q, got %q", DigestExpired, status)
|
||||
}
|
||||
|
||||
// idempotent: a second sweep finds nothing new.
|
||||
n2, err := s.ExpireStaleDigestEntries(ctx, afterExpiry.Add(time.Hour))
|
||||
if err != nil {
|
||||
t.Fatalf("second sweep: %v", err)
|
||||
}
|
||||
if n2 != 0 {
|
||||
t.Fatalf("second sweep should find nothing, got %d", n2)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDigestEntryDrainMarksDrainedNotDeleted — draining is bookkeeping, not
|
||||
// deletion: the row survives as an audit trail of what she actually said.
|
||||
func TestDigestEntryDrainMarksDrainedNotDeleted(t *testing.T) {
|
||||
s := newTestStore(t)
|
||||
ctx := context.Background()
|
||||
now := time.Now()
|
||||
|
||||
id1, _, err := s.EnqueueDigestEntry(ctx, "break", 2, "перерыв", now, now.Add(24*time.Hour))
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue 1: %v", err)
|
||||
}
|
||||
id2, _, err := s.EnqueueDigestEntry(ctx, "break2", 2, "другое", now, now.Add(24*time.Hour))
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue 2: %v", err)
|
||||
}
|
||||
|
||||
if err := s.DrainDigestEntries(ctx, []int64{id1, id2}, now.Add(time.Hour)); err != nil {
|
||||
t.Fatalf("drain: %v", err)
|
||||
}
|
||||
|
||||
entries, err := s.PendingDigestEntries(ctx, now.Add(time.Hour))
|
||||
if err != nil {
|
||||
t.Fatalf("pending: %v", err)
|
||||
}
|
||||
if len(entries) != 0 {
|
||||
t.Fatalf("drained entries must not still be pending, got %+v", entries)
|
||||
}
|
||||
|
||||
for _, id := range []int64{id1, id2} {
|
||||
var status string
|
||||
if err := s.db.QueryRowContext(ctx, `SELECT status FROM digest_entries WHERE id = ?`, id).Scan(&status); err != nil {
|
||||
t.Fatalf("read back %d: %v", id, err)
|
||||
}
|
||||
if status != DigestDrained {
|
||||
t.Fatalf("entry %d status: want %q, got %q", id, DigestDrained, status)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -35,6 +35,36 @@ func (s *Store) CreateEvent(ctx context.Context, factID int64, action, object st
|
||||
return id, nil
|
||||
}
|
||||
|
||||
// EventPair identifies one action+object grouping in the events table — the
|
||||
// unit the pattern detector reasons about.
|
||||
type EventPair struct {
|
||||
Action string
|
||||
Object string
|
||||
}
|
||||
|
||||
// DistinctEventPairs returns every distinct action+object pair that has at
|
||||
// least one event, in no particular order. This is what lets the proactive
|
||||
// digestion tick run the pattern detector over everything accumulated so far
|
||||
// instead of only the pair touched by the utterance that just landed
|
||||
// (Vikunja #43) — the tick has no "current utterance," so it has to ask the
|
||||
// store what to look at.
|
||||
func (s *Store) DistinctEventPairs(ctx context.Context) ([]EventPair, error) {
|
||||
rows, err := s.db.QueryContext(ctx, `SELECT DISTINCT action, object FROM events`)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("distinct event pairs: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
var out []EventPair
|
||||
for rows.Next() {
|
||||
var p EventPair
|
||||
if err := rows.Scan(&p.Action, &p.Object); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, p)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// EventsFor returns all events matching action+object, ordered by ts ascending
|
||||
// (oldest first — the order the pattern detector needs for interval computation).
|
||||
func (s *Store) EventsFor(ctx context.Context, action, object string) ([]Event, error) {
|
||||
|
||||
@@ -112,6 +112,25 @@ ALTER TABLE reminders ADD COLUMN next_fire_ts INTEGER;`, // #2
|
||||
DROP TABLE delivery_attempts;
|
||||
ALTER TABLE delivery_attempts_v12 RENAME TO delivery_attempts;
|
||||
CREATE INDEX IF NOT EXISTS idx_delivery_attempts_status ON delivery_attempts (status);`,
|
||||
|
||||
// #13 — durable digest outbox (Vikunja #281). A care nudge the restraint
|
||||
// gate suppresses (quiet hours / away / calendar-busy) is not necessarily
|
||||
// lost: if it's worth resurfacing, it lands here instead, and gets spoken
|
||||
// as one bundle at the next moment speaking is appropriate. body_hash
|
||||
// dedupes repeat suppressions of the "same" nudge; expires_ts bounds how
|
||||
// stale an entry may get before it's worthless and must be dropped rather
|
||||
// than delivered late.
|
||||
`CREATE TABLE IF NOT EXISTS digest_entries (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
rule TEXT NOT NULL,
|
||||
severity INTEGER NOT NULL,
|
||||
body TEXT NOT NULL,
|
||||
body_hash TEXT NOT NULL,
|
||||
status TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending','drained','expired')),
|
||||
created_ts INTEGER NOT NULL,
|
||||
expires_ts INTEGER NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_digest_entries_status ON digest_entries (status);`,
|
||||
}
|
||||
|
||||
// migrate applies every migration with a number greater than the DB's current
|
||||
|
||||
@@ -51,9 +51,10 @@ var (
|
||||
// keep finding the pattern, and every re-propose is refused here. Maven is not
|
||||
// a nag.
|
||||
//
|
||||
// TODO(vikunja#46): the detector currently only writes here from the voice
|
||||
// path. Once digestion runs the detector on its own tick, that tick should
|
||||
// call this too, so a pattern gets noticed even with nobody at the mic.
|
||||
// Vikunja #43: this is called both from the voice fact-write path (for the
|
||||
// immediate spoken confirmation) and from the digestion tick's proactive
|
||||
// scan (cmd/mavend/tick.go's detectPatterns, via patterns.go's
|
||||
// detectAndPropose), so a pattern gets noticed even with nobody at the mic.
|
||||
func (s *Store) CreateProposedRoutine(ctx context.Context, action, object string, intervalDays float64, ts time.Time) (int64, error) {
|
||||
res, err := s.db.ExecContext(ctx,
|
||||
`INSERT INTO proposed_routines (action, object, interval_days, status, created_ts)
|
||||
|
||||
Reference in New Issue
Block a user