Compare commits
12 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 2e9b9ec1cf | |||
| d1f6f6355f | |||
| 92ecb691de | |||
| 4282f6b9a9 | |||
| 7bb9f9be06 | |||
| 1e47eaca5a | |||
| 892330eb84 | |||
| 9a3bcd7c46 | |||
| 98ee701e03 | |||
| 04c1088088 | |||
| 07c191d8b8 | |||
| c668310b3e |
@@ -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
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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"}
|
||||
]
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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(¬es); 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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user