Compare commits

..

8 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 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
17 changed files with 939 additions and 27 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
+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)
}
})
+101 -10
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
@@ -786,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
@@ -802,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 == "" {
@@ -1745,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)
}
}
+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")
}
}
+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)
}
}
+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)
}
}
}
+5
View File
@@ -83,6 +83,11 @@ ALTER TABLE reminders ADD COLUMN next_fire_ts INTEGER;`, // #2
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