Compare commits

...

32 Commits

Author SHA1 Message Date
kami d29e7ba813 Gate POST /api/chat on the same step-up as /tools (#317)
/api/chat reaches the router, the LLM and, through applyAction, the whole
act path, so it is the widest state-changing surface mavweb serves. It was
the only one with no gate. It now goes through stepUpOK like POST /tools,
POST /routines and POST /api/revert: unchanged in the default deploy
(WebAuthn unconfigured, fail-open behind wg+nginx), 403 under
-require-stepup or an unasserted passkey session.

The route table now carries an explicit enumeration of every state-changing
route and its gate, and the two startup SECURITY log lines name /routines
and /api/chat alongside /tools and /api/revert.

The loopback -addr default the task also asked for landed earlier in
d12de58; the compose already publishes mavweb on 127.0.0.1 only.
2026-08-01 01:24:42 +04:00
kami f7e1187823 Match quiet-mode toggles on whole words, and resolve OFF first 2026-08-01 01:02:10 +04:00
kami ed48c59ba7 Merge branch 'refactor/query-sources' into integration/small-batch 2026-08-01 00:54:11 +04:00
kami b09967f9e6 Split actions.go into per-intent files
Pure move: actionFact, actionReminder, actionAct and actionNote each get
their own actions_<intent>.go. The two small ones (chat, system) and the
actionHandlers table stay in actions.go, which is now just the dispatch
layer and the notes about what does not belong in it. No behaviour
change — only the file a handler is read in.
2026-08-01 00:53:37 +04:00
kami b4a3867479 Turn actionQuery into a chain of query sources
The six answer sources were hand-unrolled inside one 127-line function.
The intent table is a closed set of 7, but this list is open-ended —
Kiwix (#286), RSS (#258), the crawler (#259) and email (#246) each add
one. Each is now a registry entry: a name plus a method on the handler,
walked in order until one claims the question.

Order is unchanged and still load-bearing (memory before the notes-only
pass, #373), the confidence gate keeps its position and semantics, and
every reply string, log line and best-effort failure is verbatim.
2026-08-01 00:51:44 +04:00
kami 88d07b5175 Unify the voice and text turn pipelines into runTurn
HandlePushToTalk and handleText hand-wrote the same eight-step turn
sequence twice, comments in the latter saying "same as HandlePushToTalk"
four times. Extract it into runTurn(ctx, text) string: the voice path
wraps it in stt/tts, the text path returns it directly.

The two had drifted. The text path was missing the quiet-hours toggle
check entirely, so "тихий режим" over IPC/telegram fell through to the
classifier; unifying gives it the check. It also logged the route result
and applyAction return where the voice path did not — both logs are kept
for both paths.
2026-08-01 00:48:50 +04:00
kami c00e3003bf Merge branch 'refactor/praxis-capability-registry' into integration/small-batch 2026-08-01 00:43:54 +04:00
kami c0f9834528 Turn the Praxis act dispatch into a capability registry 2026-08-01 00:43:16 +04:00
kami ad5eb2d1cf Walk a chain of confirm resolvers instead of three copied blocks 2026-08-01 00:42:23 +04:00
kami 5253123d99 Merge branch 'refactor/voice-wiring' into integration/small-batch
# Conflicts:
#	cmd/mavend/voice.go
2026-07-31 23:54:11 +04:00
kami 2abf98dea6 Move the voice daemon wiring and startup out of voice.go 2026-07-31 23:52:51 +04:00
kami f5c71b87f3 Move the confirm/park gate out of voice.go 2026-07-31 23:51:33 +04:00
kami 5934110fa8 Move the Praxis/Hexis act handling out of voice.go
voice.go is still the biggest file in cmd/mavend and most of what is left
has nothing to do with the audio path. The ecosystem integration is one
such lump: it talks to Nexus, Praxis and Hexis over HTTP and only touches
the handler for its store and clock. Lifting it into ecosystem_acts.go
puts it next to ecosystem.go, where the clients it drives already live.

Move-only: handlePraxisAct, recordPraxisTrace (called from nowhere else),
handleHexisAct and execHexis verbatim, plus the two imports that became
unused in voice.go.
2026-07-31 23:43:47 +04:00
kami 6e47a3d736 Merge branch 'worktree-agent-a88193d1d84b04a5b' into integration/small-batch 2026-07-31 23:38:00 +04:00
kami 67a5eb3805 Split applyAction's 300-line switch into a per-intent handler table
applyAction (cmd/mavend/voice.go) dispatched all 7 intents from one giant
switch. Extract each case body verbatim into its own actionXxx method in
new cmd/mavend/actions.go, dispatched from an actionHandlers table keyed by
router.Intent. applyAction itself is now just the dec.Clarify guard plus a
table lookup.

No behaviour change: same reply strings, same side-effect order, comments
moved verbatim. The destructive-act confirm gate and the enabled-tool
allowlist stay entirely inside actionAct, exactly where they lived in the
old switch's IntentAct case — they're act-specific, not cross-cutting, so
they don't move to a separate layer. dec.Clarify short-circuit, dialogue
bookkeeping and detectPattern stay outside the table since they run
regardless of intent.

voice.go: 1638 -> 1344 lines. New actions.go: 362 lines.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 23:37:33 +04:00
kami 029449eefa Table-drive the IPC dispatcher instead of a 42-arm switch
dispatch() replaces the hand-written switch with a package-level
map[Method]handlerFunc built once at init. Each entry is one
withParams/withParamsVoid/withoutParams call closing only over the
CoreAPI method it invokes — adding a method is now one table line
instead of a new arm.

Check still runs once at the top before any unmarshal, unchanged. The
three non-CoreAPI methods (assert_stepup, store_encryption_key, unlock)
are special-cased before the table lookup since they drive Server
fields (StepUp/WrapKeyFn/UnlockFn), not store state. The current
CoreAPI is loaded once per dispatch and passed into the handler as an
argument, so SetAPI's runtime swap (the unlock transition) still takes
effect on the next request — the table itself never captures an api
value. No wire-format change; existing round-trip and unknown-method
tests in ipc_test.go pass unmodified.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 23:36:46 +04:00
kami ac36216f5d Merge branch 'worktree-agent-a517419e93e6a219f' into integration/small-batch 2026-07-31 23:31:24 +04:00
kami db17cfcc65 Delete the dead lockedAPI, add UnimplementedCoreAPI for the doubles 2026-07-31 23:31:24 +04:00
kami 7d676eb941 Stop tracking the mavwaked build artifact 2026-07-31 23:30:52 +04:00
kami fe3a4e9514 Merge branch 'worktree-agent-a9e5cef90b263a5e5' into integration/small-batch 2026-07-31 23:27:36 +04:00
kami a906f2afad Extract the pure RU/string/weather helpers out of voice.go 2026-07-31 23:27:08 +04:00
kami a2031a31d1 Record the measured confidence-gate numbers 2026-07-31 23:23:37 +04:00
kami 9b8bdf73cc Merge the five small-task branches 2026-07-31 23:11:26 +04:00
kami 7ad3c9a408 Merge branch 'worktree-agent-af88d63f65f30896b' into integration/small-batch 2026-07-31 23:11:16 +04:00
kami 84e1478823 Merge branch 'worktree-agent-af0fd9507d3e2ee46' into integration/small-batch 2026-07-31 23:11:16 +04:00
kami 9c8d0baffe Merge branch 'worktree-agent-a4cef2a815e32ebbf' into integration/small-batch 2026-07-31 23:11:16 +04:00
kami b0f5a16ec9 Add digest as a real outcome: suppressed care nudges get resurfaced, not lost
Vikunja #281. The interruption policy promised four outcomes — deliver_now,
queue, digest, drop — but only three existed: a care candidate the restraint
gate suppressed for quiet hours / away / calendar-busy simply vanished in
loop.Tick's `continue`, with only the trace remembering why.

internal/morning turned out not to be the natural drain: it's a fixed
Item/FactKey checklist engine, not a generic message bundler, so gate-
suppressed nudge text has nowhere to plug into its evidence model. Built a
parallel (but small, reusing the outbox's shape) durable digest instead:

- internal/store: digest_entries table + EnqueueDigestEntry (dedupes by
  rule+body, mirroring the delivery outbox's bodyHash), PendingDigestEntries,
  ExpireStaleDigestEntries, DrainDigestEntries (mark, never delete — an
  audit trail of what she actually said).
- internal/loop: DigestEligible(severity, blockedBy) is the pure boundary —
  only genuine restraint blocks (quiet_hours/calendar_busy/presence) even
  qualify (cooldown/snooze are not "suppression"); within care, Sev2 (break)
  digests, Sev1 (water/meal — stale by the time anyone could resurface them)
  drops. High severity never digests; alarms bypass the gate and deliver
  unchanged, on purpose.
- cmd/mavend/tick.go: each tick scans ExplainTick's trace for eligible
  blocked candidates, enqueues them, sweeps stale entries (24h expiry — the
  care rules are daily-cadence, so anything older is describing a day
  that's over), and drains the bundle only once the suppression reason has
  actually cleared, capped at 3 spoken items plus a trailing count so a
  digest can't turn into the exact nagging it was built to avoid.

Tests: store-level round-trip/restart-survival/dedupe/expiry/drain, loop-
level severity-boundary unit tests, and tick-level integration tests for
the drain-only-when-clear and never-digest-high-severity behavior.
2026-07-31 23:09:22 +04:00
kami 67563ed1f6 Run pattern detection from the digestion tick, not just voice (#43)
detectPattern only ever fired as a side effect of a voice fact-write, so a
recurring pattern already sitting in history went unnoticed until he
happened to mention it again by voice — the opposite of proactive.

Split the pipeline: extraction (fact -> normalized event) stays where a fact
is written, in voice.go, since it's tied to that write regardless of who's
talking. Detection (events -> stable pattern -> proposed_routines row) moves
into shared code (patterns.go's detectAndPropose) that both the voice path
and the new tick.go:detectPatterns call. The tick runs it every cycle over
every action+object pair on record (store.DistinctEventPairs, added), so a
pattern gets noticed on the daemon's own schedule.

Idempotence and the dismiss-must-stick requirement turned out to already be
handled by the store, not something the tick needs to reinvent:
proposed_routines has UNIQUE(action, object) and CreateProposedRoutine does
ON CONFLICT DO NOTHING, and DismissProposedRoutine flips status in place
without deleting the row. So a pair already proposed, accepted, OR
dismissed is a silent no-op on every later tick — a dismissed pattern can
never resurface, and re-running the scan never spams the /routines page.
Kept the voice-path call (immediate spoken confirmation is a nice feature
UX-wise and is now redundant-but-harmless with the tick, since both paths
share the same guarded detectAndPropose).

Tick-side detection only ever writes a row; it does not notify, ring, or
speak, keeping Maven "not a nag, not autonomous" — the /routines page is
still the only place a proposal becomes visible, and only accepting it
starts producing nudges (fireAcceptedRoutines).

Also fixed the stale vikunja#46 reference in proposed_routines.go — the
TODO it named is what this commit does.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 23:07:33 +04:00
kami f0f7ebc9b2 Give LLM-routed decisions a real confidence so clarify can fire (#359)
Confidence was hardcoded to 1.0 for every LLM decision, and the LLM branch
in Router.Route returned straight from fillSlots without ever touching the
stage-3 threshold gate — so the LLM path could not produce a Clarify no
matter what confidence a model reported. That is why all 6 want_clarify
cases in the 77-case RU fixture were missed by every model in the bake-off.

Fix reads structural signal instead of changing the (parity-locked) router
prompt: a single-token utterance ("вода", "бэкап") is flagged thin evidence
in llmrouter.go; a fact left keyless or an act that never resolves to an
allowlisted fn, checked after fillSlots so the deterministic parsers get
first crack, is flagged in router.go's new gateLLMDecision. Anything below
config.DefaultRouterThreshold (0.55) now sets Clarify=true through the same
path the classifier already uses.

Added unit tests with a stubbed Completer proving both directions: thin
cases clarify, clean multi-word/resolved-slot cases stay confident. The
77-case fixture re-run against a live llama-server is still needed to
confirm the 6/6 moves — not done here, no llama-server on this box.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 23:07:32 +04:00
kami d12de589a2 mavweb: default -addr to loopback, not all interfaces
PR #47 added two state-changing routes (POST /chat, POST /routines)
behind the -addr flag, which defaulted to ":9200" (all interfaces).
Default now binds 127.0.0.1:9200; anyone who wants LAN/wider exposure
still passes an explicit bind (as deploy/docker-compose.yml already
does with "-addr :9201" inside the container, unaffected by this
default change).

Vikunja #317.
2026-07-31 23:03:45 +04:00
kami 50cc17f33a Lock down deploy/ecosystem/nginx.conf template to match the live host
The template said "drop into your nginx sites" but listened on the
wildcard `listen 80;` with no allow/deny ACL, unlike the actual deployed
hexis.kvmx.ru config which binds only to the WireGuard (10.42.0.1) and
LAN (192.168.1.104) addresses with allow/deny all. Anyone following the
template as written would expose these unauthenticated admin UIs to the
open internet.

Bind explicitly to those two addresses and add the matching ACL block,
mirroring cmd/mavweb/nginx.conf which already does this correctly.
Added a comment naming both addresses as host-specific so a deploy on a
different box swaps the IPs instead of reverting to `listen 80` when the
bind fails.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 23:01:58 +04:00
kami 73d13f1ea6 Merge pull request 'Stop the docs claiming the LLM router is off' (#49) from docs/fix-drift into master 2026-07-31 20:46:57 +02:00
40 changed files with 3777 additions and 2156 deletions
+1
View File
@@ -6,6 +6,7 @@
/mavweb
/mavpoll
/mavcaldav
/mavwaked
# Certs (private keys, don't commit)
certs/
+13 -2
View File
@@ -90,8 +90,19 @@ fallback. Any LLM error falls through to the classifier so a turn never breaks o
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. Still open: `Confidence: 1.0` is hardcoded in `llmrouter.go`, so the
LLM path never asks for clarification (6/6 refusal cases missed) — Vikunja #359.
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
+1 -1
View File
@@ -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
+71
View File
@@ -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)
}
+66
View File
@@ -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 "готово."
}
+64
View File
@@ -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
}
+41
View File
@@ -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
}
+226
View File
@@ -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
}
+33
View File
@@ -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 ""
}
+216
View File
@@ -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, " ")
}
+158
View File
@@ -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)
}
}
+312
View File
@@ -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 + "."
}
+19 -94
View File
@@ -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)")
@@ -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 {
@@ -516,7 +441,7 @@ func run(args []string) error {
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines))
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,
+74
View File
@@ -0,0 +1,74 @@
// mavend/patterns.go — the shared detect+propose step of pattern inference
// (Vikunja #43). Event *extraction* (fact -> action/object) happens at fact-
// write time in voice.go's detectPattern, 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"
"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
}
+122
View File
@@ -0,0 +1,122 @@
package main
import (
"context"
"database/sql"
"testing"
"time"
"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). Three 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, 3)
tl := newTestTickLoop(t, st, &fakeSink{}, nil)
tl.detectPatterns(ctx, now)
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, 3)
tl := newTestTickLoop(t, st, &fakeSink{}, nil)
tl.detectPatterns(ctx, now)
tl.detectPatterns(ctx, now.Add(time.Hour))
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, 3)
tl := newTestTickLoop(t, st, &fakeSink{}, nil)
tl.detectPatterns(ctx, now)
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), 3)
tl.detectPatterns(ctx, now.Add(60*24*time.Hour))
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)
}
}
+114
View File
@@ -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)
}
})
}
}
+180
View File
@@ -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")
}
}
+85
View File
@@ -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)
}
+189
View File
@@ -180,6 +180,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 +206,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)
// 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 +361,176 @@ 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.
//
// This only ever creates a row for the /routines page to show. It does not
// notify, ring, or speak — Maven is "not a nag, not autonomous" (CLAUDE.md),
// and detection is not the same act as disturbing him about it. A proposal
// only starts producing nudges once he accepts it (fireAcceptedRoutines).
func (t *tickLoop) detectPatterns(ctx context.Context, now time.Time) {
pairs, err := t.store.DistinctEventPairs(ctx)
if err != nil {
log.Printf("tick: distinct event pairs: %v", err)
return
}
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)
}
}
// 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.
+158 -1562
View File
File diff suppressed because it is too large Load Diff
+434
View File
@@ -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)
}
}
+50
View File
@@ -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"
}
+77 -3
View File
@@ -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
View File
@@ -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)
+28 -3
View File
@@ -4,10 +4,23 @@
#
# 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.
server {
listen 80;
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 +31,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 +49,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;
+5 -84
View File
@@ -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
+5 -88
View File
@@ -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
View File
@@ -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 {
+118
View File
@@ -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
}
+48
View File
@@ -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
+30
View File
@@ -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 (Sev12), 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.
+42 -1
View File
@@ -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
+97
View File
@@ -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)
}
}
+30
View File
@@ -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
+142
View File
@@ -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
}
+202
View File
@@ -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)
}
}
}
+30
View File
@@ -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) {
+19
View File
@@ -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
+4 -3
View File
@@ -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)
BIN
View File
Binary file not shown.