Compare commits

...

12 Commits

Author SHA1 Message Date
kami 2e9b9ec1cf Warn separately when a row has no text to re-embed 2026-07-31 13:52:56 +04:00
kami d1f6f6355f Merge the vector backfill 2026-07-31 13:51:25 +04:00
kami 92ecb691de Re-embed stored notes and facts after an embedder swap (#378)
The embedder swap left every stored vector in the old model's space, so cosine against a new query vector is noise. Add the one-shot backfill: store.ReembedAll re-embeds every note and fact text with the currently configured embedder (the passage side, which is the side stored text was written with) and rewrites both places a vector lives — the notes table embedding column and the memory_vectors rows.

All of it plus the embedder marker happens in one transaction, so a failure partway changes nothing and writes no marker: re-run it. A run against a DB whose marker already names the current embedder does nothing.

Triggered explicitly with `mavend -reembed`, not automatically on mismatch: ONNX on the laptop CPU makes this minutes of work, and a silent multi-minute stall on boot would look like a hang. The mismatch warning now tells the user to run it.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 13:50:43 +04:00
kami 4282f6b9a9 Warn about vectors written before the marker existed 2026-07-31 13:44:29 +04:00
kami 7bb9f9be06 Merge the embedder marker 2026-07-31 13:42:40 +04:00
kami 1e47eaca5a Record which embedder wrote the stored vectors and warn on a swap (#378)
The embedder moved from paraphrase-multilingual-MiniLM-L12-v2 to
multilingual-e5-small. Both are 384-dimensional, so nothing in the code
noticed: cosine between an old stored vector and a new query vector is
noise, and recall degrades silently.

So the DB now records the embedder that wrote its vectors. One value for
the whole DB (migration #11, a small `meta` key/value table) rather than a
column on every vector row: the backfill re-embeds every note and fact in
one pass, so a per-row marker would hold the same string in every row and
cost a column on two tables for nothing.

The identity comes from the embedder itself via a new optional ID() method
("multilingual-e5-small@384", model file name plus dimension), so pointing
the config at another model changes the string without anyone editing a
constant. mavend logs a loud WARNING at startup naming both the stored and
the configured embedder when they differ.

Detection only — recall behaviour is unchanged. TODO(#378) in
store.CheckEmbedder marks where the backfill will hook in.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 13:42:07 +04:00
kami 892330eb84 Merge the note recall fix 2026-07-31 13:33:14 +04:00
kami 9a3bcd7c46 Merge the thinking-off measurement 2026-07-31 13:31:31 +04:00
kami 98ee701e03 Let a note win a recall, not only a fact (#373)
The memory pass ran only after the notes-only gate had already rejected
the same note at the same score. Notes and facts share one vector index,
so a note that failed there failed again — the branch could only ever
return a fact.

Now the memory pass runs first: one search over everything Maven
remembers, one gate, and the memory that clearly matches best answers
(a note gets phrased, a fact is read back). The notes-only pass stays
behind it for notes the vector index does not hold. No threshold moved,
so the set of questions answered is unchanged — only which memory
answers them.

Fixture gained two mixed note+fact cases, so the answerable count goes
25 -> 27: hash recall@1 36.0% -> 37.0% (ratchet 0.32 unchanged, comment
updated), e5 recall@1 72.0% -> 70.4%, false recall still 1/5.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 13:30:38 +04:00
kami 04c1088088 Measure thinking off on routing properly — it does not win (#376)
The 67.1% "thinking off" column in ROUTING-EVAL-31-07-2026.md was an
artefact. It came from a hand-rolled HTTP client in the eval test that
did not send repeat_penalty, so it differed from the reference run on two
axes and the penalty was the one that mattered.

Re-scored back to back on an idle box with everything else held equal:
thinking off is identical to thinking on, case for case, same confusion
matrix, same three unparseable replies. A direct probe of the running
llama-server shows enable_thinking, thinking and reasoning_budget are all
ignored for this model on this build, so there was nothing to turn off.

No defaults changed. The misleading third configuration is removed from
internal/router/eval so its table cannot be quoted again.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 13:30:22 +04:00
kami 07c191d8b8 Merge dialogue session persistence 2026-07-31 13:21:03 +04:00
kami c668310b3e Persist the dialogue session so a restart keeps the conversation
Vikunja #363. The follow-up session was a plain in-memory map, so any
mavend restart dropped the thread. It now mirrors to a small TTL-pruned
sqlite table and is loaded on startup; expired sessions are deleted on
load, not revived. Clarify's pending question is untouched.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 13:20:30 +04:00
22 changed files with 1299 additions and 134 deletions
+8
View File
@@ -75,6 +75,14 @@ same vector. A note is indexed in both places with the same embedding, so if it
`QueryNotes` it fails again here — the branch can only ever return a **fact**. Its comment calls it
"additive"; for notes it is not.
**Fixed (Vikunja #373).** The memory pass now runs *first*, as one search over notes and facts with
one gate, so whichever memory is clearly the best match answers — note or fact. The notes-only pass
stays behind it for notes the vector index does not hold. No threshold changed, so the set of
questions Maven answers is the same; only which memory answers them. The fixture gained two mixed
note+fact cases (`ru-mixed-031`, `ru-mixed-032`), which is why the counts below are out of 27
answerable cases and not 25: hash recall@1 36.0% (9/25) → 37.0% (10/27), e5 recall@1 72.0% (18/25) →
70.4% (19/27) with answered-after-gate 68.0% → 66.7% and false recall unchanged at 1/5.
### 5. Ranking has no recency or type signal, and the store is not the bottleneck
`internal/store/notes.go:67` sorts by cosine and uses `ts` only to break an exact float tie, which
+59 -10
View File
@@ -66,15 +66,64 @@ Three things this run settles:
`запиши что…` phrasings toward fact, and that suspicion stands — all five `ru-note-*`
cases now land on fact. Tracked as Vikunja #375.
**Thinking off is the best configuration measured so far**, on both accuracy and latency
(Vikunja #376). That is worth understanding before flipping: routing is a short
classification into a fixed enum with grammar-constrained output, so there is little to
reason about, and the thinking trace mostly gives a small model room to talk itself out of
the right answer. Phrasing is a different job and needs measuring separately.
The `thinking off` column above read as the best configuration measured so far (Vikunja #376).
**It was wrong** — see the controlled re-run below. Ignore that column.
Still `6 / 6` missed clarify — the router has no way to say "I don't know" (Vikunja #359).
That is unchanged by anything here.
## Thinking off — 31-07-2026, controlled re-run (Vikunja #376)
The "thinking off wins by 6 points" observation above **does not hold**. It was a measurement
artefact, and the earlier table's `thinking off` column should be ignored.
The thinking-off variant was scored by a hand-rolled HTTP client living in the test file
instead of `llm.Client`. That copy did not send `repeat_penalty`, which the real router does
send (`routeRepeatPenalty = 1.15`). So the two columns differed on two axes at once, and the
one that mattered was the penalty, not the thinking mode.
Re-measured with everything else held equal — same fixture, same prompt, same grammar, same
sampling, same idle box, the three configurations run back to back and never concurrently:
| | llm-only, thinking on | llm-only, thinking off | cascade+llm |
|---|---|---|---|
| intent-only accuracy | 59.2% (45/76) | 59.2% (45/76) | 61.8% (47/76) |
| full accuracy (intent+slots+gate) | 38.2% (29/76) | 38.2% (29/76) | 57.9% (44/76) |
| route errors | 3 | 3 | 0 |
| grammar violations | 3 (all 3 route errors) | 3 (same 3 cases) | 0 |
| missed clarify | 5 / 6 | 5 / 6 | 5 / 6 |
| p50 latency | 836ms | 920ms | 810ms |
| p95 latency | 1.41s | 2.00s | 1.31s |
Thinking off is not just a tie on the headline numbers — it is identical case for case, with
the same confusion matrix and the same three unparseable replies. The latency difference is
run-to-run noise on one box, and it points the wrong way here.
The reason is simpler than any accuracy argument: **this llama-server build ignores the
request-level thinking switch for this model.** Probed directly against the running server
with `chat_template_kwargs.enable_thinking = false`, `chat_template_kwargs.thinking = false`
and top-level `reasoning_budget = 0` — all three return a byte-identical answer with the
thinking trace still in `reasoning_content`, and the server reports the prompt prefix as
cached, meaning the rendered template did not change. There was never anything being turned
off, which is also why the numbers match exactly.
Nothing was defaulted. `internal/llm` still has no `chat_template_kwargs` field, `VoiceConfig`
has no thinking flag, and `deploy/mavend.json` is unchanged. The misleading third
configuration is removed from `internal/router/eval` so the table it produced cannot be quoted
again.
Two caveats worth saying out loud:
- **The fixture is 76 cases.** A 6-point difference on 76 cases is roughly 4-5 cases and would
not have been worth trusting even if it had reproduced. This one was exactly 0 cases, which
is a much easier call.
- **This is one server build and one checkpoint** (`b9351`, Qwen3.5-0.8B Q4_K_M). If the
#122 checkpoint or a newer llama.cpp does honour the switch, the question reopens — but it
reopens as an unmeasured question, not as a 6-point win.
Phrasing was **not** measured. Whether thinking helps there is still open, and now also blocked
on the same "can we even turn it off" question.
## Findings
### 1. The resident model does route better — 50.0% vs 36.8%
@@ -145,11 +194,11 @@ Note the grammar's `string ::= "\"" ([^"\\] | "\\" .)* "\""` is unbounded, so no
### 7. Two hypotheses tested and closed
- **Thinking mode is a non-issue.** Qwen3.5's template defaults `thinking = 1`, so
grammar-constrained JSON lands in `reasoning_content` with `content` empty —
`llm.Client`'s fallback handles it. A `thinking off` run scored *identically* (18/76,
48.7%, same p50). `internal/llm` deliberately does **not** grow a `chat_template_kwargs`
field.
- **Thinking mode is a non-issue.** Confirmed twice now, the second time properly — see the
controlled re-run section. Grammar-constrained JSON lands in `reasoning_content` with
`content` empty and `llm.Client`'s fallback handles it; the request-level switch does
nothing on this build. `internal/llm` deliberately does **not** grow a
`chat_template_kwargs` field.
- **Runaway array repetition does not reproduce.** An isolated smoke test with a stripped
grammar emitted `{"intent":"reminder"}` until `MaxTokens`; under the real `routeSystem`
prompt the few-shot examples anchor it to one object. 2 errors in 76, not 76.
+2
View File
@@ -189,7 +189,9 @@ func (l *lockedAPI) MorningStatus(ctx context.Context) ([]ipc.MorningRoutineStat
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)")
reembed := flag.Bool("reembed", false, "re-embed every stored note and fact with the configured embedder, then serve normally (run once after an embedder swap)")
flag.CommandLine.Parse(args)
reembedOnStart = *reembed
cfg, err := config.Load(*cfgPath)
if err != nil {
return err
+158
View File
@@ -0,0 +1,158 @@
package main
import (
"context"
"fmt"
"math"
"testing"
"time"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/memory"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/voice"
)
// fixedEmbedder hands back a vector chosen per text, so a test can say exactly
// how close each stored memory is to the question. The real embedders make
// scores that are realistic but not controllable, and this test is about the
// gate, not about the embedder.
type fixedEmbedder struct{ vecs map[string][]float32 }
func (f *fixedEmbedder) Dim() int { return 4 }
func (f *fixedEmbedder) Close() error { return nil }
func (f *fixedEmbedder) Embed(_ context.Context, text string) ([]float32, error) {
v, ok := f.vecs[text]
if !ok {
return nil, fmt.Errorf("fixedEmbedder: no vector for %q", text)
}
return v, nil
}
// scoreVec builds a unit vector whose cosine against the query vector
// (1,0,0,0) is exactly score.
func scoreVec(score float64) []float32 {
rest := math.Sqrt(1 - score*score)
return []float32{float32(score), float32(rest), 0, 0}
}
// recordingPhraser remembers what the query path handed it to phrase, which is
// how the test can tell which pass produced the answer.
type recordingPhraser struct {
*phraser.Stub
notes []string
}
func (r *recordingPhraser) PhraseQuery(ctx context.Context, utterance string, notes []string) (string, error) {
r.notes = notes
return r.Stub.PhraseQuery(ctx, utterance, notes)
}
// recallCase — one stored memory: its text, how close it is to the question,
// whether it is a note or a fact, and whether the notes table holds it too.
type recallCase struct {
text string
score float64
kind string
}
// buildRecallHandler stores the given memories and returns a handler whose
// query path can be run directly. Notes go into BOTH the notes table and the
// vector index, which is what the daemon does (voice.go's IntentNote).
func buildRecallHandler(t *testing.T, question string, mems []recallCase) (*reactiveHandler, *recordingPhraser) {
t.Helper()
ctx := context.Background()
st := newTestStore(t)
emb := &fixedEmbedder{vecs: map[string][]float32{question: {1, 0, 0, 0}}}
mem := memory.NewInMemoryStore()
now := time.Now()
for i, m := range mems {
vec := scoreVec(m.score)
emb.vecs[m.text] = vec
id := fmt.Sprintf("%s:%d", m.kind, i)
if m.kind == "note" {
if _, err := st.WriteNote(ctx, now, m.text, vec, "tap:voice"); err != nil {
t.Fatalf("WriteNote: %v", err)
}
}
if err := mem.Insert(ctx, id, vec, map[string]string{"text": m.text, "type": m.kind}); err != nil {
t.Fatalf("memory insert: %v", err)
}
}
phr := &recordingPhraser{Stub: phraser.NewStub()}
h := &reactiveHandler{
api: ipc.NewStoreAPI(st),
embedder: emb,
replier: voice.NewStubReplier(),
phraser: phr,
now: func() time.Time { return now },
memStore: mem,
dataStore: st,
queryMinScore: 0.55,
queryMinMargin: 0.008,
weatherProvider: nil,
}
return h, phr
}
func askQuery(t *testing.T, h *reactiveHandler, question string) string {
t.Helper()
return h.applyAction(context.Background(), router.Decision{
Intent: router.IntentQuery,
Utterance: question,
})
}
// TestQueryRecallNoteCanWin — the note-recall regression (Vikunja #373). Notes
// and facts share one vector index, and a note that clearly beats everything
// else must be the answer. Before the fix the memory pass only ran after the
// notes-only gate had already rejected the same note at the same score, so only
// a fact could ever come back from it.
func TestQueryRecallNoteCanWin(t *testing.T) {
const q = "где молоко"
t.Run("a clearly best note answers", func(t *testing.T) {
h, phr := buildRecallHandler(t, q, []recallCase{
{text: "молоко стоит в холодильнике", score: 0.90, kind: "note"},
{text: "выучил пару аккордов", score: 0.50, kind: "note"},
})
reply := askQuery(t, h, q)
if want := "вот что я нашла: молоко стоит в холодильнике"; reply != want {
t.Errorf("reply %q, want %q", reply, want)
}
// One text, the winning memory's — the answer came from the memory
// pass, not from handing the phraser every note in the table.
if len(phr.notes) != 1 || phr.notes[0] != "молоко стоит в холодильнике" {
t.Errorf("phraser got %q, want just the recalled note", phr.notes)
}
})
// The other half of "one gate over everything": a fact that matches better
// than the best note now answers, instead of losing to a note that only had
// to beat other notes.
t.Run("the better-matching fact answers", func(t *testing.T) {
h, _ := buildRecallHandler(t, q, []recallCase{
{text: "молоко стоит в холодильнике", score: 0.80, kind: "note"},
{text: "купил молоко в среду", score: 0.95, kind: "fact"},
})
if reply := askQuery(t, h, q); reply != "купил молоко в среду" {
t.Errorf("reply %q, want the fact read back", reply)
}
})
// The gate is untouched: two memories this close mean the embedder cannot
// tell them apart, and silence still beats a coin flip.
t.Run("no clear best stays silent", func(t *testing.T) {
h, _ := buildRecallHandler(t, q, []recallCase{
{text: "молоко стоит в холодильнике", score: 0.860, kind: "note"},
{text: "молоко закончилось", score: 0.858, kind: "note"},
})
if reply := askQuery(t, h, q); reply != "не знаю." {
t.Errorf("reply %q, want silence", reply)
}
})
}
+17 -14
View File
@@ -2,21 +2,24 @@ package main
import "github.com/kami/maven/internal/memory"
// bestRecall is the read side of the long-term memory store: the top hit's
// stored text when it clears the confidence gate. This recalls across BOTH
// notes and facts (facts aren't in the notes table, so this is the only path
// that can answer "when did I last …?" from a captured fact). A note hit here
// is redundant with the notes-RAG path — by design; the two indexes can diverge
// once the backend is swapped for a persistent/external store. ok=false when
// the hit fails the confidence gate (see memory.Confident: an absolute floor
// plus a margin over the runner-up) or carries no text.
func bestRecall(results []memory.Result, minScore, minMargin float64) (string, bool) {
// bestRecall is the read side of the long-term memory store: the top hit when
// it clears the confidence gate. The index holds BOTH notes and facts, and
// either can win — the caller looks at the returned hit's meta["type"] to see
// which. Facts aren't in the notes table, so this is the only path that can
// answer "when did I last …?" from a captured fact.
//
// The whole hit is returned, not just its text, because "which memory answered"
// decides how the answer is said: a note gets phrased in Maven's voice, a fact
// is read back as stored.
//
// ok=false when the hit fails the confidence gate (see memory.Confident: an
// absolute floor plus a margin over the runner-up) or carries no text.
func bestRecall(results []memory.Result, minScore, minMargin float64) (memory.Result, bool) {
if !memory.Confident(results, minScore, minMargin) {
return "", false
return memory.Result{}, false
}
text := results[0].Meta["text"]
if text == "" {
return "", false
if results[0].Meta["text"] == "" {
return memory.Result{}, false
}
return text, true
return results[0], true
}
+21 -2
View File
@@ -39,8 +39,27 @@ func TestBestRecall(t *testing.T) {
if !ok {
t.Fatal("clearing hit not returned")
}
if got != "выпил воды в три часа" {
t.Errorf("wrong text: %q", got)
if got.Meta["text"] != "выпил воды в три часа" {
t.Errorf("wrong text: %q", got.Meta["text"])
}
if got.Meta["type"] != "fact" {
t.Errorf("kind lost: %q", got.Meta["type"])
}
})
// The index holds notes and facts together, so a note has to be able to win
// it — for a long time it could not (Vikunja #373).
t.Run("a note can win", func(t *testing.T) {
res := []memory.Result{
{Score: 0.86, Meta: map[string]string{"text": "молоко в холодильнике", "type": "note"}},
{Score: 0.61, Meta: map[string]string{"text": "выпил воды", "type": "fact"}},
}
got, ok := bestRecall(res, min, margin)
if !ok {
t.Fatal("clearly-best note not returned")
}
if got.Meta["type"] != "note" || got.Meta["text"] != "молоко в холодильнике" {
t.Errorf("got %v, want the note", got.Meta)
}
})
+113 -11
View File
@@ -171,6 +171,7 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
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
@@ -227,7 +228,18 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
}
// ----- dialogue (multi-turn slot carry-over; 2-min follow-up window) -----
dialogueSessions := dialogue.NewSessionStore(2 * time.Minute)
// 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()
@@ -775,11 +787,43 @@ func (h *reactiveHandler) applyAction(ctx context.Context, dec router.Decision)
log.Printf("voice: embed query: %v", err)
return "не получилось найти ответ."
}
// Long-term memory first: ONE search over everything Maven remembers
// (notes and facts share this index) and ONE confidence gate, so the
// memory that is clearly the best match answers — a note just as much
// as a fact.
//
// This used to run only after the notes-only gate below had already
// rejected the same note at the same score, which no note could ever
// survive a second time: the branch could only return a fact (#373).
// Order, not the gate, was the bug — the set of questions Maven answers
// is unchanged, only which memory gets to answer them.
if h.memStore != nil {
if hits, herr := h.memStore.Search(ctx, vec, 3); herr == nil {
if hit, ok := bestRecall(hits, h.queryMinScore, h.queryMinMargin); ok {
text := hit.Meta["text"]
// A note is phrased in Maven's voice; a fact is read back
// as it was stored.
if hit.Meta["type"] == "note" {
if reply, perr := h.phraser.PhraseQuery(ctx, dec.Utterance, []string{text}); perr == nil && reply != "" {
return reply
}
}
return text
}
} else {
log.Printf("voice: memory search: %v", herr)
}
}
notes, err := h.api.QueryNotes(ctx, vec, 5)
if err != nil {
log.Printf("voice: query notes: %v", err)
return "не получилось найти ответ."
}
// Notes-only pass, for notes the vector index above does not hold (an
// older note written before it existed). Same gate, notes-only
// candidates.
//
// Confidence gate: below it, say "I don't know" rather than read back
// the least-unrelated note — a confident wrong recall is worse than a
// gap (spec's "not a guesser-of-truth"). Same instinct as the loop's
@@ -791,16 +835,6 @@ func (h *reactiveHandler) applyAction(ctx context.Context, dec router.Decision)
noteScores[i] = n.Score
}
if !memory.ConfidentScores(noteScores, h.queryMinScore, h.queryMinMargin) {
// Long-term memory recall (notes + facts) before general knowledge:
// the notes table can't answer fact questions, but the memory store
// indexes both. Only runs when notes-RAG already gave up → additive.
if h.memStore != nil {
if hits, herr := h.memStore.Search(ctx, vec, 3); herr == nil {
if text, ok := bestRecall(hits, h.queryMinScore, h.queryMinMargin); ok {
return text
}
}
}
// Try general knowledge from the phraser before giving up
reply, err := h.phraser.PhraseQuery(ctx, dec.Utterance, nil)
if err != nil || reply == "" {
@@ -1734,3 +1768,71 @@ func jsonStringImpl(s string) string {
b = append(b, '"')
return string(b)
}
// reembedOnStart is the -reembed flag (set in run()). Opt-in on purpose: see
// runReembed.
var reembedOnStart bool
// checkStoredEmbedder compares the embedder we just loaded with the one that
// wrote the vectors already in the DB (Vikunja #378).
//
// The two models we have both make 384-dim vectors, so a size check catches
// nothing: after a swap, recall silently compares vectors from different
// spaces and the scores are noise. So we say it out loud. Recall itself is not
// changed here — the fix is `mavend -reembed`.
func checkStoredEmbedder(dataStore *store.Store, emb router.Embedder) {
if dataStore == nil {
return
}
current := router.EmbedderID(emb)
if reembedOnStart {
runReembed(dataStore, emb, current)
return
}
stored, mismatch, err := dataStore.CheckEmbedder(context.Background(), current)
if err != nil {
log.Printf("voice: embedder marker check failed: %v", err)
return
}
if mismatch {
log.Printf("voice: WARNING embedder MISMATCH — stored vectors were written by %q but the configured embedder is %q; recall scores are noise until the notes and facts are re-embedded — run `mavend -reembed` once (Vikunja #378)", stored, current)
return
}
log.Printf("voice: embedder marker ok (%s)", current)
}
// runReembed is the one-shot backfill behind -reembed.
//
// Why a flag and not automatic on mismatch: the embedder is ONNX on the
// laptop's CPU, so a few thousand notes is minutes of work. Doing that silently
// inside a normal start would look like the daemon hanging on boot. So the user
// runs it once, deliberately, after an embedder swap; the mismatch warning
// above tells them to. It re-embeds, logs what it did, and then the daemon
// carries on serving as usual — no separate binary, no second start needed.
func runReembed(dataStore *store.Store, emb router.Embedder, current string) {
log.Printf("voice: re-embedding stored notes and facts with %s — this can take a few minutes, do not interrupt", current)
res, err := dataStore.ReembedAll(context.Background(), current,
// EmbedPassage, not EmbedQuery: these are stored texts being searched
// FOR, which is the side they were written with.
func(ctx context.Context, text string) ([]float32, error) {
return router.EmbedPassage(ctx, emb, text)
})
if err != nil {
log.Printf("voice: re-embed FAILED, nothing was changed and no marker was written — safe to run again: %v", err)
return
}
if res.Skipped {
log.Printf("voice: re-embed skipped — the stored vectors were already written by %s", current)
return
}
log.Printf("voice: re-embed done — %d notes in the notes table, %d notes and %d facts in the memory index, took %s; stored vectors now belong to %s",
res.Notes, res.MemNotes, res.Facts, res.Took.Round(time.Second), current)
// A row with no text cannot be re-embedded, so its vector is still the old
// model's noise while the marker now says everything is current. Both write
// paths always store the text, so this should be zero — say it loudly
// rather than bury it in the line above if it ever isn't.
if res.NoText > 0 {
log.Printf("voice: WARNING %d stored rows had no text, so their vectors could not be re-embedded and are still noise; they will never match anything useful (Vikunja #378)", res.NoText)
}
}
+74
View File
@@ -1,8 +1,12 @@
package dialogue
import (
"context"
"encoding/json"
"sync"
"time"
"github.com/kami/maven/internal/store"
)
type Intent string
@@ -50,10 +54,22 @@ func (s *Session) IsExpired(now time.Time) bool {
return now.After(s.Timestamp.Add(s.TTL))
}
// SessionPersister — the bit of the store the session needs, as an interface
// so tests can swap it out. Data is an opaque blob: the store never looks
// inside, we encode the session as JSON here.
type SessionPersister interface {
SaveDialogueSession(ctx context.Context, id string, data []byte, ts time.Time, ttl time.Duration) error
DeleteDialogueSession(ctx context.Context, id string) error
LoadDialogueSessions(ctx context.Context, now time.Time) ([]store.DialogueSessionRow, error)
}
// SessionStore keeps the live sessions in a map (the fast path) and mirrors
// every write to the persister, so a daemon restart can load them back.
type SessionStore struct {
mu sync.RWMutex
sessions map[string]*Session
defaultTTL time.Duration
persist SessionPersister // may be nil: memory only (tests, no-store paths)
}
func NewSessionStore(defaultTTL time.Duration) *SessionStore {
@@ -66,6 +82,43 @@ func NewSessionStore(defaultTTL time.Duration) *SessionStore {
}
}
// NewPersistentSessionStore — same store, but writes also go to the DB.
// Call Load once after this to bring back sessions from a previous run.
func NewPersistentSessionStore(defaultTTL time.Duration, p SessionPersister) *SessionStore {
s := NewSessionStore(defaultTTL)
s.persist = p
return s
}
// Load — read the saved sessions back into memory. Anything past its TTL is
// dropped (and deleted from the DB by the store), never revived.
func (s *SessionStore) Load(ctx context.Context, now time.Time) error {
if s.persist == nil {
return nil
}
rows, err := s.persist.LoadDialogueSessions(ctx, now)
if err != nil {
return err
}
s.mu.Lock()
defer s.mu.Unlock()
for _, r := range rows {
var sess Session
if err := json.Unmarshal(r.Data, &sess); err != nil {
// A blob we can't read is not worth failing a startup over.
continue
}
if sess.TTL <= 0 {
sess.TTL = r.TTL
}
if sess.IsExpired(now) {
continue
}
s.sessions[r.ID] = &sess
}
return nil
}
func (s *SessionStore) Get(id string, now time.Time) *Session {
s.mu.RLock()
sess, ok := s.sessions[id]
@@ -87,12 +140,33 @@ func (s *SessionStore) Put(id string, sess *Session) {
s.mu.Lock()
s.sessions[id] = sess
s.mu.Unlock()
s.save(id, sess)
}
func (s *SessionStore) Delete(id string) {
s.mu.Lock()
delete(s.sessions, id)
s.mu.Unlock()
if s.persist != nil {
_ = s.persist.DeleteDialogueSession(context.Background(), id)
}
}
// save — mirror one session to the DB. Best effort: memory already has it, so
// a write error costs us the restart safety net, not the current turn.
func (s *SessionStore) save(id string, sess *Session) {
if s.persist == nil {
return
}
data, err := json.Marshal(sess)
if err != nil {
return
}
ts := sess.Timestamp
if ts.IsZero() {
ts = time.Now()
}
_ = s.persist.SaveDialogueSession(context.Background(), id, data, ts, sess.TTL)
}
func InheritSlots(prev, cur Slots) Slots {
+118
View File
@@ -0,0 +1,118 @@
package dialogue
import (
"context"
"path/filepath"
"testing"
"time"
"github.com/kami/maven/internal/store"
)
// openStore — a store on disk, so a second handle can reopen the same file.
func openStore(t *testing.T, path string) *store.Store {
t.Helper()
s, err := store.Open(context.Background(), path)
if err != nil {
t.Fatalf("open store: %v", err)
}
t.Cleanup(func() { _ = s.Close() })
return s
}
// A session written before a restart comes back and still merges a follow-up.
func TestSessionSurvivesRestart(t *testing.T) {
ctx := context.Background()
path := filepath.Join(t.TempDir(), "maven_test.db")
now := time.Now().UTC().Truncate(time.Millisecond)
first := openStore(t, path)
before := NewPersistentSessionStore(2*time.Minute, first)
before.Put("voice", &Session{
Intent: IntentReminder,
Slots: Slots{Text: "полить цветы", Time: now.Add(time.Hour), HasTime: true},
Timestamp: now,
TTL: 2 * time.Minute,
})
if err := first.Close(); err != nil {
t.Fatalf("close: %v", err)
}
// fresh handle, fresh in-memory map — as after a daemon restart
after := openStore(t, path)
reloaded := NewPersistentSessionStore(2*time.Minute, after)
if err := reloaded.Load(ctx, now.Add(10*time.Second)); err != nil {
t.Fatalf("Load: %v", err)
}
sess := reloaded.Get("voice", now.Add(10*time.Second))
if sess == nil {
t.Fatal("session did not survive the restart")
}
if sess.Intent != IntentReminder {
t.Fatalf("intent = %q, want reminder", sess.Intent)
}
// the follow-up carries no text of its own; it must inherit the old one
merged := InheritSlots(sess.Slots, Slots{Time: now.Add(2 * time.Hour), HasTime: true})
if merged.Text != "полить цветы" {
t.Fatalf("merged text = %q, want the earlier turn's text", merged.Text)
}
if !merged.Time.Equal(now.Add(2 * time.Hour)) {
t.Fatalf("merged time = %v, want the follow-up's time", merged.Time)
}
}
// A session past its TTL is dead: a restart must not bring it back.
func TestExpiredSessionNotResurrected(t *testing.T) {
ctx := context.Background()
path := filepath.Join(t.TempDir(), "maven_test.db")
now := time.Now().UTC().Truncate(time.Millisecond)
first := openStore(t, path)
before := NewPersistentSessionStore(2*time.Minute, first)
before.Put("voice", &Session{
Intent: IntentReminder,
Slots: Slots{Text: "полить цветы"},
Timestamp: now,
TTL: time.Minute,
})
if err := first.Close(); err != nil {
t.Fatalf("close: %v", err)
}
after := openStore(t, path)
reloaded := NewPersistentSessionStore(2*time.Minute, after)
later := now.Add(5 * time.Minute) // well past the 1-min TTL
if err := reloaded.Load(ctx, later); err != nil {
t.Fatalf("Load: %v", err)
}
if sess := reloaded.Get("voice", later); sess != nil {
t.Fatalf("expired session came back: %+v", sess)
}
// and it is gone from the DB too, not just from memory
rows, err := after.LoadDialogueSessions(ctx, later)
if err != nil {
t.Fatalf("LoadDialogueSessions: %v", err)
}
if len(rows) != 0 {
t.Fatalf("expired row still in the DB: %+v", rows)
}
}
// Delete removes the row as well, so an ended conversation stays ended.
func TestDeleteRemovesPersistedSession(t *testing.T) {
ctx := context.Background()
path := filepath.Join(t.TempDir(), "maven_test.db")
now := time.Now().UTC().Truncate(time.Millisecond)
s := openStore(t, path)
ss := NewPersistentSessionStore(2*time.Minute, s)
ss.Put("voice", &Session{Intent: IntentChat, Timestamp: now, TTL: time.Minute})
ss.Delete("voice")
rows, err := s.LoadDialogueSessions(ctx, now)
if err != nil {
t.Fatalf("LoadDialogueSessions: %v", err)
}
if len(rows) != 0 {
t.Fatalf("row survived Delete: %+v", rows)
}
}
+2
View File
@@ -422,6 +422,8 @@ func rankNote(inTop3 bool) string {
// bestRecall mirrors cmd/mavend/recall.go — the gate the daemon actually
// applies to a memory hit. Duplicated rather than imported because package main
// is not importable; recalleval_test.go asserts the two agree in behaviour.
// The daemon returns the whole hit (a note and a fact are said differently);
// the harness only scores what came back, so it keeps returning the text.
func bestRecall(results []memory.Result, minScore, minMargin float64) string {
if !memory.Confident(results, minScore, minMargin) {
return ""
@@ -197,7 +197,8 @@ func TestHashRecallBaseline(t *testing.T) {
t.Log("\n" + rep.String() + rep.Failures())
t.Log("\ngate sweep:\n" + sweep(t, router.NewHashEmbedder(hashDim), f))
// 0.32 sits under the observed 0.360 recall@1.
// 0.32 sits under the observed 0.370 recall@1 (was 0.360 over 25 answerable
// cases; the two mixed note+fact cases added with #373 make it 27).
const floorRecall1 = 0.32
if rep.Recall1() < floorRecall1 {
t.Errorf("recall@1 %.3f below ratchet %.2f — note recall regressed", rep.Recall1(), floorRecall1)
@@ -387,6 +387,32 @@
{"id": "n2", "text": "wifi channel is 6", "kind": "note"},
{"id": "n3", "text": "the guest network is off", "kind": "note"}
]
},
{
"id": "ru-mixed-031",
"lang": "ru",
"tags": ["mixed", "paraphrase", "hard"],
"note": "notes and facts in one store and the note is the answer — the daemon indexes both (Vikunja #373)",
"query": "куда я спрятал второй ключ от квартиры",
"want": "n1",
"notes": [
{"id": "n1", "text": "запасной ключ от квартиры лежит в синей коробке на полке", "kind": "note"},
{"id": "x1", "text": "поменял замок в двери двадцатого июня", "kind": "fact"},
{"id": "x2", "text": "отдал ключ соседке в мае", "kind": "fact"}
]
},
{
"id": "ru-mixed-032",
"lang": "ru",
"tags": ["mixed", "distractor"],
"note": "the mirror of ru-mixed-031: the fact answers and the notes are the distractors",
"query": "когда я в последний раз заливал бензин",
"want": "x1",
"notes": [
{"id": "x1", "text": "залил полный бак в четверг вечером", "kind": "fact"},
{"id": "n1", "text": "на заправке у моста дешевле бензин", "kind": "note"},
{"id": "n2", "text": "надо поменять зимние шины", "kind": "note"}
]
}
]
}
+23
View File
@@ -2,6 +2,7 @@ package router
import (
"context"
"fmt"
"math"
"unicode"
)
@@ -19,6 +20,24 @@ type Embedder interface {
Close() error
}
// IdentifiedEmbedder — an embedder that can name itself. The name goes into
// the DB next to the vectors it wrote, so a later model swap is caught instead
// of silently returning nonsense scores (Vikunja #378).
type IdentifiedEmbedder interface {
Embedder
ID() string
}
// EmbedderID is the stable string stored alongside the vectors. It comes from
// the embedder itself — nobody hand-types a model name twice — and changes
// whenever the model or its dimension changes.
func EmbedderID(e Embedder) string {
if i, ok := e.(IdentifiedEmbedder); ok {
return i.ID()
}
return fmt.Sprintf("unknown@%d", e.Dim())
}
// AsymmetricEmbedder — an embedder that wants to know whether a text is a
// search query or a stored passage. Recall is asymmetric: a short question
// goes in, a longer note comes out. The e5 family is trained for exactly that
@@ -70,6 +89,10 @@ func NewHashEmbedder(dim int) *HashEmbedder {
func (h *HashEmbedder) Dim() int { return h.dim }
// ID names this embedder for the DB marker. The dimension is part of it
// because a HashEmbedder of another width is a different vector space.
func (h *HashEmbedder) ID() string { return fmt.Sprintf("hash@%d", h.dim) }
func (h *HashEmbedder) Close() error { return nil }
func (h *HashEmbedder) Embed(_ context.Context, text string) ([]float32, error) {
+24
View File
@@ -0,0 +1,24 @@
package router
import "testing"
func TestEmbedderIDFromModelPath(t *testing.T) {
got := modelIDFromPath("/opt/maven/models/embedder/multilingual-e5-small.onnx")
if got != "multilingual-e5-small@384" {
t.Fatalf("modelIDFromPath = %q", got)
}
// A different model file must produce a different id, even at 384 dim.
old := modelIDFromPath("/opt/maven/models/embedder/paraphrase-multilingual-MiniLM-L12-v2.onnx")
if old == got {
t.Fatal("two different models share one id")
}
}
func TestEmbedderIDIncludesDim(t *testing.T) {
if id := EmbedderID(NewHashEmbedder(1024)); id != "hash@1024" {
t.Fatalf("EmbedderID = %q", id)
}
if EmbedderID(NewHashEmbedder(1024)) == EmbedderID(NewHashEmbedder(384)) {
t.Fatal("dimension not part of the id")
}
}
+14 -96
View File
@@ -1,11 +1,8 @@
package eval
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
"os"
"strings"
"testing"
@@ -29,13 +26,20 @@ import (
// a bake-off across checkpoints (#278, #250) produces tables you can tell
// apart. Point the variable at one server at a time.
//
// Three configurations, because "the LLM router" is ambiguous and the three
// numbers answer different questions:
// Two configurations, because "the LLM router" is ambiguous and the two numbers
// answer different questions:
//
// llm-only — the model alone. Measures the prompt + grammar contract.
// cascade+llm — what #320 would actually ship: stage-0 grammar, then the
// model, then the classifier as the failure floor.
// llm-no-thinking — diagnostic only, not a shippable path (see below).
// llm-only — the model alone. Measures the prompt + grammar contract.
// cascade+llm — what #320 would actually ship: stage-0 grammar, then the
// model, then the classifier as the failure floor.
//
// There used to be a third, "thinking off", which looked 6 points better. It is
// gone: it was measured with a hand-rolled HTTP client that quietly dropped
// repeat_penalty, so the gap was the missing penalty and not the thinking mode.
// Re-measured with everything else held equal, thinking off scores exactly the
// same, case for case — and a direct probe shows this llama-server build ignores
// enable_thinking / reasoning_budget for this model anyway, so there was nothing
// to turn off. Full write-up in ROUTING-EVAL-31-07-2026.md (Vikunja #376).
func TestLLMRouterBaseline(t *testing.T) {
base := os.Getenv("MAVEN_LLM_URL")
if base == "" {
@@ -96,104 +100,18 @@ func TestLLMRouterBaseline(t *testing.T) {
}
t.Log("\n" + repCascade.String() + repCascade.Failures())
// llm-no-thinking: same prompt and grammar with the chat template's
// thinking mode off. Qwen3.5's template defaults thinking=1, so under a
// grammar the constrained JSON lands in reasoning_content with content
// empty — llm.Client's ReasoningContent fallback is what makes the router
// work at all today, by accident rather than design.
//
// MEASURED 2026-07-31: this variant scores identically to as-deployed
// (18/76, 48.7% intent-only, 2 errors, same p50). Thinking mode is a
// non-issue under a grammar — llama.cpp constrains the same token stream
// either way. Kept so the question stays answered instead of being
// re-asked, and so internal/llm does NOT grow a chat_template_kwargs field
// for a problem that does not exist.
repNoThink, err := Score(ctx, "llm-only ("+model+", thinking off) [diagnostic]",
RouterFunc(func(ctx context.Context, u string, now time.Time) (router.Decision, error) {
d, ok, err := router.NewLLMRouter(&noThinkCompleter{base: base, http: &http.Client{Timeout: 60 * time.Second}}).Route(ctx, u, now)
if err != nil {
return d, err
}
if !ok {
return d, fmt.Errorf("llm router declined without an error")
}
return d, nil
}), f)
if err != nil {
t.Fatalf("Score no-thinking: %v", err)
}
t.Log("\n" + repNoThink.String() + repNoThink.Failures())
// Reports rather than asserts — the numbers are inputs to the #320
// decision, and an assertion here would be this test inventing the bar.
// The one thing worth failing on is a harness fault: if every single case
// errors, the run measured infrastructure, not routing, and the report
// must not be mistaken for a score.
for _, rep := range []Report{repLLM, repCascade, repNoThink} {
for _, rep := range []Report{repLLM, repCascade} {
if rep.Errors == rep.Total {
t.Errorf("%s: all %d cases errored — harness fault, not a measurement", rep.Name, rep.Total)
}
}
}
// noThinkCompleter — llm.Client with chat_template_kwargs.enable_thinking
// false. A test-local copy rather than a change to internal/llm: whether the
// daemon should send it is the open question, and answering it here by adding
// the field would prejudge #320.
type noThinkCompleter struct {
base string
http *http.Client
}
func (c *noThinkCompleter) Complete(ctx context.Context, r llm.Req) (string, error) {
payload := map[string]any{
"messages": []map[string]string{
{"role": "system", "content": r.System},
{"role": "user", "content": r.User},
},
"max_tokens": r.MaxTokens,
"temperature": 0,
"grammar": r.Grammar,
"chat_template_kwargs": map[string]any{"enable_thinking": false},
}
b, err := json.Marshal(payload)
if err != nil {
return "", err
}
req, err := http.NewRequestWithContext(ctx, "POST", c.base+"/v1/chat/completions", bytes.NewReader(b))
if err != nil {
return "", err
}
req.Header.Set("Content-Type", "application/json")
resp, err := c.http.Do(req)
if err != nil {
return "", err
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
return "", fmt.Errorf("status %d", resp.StatusCode)
}
var out struct {
Choices []struct {
Message struct {
Content string `json:"content"`
ReasoningContent string `json:"reasoning_content"`
} `json:"message"`
} `json:"choices"`
}
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
return "", err
}
if len(out.Choices) == 0 {
return "", fmt.Errorf("no choices")
}
m := out.Choices[0].Message
if m.Content != "" {
return m.Content, nil
}
return m.ReasoningContent, nil
}
func ping(ctx context.Context, c *llm.Client) error {
ctx, cancel := context.WithTimeout(ctx, 90*time.Second)
defer cancel()
+21
View File
@@ -33,6 +33,7 @@ const (
type onnxEmbedder struct {
tokenizer *unigramTokenizer
session *ort.DynamicSession[int64, float32]
id string
}
func NewONNXEmbedder(modelPath, tokenizerPath, libPath string) (*onnxEmbedder, error) {
@@ -58,11 +59,31 @@ func NewONNXEmbedder(modelPath, tokenizerPath, libPath string) (*onnxEmbedder, e
return &onnxEmbedder{
tokenizer: tok,
session: session,
id: modelIDFromPath(modelPath),
}, nil
}
func (e *onnxEmbedder) Dim() int { return embedDim }
// ID names the loaded model for the DB marker (Vikunja #378): the model file's
// own name plus the dimension, so pointing the config at another model changes
// the string on its own.
func (e *onnxEmbedder) ID() string { return e.id }
// modelIDFromPath turns /opt/.../multilingual-e5-small.onnx into
// "multilingual-e5-small@384".
func modelIDFromPath(modelPath string) string {
name := modelPath
if i := strings.LastIndexAny(name, "/\\"); i >= 0 {
name = name[i+1:]
}
name = strings.TrimSuffix(name, ".onnx")
if name == "" {
name = "onnx"
}
return fmt.Sprintf("%s@%d", name, embedDim)
}
// Embed treats the text as a query. The classifier compares one short
// utterance to another short seed phrase, so both sides get the same prefix
// and the comparison stays fair. The recall path must call EmbedQuery and
+161
View File
@@ -0,0 +1,161 @@
package store
import (
"context"
"encoding/json"
"fmt"
"time"
)
// EmbedFunc embeds one piece of stored text. The caller passes
// router.EmbedPassage — the STORE side of the query/passage asymmetry, which is
// the side every vector in the DB was written with. (Passing the query side
// would put the stored vectors in the wrong half of the space and quietly halve
// recall.) A func instead of an interface keeps this package free of any
// dependency on internal/router.
type EmbedFunc func(ctx context.Context, text string) ([]float32, error)
// BackfillResult is what the re-embed run did, for logging.
type BackfillResult struct {
Skipped bool // marker already matched — nothing to do
Notes int // rows rewritten in the notes table
Facts int // fact rows rewritten in memory_vectors
MemNotes int // note rows rewritten in memory_vectors
NoText int // memory_vectors rows with no text in their meta, left alone
Took time.Duration
}
// ReembedAll rewrites every stored vector with the currently configured
// embedder and then records that embedder as the one that owns the DB.
//
// Both places a vector lives are rewritten in the same pass: the `notes` table
// `embedding` column and the `memory_vectors` rows (notes AND facts). Doing
// only one would leave the two indexes disagreeing, which is worse than leaving
// both stale.
//
// Safe to re-run: if the marker already names the current embedder there is
// nothing to fix, so it returns immediately with Skipped set.
//
// Crash safety: everything — every vector and the marker — happens inside one
// transaction. If anything fails or the process dies partway, the transaction
// rolls back: no vectors changed and no marker written, so the next run does
// the whole job again. The marker is never set unless the full rewrite
// committed.
func (s *Store) ReembedAll(ctx context.Context, currentID string, embed EmbedFunc) (BackfillResult, error) {
start := time.Now()
var res BackfillResult
stored, err := s.Meta(ctx, metaKeyEmbedderID)
if err != nil {
return res, err
}
if stored == currentID {
res.Skipped = true
res.Took = time.Since(start)
return res, nil
}
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return res, fmt.Errorf("reembed: begin: %w", err)
}
defer tx.Rollback() // no-op once committed
// ----- notes table -----
type noteRow struct {
id int64
text string
}
var notes []noteRow
rows, err := tx.QueryContext(ctx, `SELECT id, text FROM notes WHERE text != ''`)
if err != nil {
return res, fmt.Errorf("reembed: read notes: %w", err)
}
for rows.Next() {
var n noteRow
if err := rows.Scan(&n.id, &n.text); err != nil {
rows.Close()
return res, fmt.Errorf("reembed: note row: %w", err)
}
notes = append(notes, n)
}
rows.Close()
if err := rows.Err(); err != nil {
return res, fmt.Errorf("reembed: notes: %w", err)
}
for _, n := range notes {
vec, err := embed(ctx, n.text)
if err != nil {
return res, fmt.Errorf("reembed: embed note %d: %w", n.id, err)
}
if _, err := tx.ExecContext(ctx,
`UPDATE notes SET embedding = ? WHERE id = ?`, floatsToBlob(vec), n.id); err != nil {
return res, fmt.Errorf("reembed: write note %d: %w", n.id, err)
}
res.Notes++
}
// ----- memory_vectors (the unified index: notes AND facts) -----
// The text to re-embed is the one carried in the row's meta blob, which is
// exactly the text that was embedded when the row was written.
type vecRow struct {
id, text, kind string
}
var vecs []vecRow
rows, err = tx.QueryContext(ctx, `SELECT id, meta FROM memory_vectors`)
if err != nil {
return res, fmt.Errorf("reembed: read memory vectors: %w", err)
}
for rows.Next() {
var id, metaJSON string
if err := rows.Scan(&id, &metaJSON); err != nil {
rows.Close()
return res, fmt.Errorf("reembed: memory row: %w", err)
}
meta := map[string]string{}
if err := json.Unmarshal([]byte(metaJSON), &meta); err != nil {
rows.Close()
return res, fmt.Errorf("reembed: meta for %q: %w", id, err)
}
if meta["text"] == "" {
res.NoText++
continue
}
vecs = append(vecs, vecRow{id: id, text: meta["text"], kind: meta["type"]})
}
rows.Close()
if err := rows.Err(); err != nil {
return res, fmt.Errorf("reembed: memory vectors: %w", err)
}
for _, v := range vecs {
vec, err := embed(ctx, v.text)
if err != nil {
return res, fmt.Errorf("reembed: embed %q: %w", v.id, err)
}
if _, err := tx.ExecContext(ctx,
`UPDATE memory_vectors SET vec = ? WHERE id = ?`, encodeVec(vec), v.id); err != nil {
return res, fmt.Errorf("reembed: write %q: %w", v.id, err)
}
if v.kind == "fact" {
res.Facts++
} else {
res.MemNotes++
}
}
// Same transaction as the rewrite, on purpose: the marker can only exist if
// every vector above was written.
if _, err := tx.ExecContext(ctx,
`INSERT INTO meta (key, value) VALUES (?,?)
ON CONFLICT(key) DO UPDATE SET value = excluded.value`,
metaKeyEmbedderID, currentID); err != nil {
return res, fmt.Errorf("reembed: write marker: %w", err)
}
if err := tx.Commit(); err != nil {
return res, fmt.Errorf("reembed: commit: %w", err)
}
res.Took = time.Since(start)
return res, nil
}
+171
View File
@@ -0,0 +1,171 @@
package store
import (
"context"
"errors"
"testing"
"time"
)
// markerVec is a recognisable vector: nothing in these tests writes it except
// the backfill, so finding it proves the row really was rewritten.
var markerVec = []float32{9, 9, 9}
func newEmbedder(calls *int) EmbedFunc {
return func(_ context.Context, _ string) ([]float32, error) {
*calls++
return markerVec, nil
}
}
// seedOldVectors puts one note (notes table + unified index) and one fact
// (unified index only) in the DB, both carrying obviously-old vectors.
func seedOldVectors(t *testing.T, s *Store) {
t.Helper()
ctx := context.Background()
old := []float32{0.1, 0.2, 0.3}
id, err := s.WriteNote(ctx, time.Now(), "молоко в холодильнике", old, "voice")
if err != nil {
t.Fatalf("WriteNote: %v", err)
}
mem := s.VectorMemory()
if err := mem.Insert(ctx, "note:1", old, map[string]string{
"type": "note", "text": "молоко в холодильнике",
}); err != nil {
t.Fatalf("Insert note vector: %v", err)
}
if err := mem.Insert(ctx, "fact:water:1", old, map[string]string{
"type": "fact", "text": "я пил воду",
}); err != nil {
t.Fatalf("Insert fact vector: %v", err)
}
_ = id
}
func noteVec(t *testing.T, s *Store) []float32 {
t.Helper()
var blob []byte
if err := s.db.QueryRow(`SELECT embedding FROM notes LIMIT 1`).Scan(&blob); err != nil {
t.Fatalf("read note embedding: %v", err)
}
return blobToFloats(blob)
}
func memVec(t *testing.T, s *Store, id string) []float32 {
t.Helper()
var blob []byte
if err := s.db.QueryRow(`SELECT vec FROM memory_vectors WHERE id = ?`, id).Scan(&blob); err != nil {
t.Fatalf("read memory vector %s: %v", id, err)
}
return decodeVec(blob)
}
func sameVec(a, b []float32) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i] != b[i] {
return false
}
}
return true
}
// The deployed case: old vectors everywhere, no marker. Every vector in both
// places must be rewritten and the marker recorded.
func TestReembedAllRewritesEveryVector(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
seedOldVectors(t, s)
calls := 0
res, err := s.ReembedAll(ctx, "multilingual-e5-small@384", newEmbedder(&calls))
if err != nil {
t.Fatalf("ReembedAll: %v", err)
}
if res.Skipped {
t.Fatal("first run should not skip")
}
if res.Notes != 1 || res.MemNotes != 1 || res.Facts != 1 {
t.Fatalf("counts: notes=%d memNotes=%d facts=%d", res.Notes, res.MemNotes, res.Facts)
}
if calls != 3 {
t.Fatalf("embedder called %d times, want 3", calls)
}
if !sameVec(noteVec(t, s), markerVec) {
t.Fatalf("notes table not rewritten: %v", noteVec(t, s))
}
if !sameVec(memVec(t, s, "note:1"), markerVec) {
t.Fatal("unified index note row not rewritten")
}
if !sameVec(memVec(t, s, "fact:water:1"), markerVec) {
t.Fatal("unified index fact row not rewritten")
}
got, err := s.Meta(ctx, metaKeyEmbedderID)
if err != nil {
t.Fatalf("Meta: %v", err)
}
if got != "multilingual-e5-small@384" {
t.Fatalf("marker = %q", got)
}
}
// Re-running must do nothing at all — not a second pass over the same rows.
func TestReembedAllSecondRunIsNoop(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
seedOldVectors(t, s)
calls := 0
if _, err := s.ReembedAll(ctx, "e5@384", newEmbedder(&calls)); err != nil {
t.Fatalf("first run: %v", err)
}
first := calls
res, err := s.ReembedAll(ctx, "e5@384", newEmbedder(&calls))
if err != nil {
t.Fatalf("second run: %v", err)
}
if !res.Skipped {
t.Fatal("second run should report Skipped")
}
if calls != first {
t.Fatalf("second run embedded %d more rows, want 0", calls-first)
}
}
// A failure partway must leave the DB exactly as it was: no marker, and the old
// vectors still in place (one transaction, rolled back).
func TestReembedAllPartialFailureLeavesMarkerUnset(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
seedOldVectors(t, s)
before := noteVec(t, s)
calls := 0
boom := func(_ context.Context, _ string) ([]float32, error) {
calls++
if calls == 2 {
return nil, errors.New("onnx blew up")
}
return markerVec, nil
}
if _, err := s.ReembedAll(ctx, "e5@384", boom); err == nil {
t.Fatal("expected an error")
}
got, err := s.Meta(ctx, metaKeyEmbedderID)
if err != nil {
t.Fatalf("Meta: %v", err)
}
if got != "" {
t.Fatalf("marker was set to %q after a failed run", got)
}
if !sameVec(noteVec(t, s), before) {
t.Fatal("a failed run left a partially rewritten notes table")
}
// And the mismatch warning must still fire, so the user knows to re-run.
if _, mismatch, err := s.CheckEmbedder(ctx, "e5@384"); err != nil || !mismatch {
t.Fatalf("CheckEmbedder after failed backfill: mismatch=%v err=%v", mismatch, err)
}
}
+74
View File
@@ -0,0 +1,74 @@
package store
import (
"context"
"fmt"
"time"
)
// DialogueSessionRow — one saved follow-up session. Data is the session
// encoded by the dialogue package; the store does not look inside it.
type DialogueSessionRow struct {
ID string
Data []byte
Ts time.Time
TTL time.Duration
Expires time.Time
}
// SaveDialogueSession — write (or replace) the session for one dialogue id.
// One row per id: a newer turn overwrites the older state.
func (s *Store) SaveDialogueSession(ctx context.Context, id string, data []byte, ts time.Time, ttl time.Duration) error {
expires := ts.Add(ttl)
_, err := s.db.ExecContext(ctx, `
INSERT INTO dialogue_sessions (id, data, ts, ttl_ms, expires_ts) VALUES (?, ?, ?, ?, ?)
ON CONFLICT(id) DO UPDATE SET data = excluded.data,
ts = excluded.ts,
ttl_ms = excluded.ttl_ms,
expires_ts = excluded.expires_ts`,
id, data, ts.UnixMilli(), ttl.Milliseconds(), expires.UnixMilli())
if err != nil {
return fmt.Errorf("save dialogue session: %w", err)
}
return nil
}
// DeleteDialogueSession — drop one session (ended, or expired).
func (s *Store) DeleteDialogueSession(ctx context.Context, id string) error {
if _, err := s.db.ExecContext(ctx, `DELETE FROM dialogue_sessions WHERE id = ?`, id); err != nil {
return fmt.Errorf("delete dialogue session: %w", err)
}
return nil
}
// LoadDialogueSessions — return the sessions still alive at `now` and delete
// the ones that already ran out. An expired session is dead: it never comes
// back after a restart.
func (s *Store) LoadDialogueSessions(ctx context.Context, now time.Time) ([]DialogueSessionRow, error) {
if _, err := s.db.ExecContext(ctx,
`DELETE FROM dialogue_sessions WHERE expires_ts <= ?`, now.UnixMilli()); err != nil {
return nil, fmt.Errorf("prune dialogue sessions: %w", err)
}
rows, err := s.db.QueryContext(ctx,
`SELECT id, data, ts, ttl_ms, expires_ts FROM dialogue_sessions ORDER BY id`)
if err != nil {
return nil, fmt.Errorf("load dialogue sessions: %w", err)
}
defer rows.Close()
var out []DialogueSessionRow
for rows.Next() {
var r DialogueSessionRow
var tsMilli, ttlMilli, expMilli int64
if err := rows.Scan(&r.ID, &r.Data, &tsMilli, &ttlMilli, &expMilli); err != nil {
return nil, fmt.Errorf("scan dialogue session: %w", err)
}
r.Ts = time.UnixMilli(tsMilli).UTC()
r.TTL = time.Duration(ttlMilli) * time.Millisecond
r.Expires = time.UnixMilli(expMilli).UTC()
out = append(out, r)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("load dialogue sessions: %w", err)
}
return out, nil
}
+97
View File
@@ -0,0 +1,97 @@
package store
import (
"context"
"database/sql"
"errors"
"fmt"
)
// metaKeyEmbedderID names the embedder that wrote the stored vectors.
//
// Why one value for the whole DB and not a column on every vector row: the
// vectors are only ever rewritten all at once (one backfill re-embeds every
// note and fact together), so a per-row marker would hold the same string in
// every row and cost a column on two tables for nothing.
const metaKeyEmbedderID = "embedder_id"
// Meta reads a single value from the meta table. Missing key ⇒ empty string.
func (s *Store) Meta(ctx context.Context, key string) (string, error) {
var v string
err := s.db.QueryRowContext(ctx, `SELECT value FROM meta WHERE key = ?`, key).Scan(&v)
if errors.Is(err, sql.ErrNoRows) {
return "", nil
}
if err != nil {
return "", fmt.Errorf("read meta %s: %w", key, err)
}
return v, nil
}
// SetMeta writes (or overwrites) a single meta value.
func (s *Store) SetMeta(ctx context.Context, key, value string) error {
_, err := s.db.ExecContext(ctx,
`INSERT INTO meta (key, value) VALUES (?,?)
ON CONFLICT(key) DO UPDATE SET value = excluded.value`, key, value)
if err != nil {
return fmt.Errorf("write meta %s: %w", key, err)
}
return nil
}
// EmbedderUnknown is the stored id reported for a DB that already holds
// vectors but never recorded who wrote them.
const EmbedderUnknown = "unknown (written before this marker existed)"
// CheckEmbedder compares the embedder now configured against the one that
// wrote the stored vectors. Returns the stored id and whether it differs.
//
// Vectors from two different models live in different spaces, so cosine
// between them is noise rather than a low score — and both of our models are
// 384-dimensional, so nothing else catches it.
//
// Three cases, and the middle one is the one that actually matters:
//
// - marker present ⇒ compare the two ids.
// - marker absent but vectors already stored ⇒ this is a DB from before the
// marker, so we cannot know who wrote them. Report a mismatch. This is the
// real case on the deployed box: those vectors came from the old embedder,
// and claiming them for the current one would hide the exact problem the
// marker was added to catch.
// - marker absent and no vectors ⇒ fresh DB, claim it, nothing to fix.
//
// On a mismatch the fix is ReembedAll (backfill.go), run explicitly with
// `mavend -reembed`. Nothing is re-embedded here: that work is minutes of CPU
// on the laptop and must not stall a normal start.
func (s *Store) CheckEmbedder(ctx context.Context, currentID string) (stored string, mismatch bool, err error) {
stored, err = s.Meta(ctx, metaKeyEmbedderID)
if err != nil {
return "", false, err
}
if stored != "" {
return stored, stored != currentID, nil
}
n, err := s.countVectors(ctx)
if err != nil {
return "", false, err
}
if n > 0 {
return EmbedderUnknown, true, nil
}
return currentID, false, s.SetMeta(ctx, metaKeyEmbedderID, currentID)
}
// countVectors — how many stored rows carry an embedding. Used only to tell a
// fresh DB apart from one that predates the marker.
func (s *Store) countVectors(ctx context.Context) (int, error) {
var notes, vecs int
if err := s.db.QueryRowContext(ctx,
`SELECT count(*) FROM notes WHERE embedding IS NOT NULL`).Scan(&notes); err != nil {
return 0, fmt.Errorf("count note vectors: %w", err)
}
if err := s.db.QueryRowContext(ctx,
`SELECT count(*) FROM memory_vectors`).Scan(&vecs); err != nil {
return 0, fmt.Errorf("count memory vectors: %w", err)
}
return notes + vecs, nil
}
+100
View File
@@ -0,0 +1,100 @@
package store
import (
"context"
"testing"
"time"
)
// A fresh DB has no marker yet, so the current embedder is recorded and
// nothing is flagged.
func TestCheckEmbedderFreshDBRecords(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
stored, mismatch, err := s.CheckEmbedder(ctx, "multilingual-e5-small@384")
if err != nil {
t.Fatalf("CheckEmbedder: %v", err)
}
if mismatch {
t.Fatal("fresh DB reported a mismatch")
}
if stored != "multilingual-e5-small@384" {
t.Fatalf("stored = %q", stored)
}
got, err := s.Meta(ctx, metaKeyEmbedderID)
if err != nil {
t.Fatalf("Meta: %v", err)
}
if got != "multilingual-e5-small@384" {
t.Fatalf("marker not persisted, got %q", got)
}
}
// The deployed box: notes were written by the old embedder, before the marker
// existed. Claiming them for the current one would hide exactly the problem
// the marker is for, so an unmarked DB that already holds vectors is a
// mismatch.
func TestCheckEmbedderUnmarkedDBWithVectorsIsMismatch(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
if _, err := s.WriteNote(ctx, time.Now(), "молоко в холодильнике", []float32{0.1, 0.2}, "voice"); err != nil {
t.Fatalf("WriteNote: %v", err)
}
stored, mismatch, err := s.CheckEmbedder(ctx, "multilingual-e5-small@384")
if err != nil {
t.Fatalf("CheckEmbedder: %v", err)
}
if !mismatch {
t.Fatal("an unmarked DB with stored vectors should report a mismatch")
}
if stored != EmbedderUnknown {
t.Fatalf("stored = %q, want %q", stored, EmbedderUnknown)
}
// It must NOT claim the DB — that would silence the warning on restart.
got, err := s.Meta(ctx, metaKeyEmbedderID)
if err != nil {
t.Fatalf("Meta: %v", err)
}
if got != "" {
t.Fatalf("marker written despite unknown provenance: %q", got)
}
}
// Both models are 384-dim, so this is the only thing that catches the swap.
func TestCheckEmbedderDifferentModelMismatch(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
if err := s.SetMeta(ctx, metaKeyEmbedderID, "paraphrase-multilingual-MiniLM-L12-v2@384"); err != nil {
t.Fatalf("SetMeta: %v", err)
}
stored, mismatch, err := s.CheckEmbedder(ctx, "multilingual-e5-small@384")
if err != nil {
t.Fatalf("CheckEmbedder: %v", err)
}
if !mismatch {
t.Fatal("different embedder not detected")
}
if stored != "paraphrase-multilingual-MiniLM-L12-v2@384" {
t.Fatalf("stored = %q", stored)
}
}
// The same embedder must never raise a false alarm, including on re-check.
func TestCheckEmbedderSameModelNoAlarm(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
for i := 0; i < 2; i++ {
_, mismatch, err := s.CheckEmbedder(ctx, "multilingual-e5-small@384")
if err != nil {
t.Fatalf("CheckEmbedder: %v", err)
}
if mismatch {
t.Fatalf("false alarm on pass %d", i)
}
}
}
+14
View File
@@ -74,6 +74,20 @@ ALTER TABLE reminders ADD COLUMN next_fire_ts INTEGER;`, // #2
`CREATE INDEX IF NOT EXISTS idx_nudges_snoozed ON nudges (outcome_ts) WHERE outcome = 'snoozed';`, // #8 — SnoozedUntil runs every tick; keep it off a full scan (Vikunja #364)
`ALTER TABLE proposed_routines ADD COLUMN accepted_ts INTEGER;
ALTER TABLE proposed_routines ADD COLUMN last_fired_ts INTEGER;`, // #9 — accepted routines keep firing (Vikunja #366): the tick loop needs to know when a routine was accepted and when it last nudged
`CREATE TABLE IF NOT EXISTS dialogue_sessions (
id TEXT PRIMARY KEY,
data BLOB NOT NULL,
ts INTEGER NOT NULL,
ttl_ms INTEGER NOT NULL,
expires_ts INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_dialogue_sessions_expires ON dialogue_sessions (expires_ts);`, // #10 — the follow-up session survives a restart (Vikunja #363); small, TTL-pruned table, not a history log
`CREATE TABLE IF NOT EXISTS meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);`, // #11 — small key/value table for facts about the DB itself; first key is embedder_id (Vikunja #378)
}
// migrate applies every migration with a number greater than the DB's current