Compare commits

..

36 Commits

Author SHA1 Message Date
kami 453919db20 Merge pull request 'Bug: the memory index stores the raw utterance as a fact's recall text, and nothing ever deletes a fact vector' (#100) from task/493-bug-the-memory-index-stores-the-raw-utte into task/470-bug-a-question-writes-invented-knowledge
Reviewed-on: #100
2026-08-03 20:41:09 +02:00
claude ad60e10e95 mavend: run the fact vector repair on start, and test what it does (V-493)
Automatic rather than a flag, unlike -reembed: only voice-tapped facts are in
this index, so it is tens of embeddings rather than thousands of notes. And
waiting for an operator to know the repair exists is the failure being fixed.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-03 22:36:54 +04:00
claude 1528697287 store: repair fact vectors against the facts they name (V-493)
Every write-path fix leaves the rows already stored wrong, and a box in that
state looks fine: recall answers with the wrong text and nothing logs an error.
That is how the original poison survived four restarts.

RepairFactVectors resolves each fact vector against the fact it names,
re-embeds the ones whose text is stale, and deletes the voided, superseded and
orphaned ones. Marker-guarded and idempotent, so it runs once per box and a run
that dies partway is simply redone.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-03 22:36:54 +04:00
claude dbdab2d570 store, mavend: a fact is indexed as the fact, not as the utterance (V-493)
queryMemory returns a fact's stored text verbatim, so the text the write path
indexed is what he hears. It was the utterance, which made recall of any
voice-tapped fact answer with the sentence he said: go_version = 1.20 was
indexed as "какая последняя версия языка Go?", and that question came back.

FactRecallText renders the fact instead, and the utterance stays in meta as
provenance. Correcting a value now drops the key's vectors the way voiding one
does, since the superseded value was still answering.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-03 22:36:34 +04:00
kami b9371dcac6 Merge pull request 'Bug: a question writes invented knowledge into memory as a self fact, and recall then serves it back for unrelated questions' (#99) from task/470-bug-a-question-writes-invented-knowledge into master
Reviewed-on: #99
2026-08-03 20:11:11 +02:00
claude 62c2e92ec0 mavend, recalleval: wire the topic veto into both recall sources (V-470)
queryMemory and queryNotes both gate on score alone, so both needed it. The
eval keeps its own copy of bestRecall — package main is not importable — and a
fixture that measures a weaker gate than the daemon runs flatters it, so the copy
moves in step and its test pins the new rule.

Measured on the held-out recall fixture with the real embedder: 17/32 cases pass
→ 22/32, false recall 1/5 → 0/5, answered after gate 18/27 → 17/27. The one true
recall lost is en-hard-024, an English question against a Russian note, where no
lexical test can help.
2026-08-03 13:51:04 +04:00
claude aec94eb2e8 memory: a world question must name what the memory mentions (V-470)
The score gate cannot separate the right note from an unrelated one: the
held-out fixture puts the right note at 0.791-0.890 and the must-be-silent cases
at 0.795-0.835, so a note about his slow network answered 'почему небо синее?'.

RecallAllowed adds a topic veto, and applies it only to a question that mentions
nothing of his. That restriction is the whole design: demanding a shared word of
every recall silenced four true recalls on the fixture to kill one false one,
because recall exists to find the note whose words he no longer remembers. A
question about his own life keeps the embedder as its only judge.
2026-08-03 13:50:54 +04:00
claude 4dfe106fe3 mavend: a question is never a fact about him (V-470)
IntentFact used to persist whatever the model invented for a question-shaped
utterance, at confidence 1.00, and index it for recall under the question's own
text. Two such rows then claimed seven unrelated world questions and silently
disabled world answering.

A question now goes down the query chain, which is what he asked for. The second
half is confidence: a value grounded in what he said stays 1.00, a value the model
supplied for words he never said drops to 0.60 and says so in the log. Same
reasoning as 'LLM output is not authorization' on the act path.
2026-08-03 13:40:33 +04:00
claude 2e0e2fd0bb router: a deterministic test for question-shaped text (V-470)
The predicate a fact write needs before it trusts a routing decision. Tokenized,
not substring: 'что' inside 'чтобы' is not a question. Capture verbs win over
every question signal, because 'запиши что я пил воду' contains an interrogative
and is still a capture.
2026-08-03 13:40:33 +04:00
claude f3fa6b353a store: voiding a fact drops its memory vectors (V-470)
Revert voided the fact row and left the vector, so recall kept serving the
voided fact's utterance and the documented repair reported success on a box that
stayed broken. There was no way to repair a poisoned box at all.

DeletePrefix covers every vector for the key, earlier rows included: their values
are superseded, and a superseded value has no business claiming a turn. It is
best-effort — the audit trail is already committed, and a fact that is voided but
still recallable beats a void that failed.
2026-08-03 13:40:13 +04:00
kami 6645f64c3e Merge pull request 'Name the gap: world questions through the workstation model, and the four remaining callers' (#98) from task/490-name-the-gap-world-questions-through-the into master
Reviewed-on: #98
2026-08-03 11:19:56 +02:00
claude f10e0068dd config, deploy: the workstation is workpc, not bugmachine (V-490)
Owner's correction. It is the same host CLAUDE.md already calls workpc, and
two names for one machine read as two machines. The dated eval file keeps the
old name: a measurement is never edited after the day it was taken.
2026-08-03 12:42:37 +04:00
claude 9b124d9194 docs: both halves of the degradation rule are wired, and which caller is which (V-490)
The offload inventory grows a column, because "seven callers of the resident
model" stopped being the useful fact. Which of them is offloaded, and under
which half of the rule, is. Three are resident-only on purpose and the table
now says why rather than leaving it to be rediscovered.

The three-outcome table is the part that was not obvious from the rule as
written. A configured-and-asleep workstation names the gap; a box with no
workstation block does not, because naming a gap requires a gap.
2026-08-03 12:29:12 +04:00
claude 12530c8a95 mavend: world questions ask the workstation, and name the gap when it is asleep (V-490)
queryGeneral has nothing fetched to fall back on, so it is the sharp case:
with a workstation configured and asleep he is told that, rather than told
something false in a confident voice. The 1.7B answering a world question is
where "Война и мир" got Левитан as its author.

The sources that already hold a passage — a live search, a ZIM article, a
page he named — go through the world model too, but read the passage back
when it is not there instead of naming a gap. A real quote beats "не могу
сейчас", and nothing is invented on either path.

The Stub and every test double keep the Phraser interface they have.
PhraseWorld is reached by assertion, and a phraser without it is the
no-workstation case.
2026-08-03 12:27:40 +04:00
claude 51256c4c9a phraser: test the three outcomes of a world question, and prompt parity (V-490)
The middle outcome is the whole task: a workstation that is configured and
asleep produces a gap, and the resident model is never asked. The parity
test compares the bytes PhraseWorld sends the workstation against the bytes
PhraseQuery sends the resident model, so the fixtures and the daemon cannot
measure two different prompts.

The nudge tests cover the silent half from both sides, including the
temperature, which is how the workstation would otherwise change how she
sounds without anyone deciding to.
2026-08-03 12:27:30 +04:00
claude 76481c2736 phraser: a world model seam, so a gap can be named instead of invented (V-490)
The naming half of the degradation rule in docs/offload.md. PhraseWorld has
three outcomes: no workstation configured means the resident model answers
exactly as today, a workstation that is taking work answers, and one that is
asleep returns ErrNoWorldModel so the caller can say so. Naming a gap
requires a gap — on a box that never had a second model, refusing every
world question would remove a capability he has now.

Both prompts move into knowledgePrompt and evidencePrompt, shared by
PhraseQuery and PhraseWorld, because prompt parity across two models stops
holding the moment there are two copies of a prompt.

The silent half comes with it: chatWithSystem and chatWithMessages prefer
the workstation when it will take work, at the same 0.7 the resident
transport samples at, and say nothing when it will not. That covers the
digestion worker's nudge and reminder phrasing without touching tick.go.

Only Available and CompleteRemote are in the Remote interface. Pair.Complete
has its own floor and the phraser already owns one; two floors under a
single call is one too many.
2026-08-03 12:27:30 +04:00
claude bcc2305cd0 llm: let a caller name its sampling temperature (V-490)
The phraser's own transport has always sampled at 0.7 and this client has
always been greedy. Routing a phrasing call through the client must not
change how it decodes, so Req carries the temperature and 0 — the zero
value, and what every existing caller wanted — is still greedy.
2026-08-03 12:27:04 +04:00
kami 0ceeac8df4 Merge pull request 'Point Maven at the workstation model: a workstation block, and routing plus replies through llm.Pair' (#97) from task/485-run-the-big-model-on-the-workstation-wit into master
Reviewed-on: #97
2026-08-03 10:13:09 +02:00
claude 4fae13af75 docs: record the workstation routing numbers where the router is documented (V-485)
CLAUDE.md carried only the homesrv figures, which now read as the whole story.
Also points offload.md's order at #490 for the naming half.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NJYcaBiuny9UGSpFweQVQ1
2026-08-03 10:12:21 +02:00
claude 774217199e docs: measure gemma-4-12b on the workstation against the resident model (V-485)
Both fixtures, run from homesrv across the LAN with the proxy env stripped.
Routing: 84.4% full / 93.5% intent-only at p50 329ms through the cascade, against
72.7% / 77.9% at p50 0.80-1.04s for Qwen3-1.7B. Talk: 25/27 against 20/27, with
knowledge 9/9. Nudges 15/15. Settles #485's first assumption by measurement.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NJYcaBiuny9UGSpFweQVQ1
2026-08-03 10:12:21 +02:00
claude 2db59d52a7 deploy, docs: point homesrv at bugmachine and say what is still unwired (V-485) 2026-08-03 10:12:21 +02:00
claude 92d5fd580c mavend: route and reply through the workstation when its card is free (V-485)
modelSeam builds an llm.Pair when a workstation is configured and hands it to
the router and the replier. Both are the silent half of the degradation rule:
the big model is only better there, and he is never told which model answered.
No block, no probe, and the box behaves exactly as it did.
2026-08-03 10:12:21 +02:00
claude edeef19ff0 config: a workstation block, dropped when it names no address (V-485)
Health defaults to the supervisor's /health rather than llama-server's,
because mavgpud is what answers 503 while the card is held.
2026-08-03 10:12:21 +02:00
kami 018f7a6f47 Merge pull request 'Run the big model on the workstation, with admission control and the 1.7B as the floor' (#96) from task/489-workstation-deploy-mavgpud-on-workpc-and into master
Reviewed-on: #96
2026-08-03 10:12:13 +02:00
claude eca41798bd mavgpud: turn gemma's thinking off in the chat template (V-489)
Owner's call, 02-08-2026. Without it the 12B spends the reply budget on
reasoning tokens and answers empty at low max_tokens. Verified on the box:
"Столица Франции?" now answers "Париж" with no reasoning_content.
2026-08-02 22:44:10 +04:00
claude cc423567e7 docs: record that contention is KFD presence, not a VRAM threshold (V-489) 2026-08-02 22:29:58 +04:00
claude 8088ef9e00 mavgpud: build it with the rest, and ship the workstation config and unit (V-489)
make build now catches a broken supervisor on homesrv. deploy/mavgpud.json
carries the owner's gemma-4-12b line with the MTP draft model, passed to
llama-server untouched. The unit is a systemd user unit because sudo on the
workstation wants a password; lingering is the one command left to the owner.
2026-08-02 22:29:57 +04:00
kami 666b924d29 Merge pull request 'Run the big model on the workstation, with admission control and the 1.7B as the floor' (#95) from task/488-workstation-a-supervisor-that-keeps-llam into master
Reviewed-on: #95
2026-08-02 17:04:05 +02:00
claude e52c616592 mavgpud: test the probe against the sysfs the workstation actually has (V-488)
The fixtures are the live numbers sampled from the box on 02-08-2026, where the
CPT run held 12.8GB of 16 as proc/478104/vram_35881.

The cases that matter are the ones where a mistake is silent: our own
llama-server counting as a contender, an unreadable card reading as free, and
/health hanging or proxying into a closed port instead of answering 503.
2026-08-02 17:03:56 +02:00
claude 2b97bac51e mavgpud: keep the model loaded while the card is free, yield when it is not (V-488)
The lifecycle rule from Vikunja #488. Not on demand, because a 7-14B takes tens
of seconds to load and a world question would meet a gap every time the card
had been quiet. Not always on, because that is what holds the card.

/health is answered locally and always, so Maven's prober costs nothing and
works while the model is down. Everything else is reverse-proxied to
llama-server, which is what makes the idle window measurable at all.

Yielding is checked before starting, and both transitions are damped by a poll
streak so a short-lived rocm process cannot evict the model.
2026-08-02 17:03:56 +02:00
claude ab42db2b87 mavgpud: read the card from sysfs and own llama-server's lifecycle (V-488)
The workstation cannot keep a 7-14B resident: it would hold 16GB against the
owner's CPT runs, Correx and the manga-recap pipeline. So the process that
stays up costs no VRAM and the model comes and goes under it.

Contention is detected by presence on the KFD, not by a VRAM threshold. A ROCm
process registers under /sys/class/kfd/kfd/proc when it initialises HIP, well
before it allocates, so we see a contender during its startup instead of after
it has already lost an allocation race. rocm-smi is not installed on that box
and a per-second subprocess would get tuned down until useless, so this reads
sysfs and forks nothing.

Free VRAM is read only to decide whether to start. It is never a reason to
stop: by the time free VRAM has dropped, the other job has already failed.
2026-08-02 17:03:56 +02:00
kami 94d553570d Merge pull request 'Run the big model on the workstation, with admission control and the 1.7B as the floor' (#94) from task/485-run-the-big-model-on-the-workstation-wit into master
Reviewed-on: #94
2026-08-02 17:03:25 +02:00
claude 2e97b905b4 docs: the workstation supervisor owns llama-server's lifecycle (V-485)
The remote model cannot be a llama-server that is simply left running: a
resident 7-14B holds 16GB against the CPT runs the card is for. So what
is always up on the workstation is a supervisor, and llama-server is
loaded while the card is free.

Still not a scheduler. It arbitrates nothing between callers, and Maven
never asks it to start anything.
2026-08-02 18:14:31 +04:00
claude fbcca449be llm: pin that a down workstation is invisible (V-485)
Seven cases. The load-bearing ones are the constraint from 483: an
unconfigured deploy never probes and always reaches the floor, a busy
card degrades silently with the remote untouched, and a remote that dies
between probes still completes the turn and corrects the cached answer on
its way out.

CompleteRemote is pinned not to fall back, because a named gap that
quietly became a 1.7B guess is the failure this whole split exists to
prevent. And 1000 Available calls are pinned to make zero probes.
2026-08-02 17:19:53 +04:00
claude 2076e4a788 llm: prefer the workstation model, floor on the resident one (V-485)
Pair holds both models and decides which answers. A prober asks the
remote whether it will take work and caches the answer, so a request
reads an atomic bool rather than paying for a health check. Routing sits
at p50 825ms on the hot path and must never wait on a machine that may be
asleep.

The two methods are the two halves of the degradation rule in
docs/offload.md. Complete falls back silently, for routing, replies and
nudge phrasing, where the big model is only better. CompleteRemote
returns ErrRemoteUnavailable instead, for a world question, where the
1.7B does not answer worse but invents.

A nil remote is the unconfigured deploy: nothing probes, everything goes
to the floor, and the box behaves exactly as it does today.
2026-08-02 17:19:53 +04:00
kami 30eb6add1b Merge pull request 'Docs: refresh the QA plan against the live task list' (#93) from task/483-docs-offload-design into master
Reviewed-on: #93
2026-08-02 15:08:42 +02:00
40 changed files with 3142 additions and 85 deletions
+1
View File
@@ -9,6 +9,7 @@
/mavwaked
/mavmaild
/mavupdate
/mavgpud
# Certs (private keys, don't commit)
certs/
+16 -1
View File
@@ -36,6 +36,13 @@ stays on homesrv permanently, because it backs that floor. Read `docs/offload.md
touching a daemon seam or adding a model caller. Vikunja #483 is the umbrella, #484 to #487
are the work.
Both halves are wired as of 2026-08-03. Routing and replies prefer the workstation silently
through `modelSeam`; nudge and reminder phrasing prefer it silently inside the phraser. A
world question goes through `LLMPhraser.PhraseWorld` and names the gap when the card is not
free — `worldGap` in `cmd/mavend/worldmodel.go`, which he hears instead of an invented
answer. A box with no `workstation` block behaves exactly as it did before the seam: naming
a gap requires a gap. The offload table in `docs/offload.md` says which caller is which.
## Build & test
CGO daemons (`mavend`, `mavsttd`, `mavttsd`, `mavenclient`) need the vendored toolchain
@@ -143,7 +150,15 @@ seed additions, both of which now score inside the classifier baseline. Qwen3-1.
not a doubling, and the trade is worth re-arguing rather than assuming. **The ≈2.7s figure
that stood here until 2026-08-02 was contention, not the model.** See `docs/evals/2026-07-31-routing.md` line 61, which measures the LLM router at
p50 825ms / p95 1.2s / max 3.0s and the full cascade at p50 0.80-1.04s. Do not plan latency
work off the bakeoff table. `Confidence: 1.0` used to be hardcoded in `llmrouter.go`, so the LLM
work off the bakeoff table.
**The numbers above are the homesrv floor, not the ceiling.** With the workstation up, routing
completes through `llm.Pair` against gemma-4-12b and scores **84.4% full / 93.5% intent-only at
p50 329ms** — better than the resident model and about 2.5× faster (`docs/evals/2026-08-02-workstation-gemma4-12b.md`,
Vikunja #485). The workstation is never assumed up, so both sets of numbers are live. Judge a
routing change against the classifier and the resident model, since those are what always answer.
`Confidence: 1.0` used to be hardcoded in `llmrouter.go`, so the LLM
path could never ask for clarification (6/6 refusal cases missed on the fixture) — Vikunja
#359. Fixed 31-07-2026 with structural signal (single-token utterance, keyless fact, act with
no allowlisted fn) feeding the same stage-3 gate the classifier path already had — see
+8 -2
View File
@@ -16,11 +16,11 @@ PIPER_BIN := $(shell pwd)/deps/piper/piper
PIPER_MODEL := $(shell pwd)/models/tts/ru_RU-irina-medium.onnx
PIPER_ESPEAK := $(shell pwd)/deps/piper/espeak-ng-data
.PHONY: simulate stt-fixtures test-stt-golden all build build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav clean test fmt-check vet run-stt run-tts run-web download-embedder deps-go eval-router eval-recall eval-phrasing eval-models
.PHONY: simulate stt-fixtures test-stt-golden all build build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav clean test fmt-check vet run-stt run-tts run-web download-embedder deps-go eval-router eval-recall eval-phrasing eval-models build-gpud
all: build
build: build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav build-mail build-update
build: build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav build-mail build-update build-gpud
build-stt:
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
@@ -59,6 +59,12 @@ build-mail:
build-update:
$(GO) build $(GOFLAGS) -o mavupdate ./cmd/mavupdate/
# mavgpud runs on the workstation, not here. It is built with the rest so a
# broken supervisor is caught by `make build` on homesrv rather than by the
# workstation refusing to serve. Copy the binary over, do not `make deploy` it.
build-gpud:
$(GO) build $(GOFLAGS) -o mavgpud ./cmd/mavgpud/
run-web: build-web
./mavweb -addr :9200 -voice 127.0.0.1:9100
+50 -14
View File
@@ -7,6 +7,7 @@ import (
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/store"
)
// actionFact handles router.IntentFact: persist a tapped self-fact, index
@@ -15,14 +16,41 @@ func (h *reactiveHandler) actionFact(ctx context.Context, dec router.Decision) s
if !dec.Slots.HasKey {
return "не разобрала, что записать — попробуй иначе."
}
// A question is never a fact about him (#470). "какая последняя версия
// языка Go?" used to land here, and the value stored was whatever the
// model invented for it, at confidence 1.00, indexed for recall under the
// question's own text. Two such rows then claimed seven unrelated world
// questions through recall and silently disabled world answering.
//
// The routing error itself is not fixed here — the answer is to answer.
// Sending the turn down the query chain is what he asked for anyway, and
// it costs a mis-routed capture nothing: an explicit "запиши ..." is not
// question-shaped, so it never takes this branch.
if router.IsQuestionShaped(dec.Utterance) {
log.Printf("voice: fact write refused, utterance is a question: %q (key %q) — answering as a query",
dec.Utterance, dec.Slots.Key)
q := dec
q.Intent = router.IntentQuery
// The key the model extracted is its guess at what to store, not a
// fact he has. Left in place, queryFactByKey would read it back and
// claim the turn before any real source ran.
q.Slots.Key, q.Slots.HasKey = "", false
q.Slots.Value = ""
return h.actionQuery(ctx, q)
}
now := h.now()
req := ipc.WriteFactReq{
Ts: now,
Kind: "self",
Key: dec.Slots.Key,
Value: dec.Slots.Value,
Source: "tap:voice",
Confidence: 1.0,
Ts: now,
Kind: "self",
Key: dec.Slots.Key,
Value: dec.Slots.Value,
Source: "tap:voice",
// Not 1.00 unconditionally any more (#470). A value he said is
// evidence; a value the model supplied for words he never said is a
// guess, and writing a guess at full confidence is the same mistake
// the act path already refuses under "LLM output is not
// authorization".
Confidence: factConfidence(dec.Utterance, dec.Slots.Value),
// Subject: the key doubles as the entity-resolution candidate —
// a voice-tapped fact's key is usually the thing/person it's
// about ("espresso_machine", "kate"), so queueing it for Nexus
@@ -36,17 +64,25 @@ func (h *reactiveHandler) actionFact(ctx context.Context, dec router.Decision) s
log.Printf("voice: write fact: %v", err)
return "не получилось сохранить факт."
}
// Index the fact utterance in long-term memory (best-effort, must not
// fail the fact write). Facts aren't in the notes table, so this is the
// only recall path for them — "когда я пил воду?" reads back from here.
// Index the fact in long-term memory (best-effort, must not fail the fact
// write). Facts aren't in the notes table, so this is the only recall path
// for them — "когда я пил воду?" reads back from here.
//
// The indexed text is the fact, not the utterance (#493). queryMemory
// returns a fact's stored text verbatim, so what goes in here is what he
// hears; storing the utterance meant recall answered with his own sentence
// rather than the value. The utterance stays alongside as provenance —
// readable on /trace, never the answer and never embedded.
if h.memStore != nil {
if vec, err := router.EmbedPassage(ctx, h.embedder, dec.Utterance); err != nil {
text := store.FactRecallText(dec.Slots.Key, dec.Slots.Value)
if vec, err := router.EmbedPassage(ctx, h.embedder, text); err != nil {
log.Printf("voice: embed fact for memory: %v", err)
} else if err := h.memStore.Insert(ctx, "fact:"+dec.Slots.Key+":"+strconv.FormatInt(now.Unix(), 10), vec, map[string]string{
"source": "voice",
"type": "fact",
"text": dec.Utterance,
"ts": strconv.FormatInt(now.Unix(), 10),
"source": "voice",
"type": "fact",
"text": text,
"utterance": dec.Utterance,
"ts": strconv.FormatInt(now.Unix(), 10),
}); err != nil {
log.Printf("voice: memory insert fact: %v", err)
}
+39 -24
View File
@@ -13,6 +13,7 @@ import (
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/memory"
"github.com/kami/maven/internal/morning"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/rss"
"github.com/kami/maven/internal/store"
@@ -433,6 +434,14 @@ func (h *reactiveHandler) queryMemory(ctx context.Context, t *queryTurn) (string
return "", false
}
text := hit.Meta["text"]
// The score cleared the gate and the topic still has to match (#470). A
// note about his slow network scored high enough to answer "почему небо
// синее?", because the right-note and must-be-silent score ranges overlap
// and no threshold sits between them.
if !memory.RecallAllowed(t.dec.Utterance, text) {
log.Printf("voice: recall %q rejected for %q: a world question and no shared topic word", text, t.dec.Utterance)
return "", false
}
// A note is phrased in Maven's voice; a fact is read back as it was
// stored.
if hit.Meta["type"] == "note" {
@@ -468,6 +477,12 @@ func (h *reactiveHandler) queryNotes(ctx context.Context, t *queryTurn) (string,
if !memory.ConfidentScores(noteScores, h.queryMinScore, h.queryMinMargin) {
return "", false
}
// Same topic veto as queryMemory above: the best note must be about what
// he asked, not merely the nearest vector in the index.
if !memory.RecallAllowed(t.dec.Utterance, notes[0].Text) {
log.Printf("voice: note %q rejected for %q: a world question and no shared topic word", notes[0].Text, t.dec.Utterance)
return "", false
}
texts := make([]string, len(notes))
for i, n := range notes {
texts[i] = n.Text
@@ -523,10 +538,7 @@ func (h *reactiveHandler) queryWeb(ctx context.Context, t *queryTurn) (string, b
// the question he actually asked. She answers the question, she does not
// recite the page.
snippet := page.Title + "\n" + crawl.TrimRunes(page.Text, webPageContextRunes)
reply, perr := h.phraser.PhraseQuery(ctx, t.dec.Utterance, []string{snippet})
if perr != nil {
log.Printf("voice: web: phrase: %v", perr)
}
reply := h.phraseSource(ctx, "web", t.dec.Utterance, []string{snippet})
if reply == "" {
// No phraser (or it failed): read back the top of the page rather than
// pretend the fetch did not happen.
@@ -592,14 +604,7 @@ func (h *reactiveHandler) querySearch(ctx context.Context, t *queryTurn) (string
// question he asked, not something to recite. The trim is one budget over the
// joined block, so a long first snippet cannot crowd out the rest.
evidence := crawl.TrimRunes(strings.Join(resp.Snippets(), "\n"), h.search.runes)
var reply string
if h.phraser != nil {
var perr error
reply, perr = h.phraser.PhraseQuery(ctx, t.dec.Utterance, []string{evidence})
if perr != nil {
log.Printf("voice: search: phrase: %v", perr)
}
}
reply := h.phraseSource(ctx, "search", t.dec.Utterance, []string{evidence})
if reply == "" {
// No phraser, or it failed. Read back the best evidence rather than
// pretend the search did not happen.
@@ -680,14 +685,7 @@ func (h *reactiveHandler) queryKiwix(ctx context.Context, t *queryTurn) (string,
// Handed over the same way a note or a page is: context for the question he
// asked, not something to recite.
snippet := top.Title + "\n" + crawl.TrimRunes(page.Text, h.kiwix.runes)
var reply string
if h.phraser != nil {
var perr error
reply, perr = h.phraser.PhraseQuery(ctx, t.dec.Utterance, []string{snippet})
if perr != nil {
log.Printf("voice: kiwix: phrase: %v", perr)
}
}
reply := h.phraseSource(ctx, "kiwix", t.dec.Utterance, []string{snippet})
if reply == "" {
// No phraser, or it failed. Read back the best hit rather than pretend
// the search did not happen.
@@ -755,11 +753,28 @@ func isPersonalQuery(utterance string) bool {
return false
}
// queryGeneral — general knowledge from the phraser, the last source before
// giving up. It always claims: either the model answers or Maven says she
// doesn't know.
// queryGeneral — general knowledge, the last source before giving up. It always
// claims: either a model answers, or Maven names the gap, or she says she does
// not know.
//
// This is the sharpest case for the naming half. Nothing has been fetched, so
// there is no passage to fall back on and no floor under the answer except the
// model's weights — and a 1.7B's weights are where the invented answers come
// from. With a workstation configured and asleep he is told that, rather than
// told something false in a confident voice. With no workstation configured at
// all the resident model answers exactly as it does today: naming a gap requires
// a gap, and on that box the 1.7B is the whole product.
func (h *reactiveHandler) queryGeneral(ctx context.Context, t *queryTurn) (string, bool) {
reply, err := h.phraser.PhraseQuery(ctx, t.dec.Utterance, nil)
if h.phraser == nil {
// No model of any size. That is not the workstation being asleep, so it
// is not that gap: it is simply not knowing.
return "не знаю.", true
}
reply, err := h.phraseWorld(ctx, t.dec.Utterance, nil)
if errors.Is(err, phraser.ErrNoWorldModel) {
log.Printf("voice: %q needs the world model and it is not available", t.dec.Utterance)
return worldGap, true
}
if err != nil || reply == "" {
return "не знаю.", true
}
+82
View File
@@ -0,0 +1,82 @@
package main
import (
"log"
"strings"
"unicode"
)
// ungroundedConfidence — what a self fact is worth when its value appears
// nowhere in what he said. Below `query_min_score` is not the point (recall
// gates on vector distance, not on this number); the point is that
// `/history` and every future reader can tell a value he said from a value
// the model supplied.
const ungroundedConfidence = 0.6
// factConfidence scores a self fact by whether its value is grounded in the
// utterance it came from. Grounded stays 1.00, which is what a tapped fact
// has always been worth. Ungrounded drops, and says so in the log.
//
// An empty value is grounded by definition: the key alone carries the fact
// ("поужинал"), and there is nothing for the model to have invented.
func factConfidence(utterance, value string) float64 {
if strings.TrimSpace(value) == "" {
return 1.0
}
if valueGrounded(utterance, value) {
return 1.0
}
log.Printf("voice: fact value %q is not in %q — writing at confidence %.2f",
value, utterance, ungroundedConfidence)
return ungroundedConfidence
}
// valueGrounded reports whether every word of value traces back to a word he
// actually said. The comparison is on a 4-rune prefix, so the model's
// normalization survives ("пил воду" → "вода") while an invented value
// ("1.20" for a question about Go) does not.
func valueGrounded(utterance, value string) bool {
said := factTokens(utterance)
words := factTokens(value)
if len(words) == 0 {
return true
}
for _, w := range words {
if !anyTokenMatches(said, w) {
return false
}
}
return true
}
func anyTokenMatches(said []string, w string) bool {
for _, s := range said {
if s == w || sameStem(s, w) {
return true
}
}
return false
}
// sameStem is inflection tolerance and nothing more: it compares all but the
// last rune of the shorter word, and never fewer than three. Russian marks
// case on the ending, so "пил воду" and the stored "вода" are the same word he
// said, while "1.20" and "версия" are not. A word of three runes or fewer must
// match outright, where a shorter prefix would match half the language.
func sameStem(a, b string) bool {
ar, br := []rune(a), []rune(b)
shorter := min(len(ar), len(br))
n := shorter - 1
if n < 3 || len(ar) < n || len(br) < n {
return false
}
return string(ar[:n]) == string(br[:n])
}
// factTokens lowercases and splits on everything that is not a letter or a
// digit, the same shape planTokens uses in the router.
func factTokens(s string) []string {
return strings.FieldsFunc(strings.ToLower(s), func(r rune) bool {
return !unicode.IsLetter(r) && !unicode.IsDigit(r)
})
}
+125
View File
@@ -0,0 +1,125 @@
package main
import (
"context"
"testing"
"time"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/memory"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/tool"
"github.com/kami/maven/internal/voice"
)
func newFactGateHandler(t *testing.T, now time.Time) (*reactiveHandler, ipc.CoreAPI) {
t.Helper()
st := newTestStore(t)
api := ipc.NewStoreAPI(st)
emb := router.NewHashEmbedder(1024)
h := &reactiveHandler{
api: api,
embedder: emb,
router: buildRouter(emb, tool.NewMatcher(api), 0.55, nil),
replier: voice.NewStubReplier(),
now: func() time.Time { return now },
memStore: memory.NewInMemoryStore(),
dataStore: st,
}
return h, api
}
// The write half of #470: a question routed to IntentFact must not become a
// fact about him, and must not leave a vector behind for recall to serve.
func TestActionFact_QuestionIsNotWritten(t *testing.T) {
ctx := context.Background()
h, api := newFactGateHandler(t, time.Now())
reply := h.actionFact(ctx, router.Decision{
Intent: router.IntentFact,
Utterance: "какая последняя версия языка Go?",
Slots: router.Slots{Key: "go_version", HasKey: true, Value: `"1.20"`},
})
if _, err := api.LatestFact(ctx, "go_version"); err == nil {
t.Fatal("a question was stored as a fact about him")
}
hits, err := h.memStore.Search(ctx, mustEmbedPassage(t, h, "какая последняя версия языка Go?"), 3)
if err != nil {
t.Fatalf("memory search: %v", err)
}
if len(hits) != 0 {
t.Fatalf("the question was indexed for recall: %+v", hits)
}
// It went down the query chain instead. Nothing is configured to answer a
// world question in this harness, so "не знаю." is the honest outcome —
// what matters is that the turn was answered, not stored.
if reply == "" {
t.Fatal("the turn was neither stored nor answered")
}
}
// The capture that must survive the gate: an explicit instruction to record,
// even though it contains an interrogative.
func TestActionFact_ExplicitCaptureStillWrites(t *testing.T) {
ctx := context.Background()
h, api := newFactGateHandler(t, time.Now())
h.actionFact(ctx, router.Decision{
Intent: router.IntentFact,
Utterance: "запиши что я пил воду",
Slots: router.Slots{Key: "water", HasKey: true, Value: `"вода"`},
})
f, err := api.LatestFact(ctx, "water")
if err != nil {
t.Fatalf("an explicit capture was refused: %v", err)
}
if f.Confidence != 1.0 {
t.Errorf("confidence = %v, want 1.0 for a value he said", f.Confidence)
}
// #493: what recall reads back is the fact, not the sentence he said.
// queryMemory returns a fact's text verbatim, so the utterance sitting here
// meant "запиши что я пил воду" was the answer to "когда я пил воду?".
hits, err := h.memStore.Search(ctx, mustEmbedPassage(t, h, "вода"), 3)
if err != nil {
t.Fatalf("memory search: %v", err)
}
if len(hits) != 1 {
t.Fatalf("the fact was not indexed once: %+v", hits)
}
if got := hits[0].Meta["text"]; got != "water — вода" {
t.Errorf("indexed text = %q, want the fact", got)
}
if got := hits[0].Meta["utterance"]; got != "запиши что я пил воду" {
t.Errorf("utterance provenance = %q, want it kept alongside", got)
}
}
func TestFactConfidence(t *testing.T) {
cases := []struct {
utterance, value string
want float64
}{
{"запиши что я пил воду", `"вода"`, 1.0},
{"я выпил кофе", `"кофе"`, 1.0},
{"поужинал", "", 1.0},
{"отметь что я полил кактус", `"полил кактус"`, 1.0},
{"какая последняя версия языка Go", `"1.20"`, ungroundedConfidence},
{"кто премьер Японии", `"Тонио Озаки"`, ungroundedConfidence},
}
for _, c := range cases {
if got := factConfidence(c.utterance, c.value); got != c.want {
t.Errorf("factConfidence(%q, %q) = %v, want %v", c.utterance, c.value, got, c.want)
}
}
}
func mustEmbedPassage(t *testing.T, h *reactiveHandler, text string) []float32 {
t.Helper()
vec, err := router.EmbedQuery(context.Background(), h.embedder, text)
if err != nil {
t.Fatalf("embed %q: %v", text, err)
}
return vec
}
+91
View File
@@ -0,0 +1,91 @@
package main
import (
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/llm"
)
// No `workstation` block is the shipping deploy. The seam must then be the
// resident client itself, with nothing probing anything.
func TestModelSeamUnconfiguredIsResidentOnly(t *testing.T) {
resident := llm.New("http://127.0.0.1:1", time.Second)
hot, pair := modelSeam(&config.Config{}, resident)
if pair != nil {
t.Error("built a pair with no workstation configured")
}
if hot == nil {
t.Fatal("no seam at all, so the cascade would route with the classifier")
}
}
// A workstation with no resident model behind it has no floor, and a Pair with
// no floor is a configuration mistake rather than a degraded mode.
func TestModelSeamWithoutResidentIsNil(t *testing.T) {
cfg := &config.Config{Workstation: &config.WorkstationConfig{URL: "http://127.0.0.1:1"}}
cfg.Workstation.Health = strings.TrimRight(cfg.Workstation.URL, "/") + "/health"
hot, pair := modelSeam(cfg, nil)
if hot != nil || pair != nil {
t.Errorf("built a seam with no floor: hot=%v pair=%v", hot, pair)
}
}
// The configured case: the seam is the pair, and the pair notices a workstation
// that answers /health.
func TestModelSeamPrefersAnAnsweringWorkstation(t *testing.T) {
up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
}))
defer up.Close()
cfg := &config.Config{Workstation: &config.WorkstationConfig{
URL: up.URL,
Probe: config.Duration(10 * time.Millisecond),
}}
cfg.Workstation.Health = strings.TrimRight(cfg.Workstation.URL, "/") + "/health"
hot, pair := modelSeam(cfg, llm.New("http://127.0.0.1:1", time.Second))
if pair == nil || hot == nil {
t.Fatal("no pair built for a configured workstation")
}
defer pair.Stop()
deadline := time.Now().Add(2 * time.Second)
for !pair.Available() && time.Now().Before(deadline) {
time.Sleep(5 * time.Millisecond)
}
if !pair.Available() {
t.Fatal("the pair never saw a workstation that answers /health")
}
}
// A card held by a CPT run answers 503, and that must read as unavailable
// rather than as an error a turn has to handle.
func TestModelSeamHeldCardIsUnavailable(t *testing.T) {
busy := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Error(w, "model not loaded", http.StatusServiceUnavailable)
}))
defer busy.Close()
cfg := &config.Config{Workstation: &config.WorkstationConfig{
URL: busy.URL,
Probe: config.Duration(10 * time.Millisecond),
}}
cfg.Workstation.Health = strings.TrimRight(cfg.Workstation.URL, "/") + "/health"
_, pair := modelSeam(cfg, llm.New("http://127.0.0.1:1", time.Second))
if pair == nil {
t.Fatal("no pair built for a configured workstation")
}
defer pair.Stop()
time.Sleep(50 * time.Millisecond)
if pair.Available() {
t.Error("a 503 from the supervisor read as available")
}
}
+90 -5
View File
@@ -48,7 +48,11 @@ type voiceWiring struct {
// mcp — the MCP client, nil unless the `mcp` block configures an enabled
// server (Vikunja #251). Its tools land in the same allowlist as every
// other act, so nothing else here has to know about it.
mcp *mcpWiring
// pair — the workstation model with the resident one as the floor, nil
// unless a `workstation` block names an address. Held here only so the
// prober is stopped on shutdown; callers were handed it at build time.
pair *llm.Pair
mcp *mcpWiring
// home — the Home Assistant client, nil unless the `smarthome` block is
// enabled (Vikunja #256). Its devices land in the same allowlist as every
// other act, so nothing else here has to know about it.
@@ -76,6 +80,9 @@ func (w *voiceWiring) close() {
if w.ttsClient != nil {
_ = w.ttsClient.Close()
}
if w.pair != nil {
w.pair.Stop()
}
w.mcp.close()
}
@@ -139,6 +146,7 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
emb = router.NewHashEmbedder(1024)
}
w.embedder = emb
repairFactVectors(dataStore, emb)
checkStoredEmbedder(dataStore, emb)
// ----- tool executor (the enabled act allowlist, store-backed) -----
@@ -188,6 +196,18 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
// the new llama-server when the resident model is swapped (Vikunja #250).
llmClient = llmClientFor(lp, 60*time.Second)
}
// The workstation model sits above that one when it is configured and its
// card is free. hot is what the router and the replier complete through:
// either the pair, or the resident client alone, or nothing at all.
hot, pair := modelSeam(cfg, llmClient)
w.pair = pair
// The phraser gets the same pair, which is what carries the workstation model
// into the paths that do not go through `hot`: world questions (the naming
// half), and the digestion worker's nudge and reminder phrasing (the silent
// half). Wiring, so it happens once and before the voice server listens.
if lp, ok := phr.(*phraser.LLMPhraser); ok && pair != nil {
lp.UseRemote(pair)
}
// ----- router (the cascade; floor examples seed the classifier) -----
// The act matcher's allowlist is exactly the enabled tool names — the
// router only matches acts the executor can run (one source of truth).
@@ -199,7 +219,7 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
// against the classifier's 50.0%, at about 1s a turn instead of 30ms (see
// config.VoiceConfig.LLMRouter). The classifier always stays wired as the
// fallback, so a model error never breaks a turn.
rtr := buildRouter(emb, matcher, threshold, pickLLMRouter(cfg.Voice.UseLLMRouter(), llmClient))
rtr := buildRouter(emb, matcher, threshold, pickLLMRouter(cfg.Voice.UseLLMRouter(), hot))
// ----- sessions registry (shared with voicesink) -----
sessions := voice.NewSessions()
@@ -233,8 +253,8 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
// ----- replier (LLM-backed when the engine is on, Stub floor otherwise) -----
replier := voice.Replier(voice.NewStubReplier())
if llmClient != nil {
replier = newLLMReplier(llmClient, contextBlockFn(cfg, time.Now))
if hot != nil {
replier = newLLMReplier(hot, contextBlockFn(cfg, time.Now))
}
// ----- the handler (the reactive path; closes over stt / tts / router / coreAPI / memory) -----
@@ -291,7 +311,42 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
// pickLLMRouter returns the LLM router when the operator asked for it and there
// is a llama-server to talk to, and nil otherwise. nil is safe: the cascade then
// routes with the classifier, so an unusable setting costs accuracy, not turns.
func pickLLMRouter(enabled bool, c *llm.Client) *router.LLMRouter {
// modelSeam builds the completion seam the hot paths use: routing and replies.
//
// With no `workstation` block it is the resident client and nothing probes
// anything, which is today's deploy exactly. With one, it is an llm.Pair that
// prefers the workstation and falls back to the resident model silently — the
// silent half of the degradation rule (docs/offload.md), because the big model
// is only better here and the 1.7B is today's shipping quality. He is never
// told which of the two phrased his reply.
//
// A nil resident client means the phraser is not an LLM phraser. There is then
// no floor, and a Pair with no floor is a configuration mistake rather than a
// degraded mode, so the seam is nil and the cascade routes with the classifier.
func modelSeam(cfg *config.Config, resident *llm.Client) (router.Completer, *llm.Pair) {
if resident == nil {
if cfg.Workstation != nil {
log.Printf("voice: a workstation is configured but there is no resident model to floor it with — ignoring the block")
}
return nil, nil
}
if cfg.Workstation == nil {
return resident, nil
}
ws := cfg.Workstation
pair := llm.NewPair(
llm.New(ws.URL, time.Duration(ws.Timeout)),
resident,
ws.Health,
time.Duration(ws.Probe),
)
pair.Start(context.Background())
log.Printf("voice: workstation model at %s, probed every %s, resident model as the floor",
ws.URL, time.Duration(ws.Probe))
return pair, pair
}
func pickLLMRouter(enabled bool, c router.Completer) *router.LLMRouter {
if !enabled {
return nil
}
@@ -419,6 +474,36 @@ func seedTools(api ipc.CoreAPI, tools []config.ToolConfig) {
log.Printf("voice: seeded %d act tools from config", n)
}
// repairFactVectors brings stored fact vectors in line with the facts they name
// (#493), once per box, before the embedder marker is even looked at.
//
// Automatic and not a flag, unlike -reembed: only voice-tapped facts are in
// this index, so the work is tens of embeddings rather than the thousands of
// notes that made the backfill a deliberate act. And the box that needs it is
// broken in a way nobody can see — recall answers with the wrong text and
// nothing logs an error — so waiting for an operator to know to run it is how
// the defect survived four restarts in the first place.
func repairFactVectors(dataStore *store.Store, emb router.Embedder) {
if dataStore == nil {
return
}
res, err := dataStore.RepairFactVectors(context.Background(),
// EmbedPassage, the stored side, same as every other writer of these
// vectors.
func(ctx context.Context, text string) ([]float32, error) {
return router.EmbedPassage(ctx, emb, text)
})
if err != nil {
log.Printf("voice: fact vector repair failed, no marker written and nothing half-done — retried next start: %v", err)
return
}
if res.Skipped || res.Rewritten+res.Dropped == 0 {
return
}
log.Printf("voice: fact vector repair — %d re-embedded from the fact they name, %d dropped as voided or superseded, %d already right, took %s (#493)",
res.Rewritten, res.Dropped, res.Kept, res.Took.Round(time.Millisecond))
}
// reembedOnStart is the -reembed flag (set in run()). Opt-in on purpose: see
// runReembed.
var reembedOnStart bool
+60
View File
@@ -0,0 +1,60 @@
package main
import (
"context"
"errors"
"log"
"github.com/kami/maven/internal/phraser"
)
// worldPhraser — the naming half of the degradation rule (docs/offload.md), as
// the query sources see it. Only *phraser.LLMPhraser implements it, so the
// Stub and every test double stay exactly as they are.
type worldPhraser interface {
PhraseWorld(ctx context.Context, utterance string, sources []string) (string, error)
}
// worldGap — what he hears when the question is about the world, the workstation
// model is the one configured to answer it, and that machine is not answering.
//
// It says the true thing. The resident 1.7B is not a worse answer here, it is an
// invented one: "Война и мир" came back with Левитан as its author, and a
// question about his meeting came back as a swimming competition in Nottingham.
// Naming the gap is the rule CLAUDE.md already applies to a sibling service
// being down.
const worldGap = "сейчас не могу ответить — большая модель недоступна, а придумывать не хочу."
// phraseWorld asks the world model, or reports the gap.
//
// The three outcomes come straight from LLMPhraser.PhraseWorld: no workstation
// configured means the resident model answers as it always has, a workstation
// that is up answers, and a workstation that is down returns
// phraser.ErrNoWorldModel. A phraser that has no world seam at all — the Stub,
// and the doubles in the tests — is the first of those three.
func (h *reactiveHandler) phraseWorld(ctx context.Context, utterance string, sources []string) (string, error) {
if h.phraser == nil {
return "", phraser.ErrNoWorldModel
}
if w, ok := h.phraser.(worldPhraser); ok {
return w.PhraseWorld(ctx, utterance, sources)
}
return h.phraser.PhraseQuery(ctx, utterance, sources)
}
// phraseSource asks the world model to answer from a passage someone already
// fetched — a live search result, a ZIM article, a page he named. It returns ""
// rather than the gap phrase, because these callers hold something better than a
// gap: the passage itself, which their own floor reads back to him. Nothing is
// invented either way, and a real quote beats "не могу сейчас".
func (h *reactiveHandler) phraseSource(ctx context.Context, name, utterance string, sources []string) string {
reply, err := h.phraseWorld(ctx, utterance, sources)
switch {
case errors.Is(err, phraser.ErrNoWorldModel):
log.Printf("voice: %s: no world model, reading the source back instead", name)
return ""
case err != nil:
log.Printf("voice: %s: phrase: %v", name, err)
}
return reply
}
+85
View File
@@ -0,0 +1,85 @@
package main
import (
"context"
"strings"
"testing"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/router"
)
// gapPhraser — a phraser whose world model is configured and asleep, which is
// the state the naming half exists for.
type gapPhraser struct {
*phraser.Stub
worldCalls int
}
func (g *gapPhraser) PhraseWorld(context.Context, string, []string) (string, error) {
g.worldCalls++
return "", phraser.ErrNoWorldModel
}
func worldTurn(utterance string) *queryTurn {
return &queryTurn{dec: router.Decision{Intent: router.IntentQuery, Utterance: utterance}}
}
// A world question with the workstation asleep says so. The resident model is
// not asked, because what it produces here is an invention with no signal that
// it is one.
func TestQueryGeneralNamesTheGap(t *testing.T) {
g := &gapPhraser{Stub: phraser.NewStub()}
h := &reactiveHandler{phraser: g}
reply, ok := h.queryGeneral(context.Background(), worldTurn("почему небо голубое"))
if !ok {
t.Fatal("queryGeneral passed on the last source in the chain")
}
if reply != worldGap {
t.Fatalf("reply = %q, want the named gap", reply)
}
if g.worldCalls != 1 {
t.Fatalf("PhraseWorld called %d times, want 1", g.worldCalls)
}
}
// A phraser with no world seam at all — the Stub, and every box with no
// `workstation` block — answers exactly as it did before this seam existed.
func TestQueryGeneralWithoutAWorldModelIsUnchanged(t *testing.T) {
h := &reactiveHandler{phraser: phraser.NewStub()}
reply, ok := h.queryGeneral(context.Background(), worldTurn("почему небо голубое"))
if !ok {
t.Fatal("queryGeneral passed on the last source in the chain")
}
if reply != "не знаю." {
t.Fatalf("reply = %q, want the Stub's answer", reply)
}
}
// The gap is spoken aloud by a Russian voice, so it is Russian, feminine and
// informal. "не хочу" and "не могу" are her own verbs; there is no "вы" and no
// English in it.
func TestWorldGapIsInPersona(t *testing.T) {
for _, bad := range []string{"вы", "ваш", "рад ", "дорогой", "милый"} {
if strings.Contains(worldGap, bad) {
t.Errorf("the gap phrase contains %q: %s", bad, worldGap)
}
}
if strings.ContainsAny(worldGap, "abcdefghijklmnopqrstuvwxyz") {
t.Errorf("the gap phrase has Latin letters in it: %s", worldGap)
}
}
// The sources that hold a passage read it back rather than name a gap. He gets a
// real quote instead of "не могу сейчас", and nothing is invented either way.
func TestASourceWithAPassageReadsItBackInsteadOfNamingTheGap(t *testing.T) {
g := &gapPhraser{Stub: phraser.NewStub()}
h := &reactiveHandler{phraser: g}
if got := h.phraseSource(context.Background(), "search", "почему небо голубое",
[]string{"Рэлеевское рассеяние."}); got != "" {
t.Fatalf("phraseSource = %q, want \"\" so the caller's own floor reads the passage back", got)
}
if g.worldCalls != 1 {
t.Fatalf("PhraseWorld called %d times, want 1", g.worldCalls)
}
}
+117
View File
@@ -0,0 +1,117 @@
package main
import (
"os"
"path/filepath"
"strconv"
"strings"
)
// The card is an AMD 7900 GRE with 16GB, driven by amdgpu and ROCm. Everything
// here reads sysfs and forks nothing: rocm-smi is not even installed on the
// workstation, and a poll that costs a subprocess every second is a poll that
// gets tuned down until it is useless.
// gpuProc — one process holding the compute engine.
type gpuProc struct {
PID int
Comm string
VRAM int64 // bytes, as the kernel accounts them to this process
}
// probe reads the two sysfs trees the supervisor decides from.
//
// kfdRoot is /sys/class/kfd/kfd/proc, one directory per ROCm process. The
// directory appears when the process initialises HIP, which is well before it
// allocates anything large. That is the whole reason this works: the job that
// is about to want the card announces itself while it is still starting up,
// so we see the contender rather than only the winner of an allocation race.
//
// drmDev is /sys/class/drm/cardN/device, which reports total and used VRAM for
// the card as a whole.
type probe struct {
kfdRoot string
drmDev string
}
// foreign lists every ROCm process that is not ours. selfPID is the supervisor's
// llama-server child, or 0 when it is not running.
//
// An unreadable kfd tree returns no processes and no error. That is deliberate
// and it is the safe direction only because startVRAM also has to agree before
// anything launches: a supervisor that cannot see the KFD never sees free VRAM
// either, because the CPT run holding the card shows up in the drm totals.
func (p probe) foreign(selfPID int) []gpuProc {
entries, err := os.ReadDir(p.kfdRoot)
if err != nil {
return nil
}
var out []gpuProc
for _, e := range entries {
pid, err := strconv.Atoi(e.Name())
if err != nil || pid == selfPID {
continue
}
out = append(out, gpuProc{
PID: pid,
Comm: readComm(pid),
VRAM: p.procVRAM(e.Name()),
})
}
return out
}
// procVRAM sums the per-node vram_* files under one process directory. The
// suffix is the KFD topology node id (vram_35881 on this card), so it is
// globbed rather than named, and a machine with two cards sums both.
func (p probe) procVRAM(pid string) int64 {
matches, err := filepath.Glob(filepath.Join(p.kfdRoot, pid, "vram_*"))
if err != nil {
return 0
}
var total int64
for _, m := range matches {
total += readInt(m)
}
return total
}
// freeVRAM reports the bytes the card has left. Used only to decide whether to
// start: a shortfall here means llama-server would refuse to load anyway. It is
// never used to decide to stop, because by the time free VRAM has dropped the
// other job has already failed its allocation, which is exactly the outcome
// yielding exists to prevent.
func (p probe) freeVRAM() int64 {
total := readInt(filepath.Join(p.drmDev, "mem_info_vram_total"))
used := readInt(filepath.Join(p.drmDev, "mem_info_vram_used"))
if total <= 0 {
return 0
}
if free := total - used; free > 0 {
return free
}
return 0
}
func readInt(path string) int64 {
b, err := os.ReadFile(path)
if err != nil {
return 0
}
n, err := strconv.ParseInt(strings.TrimSpace(string(b)), 10, 64)
if err != nil {
return 0
}
return n
}
// readComm names the contender for the log. The log is the instrument for the
// open question in Vikunja #488: whether a process can want this card without
// ever registering on the KFD, which a Vulkan or video-decode job would.
func readComm(pid int) string {
b, err := os.ReadFile(filepath.Join("/proc", strconv.Itoa(pid), "comm"))
if err != nil {
return "?"
}
return strings.TrimSpace(string(b))
}
+103
View File
@@ -0,0 +1,103 @@
package main
import (
"net/http"
"net/http/httptest"
"net/url"
"os"
"path/filepath"
"strconv"
"testing"
)
// fakeKFD builds the sysfs shape the workstation actually has: one directory
// per ROCm process, each holding a vram_<node> file. Sampled from the live box
// on 02-08-2026, where the CPT run appeared as proc/478104/vram_35881.
func fakeKFD(t *testing.T, vramByPID map[int]int64) string {
t.Helper()
root := t.TempDir()
for pid, vram := range vramByPID {
dir := filepath.Join(root, strconv.Itoa(pid))
if err := os.MkdirAll(dir, 0o755); err != nil {
t.Fatal(err)
}
f := filepath.Join(dir, "vram_35881")
if err := os.WriteFile(f, []byte(strconv.FormatInt(vram, 10)+"\n"), 0o644); err != nil {
t.Fatal(err)
}
}
return root
}
func TestForeignExcludesOurChild(t *testing.T) {
root := fakeKFD(t, map[int]int64{478104: 12791693312, 999: 4096})
p := probe{kfdRoot: root}
all := p.foreign(0)
if len(all) != 2 {
t.Fatalf("with no child running, both processes are foreign, got %d", len(all))
}
ours := p.foreign(999)
if len(ours) != 1 || ours[0].PID != 478104 {
t.Fatalf("our own llama-server must not count as a contender, got %+v", ours)
}
if ours[0].VRAM != 12791693312 {
t.Errorf("per-process VRAM = %d, want the value from vram_35881", ours[0].VRAM)
}
}
// An empty KFD tree is the state that permits a start, so it must read as empty
// rather than as an error the caller has to interpret.
func TestForeignEmptyAndMissing(t *testing.T) {
if got := (probe{kfdRoot: t.TempDir()}).foreign(0); len(got) != 0 {
t.Errorf("empty kfd tree: got %d processes, want 0", len(got))
}
if got := (probe{kfdRoot: "/nonexistent"}).foreign(0); got != nil {
t.Errorf("missing kfd tree: got %+v, want nil", got)
}
}
func TestFreeVRAM(t *testing.T) {
dev := t.TempDir()
write := func(name, v string) {
if err := os.WriteFile(filepath.Join(dev, name), []byte(v), 0o644); err != nil {
t.Fatal(err)
}
}
// The live numbers from the workstation while the CPT run held the card.
write("mem_info_vram_total", "17163091968\n")
write("mem_info_vram_used", "13396389888\n")
p := probe{drmDev: dev}
if got, want := p.freeVRAM(), int64(3766702080); got != want {
t.Errorf("freeVRAM = %d, want %d", got, want)
}
if got := (probe{drmDev: "/nonexistent"}).freeVRAM(); got != 0 {
t.Errorf("unreadable card reports %d free, want 0 so nothing starts", got)
}
}
// With no model loaded the supervisor must still answer, and it must answer 503
// rather than hanging or proxying into a closed port. Maven reads this endpoint
// on a timer forever, including while the workstation is busy.
func TestHealthAndProxyRefuseWhenNotReady(t *testing.T) {
s := &supervisor{run: newRunner("/bin/true", nil, "")}
h := s.handler(mustURL(t, "http://127.0.0.1:1"))
for _, path := range []string{"/health", "/v1/chat/completions"} {
w := httptest.NewRecorder()
h.ServeHTTP(w, httptest.NewRequest(http.MethodGet, path, nil))
if w.Code != http.StatusServiceUnavailable {
t.Errorf("%s with no model: got %d, want 503", path, w.Code)
}
}
}
func mustURL(t *testing.T, s string) *url.URL {
t.Helper()
u, err := url.Parse(s)
if err != nil {
t.Fatal(err)
}
return u
}
+247
View File
@@ -0,0 +1,247 @@
// mavgpud — the workstation's GPU supervisor.
//
// It runs on the workstation (an AMD 7900 GRE, 16GB), not on homesrv, and it is
// deployed separately from the Maven daemons. Maven does not participate in any
// of this and never asks for a start: it reads /health through internal/llm.Pair
// and either gets the big model or falls back to the resident 1.7B.
//
// The rule, from Vikunja #488: keep llama-server loaded whenever the card is
// free, unload it when it has been idle too long or when another process needs
// the card. Not on demand, because a 7-14B takes tens of seconds to load and a
// world question would be answered by a gap every time the card had been quiet.
// Not always on, because that holds 16GB against the owner's own jobs.
package main
import (
"context"
"encoding/json"
"flag"
"log"
"net/http"
"net/http/httputil"
"net/url"
"os"
"os/signal"
"sync/atomic"
"syscall"
"time"
)
type config struct {
Listen string `json:"listen"` // what Maven talks to
LlamaAddr string `json:"llama_addr"` // where llama-server binds
LlamaBin string `json:"llama_bin"`
// LlamaArgs must include the flags that bind LlamaAddr. They are passed
// through untouched so the model, context size and layer count stay the
// owner's business and not this daemon's schema.
LlamaArgs []string `json:"llama_args"`
KFDRoot string `json:"kfd_root"`
DRMDevice string `json:"drm_device"`
Poll duration `json:"poll"`
IdleTimeout duration `json:"idle_timeout"`
StopGrace duration `json:"stop_grace"`
MinFreeVRAM int64 `json:"min_free_vram_bytes"`
// EvictAfter and StartAfter are counted in polls, not seconds. Both exist
// to damp flapping: a one-tick blip from a short-lived rocm process must
// not evict the model, and a card that has just been released must not be
// grabbed before the previous job has finished unmapping.
EvictAfter int `json:"evict_after_polls"`
StartAfter int `json:"start_after_polls"`
}
func defaults() config {
return config{
Listen: ":8080",
LlamaAddr: "127.0.0.1:8081",
KFDRoot: "/sys/class/kfd/kfd/proc",
DRMDevice: "/sys/class/drm/card1/device",
Poll: duration(time.Second),
IdleTimeout: duration(15 * time.Minute),
StopGrace: duration(20 * time.Second),
MinFreeVRAM: 15 << 30,
EvictAfter: 2,
StartAfter: 5,
}
}
// duration lets the config file say "15m" instead of counting nanoseconds.
type duration time.Duration
func (d *duration) UnmarshalJSON(b []byte) error {
var s string
if err := json.Unmarshal(b, &s); err != nil {
return err
}
v, err := time.ParseDuration(s)
if err != nil {
return err
}
*d = duration(v)
return nil
}
func main() {
path := flag.String("config", "/etc/mavgpud.json", "config file")
flag.Parse()
cfg := defaults()
b, err := os.ReadFile(*path)
if err != nil {
log.Fatalf("mavgpud: read config: %v", err)
}
if err := json.Unmarshal(b, &cfg); err != nil {
log.Fatalf("mavgpud: parse config: %v", err)
}
if cfg.LlamaBin == "" {
log.Fatal("mavgpud: llama_bin is required")
}
base := "http://" + cfg.LlamaAddr
run := newRunner(cfg.LlamaBin, cfg.LlamaArgs, base+"/health")
sup := &supervisor{
cfg: cfg,
probe: probe{kfdRoot: cfg.KFDRoot, drmDev: cfg.DRMDevice},
run: run,
}
sup.touch()
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer cancel()
target, err := url.Parse(base)
if err != nil {
log.Fatalf("mavgpud: llama_addr: %v", err)
}
srv := &http.Server{Addr: cfg.Listen, Handler: sup.handler(target)}
go func() {
log.Printf("mavgpud: listening on %s, model %s", cfg.Listen, cfg.LlamaBin)
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Fatalf("mavgpud: listen: %v", err)
}
}()
sup.loop(ctx)
// The card must come back before we do. A supervisor that exits leaving
// llama-server holding 14GB is worse than one that never ran.
shut, done := context.WithTimeout(context.Background(), 5*time.Second)
defer done()
_ = srv.Shutdown(shut)
run.stop(time.Duration(cfg.StopGrace))
}
type supervisor struct {
cfg config
probe probe
run *runner
lastReq atomic.Int64 // unix nanos of the last request Maven sent
foreignStreak int
clearStreak int
}
func (s *supervisor) touch() { s.lastReq.Store(time.Now().UnixNano()) }
func (s *supervisor) idle() time.Duration {
return time.Since(time.Unix(0, s.lastReq.Load()))
}
// handler serves the two things the workstation exposes.
//
// /health is answered locally and always, with no GPU cost and no round trip,
// because it is the only thing Maven reads and Maven reads it on a timer
// forever. Everything else is llama-server's API, reverse-proxied. Proxying
// rather than pointing Maven straight at llama-server is what makes the idle
// window measurable: the supervisor cannot otherwise know when the model was
// last used.
func (s *supervisor) handler(target *url.URL) http.Handler {
proxy := httputil.NewSingleHostReverseProxy(target)
mux := http.NewServeMux()
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
if !s.run.isReady() {
http.Error(w, "model not loaded", http.StatusServiceUnavailable)
return
}
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"status":"ok"}`))
})
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
if !s.run.isReady() {
http.Error(w, "model not loaded", http.StatusServiceUnavailable)
return
}
s.touch()
proxy.ServeHTTP(w, r)
})
return mux
}
func (s *supervisor) loop(ctx context.Context) {
t := time.NewTicker(time.Duration(s.cfg.Poll))
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
s.tick(ctx)
}
}
}
// tick is the whole decision. Yielding is checked before starting, and presence
// on the KFD is what triggers it — not a VRAM threshold. A ROCm process
// registers under /sys/class/kfd/kfd/proc when it initialises HIP, before it
// allocates, so we see a contender during its startup rather than after it has
// already failed to get the memory it wanted.
func (s *supervisor) tick(ctx context.Context) {
others := s.probe.foreign(s.run.pid())
if len(others) > 0 {
s.foreignStreak++
s.clearStreak = 0
} else {
s.foreignStreak = 0
s.clearStreak++
}
if s.run.running() {
s.run.refreshReady(ctx)
switch {
case s.foreignStreak >= s.cfg.EvictAfter:
log.Printf("mavgpud: yielding the card to %s", describe(others))
s.run.stop(time.Duration(s.cfg.StopGrace))
case s.idle() > time.Duration(s.cfg.IdleTimeout):
log.Printf("mavgpud: idle for %s, unloading", s.idle().Round(time.Second))
s.run.stop(time.Duration(s.cfg.StopGrace))
}
return
}
if s.clearStreak < s.cfg.StartAfter {
return
}
if free := s.probe.freeVRAM(); free < s.cfg.MinFreeVRAM {
return
}
s.touch() // the idle clock starts at load, not at the last request before it
if err := s.run.start(); err != nil {
log.Printf("mavgpud: start llama-server: %v", err)
}
}
// describe names the contenders in the log. This log is the instrument for the
// open question in #488: whether polling the KFD misses a job that wants the
// card without registering there.
func describe(procs []gpuProc) string {
out := ""
for i, p := range procs {
if i > 0 {
out += ", "
}
out += p.Comm
}
return out
}
+132
View File
@@ -0,0 +1,132 @@
package main
import (
"context"
"log"
"net/http"
"os/exec"
"sync"
"syscall"
"time"
)
// runner owns one llama-server process. Owning it is the point of the daemon:
// the workstation cannot keep a 7-14B resident, because that holds 16GB against
// the owner's CPT runs, Correx and the manga-recap pipeline. So the thing that
// stays up is this, which costs no VRAM, and the model comes and goes under it.
type runner struct {
bin string
args []string
// ready is llama-server's own /health, which answers "is a model loaded".
// Loading a 7-14B takes tens of seconds, so started is not ready.
readyURL string
mu sync.Mutex
cmd *exec.Cmd
ready bool
http *http.Client
}
func newRunner(bin string, args []string, readyURL string) *runner {
return &runner{
bin: bin, args: args, readyURL: readyURL,
http: &http.Client{Timeout: 2 * time.Second},
}
}
// pid is the child's, or 0. The GPU probe needs it to tell our own model apart
// from a contender.
func (r *runner) pid() int {
r.mu.Lock()
defer r.mu.Unlock()
if r.cmd == nil || r.cmd.Process == nil {
return 0
}
return r.cmd.Process.Pid
}
func (r *runner) running() bool { return r.pid() != 0 }
// isReady reports the cached readiness. The supervisor loop refreshes it; the
// health handler only reads, so answering /health never costs a round trip.
func (r *runner) isReady() bool {
r.mu.Lock()
defer r.mu.Unlock()
return r.ready
}
// start launches llama-server. It returns as soon as the process exists, not
// when the model is loaded.
func (r *runner) start() error {
r.mu.Lock()
defer r.mu.Unlock()
if r.cmd != nil {
return nil
}
cmd := exec.Command(r.bin, r.args...)
// Own process group, so stop kills anything llama-server spawned rather
// than leaving it holding VRAM after we have declared the card yielded.
cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
if err := cmd.Start(); err != nil {
return err
}
r.cmd, r.ready = cmd, false
log.Printf("mavgpud: started llama-server pid=%d", cmd.Process.Pid)
go func() {
err := cmd.Wait()
r.mu.Lock()
r.cmd, r.ready = nil, false
r.mu.Unlock()
log.Printf("mavgpud: llama-server exited: %v", err)
}()
return nil
}
// stop ends llama-server and waits for the VRAM to come back. SIGTERM first so
// it unmaps cleanly, SIGKILL after the grace window. Returning before the
// process is gone would let the supervisor report a free card while 14GB is
// still mapped, which is the one lie that would make yielding useless.
func (r *runner) stop(grace time.Duration) {
r.mu.Lock()
cmd := r.cmd
r.ready = false
r.mu.Unlock()
if cmd == nil || cmd.Process == nil {
return
}
pgid := -cmd.Process.Pid
_ = syscall.Kill(pgid, syscall.SIGTERM)
deadline := time.Now().Add(grace)
for time.Now().Before(deadline) {
if !r.running() {
return
}
time.Sleep(100 * time.Millisecond)
}
log.Printf("mavgpud: llama-server did not exit in %s, killing", grace)
_ = syscall.Kill(pgid, syscall.SIGKILL)
}
// refreshReady asks llama-server whether the model is loaded. Called once per
// supervisor tick, never per request.
func (r *runner) refreshReady(ctx context.Context) {
if !r.running() {
return
}
ok := false
req, err := http.NewRequestWithContext(ctx, http.MethodGet, r.readyURL, nil)
if err == nil {
resp, err := r.http.Do(req)
if err == nil {
ok = resp.StatusCode == http.StatusOK
resp.Body.Close()
}
}
r.mu.Lock()
was := r.ready
r.ready = ok
r.mu.Unlock()
if ok && !was {
log.Printf("mavgpud: model ready")
}
}
+17
View File
@@ -39,6 +39,23 @@
"proxy": "socks5://192.168.240.1:10808"
},
"//workstation": [
"The big model on the desk PC (workpc, 7900 GRE 16GB), fronted by",
"mavgpud on port 8080. It runs gemma-4-12b and it is preferred over the",
"resident Qwen3-1.7B for routing and replies whenever the card is free.",
"The machine is never assumed up: it sleeps, and the card is often held by",
"a CPT run, in which case mavgpud answers 503 and Maven falls back to the",
"resident model without saying so. Deleting this block restores exactly",
"the behaviour homesrv had before it existed.",
"Addressed by LAN address, not container name: mavgpud runs on another",
"machine and there is no shared docker network to name it on."
],
"workstation": {
"url": "http://192.168.1.105:8080",
"probe": "15s",
"timeout": "90s"
},
"//search": [
"The live web, searched after his own notes and before Kiwix. Only the",
"query string leaves the box — never a note, a fact, the persona block or",
+32
View File
@@ -0,0 +1,32 @@
{
"listen": ":8080",
"llama_addr": "127.0.0.1:10000",
"llama_bin": "llama-server",
"llama_args": [
"-m", "/mnt/D/AI/gemma4/gemma-4-12B-it-qat-UD-Q4_K_XL.gguf",
"-md", "/mnt/D/AI/gemma4/mtp-gemma-4-12B-it-BF16.gguf",
"-ngl", "99",
"-fa", "on",
"-np", "1",
"--host", "127.0.0.1",
"--port", "10000",
"--ctx-size", "32768",
"--threads", "6",
"--batch-size", "2048",
"--ubatch-size", "512",
"--jinja",
"--chat-template-kwargs", "{\"enable_thinking\":false}",
"--spec-type", "draft-mtp",
"--spec-draft-n-max", "2"
],
"kfd_root": "/sys/class/kfd/kfd/proc",
"drm_device": "/sys/class/drm/card1/device",
"poll": "1s",
"idle_timeout": "15m",
"stop_grace": "20s",
"min_free_vram_bytes": 10737418240,
"evict_after_polls": 2,
"start_after_polls": 5
}
+24
View File
@@ -0,0 +1,24 @@
[Unit]
# Runs on the workstation (bugmachine), not on homesrv. Install as a systemd
# user unit and turn on lingering, so the card is supervised after a reboot
# with nobody logged in:
#
# scp mavgpud workpc:~/.local/bin/mavgpud
# scp deploy/mavgpud.json workpc:~/.config/mavgpud.json
# scp deploy/mavgpud.service workpc:~/.config/systemd/user/mavgpud.service
# ssh workpc 'systemctl --user daemon-reload && systemctl --user enable --now mavgpud'
# sudo loginctl enable-linger kami
Description=Maven GPU supervisor (holds llama-server while the card is free)
After=network.target
[Service]
ExecStart=%h/.local/bin/mavgpud -config %h/.config/mavgpud.json
Restart=always
RestartSec=5
# The card must come back when the supervisor goes down. mavgpud stops
# llama-server on SIGTERM, so give it longer than stop_grace to do that.
KillSignal=SIGTERM
TimeoutStopSec=60
[Install]
WantedBy=default.target
@@ -0,0 +1,82 @@
# gemma-4-12b on the workstation, against the resident Qwen3-1.7B
Measured 2026-08-02 on the fixtures as they stand. Dated file: it is not edited
after today, and a newer number is a new file.
Vikunja #485's first assumption was that a 7-14B measurably beats Qwen3-1.7B on
the 77-case RU routing fixture and the 27-case talk fixture. It does, on both,
and it is also faster.
## The setup
`gemma-4-12B-it-qat-UD-Q4_K_XL` with the `mtp-gemma-4-12B-it-BF16` draft model,
served by `llama-server` b10220 on bugmachine (AMD 7900 GRE, 16GB), fronted by
`mavgpud` on `192.168.1.105:8080`. Thinking is off through
`--chat-template-kwargs '{"enable_thinking":false}'`, speculative decoding is
`--spec-type draft-mtp --spec-draft-n-max 2`, context 32768. The exact line is
`deploy/mavgpud.json`.
Every number below crossed the LAN from homesrv. Note the trap: homesrv's shell
exports `HTTP_PROXY`, Go honours it, and the runs need
`env -u HTTP_PROXY -u HTTPS_PROXY -u http_proxy -u https_proxy`.
## Routing, 77-case RU fixture
| | full | intent-only | p50 | p95 |
|---|---|---|---|---|
| classifier alone (02-08) | 68.8% | — | 16.6µs | — |
| Qwen3-1.7B through the cascade (31-07, 02-08) | 72.7% | 77.9% | 0.80-1.04s | — |
| **gemma-4-12b through the cascade** | **84.4%** | **93.5%** | **329ms** | 429ms |
| gemma-4-12b alone, no cascade | 55.8% | 85.7% | 335ms | 436ms |
The workstation buys 11.7 points of full accuracy over the resident model. It
buys 15.6 points of intent-only, at a third of the latency. The router's p50 was
never the model's fault, which the 02-08 contention finding already said. A 12B
on a free 16GB card answers a routing turn in a third of a second.
Two things the table hides.
The alone-versus-cascade gap is slots, not intents. gemma reads the intent right
85.7% of the time on its own. It loses full accuracy on seven fact keys
(`вода` instead of `water`, `ужин` instead of `meal`) and on six reminder times
with no time slot. Stage 0 and the daemon's own extractor repair
both, which is why the cascade is 28 points higher. The lesson is that the
cascade earns its keep even under a much better model, not that it is scaffolding
to remove.
`errors: 6` in the alone row are declines on single-token and ambiguous
utterances, all of which the cascade caught. The remaining defects through the
cascade are three `query→fact` confusions, one `chat→query`, and one false
clarify.
## Talk, 27-case conversational fixture
| | pass | notes |
|---|---|---|
| Qwen3-1.7B (31-07) | 20/27 | 11-17/27 for the 0.8B before it |
| **gemma-4-12b** | **25/27 (92.6%)** | chat 8/9, knowledge 9/9, query 8/9 |
Knowledge is the interesting column: 9/9, in Russian, with real answers about
Rayleigh scattering, SSD versus HDD and thunder delay. That is the case the
1.7B cannot do at all and the reason the naming half of the degradation rule
exists.
Two failures, and one of them is the persona defect the CPT (#122) targets:
`query-notes-do-not-answer` wrote `заплатил` where Maven needs the feminine
form. The other is `chat-joke`, where the model told a joke without using any of
the words the check looks for. Run-to-run variance is about one case: a second
run scored 24/27 with `chat-followup-server` also off-topic.
## Nudge phrasing, 15-case fixture
15/15, every check, no errors. `mood`, `lang`, `length`, `feminine`,
`hisgender`, `address`, `cringe` and `ontopic` all clean.
## What this settles and what it does not
Settled: the size question. A 12B on the workstation beats the resident model on
every fixture we have, and it is faster. The offload argument holds.
Not settled: how often the card is free. That is #485's second assumption and
only the `mavgpud` log answers it, after a week of the owner's normal work. A
model that is better whenever it is up is worth little if it is never up.
+77 -12
View File
@@ -1,6 +1,6 @@
# Offloading model work to the workstation
*Last verified: 2026-08-02 @ 5c05163. Living doc: correct it in place, do not append.*
*Last verified: 2026-08-03 @ 12530c8. Living doc: correct it in place, do not append.*
Owner's call, 2026-08-02. Vikunja #483 is the umbrella. Tasks #484 to #487 are the
work, and this file holds the shape and the rules all four must obey.
@@ -47,6 +47,24 @@ service being down.
Nothing in between. A turn never breaks on the workstation being asleep.
Both halves are wired, 03-08-2026. `LLMPhraser.PhraseWorld`
(`internal/phraser/world.go`) is the naming half and has three outcomes, not two:
| State | What he hears |
|---|---|
| no `workstation` block | the resident model answers, exactly as before the seam existed |
| configured, card free | the workstation answers |
| configured, asleep or busy | the gap, `worldGap` in `cmd/mavend/worldmodel.go` |
The first row is the one worth stating. Naming a gap requires a gap. On a box with
no second model the 1.7B is the whole product. Refusing every world question there
would remove a capability the owner has today.
A source holding a passage is on the naming half too: a live search, a ZIM
article, a page he named. None of them says "не могу сейчас". They read the
passage back, which is what `phraseSource` returning `""` selects. A real quote
beats a gap, and neither path invents.
## Admission control, not a scheduler
There is no GPU arbiter. That is a service with its own failure modes, and nothing
@@ -58,6 +76,34 @@ jobs.
The caller must be able to ask "is this peer usable right now" without a turn
hanging on a timeout. A dead remote is a normal state, not an error state.
`internal/llm.Pair` is that check on the Maven side. A prober caches the answer,
so `Available()` is an atomic read and no turn pays for a health check.
llama-server does not stay up on the workstation. It cannot: a resident 7-14B
would hold 16GB against the owner's CPT runs. So a supervisor there owns its
lifecycle, keeps it loaded while the card is free, and unloads it on idle or
when another process needs the card (owner's call, 2026-08-02, Vikunja #488).
That supervisor is still not a scheduler, and the distinction is worth holding.
It arbitrates nothing between callers. It reports whether it can take work and
manages one process to back that answer. Maven never asks it to start anything
and never learns that it did.
Contention is decided by presence under `/sys/class/kfd/kfd/proc`, not by a VRAM
threshold. A ROCm process registers there when it initialises HIP, before it
allocates anything. So the supervisor sees a contender during that job's startup,
and yields before the job loses the memory it asked for. A
threshold reads the card too late. By the time free VRAM has dropped, the other
job has already lost the allocation race. Free VRAM is still read, but only as a
precondition for loading, never as the eviction signal. One blind spot is known.
A job can take the card without registering on the KFD, as a Vulkan or a
video-decode job would. `describe()` logs every contender's comm, and that log is
how we find out whether the blind spot is real.
`mavgpud` runs from a systemd unit on the workstation with
`deploy/mavgpud.json` as its config, and `llama_args` is passed to llama-server
untouched. The model, the context size, the layer count and the MTP flags are the
owner's business and not this daemon's schema.
## What stays on homesrv, permanently
@@ -76,17 +122,29 @@ it buys nothing. Four callers:
## Inventory: what runs a model on homesrv today
The **resident model** is one llama-server with seven callers:
The **resident model** is one llama-server with seven callers, and 03-08-2026 is
the date each of them stopped or did not stop being resident-only:
| Caller | What for |
|---|---|
| `cmd/mavend/voicewire.go` | routing |
| `cmd/mavend/replier_llm.go` | replies |
| `cmd/mavend/tick.go` | digestion worker: `PhraseNudge`, `PhraseReminder` |
| `cmd/mavend/capture.go` | capture summarisation (unreachable, see #480) |
| `cmd/mavend/mail.go` | mail extraction (off, no IMAP) |
| `cmd/mavend/kiwixwire.go` | answering from a Kiwix, search or crawl passage |
| `memoryeval.go`, `modelswap.go` | admin and evals |
| Caller | What for | Offloaded |
|---|---|---|
| `cmd/mavend/voicewire.go` | routing | silently, through `hot` |
| `cmd/mavend/replier_llm.go` | replies | silently, through `hot` |
| `cmd/mavend/tick.go` | digestion worker: `PhraseNudge`, `PhraseReminder` | silently, inside the phraser |
| `cmd/mavend/actions_query.go` | world questions, and any fetched passage | names the gap |
| `cmd/mavend/capture.go` | capture summarisation (unreachable, see #480) | no, holds its own client |
| `cmd/mavend/mail.go` | mail extraction (off, no IMAP) | no, holds its own client |
| `memoryeval.go`, `modelswap.go` | admin and evals | no, and deliberately |
The last three rows are resident-only on purpose. `memoryeval.go` and
`modelswap.go` measure and swap the resident model, so sending their work
elsewhere would measure the wrong thing. `capture.go` and `mail.go` are
background jobs that hold a gated background client (`llmBackgroundClientFor`),
and that priority has no equivalent on the remote yet. Both are also unreachable
on this deploy, so wiring them would ship an untestable path.
The `tick.go` row needs one caveat. `phraser.llm_nudges` is `false` in deploy, so
nudges come from templates and the seam under them changes nothing until that
flips. It is wired anyway: `PhraseReminder` is on the same transport and is on.
Then the embedder above, **whisper.cpp** in `mavsttd`, and **piper** in `mavttsd`.
`mavwaked` uses no model at all: an energy-threshold VAD over 30ms frames.
@@ -97,7 +155,14 @@ Then the embedder above, **whisper.cpp** in `mavsttd`, and **piper** in `mavttsd
`internal/netaddr` landed in PR #92. A seam address now carries its own scheme,
and a scheme-less one is still unix. A tcp seam requires a shared token, because
the filesystem permission that authenticated the unix socket is gone.
2. **The resident model** (#485). Biggest quality delta. A 16GB card runs a 7-14B,
2. **The resident model** (#485, #490). Wired. A `workstation` block builds an
`llm.Pair` in `modelSeam` (`cmd/mavend/voicewire.go`), routing and replies
complete through it, and the phraser holds the same pair (`UseRemote`). Both
halves of the rule are live: see the table above for which caller gets which.
Measured, `docs/evals/2026-08-02-workstation-gemma4-12b.md`: gemma-4-12b
through the cascade scores 84.4% full accuracy at p50 329ms. The resident
model scores 72.7% at p50 0.80-1.04s. On the talk fixture it is 25/27
against 20/27. Biggest quality delta. A 16GB card runs a 7-14B,
which fixes what the 1.7B gets wrong: world knowledge, and the persona the CPT
targets. The degradation path is already written and measured, since the
classifier scores 68.8% full accuracy at p50 16.6µs on its own.
+15 -6
View File
@@ -498,12 +498,21 @@ rejects `https://api.openai.com`, and forget really deletes
(`internal/store/memory.go:145` is a real `DELETE`, not a tombstone). Vision is
19/19, speaker 22/22, media 16/16.
**470 got worse.** Both poisoned facts show `voided` on `/history`, and the
defect survives. Re-measured at 15:42, after four restarts: `почему небо синее?`
still answers `какая последняя версия языка Go?` with no `search:` line. What
comes back is the question he typed, not the value the fact held. So the poison
is a vector in the memory index, and `revert` does not remove it. There is
currently no documented way to repair a poisoned box.
**470 got worse, then closed.** Both poisoned facts showed `voided` on
`/history` and the defect survived. Re-measured at 15:42, after four restarts:
`почему небо синее?` still answered `какая последняя версия языка Go?` with no
`search:` line. What came back was the question he typed, not the value the fact
held. So the poison was a vector in the memory index, and `revert` did not
remove it.
Repaired in two parts. 470 stopped the writes: a question is never a fact, and a
void drops the key's vectors. 493 fixed what the index holds. A fact is indexed
as the fact and not as the utterance, and a correction drops its superseded
vector too.
A poisoned box now repairs itself on the next start. `RepairFactVectors`
re-embeds every fact vector from the fact it names, and deletes the voided and
superseded ones. It runs once, guarded by a marker, and logs what it did.
---
+61
View File
@@ -224,6 +224,11 @@ type Config struct {
// See SearchConfig.
Search *SearchConfig `json:"search,omitempty"`
// Workstation — the big model on the owner's desktop, preferred over the
// resident one when its GPU is free. nil / absent / url empty ⇒ homesrv
// behaves exactly as it does today. See WorkstationConfig.
Workstation *WorkstationConfig `json:"workstation,omitempty"`
// Praxis — the ecosystem attention-state service. When configured, maven
// calls the Praxis HTTP tools API for attention listing and item lifecycle.
// Maven never touches Praxis's database directly (ecosystem invariant: no
@@ -1105,6 +1110,44 @@ const (
DefaultKiwixSnippetRunes = 1500
)
// WorkstationConfig — the big model on the owner's desktop (workpc, a
// 7900 GRE with 16GB), fronted by mavgpud.
//
// homesrv cannot grow a GPU, so the resident Qwen3-1.7B is the floor and this
// is the preferred model above it (owner's call, 2026-08-02, docs/offload.md).
// The workstation is never assumed up: its card is often held by a CPT run and
// the machine sleeps. No block, or an empty URL, and homesrv behaves exactly as
// it does today.
//
// Only the prompt crosses the LAN, and the workstation is not "the box". The
// rules in CLAUDE.md about what may leave still apply.
type WorkstationConfig struct {
// URL — where mavgpud listens, e.g. "http://192.168.1.105:8080". Empty ⇒
// the whole block is normalised to nil and nothing probes anything.
URL string `json:"url,omitempty"`
// Health — the admission endpoint. Empty ⇒ URL + "/health", which is what
// mavgpud serves. It answers 503 while the card is held, and that is the
// signal, so it must be the supervisor's endpoint and not llama-server's.
Health string `json:"health,omitempty"`
// Probe — how often admission is re-checked. 0 ⇒ DefaultWorkstationProbe.
// Nothing on the hot path waits for it: the answer is cached and read
// atomically, so this only sets how late Maven notices the card came back.
Probe Duration `json:"probe,omitempty"`
// Timeout — the per-request budget for a completion on the workstation.
// 0 ⇒ DefaultWorkstationTimeout. A big model on a LAN host is slower than
// the resident one, and a request that overruns falls back to the floor.
Timeout Duration `json:"timeout,omitempty"`
}
// Workstation defaults, applied in Normalise.
const (
DefaultWorkstationProbe = 15 * time.Second
DefaultWorkstationTimeout = 90 * time.Second
)
// SearchConfig — the self-hosted SearXNG instance she searches with.
//
// External search is allowed and off unless configured (CLAUDE.md). Configuring
@@ -1500,6 +1543,24 @@ func (c *Config) applyDefaults() {
}
}
// No address, no preferred model. An unconfigured workstation is the
// default deploy and must be indistinguishable from today.
if c.Workstation != nil && strings.TrimSpace(c.Workstation.URL) == "" {
c.Workstation = nil
}
if c.Workstation != nil {
w := c.Workstation
if strings.TrimSpace(w.Health) == "" {
w.Health = strings.TrimRight(w.URL, "/") + "/health"
}
if w.Probe <= 0 {
w.Probe = Duration(DefaultWorkstationProbe)
}
if w.Timeout <= 0 {
w.Timeout = Duration(DefaultWorkstationTimeout)
}
}
if c.Voice != nil {
if c.Voice.RouterThreshold <= 0 {
c.Voice.RouterThreshold = DefaultRouterThreshold
+53
View File
@@ -413,3 +413,56 @@ func TestNormaliseFillsKiwixDefaults(t *testing.T) {
t.Error("rewrite: false was not honoured")
}
}
// A workstation with no address is not a workstation. The unconfigured deploy
// must be indistinguishable from today, so the block is dropped rather than
// left to fail one probe at a time.
func TestNormaliseDropsAddresslessWorkstation(t *testing.T) {
for _, tc := range []struct {
name string
in *WorkstationConfig
}{
{"no url", &WorkstationConfig{Probe: Duration(time.Second)}},
{"blank url", &WorkstationConfig{URL: " "}},
} {
t.Run(tc.name, func(t *testing.T) {
c := &Config{Workstation: tc.in}
c.applyDefaults()
if c.Workstation != nil {
t.Errorf("kept an unusable workstation block: %+v", c.Workstation)
}
})
}
}
// The health endpoint defaults to the supervisor's, not llama-server's: mavgpud
// answers 503 while the card is held, and that refusal is the whole signal.
func TestNormaliseFillsWorkstationDefaults(t *testing.T) {
c := &Config{Workstation: &WorkstationConfig{URL: "http://192.168.1.105:8080/"}}
c.applyDefaults()
if c.Workstation == nil {
t.Fatal("dropped a usable workstation block")
}
if got, want := c.Workstation.Health, "http://192.168.1.105:8080/health"; got != want {
t.Errorf("Health = %q, want %q", got, want)
}
if time.Duration(c.Workstation.Probe) != DefaultWorkstationProbe {
t.Errorf("Probe = %s, want %s", time.Duration(c.Workstation.Probe), DefaultWorkstationProbe)
}
if time.Duration(c.Workstation.Timeout) != DefaultWorkstationTimeout {
t.Errorf("Timeout = %s, want %s", time.Duration(c.Workstation.Timeout), DefaultWorkstationTimeout)
}
}
// An explicit health URL is left alone: the supervisor may sit behind something
// that does not put /health at the root.
func TestNormaliseKeepsExplicitWorkstationHealth(t *testing.T) {
c := &Config{Workstation: &WorkstationConfig{
URL: "http://192.168.1.105:8080",
Health: "http://192.168.1.105:9000/ready",
}}
c.applyDefaults()
if got, want := c.Workstation.Health, "http://192.168.1.105:9000/ready"; got != want {
t.Errorf("Health = %q, want %q", got, want)
}
}
+6 -1
View File
@@ -129,6 +129,11 @@ type Req struct {
RepeatPenalty float64
// Stop — sequences that end generation early (e.g. newline for a one-liner).
Stop []string
// Temperature — 0 (the zero value) is greedy decoding, and greedy is what
// every caller here wanted before this field existed. It is set only by the
// phraser, whose own transport has always sampled at 0.7: routing a phrasing
// call through this client must not quietly change how it decodes.
Temperature float64
}
type msg struct {
@@ -176,7 +181,7 @@ func (c *Client) Complete(ctx context.Context, r Req) (string, error) {
Messages: []msg{{Role: "system", Content: r.System}, {Role: "user", Content: r.User}},
MaxTokens: r.MaxTokens,
Grammar: r.Grammar,
Temp: 0,
Temp: r.Temperature,
RepeatPenalty: r.RepeatPenalty,
Stop: r.Stop,
})
+193
View File
@@ -0,0 +1,193 @@
package llm
import (
"context"
"errors"
"log"
"net/http"
"sync/atomic"
"time"
)
// Pair — a preferred model on another host, with the resident one as the floor.
//
// homesrv cannot grow a GPU and the workstation has 16GB of VRAM, so the big
// model runs there and the resident Qwen3-1.7B stays here. See docs/offload.md.
// The workstation is never assumed up: its GPU is often busy with CPT runs and
// the manga-recap pipeline, and the machine sleeps. So the remote is preferred,
// never required, and Pair is what makes "preferred" mean something precise.
//
// This is admission control, not a scheduler. There is no arbiter deciding who
// gets the card. A prober asks the remote whether it will take work, caches the
// answer, and every request reads that cached answer in nanoseconds. Routing
// sits on the hot path at p50 825ms and must never wait on a machine that may
// be asleep, so no request ever pays for a health check itself.
//
// Pair satisfies nothing by itself. Callers pick a method by which half of the
// degradation rule they live under:
//
// - Complete falls back silently. For routing, replies, and nudge phrasing,
// where the big model is only better and the 1.7B is today's shipping
// quality. He is not told which model phrased his reply.
// - CompleteRemote returns ErrRemoteUnavailable instead of falling back. For
// a world question, or a long Kiwix or search passage, where a 1.7B
// confabulates rather than summarises. A named gap beats an invented
// answer.
type Pair struct {
remote *Client
floor *Client
// up — the cached admission answer, written only by the prober goroutine
// and read by every request. Atomic so the read costs nanoseconds and no
// request ever contends with the prober.
up atomic.Bool
health string
interval time.Duration
http *http.Client
stop chan struct{}
}
// ErrRemoteUnavailable — the workstation model was required and is not
// answering. Callers on the naming half of the degradation rule turn this into
// a gap in the reply ("не могу сейчас"), never into a guess from the floor.
var ErrRemoteUnavailable = errors.New("llm: workstation model unavailable")
// ErrNoFloor — a Pair was built with no resident model to fall back to. A
// configuration mistake: the floor is the whole point.
var ErrNoFloor = errors.New("llm: no floor client")
// NewPair builds the two-model arrangement. remote may be nil, which is the
// unconfigured deploy and must behave exactly as the box behaves today: every
// call goes to the floor and nothing probes anything.
//
// health is the URL the prober asks. llama-server's /health answers "is a model
// loaded and ready", which is the useful signal here, because llama-server
// refuses to load at all when VRAM is short. That makes a busy card detectable
// without any cooperation from the owner's other jobs.
func NewPair(remote, floor *Client, health string, interval time.Duration) *Pair {
p := &Pair{
remote: remote,
floor: floor,
health: health,
interval: interval,
http: &http.Client{Timeout: probeTimeout},
stop: make(chan struct{}),
}
return p
}
// probeTimeout — a remote that cannot answer /health this fast is not going to
// serve a turn either. Short on purpose: the prober runs on its own goroutine,
// but a slow probe still delays the moment Maven notices the card came back.
const probeTimeout = 2 * time.Second
// Start begins probing. It returns immediately, and the first probe runs before
// the first tick so a remote that is already up is used on the first turn
// rather than after one interval of falling back. Safe to call with a nil
// remote; it does nothing.
func (p *Pair) Start(ctx context.Context) {
if p.remote == nil || p.health == "" {
return
}
go func() {
p.probe(ctx)
t := time.NewTicker(p.interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-p.stop:
return
case <-t.C:
p.probe(ctx)
}
}
}()
}
// Stop ends the prober. Idempotent.
func (p *Pair) Stop() {
select {
case <-p.stop:
default:
close(p.stop)
}
}
// Available reports whether the workstation will take work right now. It reads
// a cached flag, so it is safe to call per turn on the hot path. A false answer
// is never stale in the direction that matters: the worst case is that Maven
// falls back for up to one probe interval after the card frees up.
func (p *Pair) Available() bool {
return p.remote != nil && p.up.Load()
}
func (p *Pair) probe(ctx context.Context) {
ctx, cancel := context.WithTimeout(ctx, probeTimeout)
defer cancel()
req, err := http.NewRequestWithContext(ctx, http.MethodGet, p.health, nil)
if err != nil {
p.set(false)
return
}
resp, err := p.http.Do(req)
if err != nil {
p.set(false)
return
}
defer resp.Body.Close()
p.set(resp.StatusCode == http.StatusOK)
}
// set records the admission answer and logs only the transitions. A machine
// that sleeps every night would otherwise write one line per interval forever.
func (p *Pair) set(up bool) {
if p.up.Swap(up) == up {
return
}
if up {
log.Printf("llm: workstation model available at %s", p.health)
} else {
log.Printf("llm: workstation model unavailable, falling back to the resident model")
}
}
// Complete runs r on the workstation when it will take work, and on the
// resident model otherwise. A remote that fails mid-request falls back too: the
// admission answer is a cache and can be one interval out of date, so an error
// here is expected rather than exceptional.
//
// This is the silent half of the degradation rule. It must be indistinguishable
// from today's behaviour when the workstation is down.
func (p *Pair) Complete(ctx context.Context, r Req) (string, error) {
if p.floor == nil {
return "", ErrNoFloor
}
if p.Available() {
out, err := p.remote.Complete(ctx, r)
if err == nil {
return out, nil
}
// The cached answer was wrong. Correct it now rather than sending the
// next request into the same hole, then fall back.
p.set(false)
}
return p.floor.Complete(ctx, r)
}
// CompleteRemote runs r on the workstation or refuses. It never falls back,
// because for a world question the resident 1.7B does not answer worse, it
// invents. Callers turn ErrRemoteUnavailable into a named gap.
func (p *Pair) CompleteRemote(ctx context.Context, r Req) (string, error) {
if !p.Available() {
return "", ErrRemoteUnavailable
}
out, err := p.remote.Complete(ctx, r)
if err != nil {
p.set(false)
return "", errors.Join(ErrRemoteUnavailable, err)
}
return out, nil
}
+210
View File
@@ -0,0 +1,210 @@
package llm
import (
"context"
"errors"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"time"
)
// completionServer stands in for a llama-server. It counts what reached it, so
// a test can say which of the two models answered.
func completionServer(t *testing.T, reply string, hits *atomic.Int64) *httptest.Server {
t.Helper()
s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
hits.Add(1)
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"choices":[{"message":{"content":"` + reply + `"}}]}`))
}))
t.Cleanup(s.Close)
return s
}
func healthServer(t *testing.T, ok *atomic.Bool) *httptest.Server {
t.Helper()
s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !ok.Load() {
w.WriteHeader(http.StatusServiceUnavailable)
return
}
w.WriteHeader(http.StatusOK)
}))
t.Cleanup(s.Close)
return s
}
// waitFor polls until cond holds or the deadline passes. The prober runs on its
// own goroutine, so a test has to wait for it rather than assume it has run.
func waitFor(t *testing.T, cond func() bool) bool {
t.Helper()
deadline := time.Now().Add(2 * time.Second)
for time.Now().Before(deadline) {
if cond() {
return true
}
time.Sleep(5 * time.Millisecond)
}
return false
}
// The unconfigured deploy. No remote, no probing, every call to the floor —
// exactly what the box does today.
func TestNoRemoteGoesToTheFloor(t *testing.T) {
var floorHits atomic.Int64
floor := completionServer(t, "floor", &floorHits)
p := NewPair(nil, New(floor.URL, time.Second), "", time.Second)
p.Start(context.Background())
defer p.Stop()
if p.Available() {
t.Fatal("a Pair with no remote reports available")
}
out, err := p.Complete(context.Background(), Req{User: "привет"})
if err != nil {
t.Fatalf("complete: %v", err)
}
if out != "floor" || floorHits.Load() != 1 {
t.Fatalf("out = %q, floor hits = %d", out, floorHits.Load())
}
}
// The workstation is up, so it answers and the resident model is not touched.
func TestAvailableRemoteAnswers(t *testing.T) {
var remoteHits, floorHits atomic.Int64
remote := completionServer(t, "remote", &remoteHits)
floor := completionServer(t, "floor", &floorHits)
up := &atomic.Bool{}
up.Store(true)
health := healthServer(t, up)
p := NewPair(New(remote.URL, time.Second), New(floor.URL, time.Second), health.URL, 20*time.Millisecond)
p.Start(context.Background())
defer p.Stop()
if !waitFor(t, p.Available) {
t.Fatal("prober never saw the remote come up")
}
out, err := p.Complete(context.Background(), Req{User: "привет"})
if err != nil {
t.Fatalf("complete: %v", err)
}
if out != "remote" || floorHits.Load() != 0 {
t.Fatalf("out = %q, floor hits = %d", out, floorHits.Load())
}
}
// The card is busy, so /health refuses and Complete degrades silently. This is
// the constraint from 483: the workstation being down is indistinguishable from
// today's behaviour.
func TestBusyCardFallsBackSilently(t *testing.T) {
var remoteHits, floorHits atomic.Int64
remote := completionServer(t, "remote", &remoteHits)
floor := completionServer(t, "floor", &floorHits)
health := healthServer(t, &atomic.Bool{}) // never ok
p := NewPair(New(remote.URL, time.Second), New(floor.URL, time.Second), health.URL, 20*time.Millisecond)
p.Start(context.Background())
defer p.Stop()
time.Sleep(60 * time.Millisecond)
out, err := p.Complete(context.Background(), Req{User: "привет"})
if err != nil {
t.Fatalf("complete: %v", err)
}
if out != "floor" || remoteHits.Load() != 0 {
t.Fatalf("out = %q, remote hits = %d", out, remoteHits.Load())
}
}
// The cached admission answer can be one interval out of date, so a remote that
// dies between probes must still not break the turn.
func TestRemoteErrorMidRequestFallsBack(t *testing.T) {
var floorHits atomic.Int64
dead := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
}))
defer dead.Close()
floor := completionServer(t, "floor", &floorHits)
up := &atomic.Bool{}
up.Store(true)
health := healthServer(t, up)
p := NewPair(New(dead.URL, time.Second), New(floor.URL, time.Second), health.URL, time.Hour)
p.Start(context.Background())
defer p.Stop()
if !waitFor(t, p.Available) {
t.Fatal("prober never saw the remote come up")
}
out, err := p.Complete(context.Background(), Req{User: "привет"})
if err != nil {
t.Fatalf("complete: %v", err)
}
if out != "floor" || floorHits.Load() != 1 {
t.Fatalf("out = %q, floor hits = %d", out, floorHits.Load())
}
// The failed request must have corrected the cached answer, so the next
// one does not walk into the same hole.
if p.Available() {
t.Fatal("a failed remote request left the admission answer up")
}
}
// The naming half of the degradation rule. A world question must not be handed
// to the resident model, because it answers by inventing.
func TestCompleteRemoteNamesTheGap(t *testing.T) {
var floorHits atomic.Int64
floor := completionServer(t, "floor", &floorHits)
health := healthServer(t, &atomic.Bool{}) // never ok
p := NewPair(New("http://127.0.0.1:1", time.Second), New(floor.URL, time.Second), health.URL, 20*time.Millisecond)
p.Start(context.Background())
defer p.Stop()
time.Sleep(60 * time.Millisecond)
if _, err := p.CompleteRemote(context.Background(), Req{User: "почему небо голубое"}); !errors.Is(err, ErrRemoteUnavailable) {
t.Fatalf("err = %v, want ErrRemoteUnavailable", err)
}
if floorHits.Load() != 0 {
t.Fatalf("CompleteRemote fell back to the floor %d times", floorHits.Load())
}
}
// Routing sits on the hot path and must never pay for a health check. Available
// reads a cached flag, so it costs no network at all.
func TestAvailableDoesNotProbe(t *testing.T) {
var probes atomic.Int64
health := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
probes.Add(1)
w.WriteHeader(http.StatusOK)
}))
defer health.Close()
p := NewPair(New("http://127.0.0.1:1", time.Second), New("http://127.0.0.1:1", time.Second), health.URL, time.Hour)
p.Start(context.Background())
defer p.Stop()
if !waitFor(t, p.Available) {
t.Fatal("prober never ran")
}
before := probes.Load()
for range 1000 {
p.Available()
}
if got := probes.Load(); got != before {
t.Fatalf("1000 Available calls made %d probes", got-before)
}
}
// A Pair with no floor is a configuration mistake, and it must say so rather
// than silently having nowhere to degrade to.
func TestNoFloorIsAnError(t *testing.T) {
p := NewPair(nil, nil, "", time.Second)
if _, err := p.Complete(context.Background(), Req{User: "привет"}); !errors.Is(err, ErrNoFloor) {
t.Fatalf("err = %v, want ErrNoFloor", err)
}
}
+126
View File
@@ -0,0 +1,126 @@
package memory
import (
"strings"
"unicode"
)
// stopwords — words that carry no topic. A question and a note that share only
// these share nothing: "почему небо синее" and "сеть какая-то медленная" both
// contain "какая"-shaped filler and are about different worlds.
var stopwords = map[string]bool{
// interrogatives and demonstratives
"что": true, "чего": true, "какой": true, "какая": true, "какое": true,
"какие": true, "каких": true, "кто": true, "кого": true, "кому": true,
"почему": true, "зачем": true, "где": true, "куда": true, "откуда": true,
"когда": true, "сколько": true, "как": true, "то": true, "это": true,
"этот": true, "тот": true, "там": true, "тут": true, "такой": true,
// pronouns — every sentence he says is about him, so "я" is not a topic
"я": true, "меня": true, "мне": true, "мой": true, "моя": true, "мои": true,
"ты": true, "тебя": true, "тебе": true, "твой": true, "он": true, "она": true,
"они": true, "мы": true, "себя": true, "свой": true,
// prepositions, conjunctions, particles, copulas
"в": true, "во": true, "на": true, "с": true, "со": true, "у": true,
"о": true, "об": true, "про": true, "за": true, "из": true, "по": true,
"до": true, "от": true, "для": true, "над": true, "под": true, "при": true,
"и": true, "а": true, "но": true, "или": true, "же": true, "ли": true,
"не": true, "ни": true, "бы": true, "был": true, "была": true, "было": true,
"быть": true, "есть": true, "был-ли": true, "уже": true, "ещё": true,
"еще": true, "так": true, "вот": true, "там-же": true,
// English filler, for the mixed utterances he does say
"the": true, "a": true, "an": true, "is": true, "are": true, "was": true,
"were": true, "be": true, "of": true, "in": true, "on": true, "at": true,
"to": true, "for": true, "about": true, "and": true, "or": true, "not": true,
"what": true, "who": true, "why": true, "when": true, "where": true,
"which": true, "how": true, "i": true, "my": true, "me": true, "it": true,
"this": true, "that": true,
}
// firstPerson — the words that make an utterance a question about his own
// life. Not possession only: "как я восстановил конфиги" owns nothing and is
// still about him.
var firstPerson = map[string]bool{
"я": true, "меня": true, "мне": true, "мной": true, "мой": true,
"моя": true, "моё": true, "мое": true, "мои": true, "моего": true,
"моей": true, "моих": true, "моим": true, "себя": true, "свой": true,
"своя": true, "свои": true, "своего": true, "мною": true,
"i": true, "me": true, "my": true, "mine": true, "myself": true,
}
// RecallAllowed is the second half of the recall gate (#470). A hit that
// cleared the score and margin gate may still be about something else
// entirely: the held-out fixture puts the right note at 0.791-0.890 and the
// must-be-silent cases at 0.795-0.835, so no threshold sits between them, and
// a note about his slow network answered "почему небо синее?".
//
// The veto applies only to a question that mentions nothing of his. That
// restriction is what keeps the fix from costing more than it saves: recall
// exists to find the note whose words he no longer remembers, and demanding a
// shared word of every recall silenced four true recalls on the fixture to
// kill one false one. A question about his own life keeps the embedder alone
// as its judge. A question about the world has to name something the memory
// actually mentions.
func RecallAllowed(query, text string) bool {
if mentionsHim(query) {
return true
}
return SharesContentWord(query, text)
}
func mentionsHim(query string) bool {
for _, w := range strings.FieldsFunc(strings.ToLower(query), func(r rune) bool {
return !unicode.IsLetter(r) && !unicode.IsDigit(r)
}) {
if firstPerson[w] {
return true
}
}
return false
}
// SharesContentWord reports whether query and text have at least one topic
// word in common, after dropping the words that carry no topic. Stems are
// compared, so the note and the question do not have to inflect alike.
func SharesContentWord(query, text string) bool {
q := contentWords(query)
if len(q) == 0 {
// Nothing to compare — a question made entirely of filler. The score
// gate is then the only judge it can have.
return true
}
t := contentWords(text)
for _, a := range q {
for _, b := range t {
if a == b || sameStem(a, b) {
return true
}
}
}
return false
}
func contentWords(s string) []string {
var out []string
for _, w := range strings.FieldsFunc(strings.ToLower(s), func(r rune) bool {
return !unicode.IsLetter(r) && !unicode.IsDigit(r)
}) {
if !stopwords[w] {
out = append(out, w)
}
}
return out
}
// sameStem is inflection and derivation tolerance: Russian marks case and
// tense on the ending, and the note and the question rarely use the same form.
// "воду" and "вода" are the same water, "кормить" and "корм" the same feeding.
// All but the last rune of the shorter word must match, and never fewer than
// three, which is what keeps "сеть" clear of "сеанс".
func sameStem(a, b string) bool {
ar, br := []rune(a), []rune(b)
n := min(len(ar), len(br)) - 1
if n < 3 {
return false
}
return string(ar[:n]) == string(br[:n])
}
+40
View File
@@ -0,0 +1,40 @@
package memory
import "testing"
func TestRecallAllowed(t *testing.T) {
cases := []struct {
name string
query, text string
want bool
}{
// The #470 shape: a world question and a note about his box.
{"world question, unrelated note", "почему небо синее", "сеть какая-то медленная", false},
{"world question, unrelated fact", "какая столица Франции", "какая последняя версия языка Go", false},
{"silent fixture case", "во сколько отходит поезд", "бэкап запускается в три ночи", false},
// A world question that does name the topic keeps its answer.
{"world question, same topic", "какой поезд идёт в Минск", "поезда в Минск ходят утром", true},
// A question about his own life is judged by the embedder alone,
// because recall exists for words he no longer remembers.
{"about him, no shared word", "во сколько я обычно засыпаю", "ложусь около одиннадцати", true},
{"about him, english", "which colour scheme do i like", "тёмная тема везде", true},
// Inflection must not break a match.
{"inflected", "чем кормить кота", "корм для кота в шкафу", true},
}
for _, c := range cases {
if got := RecallAllowed(c.query, c.text); got != c.want {
t.Errorf("%s: RecallAllowed(%q, %q) = %v, want %v", c.name, c.query, c.text, got, c.want)
}
}
}
// A question made only of filler has no topic word to match on, and the score
// gate is then the only judge it can have.
func TestRecallAllowedFallsBackWhenNothingToCompare(t *testing.T) {
if !RecallAllowed("что это", "сеть какая-то медленная") {
t.Error("a question with no content word must not be vetoed")
}
}
+11 -3
View File
@@ -378,7 +378,7 @@ func scoreCase(ctx context.Context, emb router.Embedder, newStore NewStore, minS
if len(hits) > 1 {
o.Margin = hits[0].Score - hits[1].Score
}
o.Recalled = bestRecall(hits, minScore, minMargin)
o.Recalled = bestRecall(c.Query, hits, minScore, minMargin)
}
for i, h := range hits {
if h.ID != c.Want {
@@ -424,11 +424,19 @@ func rankNote(inTop3 bool) string {
// 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 {
// bestRecall mirrors the daemon's gate in cmd/mavend/recall.go, including the
// topic veto added for #470: a score that clears the gate still has to be
// about what he asked. Keep the two in step — a fixture that measures a
// weaker gate than the daemon runs flatters it.
func bestRecall(query string, results []memory.Result, minScore, minMargin float64) string {
if !memory.Confident(results, minScore, minMargin) {
return ""
}
return results[0].Meta["text"]
text := results[0].Meta["text"]
if !memory.RecallAllowed(query, text) {
return ""
}
return text
}
func bump(m map[string]TagStat, key string, pass bool) {
+15 -8
View File
@@ -140,21 +140,22 @@ func words(s string) []string {
// TestBestRecallMatchesDaemon — the harness duplicates bestRecall from
// cmd/mavend/recall.go (package main is not importable). This pins the copy to
// the original's three rules: no hits, below the gate, or no text ⇒ silence.
// the original's rules: no hits, below the gate, no text, or no shared topic
// word ⇒ silence.
func TestBestRecallMatchesDaemon(t *testing.T) {
if got := bestRecall(nil, 0.55, 0); got != "" {
if got := bestRecall("чай", nil, 0.55, 0); got != "" {
t.Errorf("no hits: got %q, want silence", got)
}
low := []memory.Result{{ID: "a", Score: 0.4, Meta: map[string]string{"text": "чай"}}}
if got := bestRecall(low, 0.55, 0); got != "" {
if got := bestRecall("чай", low, 0.55, 0); got != "" {
t.Errorf("below gate: got %q, want silence", got)
}
noText := []memory.Result{{ID: "a", Score: 0.9, Meta: map[string]string{}}}
if got := bestRecall(noText, 0.55, 0); got != "" {
if got := bestRecall("чай", noText, 0.55, 0); got != "" {
t.Errorf("no text: got %q, want silence", got)
}
ok := []memory.Result{{ID: "a", Score: 0.9, Meta: map[string]string{"text": "чай"}}}
if got := bestRecall(ok, 0.55, 0); got != "чай" {
if got := bestRecall("чай", ok, 0.55, 0); got != "чай" {
t.Errorf("above gate: got %q, want %q", got, "чай")
}
// Margin: a close runner-up means the embedder cannot tell the two apart,
@@ -163,17 +164,23 @@ func TestBestRecallMatchesDaemon(t *testing.T) {
{ID: "a", Score: 0.86, Meta: map[string]string{"text": "чай"}},
{ID: "b", Score: 0.85, Meta: map[string]string{"text": "кофе"}},
}
if got := bestRecall(close, 0.55, 0.03); got != "" {
if got := bestRecall("чай", close, 0.55, 0.03); got != "" {
t.Errorf("thin margin: got %q, want silence", got)
}
if got := bestRecall(close, 0.55, 0); got != "чай" {
if got := bestRecall("чай", close, 0.55, 0); got != "чай" {
t.Errorf("margin off: got %q, want %q", got, "чай")
}
// The topic veto (#470): the score is fine and the note is about
// something else.
offTopic := []memory.Result{{ID: "a", Score: 0.9, Meta: map[string]string{"text": "сеть какая-то медленная"}}}
if got := bestRecall("почему небо синее", offTopic, 0.55, 0); got != "" {
t.Errorf("off topic: got %q, want silence", got)
}
clear := []memory.Result{
{ID: "a", Score: 0.86, Meta: map[string]string{"text": "чай"}},
{ID: "b", Score: 0.70, Meta: map[string]string{"text": "кофе"}},
}
if got := bestRecall(clear, 0.55, 0.03); got != "чай" {
if got := bestRecall("чай", clear, 0.55, 0.03); got != "чай" {
t.Errorf("wide margin: got %q, want %q", got, "чай")
}
}
+46 -9
View File
@@ -46,6 +46,12 @@ type LLMPhraser struct {
launch func(ctx context.Context, cfg Config) (backend, error)
probe func(ctx context.Context, base string) (string, error)
// remote — the workstation model, when one is configured. Set once at wiring
// time by UseRemote and read on every phrasing call. nil ⇒ every call goes to
// the resident llama-server this phraser owns, which is the whole deploy
// before a `workstation` block exists. See world.go.
remote Remote
// swapMu — single-flight around Swap. Held for the whole swap, including the
// model load, so two concurrent swap requests can never both be loading.
swapMu sync.Mutex
@@ -363,10 +369,7 @@ func (p *LLMPhraser) PhraseQuery(ctx context.Context, utterance string, notes []
// prompt guaranteed to make a small model fill the gap from memory.
notes = nonEmpty(notes)
if len(notes) == 0 {
// General knowledge — no notes to ground the answer. The system
// prompt is the single tested source in router.KnowledgePrompt.
sys := persona.Prepend(p.cfg.ContextBlock, router.KnowledgePrompt())
prompt := fmt.Sprintf("Пользователь спрашивает: \"%s\".", utterance)
sys, prompt := p.knowledgePrompt(utterance)
resp, err := p.chatWithSystem(ctx, sys, prompt, 768)
if err != nil || resp == "" {
return "не знаю.", nil
@@ -381,11 +384,7 @@ func (p *LLMPhraser) PhraseQuery(ctx context.Context, utterance string, notes []
}
return resp, nil
}
sys := p.querySystemPrompt()
prompt := fmt.Sprintf(
"Он спрашивает: \"%s\"\n\nИсточники:\n%s\nОтветь ему коротко и своими словами, опираясь только на эти источники. Если ответа в них нет — так и скажи.",
utterance, evidenceBlock(notes),
)
sys, prompt := p.evidencePrompt(utterance, notes)
resp, err := p.chatWithSystem(ctx, sys, prompt, 768)
text, _, perr := parseResponseMood(resp)
if err != nil || perr != nil {
@@ -481,6 +480,17 @@ func chatSystemPrompt(block func() string) string {
// the LLM completion endpoint. Like chatWithSystem but for an arbitrary message
// slice — the caller owns the system prompt placement.
func (p *LLMPhraser) chatWithMessages(ctx context.Context, msgs []chatMsg, maxTokens int) (string, error) {
// Same silent preference as chatWithSystem, when the array is the shape
// llm.Req can carry: one system turn and one user turn. PhraseChat already
// folds the history into a single user message (some chat templates reject
// consecutive user turns), so today that is every call. A longer array goes
// to the resident model rather than get flattened here, because flattening a
// conversation is a decision its owner should make.
if len(msgs) == 2 && msgs[0].Role == "system" && msgs[1].Role == "user" {
if out, ok := p.remoteChat(ctx, msgs[0].Content, msgs[1].Content, maxTokens); ok {
return out, nil
}
}
base, release, err := p.acquire()
if err != nil {
return "", err
@@ -640,6 +650,12 @@ func (p *LLMPhraser) chat(ctx context.Context, userPrompt string) (string, error
}
func (p *LLMPhraser) chatWithSystem(ctx context.Context, system, user string, maxTokens int) (string, error) {
// The workstation model first when it will take work, and silently: every
// caller of this helper is on the silent half of the degradation rule. It
// answering is not news, and it being asleep is not news either.
if out, ok := p.remoteChat(ctx, system, user, maxTokens); ok {
return out, nil
}
base, release, err := p.acquire()
if err != nil {
return "", err
@@ -733,6 +749,27 @@ func (p *LLMPhraser) systemPrompt() string {
return persona.Prepend(p.cfg.ContextBlock, nudgeSystem)
}
// knowledgePrompt — the no-sources branch: a world question, answered from
// weights alone. The system prompt is the single tested source in
// router.KnowledgePrompt.
//
// Split out of PhraseQuery so PhraseWorld sends the workstation model the same
// bytes the resident model gets. Prompt parity across two models is a stated
// constraint (CLAUDE.md), and two copies of a prompt is how it stops holding.
func (p *LLMPhraser) knowledgePrompt(utterance string) (sys, user string) {
return persona.Prepend(p.cfg.ContextBlock, router.KnowledgePrompt()),
fmt.Sprintf("Пользователь спрашивает: \"%s\".", utterance)
}
// evidencePrompt — the sources branch: read these, add nothing. Shared with
// PhraseWorld for the same reason as knowledgePrompt.
func (p *LLMPhraser) evidencePrompt(utterance string, notes []string) (sys, user string) {
return p.querySystemPrompt(), fmt.Sprintf(
"Он спрашивает: \"%s\"\n\nИсточники:\n%s\nОтветь ему коротко и своими словами, опираясь только на эти источники. Если ответа в них нет — так и скажи.",
utterance, evidenceBlock(notes),
)
}
// querySystemPrompt returns the system prompt for the evidence branch of
// PhraseQuery. Prepends the configured persona when set.
//
+137
View File
@@ -0,0 +1,137 @@
package phraser
import (
"context"
"errors"
"log"
"github.com/kami/maven/internal/llm"
)
// Remote — the workstation model, seen from the phraser. `*llm.Pair` satisfies
// it, and a test fake satisfies it in three lines.
//
// Only the refusing half of Pair is here on purpose. Pair.Complete falls back to
// its own floor client, and the phraser already owns a floor: the llama-server it
// spawned. Two floors under one call is one too many, so the phraser asks whether
// the remote will take work, uses it when it will, and otherwise does exactly
// what it did before this file existed.
type Remote interface {
// Available is an atomic read of a cached probe, so it is free to call per
// turn. See llm.Pair.
Available() bool
// CompleteRemote runs on the workstation or returns ErrRemoteUnavailable. It
// never falls back.
CompleteRemote(ctx context.Context, r llm.Req) (string, error)
}
// ErrNoWorldModel — a world question was asked, a workstation model is
// configured to answer it, and that machine is not answering. The caller turns
// this into a gap he is told about ("не могу сейчас"), never into an answer from
// the resident model.
//
// This is the naming half of the degradation rule in docs/offload.md. The
// resident Qwen3-1.7B does not answer a world question worse than the 12B, it
// invents: measured, the workstation model scores knowledge 9/9 on the talk
// fixture against the resident model's confabulations
// (docs/evals/2026-08-02-workstation-gemma4-12b.md).
var ErrNoWorldModel = errors.New("phraser: no world model available")
// chatTemperature — what the phraser's own transport has always sampled at.
// Named so the remote path cannot drift from it silently. Whether 0.7 is right
// at all is Vikunja #402, and answering that here would hide a phrasing change
// inside a routing change.
const chatTemperature = 0.7
// UseRemote points the phraser at the workstation model. Wiring time only, once,
// before anything phrases: the field is read without a lock on every call
// because a per-turn lock to answer a question that changes at deploy time is
// not worth paying for.
//
// A nil remote is the normal state of a box with no `workstation` block, and it
// must behave exactly as the box behaved before this seam existed.
func (p *LLMPhraser) UseRemote(r Remote) {
p.remote = r
}
// PhraseWorld answers a question about the world — either from the model's own
// knowledge (no sources) or from a passage someone fetched (a live search, a ZIM
// article, a page he named). Three outcomes, and the middle one is the point:
//
// - No workstation configured. The resident model answers, exactly as it does
// today. Naming a gap needs a gap: on a box that never had a second model,
// refusing every world question would remove a capability he has now.
// - Workstation configured and taking work. It answers.
// - Workstation configured and down. ErrNoWorldModel, and the caller says so.
//
// The prompts are the ones PhraseQuery uses, built by the same two functions, so
// the two models are asked the same question in the same words.
func (p *LLMPhraser) PhraseWorld(ctx context.Context, utterance string, sources []string) (string, error) {
sources = nonEmpty(sources)
if p.remote == nil {
return p.PhraseQuery(ctx, utterance, sources)
}
var sys, user string
if len(sources) == 0 {
sys, user = p.knowledgePrompt(utterance)
} else {
sys, user = p.evidencePrompt(utterance, sources)
}
if !p.remote.Available() {
return "", ErrNoWorldModel
}
resp, err := p.remote.CompleteRemote(ctx, llm.Req{
System: sys,
User: user,
Grammar: p.grammar(),
MaxTokens: 768,
Temperature: chatTemperature,
})
if err != nil {
// The cached probe was one interval stale, or the card went away
// mid-request. Either way this is the gap, not an error to log and
// paper over with the smaller model.
log.Printf("phraser: world model: %v", err)
return "", errors.Join(ErrNoWorldModel, err)
}
resp = stripThink(resp)
text, _, perr := parseResponseMood(resp)
if perr != nil {
log.Printf("phraser: PhraseWorld: %v", perr)
return "", errors.Join(ErrNoWorldModel, perr)
}
if text != "" {
return text, nil
}
if resp == "" {
return "", ErrNoWorldModel
}
return resp, nil
}
// remoteChat is the silent half, for the phrasing paths where the workstation
// model is only better: a nudge, a reminder, a reply, a question answered from
// his own notes. It reports whether it answered; it never reports why not,
// because the caller's next move is the resident model either way.
//
// He is not told which of the two models phrased his reply. That is the rule.
func (p *LLMPhraser) remoteChat(ctx context.Context, system, user string, maxTokens int) (string, bool) {
if p.remote == nil || !p.remote.Available() {
return "", false
}
out, err := p.remote.CompleteRemote(ctx, llm.Req{
System: system,
User: user,
Grammar: p.grammar(),
MaxTokens: maxTokens,
Temperature: chatTemperature,
})
if err != nil {
log.Printf("phraser: workstation model declined, phrasing here instead: %v", err)
return "", false
}
if out = stripThink(out); out == "" {
return "", false
}
return out, true
}
+164
View File
@@ -0,0 +1,164 @@
package phraser
import (
"context"
"errors"
"strings"
"testing"
"github.com/kami/maven/internal/llm"
"github.com/kami/maven/internal/loop"
)
// fakeRemote — a workstation model that is up or down on command, and records
// what it was asked.
type fakeRemote struct {
up bool
reply string
err error
got []llm.Req
}
func (f *fakeRemote) Available() bool { return f.up }
func (f *fakeRemote) CompleteRemote(_ context.Context, r llm.Req) (string, error) {
f.got = append(f.got, r)
if f.err != nil {
return "", f.err
}
return f.reply, nil
}
// The three outcomes of the naming half, in one place. The middle one is the
// whole task: a gap he is told about, not an answer from the smaller model.
func TestPhraseWorldNamesTheGapOnlyWhenThereIsOne(t *testing.T) {
answer := `{"response": "Небо голубое из-за рэлеевского рассеяния.", "mood": "neutral"}`
t.Run("no workstation configured: the resident model answers as today", func(t *testing.T) {
spy := newPromptSpy(t)
p := NewLLMPhraserAt(spy.srv.URL, Config{})
got, err := p.PhraseWorld(context.Background(), "почему небо голубое", nil)
if err != nil {
t.Fatalf("PhraseWorld: %v", err)
}
if got == "" {
t.Fatal("no reply from the resident model")
}
if len(spy.user) != 1 {
t.Fatalf("resident model saw %d requests, want 1", len(spy.user))
}
})
t.Run("workstation up: it answers and the resident model is not asked", func(t *testing.T) {
spy := newPromptSpy(t)
p := NewLLMPhraserAt(spy.srv.URL, Config{})
remote := &fakeRemote{up: true, reply: answer}
p.UseRemote(remote)
got, err := p.PhraseWorld(context.Background(), "почему небо голубое", nil)
if err != nil {
t.Fatalf("PhraseWorld: %v", err)
}
if !strings.Contains(got, "рассеяния") {
t.Errorf("reply is not the workstation's: %q", got)
}
if len(spy.user) != 0 {
t.Errorf("the resident model was asked %d times, want 0", len(spy.user))
}
})
t.Run("workstation down: the gap, and nothing invented", func(t *testing.T) {
spy := newPromptSpy(t)
p := NewLLMPhraserAt(spy.srv.URL, Config{})
p.UseRemote(&fakeRemote{up: false})
got, err := p.PhraseWorld(context.Background(), "почему небо голубое", nil)
if !errors.Is(err, ErrNoWorldModel) {
t.Fatalf("err = %v, want ErrNoWorldModel", err)
}
if got != "" {
t.Errorf("got a reply %q with no world model", got)
}
if len(spy.user) != 0 {
t.Errorf("the resident model answered a world question %d times, want 0", len(spy.user))
}
})
t.Run("workstation errors mid-request: still the gap", func(t *testing.T) {
spy := newPromptSpy(t)
p := NewLLMPhraserAt(spy.srv.URL, Config{})
p.UseRemote(&fakeRemote{up: true, err: errors.New("connection refused")})
if _, err := p.PhraseWorld(context.Background(), "почему небо голубое", nil); !errors.Is(err, ErrNoWorldModel) {
t.Fatalf("err = %v, want ErrNoWorldModel", err)
}
if len(spy.user) != 0 {
t.Errorf("the resident model answered a world question %d times, want 0", len(spy.user))
}
})
}
// Prompt parity: the workstation model is asked the same question in the same
// words, or the fixtures measure one thing and the daemon ships another.
func TestPhraseWorldSendsTheSamePromptsAsPhraseQuery(t *testing.T) {
spy := newPromptSpy(t)
resident := NewLLMPhraserAt(spy.srv.URL, Config{})
if _, err := resident.PhraseQuery(context.Background(), "кто написал войну и мир", []string{"Лев Толстой"}); err != nil {
t.Fatal(err)
}
remote := &fakeRemote{up: true, reply: `{"response": "Толстой.", "mood": "neutral"}`}
offloaded := NewLLMPhraserAt(spy.srv.URL, Config{})
offloaded.UseRemote(remote)
if _, err := offloaded.PhraseWorld(context.Background(), "кто написал войну и мир", []string{"Лев Толстой"}); err != nil {
t.Fatal(err)
}
if len(remote.got) != 1 {
t.Fatalf("the workstation saw %d requests, want 1", len(remote.got))
}
if remote.got[0].System != spy.system[0] {
t.Errorf("system prompts differ:\nremote: %q\nresident: %q", remote.got[0].System, spy.system[0])
}
if remote.got[0].User != spy.user[0] {
t.Errorf("user prompts differ:\nremote: %q\nresident: %q", remote.got[0].User, spy.user[0])
}
}
// The silent half. A nudge phrased on the workstation is not news, and one
// phrased here because the card is busy is not news either — but it must be
// sampled the same way, or the workstation quietly changes how she sounds.
func TestNudgePhrasingPrefersTheWorkstationSilently(t *testing.T) {
spy := newPromptSpy(t)
p := NewLLMPhraserAt(spy.srv.URL, Config{LLMNudges: true})
remote := &fakeRemote{up: true, reply: `{"response": "Выпей воды.", "mood": "neutral"}`}
p.UseRemote(remote)
pn, err := p.PhraseNudge(context.Background(), loop.Candidate{Rule: loop.WaterRule(), Severity: loop.Sev1})
if err != nil {
t.Fatalf("PhraseNudge: %v", err)
}
if pn.Body != "Выпей воды." {
t.Errorf("body = %q, want the workstation's wording", pn.Body)
}
if len(remote.got) != 1 {
t.Fatalf("the workstation saw %d requests, want 1", len(remote.got))
}
if remote.got[0].Temperature != chatTemperature {
t.Errorf("temperature = %v, want %v (what the resident transport samples at)",
remote.got[0].Temperature, chatTemperature)
}
if len(spy.user) != 0 {
t.Errorf("the resident model phrased %d nudges, want 0", len(spy.user))
}
}
func TestNudgePhrasingFallsBackWhenTheCardIsBusy(t *testing.T) {
spy := newPromptSpy(t)
p := NewLLMPhraserAt(spy.srv.URL, Config{LLMNudges: true})
p.UseRemote(&fakeRemote{up: false})
if _, err := p.PhraseNudge(context.Background(), loop.Candidate{Rule: loop.WaterRule(), Severity: loop.Sev1}); err != nil {
t.Fatalf("PhraseNudge: %v", err)
}
if len(spy.user) != 1 {
t.Fatalf("the resident model phrased %d nudges, want 1", len(spy.user))
}
}
+64
View File
@@ -0,0 +1,64 @@
package router
import "strings"
// interrogatives — the question words that mark an utterance as asking rather
// than telling. Tokenized, never substring: "что" inside "чтобы" and "как"
// inside "какао" are not questions.
var interrogatives = []string{
"что", "чего", "какой", "какая", "какое", "какие", "каких",
"кто", "кого", "кому", "чей", "почему", "зачем", "отчего",
"где", "куда", "откуда", "когда", "сколько", "как",
"what", "who", "whom", "why", "when", "where", "which", "how",
}
// narrativeRequests — "tell me about X" asks for knowledge Maven does not
// hold about him. It carries no question mark and no interrogative, which is
// how "расскажи про битву при Ватерлоо" reached the fact store (#470).
var narrativeRequests = []string{
"расскажи", "объясни", "опиши", "перечисли",
"tell", "explain", "describe",
}
// captureVerbs — an explicit instruction to record something. These win over
// every test below, because "запиши что я пил воду" contains an interrogative
// and is still a capture: the word he said is "запиши".
var captureVerbs = []string{
"запиши", "запомни", "отметь", "заметь", "добавь", "сохрани",
"note", "remember", "log", "save",
}
// IsQuestionShaped reports whether text asks for something rather than
// records it. It is a deterministic offline test over tokens, so it costs
// nothing and never depends on the model that produced the routing decision.
//
// It exists because a mis-routed question used to be persisted as a fact
// about the owner, with the model's invented answer as the value (#470). The
// predicate is deliberately blunt: refusing to store a question is cheap and
// reversible, storing an invented fact about him is neither.
func IsQuestionShaped(text string) bool {
t := strings.TrimSpace(text)
if t == "" {
return false
}
toks := planTokens(t)
for _, v := range captureVerbs {
if hasTok(toks, v) {
return false
}
}
if strings.HasSuffix(t, "?") {
return true
}
for _, w := range interrogatives {
if hasTok(toks, w) {
return true
}
}
for _, w := range narrativeRequests {
if hasTok(toks, w) {
return true
}
}
return false
}
+48
View File
@@ -0,0 +1,48 @@
package router
import "testing"
func TestIsQuestionShaped(t *testing.T) {
// The seven utterances #470 recorded, plus the captures that must keep
// working. A capture misread as a question loses a fact; a question
// misread as a capture poisons recall, so the captures are the ones worth
// pinning here.
cases := []struct {
text string
want bool
}{
{"какая последняя версия языка Go?", true},
{"что дальше?", true},
{"расскажи про битву при Ватерлоо", true},
{"почему небо синее?", true},
{"какая столица Австралии?", true},
{"кто такой Никола Тесла?", true},
{"сколько стоит доллар", true},
{"who is the premier of Japan", true},
{"объясни линии Фраунгофера", true},
{"запиши что я пил воду", false},
{"запомни какая у меня машина", false},
{"отметь что я поужинал", false},
{"поужинал", false},
{"я выпил кофе", false},
{"вода", false},
{"привет", false},
{"", false},
}
for _, c := range cases {
if got := IsQuestionShaped(c.text); got != c.want {
t.Errorf("IsQuestionShaped(%q) = %v, want %v", c.text, got, c.want)
}
}
}
// Substring matching is what made the day-plan predicates wrong before, and
// this predicate gates a write, so it gets the same guard.
func TestIsQuestionShapedIsTokenized(t *testing.T) {
for _, text := range []string{"чтобы не забыть, я полил кактус", "какао выпил"} {
if IsQuestionShaped(text) {
t.Errorf("IsQuestionShaped(%q) = true; a question word inside a longer word is not a question", text)
}
}
}
+62
View File
@@ -6,6 +6,7 @@ import (
"encoding/json"
"errors"
"fmt"
"log"
"strings"
"time"
@@ -46,6 +47,39 @@ func (s *Store) WriteFact(ctx context.Context, ts time.Time, kind FactKind, key,
return id, nil
}
// FactRecallText is the text a fact is indexed under and read back as (#493).
//
// It used to be the utterance that wrote the fact, so recall of ANY
// voice-tapped fact answered with the sentence he said instead of the value
// stored: `go_version = 1.20` was indexed as "какая последняя версия языка
// Go?", and that question is what came back. The poisoned rows made the defect
// visible; the shape was wrong for legitimate facts too.
//
// The key is spoken with its underscores dropped, because a key is written for
// the store and this string is read out loud.
func FactRecallText(key, value string) string {
spoken := strings.TrimSpace(strings.ReplaceAll(key, "_", " "))
v := strings.TrimSpace(DecodeFactValue(value))
switch {
case v == "":
return spoken
case spoken == "":
return v
}
return spoken + " — " + v
}
// DecodeFactValue unwraps a stored value for reading. The column holds raw json
// when the writer serialized one (SetValue, CorrectValue) and a plain string
// when it did not (a voice tap), so a reader that wants the text handles both.
func DecodeFactValue(value string) string {
var s string
if err := json.Unmarshal([]byte(value), &s); err == nil {
return s
}
return value
}
// LatestFact returns the latest non-voided fact for key, or ErrNoFact.
// "Non-voided" = no later row has voids_id pointing at it. We resolve this by
// taking the newest row whose id is not referenced by any voids_id.
@@ -290,6 +324,20 @@ func (s *Store) CorrectValue(ctx context.Context, key, source string, value any,
if err != nil {
return 0, fmt.Errorf("last insert id: %w", err)
}
// The same repair a void needs, for the same reason (#493). A correction
// supersedes the value, and the vector still holds the old one, so recall
// kept answering with the value he had just corrected. Dropping it costs
// the key its recall vector until the fact is tapped again: this layer has
// no embedder, and a missing vector loses a question while a stale one
// answers it wrongly.
//
// Best-effort: the corrected row is committed, and a correction that lands
// beats one that fails on cleanup.
if n, derr := s.VectorMemory().DeletePrefix(ctx, "fact:"+key+":"); derr != nil {
log.Printf("store: correct %q: memory vectors survive: %v", key, derr)
} else if n > 0 {
log.Printf("store: correct %q: dropped %d superseded memory vector(s)", key, n)
}
return newID, nil
}
@@ -332,6 +380,20 @@ func (s *Store) VoidLatestFact(ctx context.Context, key, source string, ts time.
if err != nil {
return 0, 0, fmt.Errorf("void: last insert id: %w", err)
}
// The other half of the repair (#470). A fact reaches recall through a
// vector keyed `fact:<key>:<unix>`, holding the utterance that wrote it.
// Voiding the row alone left that vector answering questions, so revert
// reported success on a box that stayed broken. Deleting every vector for
// the key covers the earlier rows too: their values are superseded, and a
// superseded value has no business claiming a turn.
//
// Best-effort by design: the audit trail is already committed, and a fact
// that is voided but still recallable is better than a void that failed.
if n, derr := s.VectorMemory().DeletePrefix(ctx, "fact:"+key+":"); derr != nil {
log.Printf("store: void %q: memory vectors survive: %v", key, derr)
} else if n > 0 {
log.Printf("store: void %q: dropped %d memory vector(s)", key, n)
}
return oldID, newID, nil
}
+177
View File
@@ -0,0 +1,177 @@
package store
import (
"context"
"encoding/json"
"errors"
"fmt"
"strconv"
"strings"
"time"
)
// metaKeyFactVectorShape names the shape the stored fact vectors were written
// in. It exists so the repair below runs once per box instead of on every
// start: the rows it fixes were written by a code path that no longer exists,
// and once fixed nothing writes that shape again.
const metaKeyFactVectorShape = "fact_vector_shape"
// factVectorShapeFact is the shape FactRecallText produces. Anything else in
// the marker (including nothing, which is every box written before #493) means
// the fact vectors still hold utterances.
const factVectorShapeFact = "fact-text (#493)"
// FactVectorRepair is what one repair run did, for logging.
type FactVectorRepair struct {
Skipped bool // marker already matched — nothing to do
Rewritten int // rows re-embedded from the fact they name
Dropped int // rows deleted: voided, superseded, or naming no fact at all
Kept int // rows already holding the right text
Took time.Duration
}
// RepairFactVectors brings the fact rows of memory_vectors in line with the
// facts they name, and is the operator recovery a poisoned box had no path to
// (#470 point 4, #493).
//
// Three defects put wrong text in that index, and all three are write-path
// fixes that do nothing for rows already stored:
//
// - the indexed text was the utterance, so every fact row reads back a
// sentence rather than a value;
// - a void left its vector behind, so retracted junk kept answering;
// - a correction left its vector behind, so the superseded value did.
//
// So each fact row is resolved against the fact store and one of three things
// happens. It is dropped when the key has no fact, when the newest row for the
// key is a void marker, or when a newer vector for the same key exists — a
// superseded value has no business claiming a turn. It is re-embedded when its
// text is not what FactRecallText says the fact is. Otherwise it is left alone.
//
// Idempotent, and safe to interrupt: every step compares before writing and the
// marker is written last, so a run that dies partway is simply redone.
func (s *Store) RepairFactVectors(ctx context.Context, embed EmbedFunc) (FactVectorRepair, error) {
start := time.Now()
var res FactVectorRepair
shape, err := s.Meta(ctx, metaKeyFactVectorShape)
if err != nil {
return res, err
}
if shape == factVectorShapeFact {
res.Skipped = true
res.Took = time.Since(start)
return res, nil
}
rows, err := s.db.QueryContext(ctx, `SELECT id, meta FROM memory_vectors`)
if err != nil {
return res, fmt.Errorf("repair fact vectors: read: %w", err)
}
type factVec struct {
id, key string
meta map[string]string
ts int64
}
var vecs []factVec
newest := map[string]int64{} // key → newest ts seen for it
for rows.Next() {
var id, metaJSON string
if err := rows.Scan(&id, &metaJSON); err != nil {
rows.Close()
return res, fmt.Errorf("repair fact vectors: row: %w", err)
}
meta := map[string]string{}
if err := json.Unmarshal([]byte(metaJSON), &meta); err != nil {
rows.Close()
return res, fmt.Errorf("repair fact vectors: meta for %q: %w", id, err)
}
if meta["type"] != "fact" {
continue
}
key, ts, ok := splitFactVectorID(id)
if !ok {
continue
}
vecs = append(vecs, factVec{id: id, key: key, meta: meta, ts: ts})
if ts > newest[key] {
newest[key] = ts
}
}
rows.Close()
if err := rows.Err(); err != nil {
return res, fmt.Errorf("repair fact vectors: rows: %w", err)
}
for _, v := range vecs {
drop := v.ts < newest[v.key]
var want string
if !drop {
f, ferr := s.LatestFact(ctx, v.key)
switch {
case errors.Is(ferr, ErrNoFact):
drop = true
case ferr != nil:
return res, fmt.Errorf("repair fact vectors: fact %q: %w", v.key, ferr)
case DecodeFactValue(f.Value) == "voided":
drop = true
default:
want = FactRecallText(v.key, f.Value)
}
}
if drop {
if err := s.VectorMemory().Delete(ctx, v.id); err != nil {
return res, err
}
res.Dropped++
continue
}
if v.meta["text"] == want {
res.Kept++
continue
}
vec, err := embed(ctx, want)
if err != nil {
return res, fmt.Errorf("repair fact vectors: embed %q: %w", v.id, err)
}
// The whole meta blob is rewritten in Go rather than patched in SQL,
// because json_set needs the JSON1 extension and this store is opened
// through sqlcipher.
v.meta["text"] = want
metaJSON, err := json.Marshal(v.meta)
if err != nil {
return res, fmt.Errorf("repair fact vectors: meta %q: %w", v.id, err)
}
if _, err := s.db.ExecContext(ctx,
`UPDATE memory_vectors SET vec = ?, meta = ? WHERE id = ?`,
encodeVec(vec), string(metaJSON), v.id); err != nil {
return res, fmt.Errorf("repair fact vectors: write %q: %w", v.id, err)
}
res.Rewritten++
}
if err := s.SetMeta(ctx, metaKeyFactVectorShape, factVectorShapeFact); err != nil {
return res, err
}
res.Took = time.Since(start)
return res, nil
}
// splitFactVectorID reads the key and write time back out of a fact vector's
// id, which the write path builds as `fact:<key>:<unix>`. A key may hold a
// colon, the timestamp may not, so the split is from the right.
func splitFactVectorID(id string) (key string, ts int64, ok bool) {
rest, found := strings.CutPrefix(id, "fact:")
if !found {
return "", 0, false
}
cut := strings.LastIndex(rest, ":")
if cut <= 0 {
return "", 0, false
}
ts, err := strconv.ParseInt(rest[cut+1:], 10, 64)
if err != nil {
return "", 0, false
}
return rest[:cut], ts, true
}
+151
View File
@@ -0,0 +1,151 @@
package store
import (
"context"
"database/sql"
"testing"
"time"
)
// The write-path half of #493: recall of a fact must read back the fact, not
// the sentence he happened to say.
func TestFactRecallText(t *testing.T) {
for _, tc := range []struct {
name, key, value, want string
}{
{"json value", "go_version", `"1.20"`, "go version — 1.20"},
{"plain value", "water", "выпил", "water — выпил"},
{"no value", "shower", "", "shower"},
{"underscores are spoken as spaces", "espresso_machine", `"чистая"`, "espresso machine — чистая"},
} {
t.Run(tc.name, func(t *testing.T) {
if got := FactRecallText(tc.key, tc.value); got != tc.want {
t.Fatalf("FactRecallText(%q, %q) = %q; want %q", tc.key, tc.value, got, tc.want)
}
})
}
}
// A correction left the superseded value in the index, so recall answered with
// the value he had just corrected (#493).
func TestCorrectValueDropsMemoryVectors(t *testing.T) {
ctx := context.Background()
s := newTestStore(t)
now := time.Now()
mem := s.VectorMemory()
if _, err := s.WriteFact(ctx, now, KindSelf, "go_version", `"1.20"`, "tap:voice", 1.0, sql.NullInt64{}); err != nil {
t.Fatalf("WriteFact: %v", err)
}
if err := mem.Insert(ctx, "fact:go_version:1", []float32{1, 0, 0}, map[string]string{
"type": "fact", "text": "go version — 1.20",
}); err != nil {
t.Fatalf("Insert: %v", err)
}
if _, err := s.CorrectValue(ctx, "go_version", "feedback", "1.25", now.Add(time.Minute)); err != nil {
t.Fatalf("CorrectValue: %v", err)
}
got, err := mem.ByPrefix(ctx, "fact:")
if err != nil {
t.Fatalf("ByPrefix: %v", err)
}
if len(got) != 0 {
t.Fatalf("after the correction the index still holds %+v; the superseded value must not answer", got)
}
}
// The recovery path a poisoned box had none of (#470 point 4, #493): rows
// written before the fix hold utterances, voided junk and superseded values,
// and no write-path change reaches any of them.
func TestRepairFactVectors(t *testing.T) {
ctx := context.Background()
s := newTestStore(t)
now := time.Now()
mem := s.VectorMemory()
embed := func(ctx context.Context, text string) ([]float32, error) {
return []float32{float32(len(text)), 1, 0}, nil
}
// A live fact indexed under the question that wrote it — the defect.
if _, err := s.WriteFact(ctx, now, KindSelf, "water", `"выпил"`, "tap:voice", 1.0, sql.NullInt64{}); err != nil {
t.Fatalf("WriteFact water: %v", err)
}
if err := mem.Insert(ctx, "fact:water:100", []float32{9, 9, 9}, map[string]string{
"type": "fact", "source": "voice", "text": "запиши что я пил воду",
}); err != nil {
t.Fatalf("Insert water: %v", err)
}
// A voided fact whose vector survived the void.
if _, err := s.WriteFact(ctx, now, KindSelf, "go_version", `"1.20"`, "tap:voice", 1.0, sql.NullInt64{}); err != nil {
t.Fatalf("WriteFact go_version: %v", err)
}
if _, _, err := s.VoidLatestFact(ctx, "go_version", "feedback", now.Add(time.Minute)); err != nil {
t.Fatalf("VoidLatestFact: %v", err)
}
if err := mem.Insert(ctx, "fact:go_version:100", []float32{9, 9, 9}, map[string]string{
"type": "fact", "text": "какая последняя версия языка Go?",
}); err != nil {
t.Fatalf("Insert go_version: %v", err)
}
// A key with two vectors: only the newest may answer.
if _, err := s.WriteFact(ctx, now, KindSelf, "mood", `"устал"`, "tap:voice", 1.0, sql.NullInt64{}); err != nil {
t.Fatalf("WriteFact mood: %v", err)
}
for _, ts := range []string{"100", "200"} {
if err := mem.Insert(ctx, "fact:mood:"+ts, []float32{9, 9, 9}, map[string]string{
"type": "fact", "text": "мне грустно",
}); err != nil {
t.Fatalf("Insert mood %s: %v", ts, err)
}
}
// A note must be left entirely alone.
if err := mem.Insert(ctx, "note:7", []float32{5, 5, 5}, map[string]string{
"type": "note", "text": "сеть тормозит по вечерам",
}); err != nil {
t.Fatalf("Insert note: %v", err)
}
res, err := s.RepairFactVectors(ctx, embed)
if err != nil {
t.Fatalf("RepairFactVectors: %v", err)
}
if res.Rewritten != 2 || res.Dropped != 2 {
t.Fatalf("repair reported %+v; want 2 rewritten (water, newest mood) and 2 dropped (voided go_version, superseded mood)", res)
}
got, err := mem.ByPrefix(ctx, "fact:")
if err != nil {
t.Fatalf("ByPrefix: %v", err)
}
texts := map[string]string{}
for _, r := range got {
texts[r.ID] = r.Meta["text"]
}
if len(texts) != 2 {
t.Fatalf("the index holds %+v; want only fact:water:100 and fact:mood:200", texts)
}
if texts["fact:water:100"] != "water — выпил" {
t.Fatalf("water reads back %q; want the fact, not the utterance", texts["fact:water:100"])
}
if texts["fact:mood:200"] != "mood — устал" {
t.Fatalf("mood reads back %q", texts["fact:mood:200"])
}
// Provenance the row already carried must survive the rewrite.
for _, r := range got {
if r.ID == "fact:water:100" && r.Meta["source"] != "voice" {
t.Fatalf("water lost its source meta: %+v", r.Meta)
}
}
if notes, err := mem.ByPrefix(ctx, "note:"); err != nil || len(notes) != 1 {
t.Fatalf("the note row was touched: %+v (err %v)", notes, err)
}
// Marker written, so a second run is free and changes nothing.
again, err := s.RepairFactVectors(ctx, embed)
if err != nil {
t.Fatalf("second RepairFactVectors: %v", err)
}
if !again.Skipped {
t.Fatalf("second run did work: %+v; the marker must make it a no-op", again)
}
}
+21
View File
@@ -148,6 +148,27 @@ func (m *MemoryStore) Delete(ctx context.Context, id string) error {
return nil
}
// DeletePrefix removes every vector whose id starts with prefix and returns
// how many went. Same escaping as ByPrefix, so a key containing % or _ cannot
// widen the delete.
//
// It exists for the repair half of a revert (#470). Voiding a fact row left
// its vector in the index, so recall kept serving the voided fact's utterance
// and the documented repair did not repair.
func (m *MemoryStore) DeletePrefix(ctx context.Context, prefix string) (int64, error) {
pattern := escapeLike(prefix) + "%"
res, err := m.db.ExecContext(ctx,
`DELETE FROM memory_vectors WHERE id LIKE ? ESCAPE '\'`, pattern)
if err != nil {
return 0, fmt.Errorf("memory: delete prefix %q: %w", prefix, err)
}
n, err := res.RowsAffected()
if err != nil {
return 0, fmt.Errorf("memory: delete prefix %q: rows affected: %w", prefix, err)
}
return n, nil
}
// escapeLike neutralises the LIKE wildcards in a literal prefix.
func escapeLike(s string) string {
r := strings.NewReplacer(`\`, `\\`, `%`, `\%`, `_`, `\_`)
+64
View File
@@ -0,0 +1,64 @@
package store
import (
"context"
"database/sql"
"testing"
"time"
)
// Stage 3 of #470: reverting a fact reported success and left the vector that
// was answering questions, so the documented repair did not repair.
func TestVoidLatestFactDropsMemoryVectors(t *testing.T) {
ctx := context.Background()
s := newTestStore(t)
now := time.Now()
mem := s.VectorMemory()
if _, err := s.WriteFact(ctx, now, KindSelf, "go_version", `"1.20"`, "tap:voice", 1.0, sql.NullInt64{}); err != nil {
t.Fatalf("WriteFact: %v", err)
}
// The id shape actionFact writes: fact:<key>:<unix>.
if err := mem.Insert(ctx, "fact:go_version:1", []float32{1, 0, 0}, map[string]string{
"type": "fact", "text": "какая последняя версия языка Go?",
}); err != nil {
t.Fatalf("Insert: %v", err)
}
// A vector for another key must survive the void.
if err := mem.Insert(ctx, "fact:water:1", []float32{0, 1, 0}, map[string]string{
"type": "fact", "text": "запиши что я пил воду",
}); err != nil {
t.Fatalf("Insert: %v", err)
}
if _, _, err := s.VoidLatestFact(ctx, "go_version", "feedback", now.Add(time.Minute)); err != nil {
t.Fatalf("VoidLatestFact: %v", err)
}
got, err := mem.ByPrefix(ctx, "fact:")
if err != nil {
t.Fatalf("ByPrefix: %v", err)
}
if len(got) != 1 || got[0].ID != "fact:water:1" {
t.Fatalf("after the void the index holds %+v; want only fact:water:1", got)
}
}
func TestDeletePrefixDoesNotWidenOnWildcards(t *testing.T) {
ctx := context.Background()
s := newTestStore(t)
mem := s.VectorMemory()
for _, id := range []string{"fact:a_b:1", "fact:axb:1"} {
if err := mem.Insert(ctx, id, []float32{1, 0}, map[string]string{"type": "fact"}); err != nil {
t.Fatalf("Insert %q: %v", id, err)
}
}
n, err := mem.DeletePrefix(ctx, "fact:a_b:")
if err != nil {
t.Fatalf("DeletePrefix: %v", err)
}
if n != 1 {
t.Fatalf("deleted %d rows; the _ in the key must not match x", n)
}
}