Compare commits

..

19 Commits

Author SHA1 Message Date
claude 8aba4845bf Merge remote-tracking branch 'origin/master' into task/486-deploy-the-workstation-transcriber 2026-08-09 01:22:03 +04:00
claude 2ec92ee8bf Merge pull request 'Move STT and TTS to the workstation, where the microphone already is' (#208) from task/486-move-stt-and-tts-to-the-workstation-wher into master 2026-08-08 23:21:51 +02:00
claude 7c77a378c1 Merge pull request 'Run the routing heads in Go and route with them' (#206) from task/664-routing-heads-in-go into master 2026-08-08 23:21:40 +02:00
claude 672eabc134 Merge pull request 'Measure CrisperWhisper 2.0 turbo in Russian before wiring a runtime for it' (#207) from task/665-crisperwhisper-2-russian into master 2026-08-08 23:21:13 +02:00
claude a1a2fa3704 Swap the workstation model to gemma-4-E4B (V-486)
Owner's call. E4B is 4.2GB against 6.7GB plus a 0.86GB draft, so with CW2
resident the card holds 5.8GB of 16GB instead of 9.2GB.

Measured against a same-session 12B control on the 96-case fixture: 83.3% full
against 84.4%, 89.6% intent-only against 91.7%, destination 19/33 against
23/33, p50 294ms against 344ms. Destination is the column that moved. E4B names
nothing where the 12B names recall or calendar, which walks the whole chain
rather than answering wrong.

MTP is gone with the 12B and cannot come back. It is a separate gguf of
architecture gemma4-assistant with nextn_predict_layers=4, and the only one on
disk is trained against the 12B's hidden states. Neither target gguf carries
nextn tensors, so neither self-speculates.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 01:15:00 +04:00
claude 22a4978459 Say that the card takes one supervisor (V-486)
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 01:07:55 +04:00
claude 1456336652 The transcriber ships with the daemon that starts it (V-486)
serve.py lived only on workpc, which was fine while systemd launched it and is
not fine now that mavgpud does. Two endpoints and no framework: /health answers
503 until the model is loaded, /transcribe takes raw PCM and returns
{"text","confidence"}.

The unit carries CW2_TOKEN through EnvironmentFile and the child inherits it,
so the token is never a flag value.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 01:07:55 +04:00
claude b975716759 One owner for the card, not two neighbours (V-486)
CW2 is a ROCm process, so it registers on the KFD like any contender. Running
it as its own systemd unit made mavgpud yield llama-server to it every few
seconds. The gemma-4-12b arm was down for eight minutes on 2026-08-09 and
routing had silently fallen back to the resident model.

So mavgpud takes an `stt` block and runs the transcriber itself. `foreign` now
excludes every child rather than one pid, which is the fix. Yielding is all or
nothing, because a job that wants the card wants all of it. Idle unloading
stays llama-server's alone: CW2 holds 1.6GB and unloading it would only send
the next voice turn to the homesrv floor.

Maven still talks to the transcriber directly on 8081. There is no proxy,
because with no idle timer there is nothing for one to measure.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 01:07:46 +04:00
claude 4b1edb0617 Record which machine hears him now (V-486)
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 00:46:51 +04:00
claude 944e553669 Point this box at the workstation transcriber (V-486)
The block is inert until the code in PR #208 lands, and deleting it sends
every utterance back to mavsttd, which is what the box does today.

Port 8081 and not mavgpud's 8080, because whisper.cpp cannot load
CrisperWhisper 2.0 at all and it runs under transformers as its own service.
The token comes from deploy/telegram.env like every other secret here. It is
what stops anything on the LAN posting audio to that port.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 00:46:23 +04:00
claude cc32c2c4ab Wire the transcription seam beside the model seam (V-486)
sttSeam is modelSeam for audio and sits at the same place in wireVoice, so
the voice path and the meeting recorder share one transcriber as they
always have.

A box with no workstation.stt block behaves byte-for-byte as it did before
this existed: the floor is handed back untouched and nothing probes. An
empty URL is normalised to no block at all, the way the model block already
works.

Health defaults to the URL's origin rather than the URL itself, because the
transcribe endpoint names a path and appending would ask for
/transcribe/health. A block with no token logs once that anything on the
LAN can post audio to that port.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 00:42:15 +04:00
claude a1e97c94ac The workstation transcribes, homesrv is the floor (V-486)
Same arrangement as llm.Pair and for the same reason. The microphone is at
workpc, the card there has 16GB, and CrisperWhisper 2.0 turbo scores 10.4%
WER in Russian against 27.5% for the ggml-small.bin homesrv loads. The
workstation is never assumed up: it sleeps, and the card is often held.

Admission is a cached atomic written only by the prober, so no voice turn
ever waits on a machine that may be asleep.

Speech-to-text has only the silent half of the degradation rule. A worse
transcript is still a turn, so there is nothing to name a gap about and
Transcribe always falls back. That is the whole difference from llm.Pair,
which also carries CompleteRemote for callers that must refuse instead. A
remote that dies mid-request corrects the cache and falls back in the same
turn, which is what TestPairFallsBackWhenRemoteFails pins.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 00:42:05 +04:00
claude c7f59e48f4 CrisperWhisper reads audio over HTTP, not a socket (V-486)
mavsttd is whisper.cpp linked into a Go daemon and reached over a unix
socket. CrisperWhisper 2.0 cannot be reached that way. whisper.cpp derives
its language count from the vocabulary size, and CW2's 51897 tokens shift
seven special token ids, so it never loads at all.

So it runs under transformers on workpc and this is the client. Same
stt.Transcriber interface and one method, a second transport rather than a
second seam. The body is the PCM itself, because a minute of 16kHz mono is
under 2MB raw and the format is fixed by audio.PCM16kMono.

Audio is the most sensitive thing that crosses this seam, so the client
carries a bearer token.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 00:41:55 +04:00
claude 7138086c3f The routing heads run in Go now, so say so (V-664)
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-08 22:33:35 +04:00
claude 83e168f326 Record what the routing heads score in Go (V-664)
Two defects were found on the way: the tokenizer read every long word
backwards, and the clarify head was discarded below the intent threshold.
Both numbers are in the doc.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-08 22:32:48 +04:00
claude a4abcdefa3 Give the daemon a heads_path and a fixture arm (V-664)
embedder.heads_path is empty by default and deploy/mavend.json sets
it. A missing or broken weights file logs and leaves the heads nil,
because refusing to start over a routing accelerator would trade a
working box for a better one.

TestONNXRoutingHeads is the same cascade TestONNXBaseline scores with
one arm added, so the two are directly comparable. It also checks the
Go tokenizer against the Python one, since the heads were trained
through transformers and are read through a hand-written tokenizer: a
mismatch shows up here as a score below what Python measured on the
same weights, and nowhere else. That is how the reversed word pieces
were found.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-08 22:24:47 +04:00
claude 68a3c85186 Wire the heads between stage 0 and the resident model (V-664)
They run before the model because they are two orders of magnitude
faster and score better on both halves of the route. They decline
rather than clarify, so a declined turn carries on to the model and
then the classifier, which is what a box with no weights file does on
every turn. Nil heads are byte-for-byte the cascade that shipped
before this.

Measured on the 96-case fixture, classifier+ONNX either way:

  intent       76.0% -> 96.9%
  destination  36.4% -> 75.8%
  false clarify   0 -> 1
  missed clarify  8 -> 1
  p50          24.5ms -> 27.9ms

That beats the gemma-4-12b cascade on both halves, 84.4% and 72.7%, at
a twelfth of its 329ms. The four remaining destination misses are all
calendar, which is the stage 0 trade V-660 flagged and the owner has
not called yet.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-08 22:24:37 +04:00
claude 88c086482e Load the routing heads and read three of the four (V-664)
The heads trained in V-661 ran nowhere. This loads the exported graph
and reads intent, destination and clarify off one forward pass. It
declines below 0.6 max softmax rather than clarifying, so a declined
turn reaches whatever is behind it.

The slot head is exported and deliberately not read: slots already
come from the stage-2 extractor, and mapping BIO tags back to text
needs character offsets the tokenizer does not keep.

The clarify head decides on its own and decides first. It answers a
different question from the intent head, so a low intent confidence is
no reason to discard it. Reading it only above the intent threshold
cost 6 of the 8 ambiguous cases on the fixture: the word for water
reads as intent act at 0.23 and clarify at 0.98.

0.6 is the knee measured on the intent fixture: every higher value up
to 0.9 drops right answers and keeps the same two wrong ones.

The body is a fine-tuned COPY of the resident embedder and must never
replace it, because memory recall depends on that file scoring what it
scored.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-08 22:24:37 +04:00
claude feabf9f350 The tokenizer read every long word backwards (V-664)
encodeWord backtracks the Viterbi path from the end of the word and
prepends each piece, which puts them back in reading order. A second
reverse after that loop undid it. So "query: вода" tokenized to
[0 12 1294 41 12489 2] where the reference tokenizer gives
[0 41 1294 12 12489 2], and every multi-piece Russian word reached the
model with its pieces in the wrong order.

Measured on the recall fixture, same 27 cases either way:

  recall@1  70.4% -> 77.8%
  recall@3  85.2% -> 96.3%
  answered after gate  63.0% -> 66.7%
  false recall  0/5 -> 1/5

The classifier barely moves, 76.0% to 75.0% on the routing fixture,
because seeds and queries were mangled the same way and cosine survived
it. Recall is where it cost, because a stored passage and a live query
are different lengths and break differently.

The embedder id now names a tokenizer revision. Stored vectors were
written under rev 1 and no longer sit in the same space as a query
embedded now, and the model file's name never moved, so nothing would
have triggered ReembedAll.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-08 22:23:56 +04:00
36 changed files with 1801 additions and 83 deletions
+68 -6
View File
@@ -47,6 +47,31 @@ free — `worldGap` in `cmd/mavend/worldmodel.go`, which the owner hears instead
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.
**Speech-to-text moved on 2026-08-09** (V-486). `sttSeam` in `cmd/mavend/voicewire.go`
builds an `stt.Pair` beside `modelSeam`, preferring CrisperWhisper 2.0 turbo on workpc
with mavsttd as the floor. It takes only the silent half of the rule. A worse
transcript is still a turn, so `stt.Pair` has no `TranscribeRemote`. The fallback is
never spoken. CW2 turbo scores **10.4% WER in Russian against 27.5%** for the `ggml-small.bin`
mavsttd loads, over 200 Golos clips
(`docs/evals/2026-08-09-crisperwhisper2-russian-wer.md`). It runs in Intended mode, not
Verbatim, though that corpus cannot separate the two.
**whisper.cpp cannot load CW2 at all.** It reads its language count off the vocabulary
size, and CW2's 51897 tokens shift seven special token ids. So it is not a second
endpoint on mavgpud. It is its own transformers service on port 8081
(`deploy/cw2/serve.py`), which Maven reaches directly. `stt.HTTPTranscriber`
posts raw PCM to it with a bearer token, because audio is the most sensitive thing that
crosses this seam. The switch is `workstation.stt` in
`deploy/mavend.json`, and deleting the block sends every utterance to mavsttd.
**mavgpud runs that service as a second child.** That is not an optimisation. CW2 is a
ROCm process on the same card, so it registers on the KFD like any contender. Under its own
systemd unit it made mavgpud evict llama-server every few seconds. That took the
gemma-4-12b arm down for eight minutes on 2026-08-09 before anyone noticed. The card needs
one owner. Any GPU service added beside this daemon has the same defect, so add it to
`cmd/mavgpud` and not to systemd. CW2 is on the yield clock and not the idle one. At 1.6GB
it denies the card to nobody, and unloading it would only send the next voice turn to the
homesrv floor.
Text-to-speech has not moved and piper on homesrv is still the only synthesizer.
## Build & test
CGO daemons (`mavend`, `mavsttd`, `mavttsd`, `mavenclient`) need the vendored toolchain
@@ -230,11 +255,21 @@ re-run it, start a **second** llama-server on a fixed host port — the resident
`--port 0` inside the container and no host process can reach it.
**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
completes through `llm.Pair` against the model mavgpud holds, which is better than the resident
model and about 2.5× faster. gemma-4-12b scored **84.4% full / 93.5% intent-only at p50 329ms**
(`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.
**The workstation runs gemma-4-E4B since 2026-08-09** (owner's call), and it is a
step down measured the same day (`docs/evals/2026-08-09-e4b-vs-12b-routing.md`).
Against a same-session 12B control it scores **83.3% full / 89.6% intent-only,
destination 19/33 against 23/33, at p50 294ms against 344ms**. So it costs four
destination cases and buys 50ms. Read destination as the finding: it names nothing
where the 12B names `recall` or `calendar`, which is safe but walks the whole chain.
It also has no MTP and cannot be given any here. The only `gemma4-assistant`
draft on disk is trained against the 12B's hidden states.
**The intended third engine is not a generative model** (owner's call, 05-08-2026, V-546,
`docs/plans/18-routing-heads-on-e5-small.md`). Routing has a bounded output space, so it is
classification, and the 118M multilingual-e5-small is already resident. Three heads on one
@@ -308,9 +343,36 @@ the possessive agenda rules claim those cases at stage 0 and name nothing, so no
label reaches the head. That is the same trade V-660 flagged and it wants the
owner's call.
**Nothing of this runs in Go.** The weights are `heads.pt` and `out/body_heads/`
on workpc. Reaching the daemon needs an ONNX export and a caller. The resident
e5-small must not be replaced by the copy, because recall depends on that file.
**The heads run in Go and route every turn, since 08-08-2026** (V-664,
`docs/evals/2026-08-08-routing-heads-in-go.md`). This section used to say
nothing of it ran. `RouterHeads` in `internal/router/heads.go` loads
`router_heads.onnx` and reads intent, destination and clarify off one forward
pass. It is stage 0b: after the grammars, **before** the resident model, and the
classifier is still behind both. Through the cascade it scores intent **96.9%**
and destination **75.8%** at p50 27.9ms. That beats the gemma-4-12b cascade,
84.4% and 72.7%, at a twelfth of its 329ms. The workstation stays the better
phraser and is no longer the better router.
Three rules around it. The **clarify head decides first**, before the intent
threshold. It answers a different question. A thin utterance scores low
intent by construction, so gating it cost 6 of 8 ambiguous cases. The
**destination head is read on `IntentQuery` only**, since no other intent
reaches `queryWalk`. And `headsThreshold` is 0.6, the measured knee: every value
to 0.85 drops right answers and keeps the same two wrong ones.
`voice.embedder.heads_path` is the whole switch. Empty, missing or unloadable
means the heads are nil and the cascade is byte-for-byte what shipped before
them. **It must never be pointed at `model_path`.** The resident e5-small must
not be replaced by the fine-tuned copy. Recall depends on that file scoring
what it scored.
**The hand-written tokenizer read every long word backwards** until this task
(`encodeWord`, `onnxembedder.go`). It cost recall@1 7.4 points and recall@3 11.1.
Nothing caught it because seeds and queries were mangled the same way, so cosine
survived. The heads found it. They are trained through transformers and read
through this. The embedder id now carries a tokenizer revision
(`@384/tok2`), so fixing the tokenizer triggers `ReembedAll` the way swapping the
model file does. Bump `tokenizerRev` on any change to what it emits.
`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
+1 -1
View File
@@ -19,7 +19,7 @@ func TestChatAnswersWithNoLlamaServer(t *testing.T) {
dead := llm.New("http://127.0.0.1:1", 500*time.Millisecond)
emb := router.NewHashEmbedder(1024)
h.recall.embedder = emb
h.router = buildRouter(emb, h.matcher, 0.55, pickLLMRouter(true, dead))
h.router = buildRouter(emb, h.matcher, 0.55, pickLLMRouter(true, dead), nil)
h.replier = newLLMReplier(dead, nil)
ctx := withDialogueID(context.Background(), dialogueIDFor(sourceText, "web"))
+2 -2
View File
@@ -317,7 +317,7 @@ func TestClarifyExpiryIsAnnouncedAndWordsStillRoute(t *testing.T) {
h, _, now := newClarifyHandler(t)
emb := router.NewHashEmbedder(1024)
h.recall.embedder = emb
h.router = buildRouter(emb, h.matcher, 0.55, nil)
h.router = buildRouter(emb, h.matcher, 0.55, nil, nil)
if _, asked := h.askClarify(ctx, clarifyDec(router.IntentReminder, router.Slots{Text: "напомни"}, "напомни")); !asked {
t.Fatal("expected a question")
@@ -671,7 +671,7 @@ func TestUnresolvedActSaysItDoesNotKnowTheCommand(t *testing.T) {
func newRoutingClarifyHandler(t *testing.T) (*reactiveHandler, *store.Store) {
t.Helper()
h, st, _ := newClarifyHandler(t)
h.router = buildRouter(router.NewHashEmbedder(1024), h.matcher, 0.55, nil)
h.router = buildRouter(router.NewHashEmbedder(1024), h.matcher, 0.55, nil, nil)
h.recall = recallWiring{embedder: router.NewHashEmbedder(1024), memStore: memory.NewInMemoryStore()}
return h, st
}
+1 -1
View File
@@ -25,7 +25,7 @@ func traceHandler(t *testing.T, ring *decision.Ring) *reactiveHandler {
return &reactiveHandler{
api: api,
recall: recallWiring{embedder: emb, memStore: memory.NewInMemoryStore()},
router: buildRouter(emb, tool.NewMatcher(api), 0.55, nil),
router: buildRouter(emb, tool.NewMatcher(api), 0.55, nil, nil),
replier: voice.NewStubReplier(),
now: func() time.Time { return now },
dataStore: st,
+1 -1
View File
@@ -171,7 +171,7 @@ func newDialogueHandler(t *testing.T) (*reactiveHandler, *store.Store, *time.Tim
// and never a coincidence (V-577, V-579). checkEnd refuses any reminder
// landing on it, and at 09:00 the row that answers "на 9" would trip that.
*now = time.Date(2026, 7, 31, 9, 17, 0, 0, time.UTC)
h.router = buildRouter(router.NewHashEmbedder(1024), h.matcher, 0.55, nil)
h.router = buildRouter(router.NewHashEmbedder(1024), h.matcher, 0.55, nil, nil)
h.recall = recallWiring{embedder: router.NewHashEmbedder(1024), memStore: memory.NewInMemoryStore()}
return h, st, now
}
+1 -1
View File
@@ -24,7 +24,7 @@ func TestApplyAction_FactCapture_QueuesEntityResolution(t *testing.T) {
emb := router.NewHashEmbedder(1024)
matcher := tool.NewMatcher(api)
rtr := buildRouter(emb, matcher, 0.55, nil)
rtr := buildRouter(emb, matcher, 0.55, nil, nil)
h := &reactiveHandler{
api: api,
+1 -1
View File
@@ -20,7 +20,7 @@ func newFactGateHandler(t *testing.T, now time.Time) (*reactiveHandler, ipc.Core
h := &reactiveHandler{
api: api,
recall: recallWiring{embedder: emb, memStore: memory.NewInMemoryStore()},
router: buildRouter(emb, tool.NewMatcher(api), 0.55, nil),
router: buildRouter(emb, tool.NewMatcher(api), 0.55, nil, nil),
replier: voice.NewStubReplier(),
now: func() time.Time { return now },
dataStore: st,
+1 -1
View File
@@ -45,7 +45,7 @@ func newNoteHandler(t *testing.T) (*reactiveHandler, *store.Store) {
h := &reactiveHandler{
api: api,
recall: recallWiring{embedder: emb, memStore: memory.NewInMemoryStore()},
router: buildRouter(emb, tool.NewMatcher(api), 0.55, nil),
router: buildRouter(emb, tool.NewMatcher(api), 0.55, nil, nil),
replier: voice.NewStubReplier(),
now: func() time.Time { return now },
dataStore: st,
+2 -2
View File
@@ -22,7 +22,7 @@ func TestReactiveNotesReminders(t *testing.T) {
emb := router.NewHashEmbedder(1024)
matcher := tool.NewMatcher(api)
rtr := buildRouter(emb, matcher, 0.55, nil)
rtr := buildRouter(emb, matcher, 0.55, nil, nil)
h := &reactiveHandler{
api: api,
@@ -104,7 +104,7 @@ func TestSpokenTaskCaptureFilesATask(t *testing.T) {
h := &reactiveHandler{
api: api,
recall: recallWiring{embedder: emb, memStore: memory.NewInMemoryStore()},
router: buildRouter(emb, matcher, 0.55, nil),
router: buildRouter(emb, matcher, 0.55, nil, nil),
replier: voice.NewStubReplier(),
now: func() time.Time { return now },
dataStore: st,
+1 -1
View File
@@ -474,7 +474,7 @@ func newSimWorld(t *testing.T, sc scenario) *simWorld {
// used to be built on a nil API, which meant any scenario that produced an
// act panicked the moment the matcher was consulted.
matcher := tool.NewMatcher(api)
rtr := buildRouter(emb, matcher, config.DefaultRouterThreshold, router.NewLLMRouter(scripted))
rtr := buildRouter(emb, matcher, config.DefaultRouterThreshold, router.NewLLMRouter(scripted), nil)
w.handler = &reactiveHandler{
stt: simTranscriber{},
+43
View File
@@ -0,0 +1,43 @@
package main
import (
"testing"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/stt"
)
// A box with no workstation.stt block transcribes exactly as it did before the
// seam existed: the floor is handed back untouched, and nothing probes.
func TestSttSeamWithNoBlockIsTheFloor(t *testing.T) {
floor := stt.NewStub()
got, pair := sttSeam(&config.Config{}, floor)
if pair != nil {
t.Fatal("no block must build no pair")
}
if got != stt.Transcriber(floor) {
t.Fatal("no block must hand back the floor itself")
}
}
func TestSttSeamPrefersTheWorkstation(t *testing.T) {
cfg := &config.Config{Workstation: &config.WorkstationConfig{
URL: "http://127.0.0.1:1",
Stt: &config.WorkstationSttConfig{
URL: "http://127.0.0.1:2/transcribe",
Health: "http://127.0.0.1:2/health",
},
}}
got, pair := sttSeam(cfg, stt.NewStub())
if pair == nil {
t.Fatal("a configured block must build a pair")
}
defer pair.Stop()
if got != stt.Transcriber(pair) {
t.Fatal("the pair is what callers must transcribe through")
}
// Nothing answers on port 2, so the seam is the floor until it does.
if pair.Available() {
t.Fatal("an unreachable workstation must not be available")
}
}
+76 -4
View File
@@ -36,7 +36,9 @@ type voiceWiring struct {
sessions *voice.Sessions
voiceSink delivery.Sink
embedder router.Embedder
handler *reactiveHandler // the reactive handler for IPC Chat
// heads — the routing heads, nil unless embedder.heads_path is set.
heads *router.RouterHeads
handler *reactiveHandler // the reactive handler for IPC Chat
// worker clients (set when configured as Remote): closed on shutdown so
// mavsttd / mavttsd don't keep a stale conn into a restarting daemon.
sttClient *worker.Client
@@ -53,7 +55,11 @@ type voiceWiring struct {
// 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
// sttPair — CrisperWhisper 2.0 on the workstation with mavsttd as the
// floor, nil unless the `workstation.stt` block names an address. Held for
// the same reason as pair: to stop its prober on shutdown.
sttPair *stt.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.
@@ -72,6 +78,9 @@ func (w *voiceWiring) close() {
if w.embedder != nil {
_ = w.embedder.Close()
}
if w.heads != nil {
_ = w.heads.Close()
}
if w.server != nil {
_ = w.server.Close()
}
@@ -84,6 +93,9 @@ func (w *voiceWiring) close() {
if w.pair != nil {
w.pair.Stop()
}
if w.sttPair != nil {
w.sttPair.Stop()
}
w.mcp.close()
}
@@ -112,6 +124,7 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
} else {
transcriber = stt.NewStub()
}
transcriber, w.sttPair = sttSeam(cfg, transcriber)
w.transcriber = transcriber
// ----- tts (Stub in-process OR Remote) -----
@@ -147,6 +160,24 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
emb = router.NewHashEmbedder(1024)
}
w.embedder = emb
// ----- router: routing heads (only when configured, and never fatal) -----
// A missing or broken weights file logs and leaves w.heads nil, which is
// byte-for-byte the cascade that shipped before V-664. Refusing to start
// over a routing accelerator would trade a working box for a better one.
if cfg.Voice.Embedder != nil && cfg.Voice.Embedder.HeadsPath != "" {
h, err := router.NewRouterHeads(
cfg.Voice.Embedder.HeadsPath,
cfg.Voice.Embedder.TokenizerPath,
)
if err != nil {
log.Printf("voice: routing heads unavailable, cascade unchanged: %v", err)
} else {
log.Printf("voice: routing heads loaded from %s", cfg.Voice.Embedder.HeadsPath)
w.heads = h
}
}
repairFactVectors(dataStore, emb)
checkStoredEmbedder(dataStore, emb)
// Retention is enforced on write, which is not enough on its own: a box that
@@ -223,7 +254,8 @@ 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(), hot))
rtr := buildRouter(emb, matcher, threshold,
pickLLMRouter(cfg.Voice.UseLLMRouter(), hot), w.heads)
// ----- sessions registry (shared with voicesink) -----
sessions := voice.NewSessions()
@@ -366,6 +398,44 @@ func modelSeam(cfg *config.Config, resident *llm.Client) (router.Completer, *llm
return pair, pair
}
// sttSeam builds the transcription seam the voice path and the meeting
// recorder share. It is modelSeam for audio and follows the same rule.
//
// With no `workstation.stt` block it hands back the floor untouched, which is
// today's deploy exactly. With one, it is an stt.Pair preferring CrisperWhisper
// 2.0 on workpc, which scores 10.4% WER in Russian against the floor's 27.5%
// (docs/evals/2026-08-09-crisperwhisper2-russian-wer.md).
//
// Only the silent half of the degradation rule applies here. A worse transcript
// is still a turn, so there is nothing to name a gap about and the fallback is
// never spoken. That is why stt.Pair has no TranscribeRemote.
func sttSeam(cfg *config.Config, floor stt.Transcriber) (stt.Transcriber, *stt.Pair) {
if cfg.Workstation == nil || cfg.Workstation.Stt == nil {
return floor, nil
}
s := cfg.Workstation.Stt
lang := ""
if cfg.Voice != nil {
lang = cfg.Voice.Lang
if cfg.Voice.Stt != nil && cfg.Voice.Stt.Lang != "" {
lang = cfg.Voice.Stt.Lang
}
}
pair := stt.NewPair(
stt.NewHTTPTranscriber(s.URL, s.Token, lang, time.Duration(s.Timeout)),
floor,
s.Health,
time.Duration(s.Probe),
)
pair.Start(context.Background())
if s.Token == "" {
log.Print("voice: the workstation transcriber has no token, so anything on the LAN can post audio to it")
}
log.Printf("voice: workstation transcriber at %s, probed every %s, mavsttd as the floor",
s.URL, time.Duration(s.Probe))
return pair, pair
}
func pickLLMRouter(enabled bool, c router.Completer) *router.LLMRouter {
if !enabled {
return nil
@@ -390,7 +460,8 @@ func pickLLMRouter(enabled bool, c router.Completer) *router.LLMRouter {
// intent from seedDir (models/seeds/<intent>.txt) — see seedClassifier
// below for the current intent list and file names.
// - Threshold is from voice.router_threshold config (default 0.55).
func buildRouter(emb router.Embedder, acts router.ActMatcher, threshold float64, llmR *router.LLMRouter) *router.Router {
func buildRouter(emb router.Embedder, acts router.ActMatcher, threshold float64,
llmR *router.LLMRouter, heads *router.RouterHeads) *router.Router {
cls := router.NewClassifier(emb)
seedClassifier(cls)
grammars := router.DefaultGrammars(acts)
@@ -442,6 +513,7 @@ func buildRouter(emb router.Embedder, acts router.ActMatcher, threshold float64,
},
Threshold: threshold,
LLM: llmR,
Heads: heads,
})
}
+13 -4
View File
@@ -34,22 +34,31 @@ type probe struct {
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.
// foreign lists every ROCm process that is not ours. self holds the pids of the
// supervisor's own children, and a child that is not running contributes 0.
//
// There is more than one child since 09-08-2026. CW2 registers on the KFD like
// any ROCm job, so a supervisor that excluded only llama-server would read its
// own transcriber as a contender, yield the card to it, and never keep a model
// loaded again.
//
// 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 {
func (p probe) foreign(self ...int) []gpuProc {
entries, err := os.ReadDir(p.kfdRoot)
if err != nil {
return nil
}
mine := make(map[int]bool, len(self))
for _, pid := range self {
mine[pid] = true
}
var out []gpuProc
for _, e := range entries {
pid, err := strconv.Atoi(e.Name())
if err != nil || pid == selfPID {
if err != nil || mine[pid] {
continue
}
out = append(out, gpuProc{
+46 -1
View File
@@ -1,6 +1,7 @@
package main
import (
"context"
"net/http"
"net/http/httptest"
"net/url"
@@ -8,6 +9,7 @@ import (
"path/filepath"
"strconv"
"testing"
"time"
)
// fakeKFD builds the sysfs shape the workstation actually has: one directory
@@ -47,6 +49,24 @@ func TestForeignExcludesOurChild(t *testing.T) {
}
}
// The transcriber is a ROCm process on the same card, so it registers on the
// KFD exactly like a contender does. Reading it as one is what happened on
// 2026-08-09 while CW2 ran under its own systemd unit: mavgpud yielded, waited
// five polls, loaded the model, yielded again, and never held it for a whole
// minute. Excluding every child is the fix and this is the test of it.
func TestForeignExcludesEveryChild(t *testing.T) {
p := probe{kfdRoot: fakeKFD(t, map[int]int64{478104: 12791693312, 999: 4096, 1001: 1717986918})}
ours := p.foreign(999, 1001)
if len(ours) != 1 || ours[0].PID != 478104 {
t.Fatalf("only the CPT run is a contender, got %+v", ours)
}
// A child that is not running reports pid 0, which must exclude nothing.
if got := p.foreign(999, 0); len(got) != 2 {
t.Errorf("a stopped child excludes nobody: got %d contenders, want 2", len(got))
}
}
// 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) {
@@ -81,7 +101,7 @@ func TestFreeVRAM(t *testing.T) {
// 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, "")}
s := &supervisor{run: newRunner("fake", "/bin/true", nil, "")}
h := s.handler(mustURL(t, "http://127.0.0.1:1"))
for _, path := range []string{"/health", "/v1/chat/completions"} {
@@ -101,3 +121,28 @@ func mustURL(t *testing.T, s string) *url.URL {
}
return u
}
// Yielding is all or nothing. A CPT run wants the whole card, so handing back
// the language model while the transcriber keeps 1.6GB mapped would leave the
// other job failing its allocation, which is the outcome yielding exists to
// prevent.
func TestYieldStopsEveryChild(t *testing.T) {
idle := "while : ; do sleep 1 ; done"
s := &supervisor{
cfg: config{EvictAfter: 1, StopGrace: duration(2 * time.Second)},
probe: probe{kfdRoot: fakeKFD(t, map[int]int64{478104: 12791693312})},
run: newRunner("llama-server", fakeServer(t, idle), nil, ""),
stt: newRunner("cw2", fakeServer(t, idle), nil, ""),
}
for _, r := range s.children() {
if err := r.start(); err != nil {
t.Fatal(err)
}
}
s.tick(context.Background())
for _, r := range s.children() {
if r.running() {
t.Errorf("%s outlived the yield", r.name)
}
}
}
+90 -17
View File
@@ -10,6 +10,11 @@
// 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.
//
// It supervises a second child since 09-08-2026, the CW2 transcriber, and for
// one reason only: it is a ROCm process on the same card. Any GPU service the
// owner leaves running beside this daemon reads as a contender and evicts the
// model, so the card needs one owner rather than two neighbours.
package main
import (
@@ -36,6 +41,10 @@ type config struct {
// owner's business and not this daemon's schema.
LlamaArgs []string `json:"llama_args"`
// Stt is optional. Without it mavgpud supervises llama-server alone, which
// is everything it did before 09-08-2026.
Stt *sttConfig `json:"stt,omitempty"`
KFDRoot string `json:"kfd_root"`
DRMDevice string `json:"drm_device"`
@@ -51,6 +60,22 @@ type config struct {
StartAfter int `json:"start_after_polls"`
}
// sttConfig is the CW2 transcriber, which mavgpud runs for one reason: it is a
// ROCm process on this card. Left to its own systemd unit it registers on the
// KFD, the supervisor reads it as a contender, and llama-server is evicted
// within two polls and restarted five polls later, forever. That thrash was
// observed on 2026-08-09 and it is what folded the service in here.
//
// Maven talks to it directly, not through this daemon. There is no proxy and no
// idle timer: at 1.6GB it denies the card to nobody, and unloading it would only
// send the next voice turn to the homesrv floor for no gain.
type sttConfig struct {
// Addr is where the service binds, and it is read only to probe /health.
Addr string `json:"addr"`
Bin string `json:"bin"`
Args []string `json:"args"`
}
func defaults() config {
return config{
Listen: ":8080",
@@ -99,12 +124,18 @@ func main() {
}
base := "http://" + cfg.LlamaAddr
run := newRunner(cfg.LlamaBin, cfg.LlamaArgs, base+"/health")
run := newRunner("llama-server", cfg.LlamaBin, cfg.LlamaArgs, base+"/health")
sup := &supervisor{
cfg: cfg,
probe: probe{kfdRoot: cfg.KFDRoot, drmDev: cfg.DRMDevice},
run: run,
}
if s := cfg.Stt; s != nil {
if s.Bin == "" || s.Addr == "" {
log.Fatal("mavgpud: stt needs both bin and addr")
}
sup.stt = newRunner("cw2", s.Bin, s.Args, "http://"+s.Addr+"/health")
}
sup.touch()
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
@@ -129,13 +160,17 @@ func main() {
shut, done := context.WithTimeout(context.Background(), 5*time.Second)
defer done()
_ = srv.Shutdown(shut)
run.stop(time.Duration(cfg.StopGrace))
for _, r := range sup.children() {
r.stop(time.Duration(cfg.StopGrace))
}
}
type supervisor struct {
cfg config
probe probe
run *runner
// stt is the CW2 transcriber, or nil when the config names none.
stt *runner
lastReq atomic.Int64 // unix nanos of the last request Maven sent
@@ -198,7 +233,11 @@ func (s *supervisor) loop(ctx context.Context) {
// 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())
var pids []int
for _, r := range s.children() {
pids = append(pids, r.pid())
}
others := s.probe.foreign(pids...)
if len(others) > 0 {
s.foreignStreak++
s.clearStreak = 0
@@ -207,31 +246,65 @@ func (s *supervisor) tick(ctx context.Context) {
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))
// Yielding is all or nothing. A CPT run wants the whole card, and handing
// back 8GB while holding 1.6GB is the shape of a failed allocation.
if s.foreignStreak >= s.cfg.EvictAfter && s.anyRunning() {
log.Printf("mavgpud: yielding the card to %s", describe(others))
for _, r := range s.children() {
r.stop(time.Duration(s.cfg.StopGrace))
}
return
}
if s.clearStreak < s.cfg.StartAfter {
clear := s.clearStreak >= s.cfg.StartAfter
if s.run.running() {
s.run.refreshReady(ctx)
if 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))
}
} else if clear && s.probe.freeVRAM() >= s.cfg.MinFreeVRAM {
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)
}
}
if s.stt == nil {
return
}
if free := s.probe.freeVRAM(); free < s.cfg.MinFreeVRAM {
if s.stt.running() {
s.stt.refreshReady(ctx)
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)
// No VRAM precondition here, unlike llama-server. That check exists because
// a 12B refuses to load when the card is short, and 1.6GB fits wherever the
// KFD is clear. Reading free VRAM would also block the transcriber for good
// once the language model was resident, since it holds more than the floor.
if clear {
if err := s.stt.start(); err != nil {
log.Printf("mavgpud: start cw2: %v", err)
}
}
}
func (s *supervisor) children() []*runner {
if s.stt == nil {
return []*runner{s.run}
}
return []*runner{s.run, s.stt}
}
func (s *supervisor) anyRunning() bool {
for _, r := range s.children() {
if r.running() {
return true
}
}
return false
}
// 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.
+17 -13
View File
@@ -10,14 +10,18 @@ import (
"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
// runner owns one GPU 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.
//
// There are two of them since 09-08-2026: llama-server and the CW2 transcriber.
// name is what the log calls this one.
type runner struct {
name string
bin string
args []string
// ready is llama-server's own /health, which answers "is a model loaded".
// ready is the child's own /health, which answers "is a model loaded".
// Loading a 7-14B takes tens of seconds, so started is not ready.
readyURL string
@@ -32,9 +36,9 @@ type runner struct {
http *http.Client
}
func newRunner(bin string, args []string, readyURL string) *runner {
func newRunner(name, bin string, args []string, readyURL string) *runner {
return &runner{
bin: bin, args: args, readyURL: readyURL,
name: name, bin: bin, args: args, readyURL: readyURL,
http: &http.Client{Timeout: 2 * time.Second},
}
}
@@ -60,7 +64,7 @@ func (r *runner) isReady() bool {
return r.ready
}
// start launches llama-server. It returns as soon as the process exists, not
// start launches the child. It returns as soon as the process exists, not
// when the model is loaded.
func (r *runner) start() error {
r.mu.Lock()
@@ -76,7 +80,7 @@ func (r *runner) start() error {
return err
}
r.cmd, r.ready, r.yielding = cmd, false, false
log.Printf("mavgpud: started llama-server pid=%d", cmd.Process.Pid)
log.Printf("mavgpud: started %s pid=%d", r.name, cmd.Process.Pid)
go func() {
err := cmd.Wait()
r.mu.Lock()
@@ -84,15 +88,15 @@ func (r *runner) start() error {
r.cmd, r.ready, r.yielding = nil, false, false
r.mu.Unlock()
if yielded {
log.Printf("mavgpud: llama-server stopped, card yielded (%v)", err)
log.Printf("mavgpud: %s stopped, card yielded (%v)", r.name, err)
return
}
log.Printf("mavgpud: llama-server exited: %v", err)
log.Printf("mavgpud: %s exited: %v", r.name, err)
}()
return nil
}
// stop ends llama-server and waits for the VRAM to come back. SIGTERM first so
// stop ends the child 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.
@@ -117,11 +121,11 @@ func (r *runner) stop(grace time.Duration) {
}
time.Sleep(100 * time.Millisecond)
}
log.Printf("mavgpud: llama-server did not exit in %s, killing", grace)
log.Printf("mavgpud: %s did not exit in %s, killing", r.name, grace)
_ = syscall.Kill(pgid, syscall.SIGKILL)
}
// refreshReady asks llama-server whether the model is loaded. Called once per
// refreshReady asks the child whether the model is loaded. Called once per
// supervisor tick, never per request.
func (r *runner) refreshReady(ctx context.Context) {
if !r.running() {
@@ -141,6 +145,6 @@ func (r *runner) refreshReady(ctx context.Context) {
r.ready = ok
r.mu.Unlock()
if ok && !was {
log.Printf("mavgpud: model ready")
log.Printf("mavgpud: %s ready", r.name)
}
}
+2 -2
View File
@@ -25,7 +25,7 @@ func fakeServer(t *testing.T, body string) string {
// status of a routine yield is identical to that of a real crash. Reading the
// mavgpud log, the two were indistinguishable (Vikunja #491).
func TestStopMarksTheExitAsAYield(t *testing.T) {
r := newRunner(fakeServer(t, "while : ; do sleep 1 ; done"), nil, "")
r := newRunner("fake", fakeServer(t, "while : ; do sleep 1 ; done"), nil, "")
if err := r.start(); err != nil {
t.Fatalf("start: %v", err)
}
@@ -49,7 +49,7 @@ func TestStopMarksTheExitAsAYield(t *testing.T) {
// Stopping when nothing is running must not arm the flag for the next child.
// The next exit after that would be a real crash logged as a yield.
func TestStopWithNoChildDoesNotArmTheFlag(t *testing.T) {
r := newRunner("/nonexistent", nil, "")
r := newRunner("fake", "/nonexistent", nil, "")
r.stop(10 * time.Millisecond)
r.mu.Lock()
defer r.mu.Unlock()
+156
View File
@@ -0,0 +1,156 @@
"""CrisperWhisper 2.0 turbo as an HTTP service, for Maven's stt.Pair.
Two endpoints and no framework.
GET /health 200 once the model is loaded, 503 while it is loading.
POST /transcribe raw 16kHz mono PCM in, {"text","confidence"} out.
The body is the PCM itself rather than JSON. A minute of 16kHz mono is under
2MB raw and about 2.6MB base64, and the format is fixed at the Maven seam, so
headers carry it more cheaply than an envelope.
Why this exists at all: whisper.cpp cannot load CW2. It derives its language
count from the vocabulary size, and CW2's 51897 tokens shift seven special
token ids. So mavsttd stays whisper.cpp on homesrv and this runs beside the
model on workpc, where it scores 10.4% WER in Russian against the floor's 27.5%
(docs/evals/2026-08-09-crisperwhisper2-russian-wer.md in the Maven repo).
Intended mode, not verbatim. The owner asked for what he meant to say, not
every stutter on the way there.
"""
import hmac
import json
import logging
import os
import sys
import threading
import time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
import numpy as np
HOST = os.environ.get("CW2_HOST", "0.0.0.0")
PORT = int(os.environ.get("CW2_PORT", "8081"))
SIZE = os.environ.get("CW2_SIZE", "turbo")
MODE = os.environ.get("CW2_MODE", "intended")
TOKEN = os.environ.get("CW2_TOKEN", "")
# 25MB is about thirteen minutes of 16kHz mono. Longer than any utterance and
# short enough that a wrong caller cannot exhaust memory.
MAX_BODY = int(os.environ.get("CW2_MAX_BODY", str(25 * 1024 * 1024)))
logging.basicConfig(
level=logging.INFO, format="%(asctime)s cw2: %(message)s", stream=sys.stderr
)
log = logging.getLogger("cw2")
_model = None
# The card holds one model and transcribes one utterance at a time. The lock is
# what makes a second caller wait rather than corrupt the first.
_lock = threading.Lock()
def load_model():
global _model
from crisperwhisper import CrisperWhisperModel
t0 = time.perf_counter()
# backend is forced. With ctranslate2 importable, "auto" picks ct2, which is
# CUDA-only and this card is AMD.
m = CrisperWhisperModel(
SIZE, backend="transformers", compute_type="float16", device="cuda"
)
_model = m
log.info("loaded %s in %.1fs, mode=%s", SIZE, time.perf_counter() - t0, MODE)
def authorised(headers):
if not TOKEN:
return True
got = headers.get("Authorization", "")
return hmac.compare_digest(got, "Bearer " + TOKEN)
class Handler(BaseHTTPRequestHandler):
protocol_version = "HTTP/1.1"
def log_message(self, fmt, *args):
log.info(fmt, *args)
def _send(self, code, payload):
body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
self.send_response(code)
self.send_header("Content-Type", "application/json; charset=utf-8")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def do_GET(self):
if self.path.rstrip("/") != "/health":
self._send(404, {"error": "not found"})
return
if _model is None:
self._send(503, {"status": "loading"})
return
self._send(200, {"status": "ok", "model": SIZE, "mode": MODE})
def do_POST(self):
if self.path.rstrip("/") != "/transcribe":
self._send(404, {"error": "not found"})
return
if not authorised(self.headers):
self._send(401, {"error": "unauthorised"})
return
if _model is None:
self._send(503, {"error": "loading"})
return
length = int(self.headers.get("Content-Length", "0"))
if length <= 0 or length > MAX_BODY:
self._send(413, {"error": "bad body length"})
return
raw = self.rfile.read(length)
rate = int(self.headers.get("X-Sample-Rate", "16000"))
channels = int(self.headers.get("X-Channels", "1"))
bits = int(self.headers.get("X-Sample-Bits", "16"))
lang = self.headers.get("X-Language", "ru") or "ru"
if channels != 1 or bits != 16:
self._send(400, {"error": "want 16-bit mono pcm"})
return
# int16 little-endian to the float32 the encoder wants.
wav = np.frombuffer(raw, dtype="<i2").astype(np.float32) / 32768.0
if wav.size == 0:
self._send(200, {"text": "", "confidence": 0.0})
return
t0 = time.perf_counter()
try:
with _lock:
res = _model.transcribe(wav, sr=rate, language=lang, mode=MODE)
except Exception as exc: # noqa: BLE001 - the caller falls back to mavsttd
log.exception("transcribe failed")
self._send(500, {"error": str(exc)})
return
elapsed = time.perf_counter() - t0
text = (res.text or "").strip()
log.info("%.2fs audio in %.2fs: %r", wav.size / rate, elapsed, text[:60])
# The model reports no calibrated score. 1.0 would be a claim, and the
# Maven side reads confidence only to log it.
self._send(200, {"text": text, "confidence": 0.0})
def main():
if not TOKEN:
log.warning("no CW2_TOKEN set: anything on the LAN can post audio here")
# Bind before loading, so a restart answers 503 rather than refusing the
# connection. Both make Maven fall back, but only one of them says why.
srv = ThreadingHTTPServer((HOST, PORT), Handler)
threading.Thread(target=load_model, daemon=True).start()
log.info("listening on %s:%d", HOST, PORT)
srv.serve_forever()
if __name__ == "__main__":
main()
+23 -2
View File
@@ -78,10 +78,30 @@
"Addressed by LAN address, not container name: mavgpud runs on another",
"machine and there is no shared docker network to name it on."
],
"//workstation.stt": [
"CrisperWhisper 2.0 turbo on the same machine, a second service on port",
"8081 and not a second endpoint on mavgpud. whisper.cpp cannot load CW2 at",
"all: it derives its language count from the vocabulary size, and CW2's",
"51897 tokens shift seven special token ids. So it runs under transformers",
"there and mavsttd stays whisper.cpp here.",
"Worth the second service: CW2 turbo scores 10.4% WER in Russian against",
"27.5% for the ggml-small.bin mavsttd loads, measured on 200 Golos clips",
"in docs/evals/2026-08-09-crisperwhisper2-russian-wer.md.",
"Deleting this block sends every utterance to mavsttd, which is what the",
"box did before it existed. A worse transcript is still a turn, so the",
"fallback is silent and Kami is never told which machine heard him.",
"The token is what stops anything on the LAN posting audio to that port."
],
"workstation": {
"url": "http://192.168.1.105:8080",
"probe": "15s",
"timeout": "90s"
"timeout": "90s",
"stt": {
"url": "http://192.168.1.105:8081/transcribe",
"token": "${MAVEN_STT_TOKEN}",
"probe": "15s",
"timeout": "10s"
}
},
"//search": [
@@ -231,7 +251,8 @@
"embedder": {
"model_path": "/opt/maven/models/embedder/multilingual-e5-small/model_quantized.onnx",
"tokenizer_path": "/opt/maven/models/embedder/multilingual-e5-small/tokenizer.json",
"lib_path": "/opt/maven/lib/libonnxruntime.so"
"lib_path": "/opt/maven/lib/libonnxruntime.so",
"heads_path": "/opt/maven/models/embedder/router-heads/router_heads.onnx"
},
"llm_router": true,
"query_min_score": 0.55,
+21 -5
View File
@@ -2,9 +2,14 @@
"listen": ":8080",
"llama_addr": "127.0.0.1:10000",
"llama_bin": "llama-server",
"//llama_args": [
"E4B carries no MTP tensors, so the speculative flags are gone with the 12B.",
"MTP on this box is a separate gguf of architecture gemma4-assistant with",
"nextn_predict_layers=4, and mtp-gemma-4-12B-it-BF16 is the only one there is.",
"Its head is trained against the 12B's hidden states, so it cannot drive E4B."
],
"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",
"-m", "/mnt/D/AI/gemma4/gemma-4-E4B-it-qat-UD-Q4_K_XL.gguf",
"-ngl", "99",
"-fa", "on",
"-np", "1",
@@ -15,11 +20,22 @@
"--batch-size", "2048",
"--ubatch-size", "512",
"--jinja",
"--chat-template-kwargs", "{\"enable_thinking\":false}",
"--spec-type", "draft-mtp",
"--spec-draft-n-max", "2"
"--chat-template-kwargs", "{\"enable_thinking\":false}"
],
"//stt": [
"CrisperWhisper 2.0 turbo, which Maven reaches directly on port 8081.",
"mavgpud runs it because it is a ROCm process on this card: under its own",
"systemd unit it registered on the KFD and the supervisor evicted",
"llama-server every few seconds. CW2_TOKEN comes from the unit's",
"EnvironmentFile and is never a flag value."
],
"stt": {
"addr": "127.0.0.1:8081",
"bin": "/home/kami/Programs/cw2-eval/.venv/bin/python",
"args": ["/home/kami/Programs/cw2-service/serve.py"]
},
"kfd_root": "/sys/class/kfd/kfd/proc",
"drm_device": "/sys/class/drm/card1/device",
+7 -2
View File
@@ -1,5 +1,5 @@
[Unit]
# Runs on the workstation (bugmachine), not on homesrv. Install as a systemd
# Runs on the workstation (workpc), 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:
#
@@ -8,10 +8,15 @@
# 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)
Description=Maven GPU supervisor (holds llama-server and CW2 while the card is free)
After=network.target
[Service]
# CW2_TOKEN for the transcriber child, which inherits this environment. The
# token is read from a file and never appears as a flag value, the rule
# mavpoll and mavmaild follow. Missing file, no transcriber auth, so keep the
# dash off: a mavgpud that cannot read it must fail loudly.
EnvironmentFile=%h/Programs/cw2-service/cw2.env
ExecStart=%h/.local/bin/mavgpud -config %h/.config/mavgpud.json
Restart=always
RestartSec=5
@@ -0,0 +1,146 @@
# The routing heads, running in Go
Date: 2026-08-08. Vikunja V-664.
Weights: `router_heads.onnx`, fp32, exported from `heads.pt` on workpc.
Fixture: `internal/router/eval/ru_routing_v1.json`, 96 cases, 33 carrying a destination.
Runner: `make t PKG=./internal/router/eval/ RUN=TestONNXRoutingHeads`.
The four heads of V-661 ran nowhere. This is the number they score through the
Go cascade. Same fixture and same grammars as `TestONNXBaseline`, and only the
middle stage varies.
## Headline
| | classifier + ONNX | heads + classifier | gemma-4-12b cascade |
|---|---|---|---|
| intent | 75.0% (72/96) | **96.9% (93/96)** | 84.4% |
| destination | 33.3% (11/33) | **75.8% (25/33)** | 72.7% |
| false clarify | 0 | 1 | 2 |
| missed clarify | 8 | 1 | 1 |
| p50 | 24.5ms | 27.9ms | 329ms |
A 118M encoder beats the 12B teacher it was distilled from. It wins on both
halves of the route, at a twelfth of the latency. The workstation stays the
better phraser and is no longer the better router.
The p50 is not the heads. Most of it is the classifier's own embedder pass on
the turns the heads decline, plus process warm-up on the first case. The heads'
own forward pass measures 7.3ms on workpc.
## Two defects were in the way, and the first was not in the heads
**The tokenizer read every long word backwards.** `encodeWord` backtracks the
Viterbi path from the end of a word and prepends each piece. That puts them back
in reading order, and a second reverse after the loop undid it. So
`query: вода` tokenized to `[0 12 1294 41 12489 2]` where the reference
tokenizer gives `[0 41 1294 12 12489 2]`.
It was found here and only here. The heads were trained through transformers and
are read through the hand-written tokenizer. So a mismatch shows up as a score
far below what Python measured on the same weights. Nothing else in the suite
compares the two.
Measured on the recall fixture, same 27 cases either way:
| | reversed | fixed |
|---|---|---|
| recall@1 | 70.4% (19/27) | **77.8% (21/27)** |
| recall@3 | 85.2% (23/27) | **96.3% (26/27)** |
| answered after gate | 63.0% | 66.7% |
| wrong note on top | 8 | 6 |
| false recall | 0/5 | 1/5 |
The classifier barely moved, 76.0% to 75.0%, and destination 36.4% to 33.3%.
Both are one case on 96 and neither is a finding. Seeds and queries were mangled
the same way, so cosine survived it. Recall is where it cost, because a stored
passage and a live query are different lengths and break differently.
The one new false recall is the honest cost and it is not being hidden. A
sharper embedder scores every candidate higher, including the ones that should
have stayed under the gate. That is the same trade `2026-08-04-recall-e5-small.md`
recorded when e5-small replaced MiniLM.
The embedder id now carries a tokenizer revision, `model_quantized@384/tok2`.
Stored vectors were written under rev 1 and no longer sit in the same space as a
query embedded now. The model file's name never moved, so nothing would have
triggered `ReembedAll`. On the box the marker fired on start, and the re-embed
rewrote 65 notes and 19 facts in 5 seconds.
**The clarify head was being thrown away.** It was read only when the intent head
cleared its own threshold. That cost 6 of the 8 ambiguous cases. `вода` reads as intent
`act` at 0.233 and clarify at 0.983. Burying that handed the turn to the
classifier, which routed it confidently and never asked. The clarify head answers
a different question, which is whether there is enough here to act on at all. So
it decides on its own and decides first.
| | intent-gated | clarify decides first |
|---|---|---|
| intent | 90.6% | 96.9% |
| missed clarify | 7 | 1 |
| false clarify | 0 | 1 |
## The threshold is measured, not chosen
Max softmax over the intent head, on the 88 cases carrying an intent:
| threshold | kept | accuracy kept | wrong kept | right dropped |
|---|---|---|---|---|
| 0.5 | 84 | 96.4% | 3 | 2 |
| **0.6** | **81** | **97.5%** | **2** | **4** |
| 0.7 | 75 | 97.3% | 2 | 10 |
| 0.8 | 64 | 96.9% | 2 | 21 |
| 0.9 | 46 | 100.0% | 0 | 37 |
0.6 is the knee. Every value from 0.7 to 0.85 drops right answers and keeps the
same two wrong ones. 0.9 is the only value that clears them, and it costs 37
correct routes to do it.
## Quantization was measured and rejected
| build | size | intent | destination | p50 |
|---|---|---|---|---|
| fp32 | 470MB | 83/88 (94.3%) | 28/33 (84.8%) | 7.3ms |
| int8 | 118MB | 79/88 (89.8%) | 26/33 (78.8%) | 4.0ms |
| fp16 | 235MB | will not load | — | — |
Python numbers, on the heads alone rather than through the cascade. int8 costs
4.5 points of intent and 6 of destination to save 3ms. The cascade around it has
a p50 over a second when the resident model answers. The fp16 graph is broken:
`convert_float_to_float16` leaves a Cast node emitting float16 where the graph
expects float, and onnxruntime refuses the session. It was not worth fixing.
The exporter also had to be told to write one file. It splits weights into a
`.onnx.data` sidecar by default. This onnxruntime resolves that path against the
process working directory rather than the model. A split graph loads from one
directory only.
## What is still wrong
**Four of the eight destination misses are calendar.** Training cannot move them.
The possessive agenda rules claim those cases at stage 0 and name nothing on
purpose. That caution was free while nothing downstream could name anything
either. It has now cost four points in three separate measurements. The call is
the owner's and it is still open.
**The slot head is exported and not read.** Slots come from the stage-2
extractor. Mapping BIO tags back to text needs character offsets the unigram
tokenizer does not keep, which is its own piece of work.
**`поужинал` is a false clarify**, which is the same defect `thinSingleToken`
was narrowed for on 2026-08-01, arriving now from a different direction.
## On the box
Deployed to homesrv the same day. `voice: routing heads loaded` on start, and
`/trace` shows `routing-heads` winning or thinning every turn. The resident model
and the classifier are both marked never asked. Live probes:
```text
что такое TCP? -> kiwix a real definition
кто такой Линус Торвальдс? -> kiwix a real answer
во сколько я лёг вчера -> personal не нашла у тебя такой записи
вода -> thinned to clarify at 0.233 / 0.983
```
A missing or broken weights file logs and leaves the heads nil, which is
byte-for-byte the cascade that shipped before this.
@@ -0,0 +1,51 @@
# gemma-4-E4B against gemma-4-12B on the routing fixture
*Measured 2026-08-09 on workpc. The owner asked for the swap. This is what it costs.*
Both arms ran the same 96-case fixture through `TestLLMRouterBaseline`, minutes
apart, against the same llama-server build and the same mavgpud. The 12B arm is a
control run and not the 2026-08-02 number. That one predates five fixture cases,
the destination labels and a llama.cpp upgrade.
| | full | intent-only | destination | p50 | p95 |
|---|---|---|---|---|---|
| gemma-4-12B-it-qat-UD-Q4_K_XL, MTP draft | 81/96 (84.4%) | 91.7% | 23/33 (69.7%) | 344ms | 471ms |
| gemma-4-E4B-it-qat-UD-Q4_K_XL | 80/96 (83.3%) | 89.6% | 19/33 (57.6%) | 294ms | 562ms |
E4B costs one case of full accuracy, two of intent and **four of destination**,
and buys 50ms at p50. Read the destination column as the finding. One case is
three points on 33. So 23 against 19 is outside the noise a single case makes,
and the other two columns are not.
Both arms produce three false clarifies and one missed clarify, and neither
errored on any case.
## What E4B loses
Four of the five destination regressions are the same shape: it names nothing
where the 12B names `recall` or `calendar`. `ru-query-015` ("сколько я прошёл
шагов") goes further and names `self`. Naming nothing is the safe direction,
because `SourceUnknown` walks the whole chain, so these turns are still answered.
They cost latency and they are what a fourth head is meant to fix (V-546).
Two Russian intent cases regress, both with the interrogative off the front.
`ru-chat-003` ("расскажи анекдот про программистов") goes to `query`.
`ru-fact-003` ("поужинал") goes to `chat`.
## MTP
E4B has none, and there is no way to give it any on this box. MTP on workpc is
a separate gguf of architecture `gemma4-assistant` carrying
`nextn_predict_layers=4`, and `mtp-gemma-4-12B-it-BF16.gguf` is the only one on
disk. Its head is trained against the 12B's hidden states, so it cannot drive an
E4B target. Scanning both target ggufs finds no `nextn` tensors in either, so
neither model self-speculates.
So the 12B arm above ran with speculative decoding and E4B ran without, and E4B
was still faster.
## Cost on the card
E4B is 4.2GB against 6.7GB plus a 0.86GB draft. With CW2 resident at 1.6GB that
is 5.8GB of 16GB against 9.2GB. Nothing in Maven needs the difference, so this is
headroom for the owner's own jobs rather than a capability.
+38 -3
View File
@@ -1,6 +1,6 @@
# Offloading model work to the workstation
*Last verified: 2026-08-05 @ b789676. Living doc: correct it in place, do not append.*
*Last verified: 2026-08-09 @ 50c6637. 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.
@@ -105,6 +105,17 @@ how we find out whether the blind spot is real.
untouched. The model, the context size, the layer count and the MTP flags are the
owner's business and not this daemon's schema.
**Every GPU service on that box belongs under this supervisor**, added to
`cmd/mavgpud` rather than to systemd beside it. The rule was learned on
2026-08-09. The CW2 transcriber ran as its own user unit and registered on the
KFD like any ROCm job. So the supervisor read its own transcriber as a
contender. It yielded the card every few seconds and the gemma-4-12b arm was
down for eight minutes before anyone looked. So the supervisor takes a `stt`
block and starts CW2 itself. Yielding is all or nothing, because a job that
wants the card wants all of it. Idle unloading is not. It applies to
llama-server, which holds 8GB. CW2 holds 1.6GB, and unloading it would cost the
next voice turn its quality for nothing.
## What stays on homesrv, permanently
The **embedder** (multilingual-e5-small, ONNX, CPU). It backs the classifier, which
@@ -149,6 +160,29 @@ 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.
Speech-to-text is wired as of 09-08-2026, and it takes only the silent half of the
rule. A worse transcript is still a turn, so there is nothing to name a gap about
and `stt.Pair` has no `TranscribeRemote`. `sttSeam` in `cmd/mavend/voicewire.go`
builds it, beside `modelSeam` and at the same place in `wireVoice`, so the voice
path and the meeting recorder still share one transcriber.
The remote is not a second endpoint on mavgpud. whisper.cpp cannot load
CrisperWhisper 2.0 at all. It reads its language count off the vocabulary
size, and CW2's 51897 tokens shift seven special token ids. So CW2 runs under
transformers as its own service on port 8081, and `stt.HTTPTranscriber` is the
second transport for the same seam. It posts raw PCM with the format in headers.
It carries a bearer token, because audio is the most sensitive thing that
crosses here.
It is a second endpoint on nothing, but it is a second **child** of mavgpud, and
that part is not optional. See the supervisor section above for why: a ROCm
service the supervisor does not own is a contender it yields to.
The margin is the reason: CW2 turbo scores 10.4% WER in Russian against 27.5% for
the `ggml-small.bin` mavsttd loads, over 200 Golos clips
(`docs/evals/2026-08-09-crisperwhisper2-russian-wer.md`). Text-to-speech has not
moved and piper on homesrv is still the only synthesizer.
Speech-to-text stays two stages when it moves. One call carrying both a clip and the router
prompt was measured on 05-08-2026. It scores 54.2% intent-only against 84.7% for whisper on
homesrv, on the same 72 cases. The model transcribes clips it then routes wrong, so a long
@@ -173,8 +207,9 @@ cleaner transcripts, not accuracy. See `docs/evals/2026-08-05-audio-in-routing.m
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.
3. **Speech-to-text and text-to-speech** (#486). They gain a real margin, but on
quality alone, and both already work.
3. **Speech-to-text and text-to-speech** (#486). Speech-to-text is wired, see
above. Text-to-speech is not, and piper is good enough that nothing argues
for moving it yet.
4. **The wake word** (#487). Independent of all of the above.
## Assumptions
+10
View File
@@ -37,6 +37,16 @@ type EmbedderConfig struct {
ModelPath string `json:"model_path,omitempty"`
TokenizerPath string `json:"tokenizer_path,omitempty"`
LibPath string `json:"lib_path,omitempty"`
// HeadsPath — the routing heads graph, which is a fine-tuned COPY of the
// model above with four linear heads on its pooled output (V-664). Empty
// means no heads, and the cascade runs exactly as it did before they
// existed. It shares LibPath and TokenizerPath, and router_heads.json is
// read from the same directory.
//
// It must never be pointed at ModelPath. Memory recall depends on the
// resident copy scoring what it scored, and the fine-tuned one does not.
HeadsPath string `json:"heads_path,omitempty"`
}
// WeatherConfig configures the weather provider for voice queries.
+75
View File
@@ -1,6 +1,7 @@
package config
import (
"net/url"
"strings"
"time"
)
@@ -35,12 +36,54 @@ type WorkstationConfig struct {
// 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"`
// Stt — CrisperWhisper 2.0 on the same machine, a separate service on its
// own port. Absent ⇒ every utterance goes to mavsttd, which is today.
Stt *WorkstationSttConfig `json:"stt,omitempty"`
}
// WorkstationSttConfig — speech-to-text on the workstation.
//
// It is a second service and not a second endpoint on mavgpud: whisper.cpp
// cannot load CrisperWhisper 2.0 at all, because it derives its language count
// from the vocabulary size and CW2's 51897 tokens shift seven special token
// ids. So CW2 runs under transformers, and this block addresses it.
//
// Worth the trouble: CW2 turbo scores 10.4% WER in Russian against 27.5% for
// the ggml-small.bin homesrv loads
// (docs/evals/2026-08-09-crisperwhisper2-russian-wer.md).
type WorkstationSttConfig struct {
// URL — the transcribe endpoint, e.g.
// "http://192.168.1.105:8081/transcribe". Empty ⇒ the block is normalised
// to nil and mavsttd takes every turn.
URL string `json:"url,omitempty"`
// Health — the admission endpoint. Empty ⇒ the URL's origin + "/health".
// It answers 503 while the card is held, and that is the signal.
Health string `json:"health,omitempty"`
// Token — the bearer token the service checks. Audio is the most sensitive
// thing that crosses this seam, so a LAN deployment should set one. Write
// it as ${MAVEN_STT_TOKEN} and keep the value in deploy/telegram.env, the
// way every other secret in this file is written.
Token string `json:"token,omitempty"`
// Probe — how often admission is re-checked. 0 ⇒ DefaultWorkstationProbe.
Probe Duration `json:"probe,omitempty"`
// Timeout — the per-request budget for one utterance. 0 ⇒
// DefaultWorkstationSttTimeout. A request that overruns falls back to
// mavsttd, which costs a worse transcript and not the turn.
Timeout Duration `json:"timeout,omitempty"`
}
// Workstation defaults, applied in normaliseWorkstation.
const (
DefaultWorkstationProbe = 15 * time.Second
DefaultWorkstationTimeout = 90 * time.Second
// One utterance, not one completion. A voice turn waits on this, so the
// budget is a few seconds and not a minute and a half.
DefaultWorkstationSttTimeout = 10 * time.Second
)
// normaliseWorkstation applies the block's defaults. No address, no preferred
@@ -63,4 +106,36 @@ func (c *Config) normaliseWorkstation() {
if w.Timeout <= 0 {
w.Timeout = Duration(DefaultWorkstationTimeout)
}
normaliseWorkstationStt(w)
}
// normaliseWorkstationStt applies the speech-to-text block's defaults. No
// address, no remote: mavsttd then takes every utterance, which is today.
func normaliseWorkstationStt(w *WorkstationConfig) {
if w.Stt != nil && strings.TrimSpace(w.Stt.URL) == "" {
w.Stt = nil
}
if w.Stt == nil {
return
}
s := w.Stt
if strings.TrimSpace(s.Health) == "" {
s.Health = healthOrigin(s.URL)
}
if s.Probe <= 0 {
s.Probe = Duration(DefaultWorkstationProbe)
}
if s.Timeout <= 0 {
s.Timeout = Duration(DefaultWorkstationSttTimeout)
}
}
// healthOrigin derives the admission endpoint from the transcribe endpoint.
// The URL names a path, so appending to it would ask for /transcribe/health.
func healthOrigin(raw string) string {
u, err := url.Parse(raw)
if err != nil || u.Host == "" {
return strings.TrimRight(raw, "/") + "/health"
}
return u.Scheme + "://" + u.Host + "/health"
}
+6 -3
View File
@@ -14,12 +14,15 @@ import (
"github.com/kami/maven/internal/decision"
)
// The two routing engines, named as claimants. They are one stage and not two,
// because only one of them ever runs: the classifier is reached when the model
// is absent or errored, never alongside it.
// The three routing engines, named as claimants. The model and the classifier
// are one stage and not two, because only one of them ever runs: the classifier
// is reached when the model is absent or errored, never alongside it. The heads
// run before both and decline on low confidence, so they can appear beside
// either one in a record.
const (
claimantLLM = "llm-router"
claimantClassifier = "classifier"
claimantHeads = "routing-heads"
)
// thinReason names which arm of gateLLMDecision cut the confidence. The gate
+12 -2
View File
@@ -1,10 +1,13 @@
package router
import "testing"
import (
"strings"
"testing"
)
func TestEmbedderIDFromModelPath(t *testing.T) {
got := modelIDFromPath("/opt/maven/models/embedder/multilingual-e5-small.onnx")
if got != "multilingual-e5-small@384" {
if got != "multilingual-e5-small@384/tok2" {
t.Fatalf("modelIDFromPath = %q", got)
}
// A different model file must produce a different id, even at 384 dim.
@@ -12,6 +15,13 @@ func TestEmbedderIDFromModelPath(t *testing.T) {
if old == got {
t.Fatal("two different models share one id")
}
// The tokenizer is half of what makes a vector, and it changes under a
// model file whose name never moves (V-664). An id that ignored it would
// leave stored passages in one space and every new query in another, with
// nothing to trigger the re-embed.
if !strings.Contains(got, "/tok") {
t.Fatalf("id %q does not name the tokenizer revision", got)
}
}
func TestEmbedderIDIncludesDim(t *testing.T) {
+86
View File
@@ -0,0 +1,86 @@
package eval
import (
"context"
"os"
"path/filepath"
"testing"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/router"
)
// TestONNXRoutingHeads — the cascade with the routing heads wired, which is
// what V-664 deploys. Opt-in via MAVEN_ONNX_LIB, same as TestONNXBaseline, and
// one TestONNX* per process.
//
// The comparison worth reading is against TestONNXBaseline, which is the same
// cascade with the same grammars and the same classifier floor and no heads.
// Only the middle arm varies.
//
// It also checks the Go unigram tokenizer against the Python one, because the
// heads were trained through transformers and are read through a hand-written
// tokenizer. A mismatch shows up here as a score below what Python measured on
// the same weights, and nowhere else.
func TestONNXRoutingHeads(t *testing.T) {
lib := os.Getenv("MAVEN_ONNX_LIB")
if lib == "" {
t.Skip("MAVEN_ONNX_LIB unset — see AGENTS.md § Embedder model for intent routing")
}
// Absolute, because onnxruntime resolves a graph's external weights file
// against the model path it was given, and a relative one lands in the
// test's working directory.
root, err := filepath.Abs("../../..")
if err != nil {
t.Fatal(err)
}
model := filepath.Join(root, "models/embedder/multilingual-e5-small/model_quantized.onnx")
tok := filepath.Join(root, "models/embedder/multilingual-e5-small/tokenizer.json")
heads := filepath.Join(root, "models/embedder/router-heads/router_heads.onnx")
for _, p := range []string{lib, model, tok, heads} {
if _, err := os.Stat(p); err != nil {
t.Skipf("missing %s: %v", p, err)
}
}
emb, err2 := router.NewONNXEmbedder(model, tok, lib)
if err2 != nil {
t.Skipf("onnx embedder unavailable: %v", err2)
}
err = nil
defer emb.Close()
h, err := router.NewRouterHeads(heads, tok)
if err != nil {
t.Skipf("routing heads unavailable: %v", err)
}
defer h.Close()
f, err := Load()
if err != nil {
t.Fatalf("Load: %v", err)
}
rep, err := Score(context.Background(), "heads+classifier", withHeads(t, emb, h), f)
if err != nil {
t.Fatalf("Score: %v", err)
}
t.Log("\n" + rep.String() + rep.Failures())
}
// withHeads mirrors newBaselineRouter and adds the one arm under test. It is a
// separate function rather than a parameter so the baseline's signature stays
// the shape every other test calls it with.
func withHeads(t *testing.T, emb router.Embedder, h *router.RouterHeads) *router.Router {
t.Helper()
acts := router.DefaultActMatcher{Fns: actFns}
return router.New(router.Config{
Grammars: baselineGrammars(acts),
Classifier: newBaselineClassifier(t, emb),
Extractor: router.Extractor{
Time: router.StubDateTimeParser{},
Acts: acts,
Facts: router.DefaultFactParser{},
},
Threshold: config.DefaultRouterThreshold,
Heads: h,
})
}
+243
View File
@@ -0,0 +1,243 @@
package router
import (
"context"
"encoding/json"
"fmt"
"math"
"os"
"path/filepath"
ort "github.com/yalue/onnxruntime_go"
)
// The routing heads (V-546, V-661, V-664). Four linear heads over one masked
// mean pool of a fine-tuned copy of multilingual-e5-small: intent,
// destination, BIO slot tags and clarify. Trained on workpc, exported to ONNX,
// and read here.
//
// Why this is not the classifier. The classifier compares one utterance to
// frozen seed phrases by cosine. A head is a softmax over the label set, so it
// cannot name a value that does not exist, and its max is a calibratable
// confidence where Confidence: 1.0 was a hardcode.
//
// Why it is not the resident model either. It answers in single-digit
// milliseconds against the model's p50 of 1.19s, and it names a destination
// the classifier arm never names at all.
//
// The body is a COPY of the embedder weights, fine-tuned. It must never
// replace models/embedder/multilingual-e5-small — memory recall depends on
// that file scoring what it scored.
//
// The slot head is exported and deliberately not read. Slots already come from
// the stage-2 extractor, and mapping BIO tags back to text needs character
// offsets the unigram tokenizer does not keep. Reading it is separate work.
const (
// headsSeq — the sequence length the heads were trained at. Padding is
// masked out of both attention and the pool, so this changes nothing but
// truncation, and truncation is what training did at 64.
headsSeq = 64
// headsThreshold — max softmax over the intent head, below which the heads
// decline and the cascade carries on to the resident model.
//
// 0.6 is the knee measured on the 88-case intent fixture
// (docs/evals/2026-08-08-routing-heads-in-go.md). It keeps 81 of 88 cases
// at 97.5% accuracy. Every higher value up to 0.9 drops right answers and
// keeps the same two wrong ones, so it buys nothing.
headsThreshold = 0.6
)
// RouterHeads runs the exported graph. Nil is a working value everywhere: a
// deployment with no weights file routes exactly as it did before this
// existed.
type RouterHeads struct {
tokenizer *unigramTokenizer
session *ort.DynamicSession[int64, float32]
intents []Intent
sources []Source
threshold float64
}
// headsMeta — router_heads.json, written beside the weights by the exporter.
// The label order is the head's output order and cannot be inferred from Go.
type headsMeta struct {
Intents []string `json:"intents"`
Sources []string `json:"sources"`
Prefix string `json:"prefix"`
}
// NewRouterHeads loads the graph and its label order. modelPath points at the
// .onnx; the external weights and router_heads.json sit beside it.
//
// It assumes the ONNX environment is already initialised, because the embedder
// does that at startup and the runtime allows it once.
func NewRouterHeads(modelPath, tokenizerPath string) (*RouterHeads, error) {
metaPath := filepath.Join(filepath.Dir(modelPath), "router_heads.json")
raw, err := os.ReadFile(metaPath)
if err != nil {
return nil, fmt.Errorf("heads: read %s: %w", metaPath, err)
}
var meta headsMeta
if err := json.Unmarshal(raw, &meta); err != nil {
return nil, fmt.Errorf("heads: parse %s: %w", metaPath, err)
}
if meta.Prefix != queryPrefix {
return nil, fmt.Errorf("heads: trained with prefix %q, this build uses %q",
meta.Prefix, queryPrefix)
}
intents := make([]Intent, len(meta.Intents))
for i, s := range meta.Intents {
intents[i] = Intent(s)
}
sources := make([]Source, len(meta.Sources))
for i, s := range meta.Sources {
// SourceUnknown is not in Sources, because it is the absence of a
// choice. It is a class the head can emit, and the one it should emit
// often, so it is allowed here and nowhere else.
if s != string(SourceUnknown) && !ValidSource(Source(s)) {
return nil, fmt.Errorf("heads: unknown destination %q in %s", s, metaPath)
}
sources[i] = Source(s)
}
tok, err := newUnigramTokenizer(tokenizerPath)
if err != nil {
return nil, fmt.Errorf("heads: tokenizer: %w", err)
}
session, err := ort.NewDynamicSession[int64, float32](
modelPath,
[]string{"input_ids", "attention_mask"},
[]string{"intent", "source", "slots", "clarify"},
)
if err != nil {
return nil, fmt.Errorf("heads: create session: %w", err)
}
return &RouterHeads{
tokenizer: tok,
session: session,
intents: intents,
sources: sources,
threshold: headsThreshold,
}, nil
}
func (h *RouterHeads) Close() error {
if h == nil {
return nil
}
h.session.Destroy()
return nil
}
// headsResult — one forward pass, read back.
type headsResult struct {
Intent Intent
Source Source
Confidence float64
Clarify bool
}
// Route runs the heads and reports whether they are confident enough to answer.
// A false second return is a decline, not an error: the cascade goes on to the
// resident model, which is what happens today.
func (h *RouterHeads) Route(ctx context.Context, utterance string) (headsResult, bool, error) {
if h == nil {
return headsResult{}, false, nil
}
ids, mask, _ := h.tokenizer.Encode(queryPrefix + utterance)
ids, mask = ids[:headsSeq], mask[:headsSeq]
// The tokenizer pads and truncates to its own length, which is longer than
// this one. Cutting the tail can cut the separator with it, so put it back.
if mask[headsSeq-1] == 1 {
ids[headsSeq-1] = sepTokenID
}
shape := ort.NewShape(1, headsSeq)
idsT, err := ort.NewTensor(shape, ids)
if err != nil {
return headsResult{}, false, fmt.Errorf("heads: ids tensor: %w", err)
}
defer idsT.Destroy()
maskT, err := ort.NewTensor(shape, mask)
if err != nil {
return headsResult{}, false, fmt.Errorf("heads: mask tensor: %w", err)
}
defer maskT.Destroy()
intentT, err := ort.NewEmptyTensor[float32](ort.NewShape(1, int64(len(h.intents))))
if err != nil {
return headsResult{}, false, fmt.Errorf("heads: intent tensor: %w", err)
}
defer intentT.Destroy()
sourceT, err := ort.NewEmptyTensor[float32](ort.NewShape(1, int64(len(h.sources))))
if err != nil {
return headsResult{}, false, fmt.Errorf("heads: source tensor: %w", err)
}
defer sourceT.Destroy()
slotsT, err := ort.NewEmptyTensor[float32](ort.NewShape(1, headsSeq, int64(numBIOTags)))
if err != nil {
return headsResult{}, false, fmt.Errorf("heads: slots tensor: %w", err)
}
defer slotsT.Destroy()
clarifyT, err := ort.NewEmptyTensor[float32](ort.NewShape(1, 2))
if err != nil {
return headsResult{}, false, fmt.Errorf("heads: clarify tensor: %w", err)
}
defer clarifyT.Destroy()
if err := h.session.Run(
[]*ort.Tensor[int64]{idsT, maskT},
[]*ort.Tensor[float32]{intentT, sourceT, slotsT, clarifyT},
); err != nil {
return headsResult{}, false, fmt.Errorf("heads: run: %w", err)
}
// The graph applies its own softmax, so these are probabilities and the max
// is the same number the eval calibrated the threshold against.
i, conf := argmax(intentT.GetData())
res := headsResult{
Intent: h.intents[i],
Confidence: conf,
}
cl := clarifyT.GetData()
res.Clarify = len(cl) == 2 && cl[1] > cl[0]
// The destination head is trained on query rows and is meaningless on any
// other intent, the same way queryWalk is never reached by one.
if res.Intent == IntentQuery {
s, _ := argmax(sourceT.GetData())
res.Source = h.sources[s]
}
// The clarify head decides on its own, and it decides first. It answers a
// different question from the intent head — not which intent, but whether
// there is enough here to act on at all — so a low intent confidence is no
// reason to discard it. It is usually the same turns: "вода" reads as
// intent act at 0.23 and clarify at 0.98, and letting the intent threshold
// bury that hands the turn to the classifier, which routes it confidently
// and never asks.
if res.Clarify {
return res, true, nil
}
if conf < h.threshold {
return res, false, nil
}
return res, true, nil
}
// numBIOTags — O plus B- and I- for each of Maven's five slots. The head is not
// read, but the graph writes it and the output tensor has to be the right size.
const numBIOTags = 11
func argmax(v []float32) (int, float64) {
best, bestV := 0, math.Inf(-1)
for i, x := range v {
if float64(x) > bestV {
best, bestV = i, float64(x)
}
}
return best, bestV
}
+18 -8
View File
@@ -66,12 +66,19 @@ func NewONNXEmbedder(modelPath, tokenizerPath, libPath string) (*onnxEmbedder, e
func (e *onnxEmbedder) Dim() int { return embedDim }
// ID names the loaded model for the DB marker (Vikunja #378): the model file's
// own name plus the dimension, so pointing the config at another model changes
// the string on its own.
// own name, the dimension, and the tokenizer revision, so pointing the config
// at another model changes the string on its own.
func (e *onnxEmbedder) ID() string { return e.id }
// tokenizerRev — bumped whenever the tokenizer changes what it emits for the
// same text, because that changes every vector while the model file's name
// stays put. Rev 2 is the fix for the reversed word pieces (V-664): stored
// passages embedded under rev 1 no longer sit in the same space as a query
// embedded now, and ReembedAll rewrites them because this string moved.
const tokenizerRev = 2
// modelIDFromPath turns /opt/.../multilingual-e5-small.onnx into
// "multilingual-e5-small@384".
// "multilingual-e5-small@384/tok2".
func modelIDFromPath(modelPath string) string {
name := modelPath
if i := strings.LastIndexAny(name, "/\\"); i >= 0 {
@@ -81,7 +88,7 @@ func modelIDFromPath(modelPath string) string {
if name == "" {
name = "onnx"
}
return fmt.Sprintf("%s@%d", name, embedDim)
return fmt.Sprintf("%s@%d/tok%d", name, embedDim, tokenizerRev)
}
// Embed treats the text as a query. The classifier compares one short
@@ -339,14 +346,17 @@ func (t *unigramTokenizer) encodeWord(word string) []int64 {
}
}
// Backtracking walks the word from its end, and prepending each piece puts
// it back in reading order. There used to be a second reverse after this
// loop, which undid it: every multi-piece word came out backwards, and
// "query: вода" tokenized to [0 12 1294 41 12489 2] where the reference
// tokenizer gives [0 41 1294 12 12489 2] (V-664). A transformer reads
// position, so the pieces of a long Russian word were being read in the
// wrong order on every turn.
var result []int64
for i := n; i > 0; i = prev[i] {
result = append([]int64{bestID[i]}, result...)
}
// Reverse
for l, r := 0, len(result)-1; l < r; l, r = l+1, r-1 {
result[l], result[r] = result[r], result[l]
}
return result
}
+62
View File
@@ -31,6 +31,12 @@ type Config struct {
// error/parse failure, falls through to the classifier (never fails the
// turn on the model).
LLM *LLMRouter
// Heads — optional routing heads over the fine-tuned embedder copy. When
// set, Route consults them after stage 0 and before the LLM router. They
// decline below their own confidence threshold, so a low-confidence turn
// reaches the model exactly as it does today. Nil is the shipped-before
// behaviour and costs nothing.
Heads *RouterHeads
}
// Router — the deterministic cascade. Route never guesses: stage 0 wins
@@ -42,6 +48,7 @@ type Router struct {
extractor Extractor
threshold float64
llm *LLMRouter
heads *RouterHeads
}
func New(cfg Config) *Router {
@@ -51,6 +58,7 @@ func New(cfg Config) *Router {
extractor: cfg.Extractor,
threshold: cfg.Threshold,
llm: cfg.LLM,
heads: cfg.Heads,
}
}
@@ -98,6 +106,60 @@ func (r *Router) Route(ctx context.Context, utterance string, now time.Time) (De
}
r.noteGrammarOutcomes(ctx, len(r.grammars), declinedBuild, "", "")
// stage 0b — routing heads (when wired). A softmax over the label set, so
// it cannot name an intent or a destination that does not exist, and its
// max is a real confidence. It runs before the model because it is three
// orders of magnitude faster and scores better on both halves of the route.
//
// It declines below its threshold rather than clarifying. A declined turn
// carries on to the model and then the classifier, which is what a box with
// no weights file does on every turn.
if r.heads != nil {
res, ok, err := r.heads.Route(ctx, utterance)
switch {
case err != nil:
log.Printf("router: heads fell through to the rest of the cascade: %v", err)
decision.Note(ctx, decision.Claim{
Stage: decision.StageRoute, Claimant: claimantHeads,
Outcome: decision.Declined, Reason: "error: " + err.Error(),
})
case !ok:
decision.Note(ctx, decision.Scored(decision.StageRoute, claimantHeads,
string(res.Intent), res.Confidence, decision.Declined,
"below the heads confidence threshold"))
default:
d := Decision{
Utterance: utterance,
Stage: 2,
Intent: res.Intent,
Confidence: res.Confidence,
Source: res.Source,
Clarify: res.Clarify,
}
r.fillSlots(ctx, &d, now)
decision.Note(ctx, decision.Claim{
Stage: decision.StageRoute, Claimant: claimantLLM,
Outcome: decision.NeverAsked, Reason: "the routing heads answered",
})
decision.Note(ctx, decision.Claim{
Stage: decision.StageRoute, Claimant: claimantClassifier,
Outcome: decision.NeverAsked, Reason: "the routing heads answered",
})
outcome, reason := decision.Won, ""
if d.Clarify {
outcome, reason = decision.Thinned, "the clarify head says there is too little here to act on"
}
decision.Note(ctx, decision.Scored(decision.StageRoute, claimantHeads,
string(d.Intent), d.Confidence, outcome, reason))
return d, nil
}
} else {
decision.Note(ctx, decision.Claim{
Stage: decision.StageRoute, Claimant: claimantHeads,
Outcome: decision.NeverAsked, Reason: "no routing heads are wired",
})
}
// stage 1a — LLM router (when wired). It reasons over the utterance instead
// of nearest-centroid guessing. On any error/parse-fail, fall through to the
// classifier cascade (never fail the turn on the model).
+92
View File
@@ -0,0 +1,92 @@
package stt
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"strconv"
"time"
"github.com/kami/maven/internal/audio"
)
// HTTPTranscriber — speech-to-text on another host, over HTTP.
//
// mavsttd is whisper.cpp linked into a Go daemon and reached over a unix
// socket. CrisperWhisper 2.0 cannot be reached that way: whisper.cpp derives
// its language count from the vocabulary size, and CW2's 51897 tokens shift
// seven special token ids. It runs under transformers instead, as a service
// beside the model on workpc. See docs/evals/2026-08-09-crisperwhisper2-russian-wer.md.
//
// So this is the second transport for the same seam, not a second seam. The
// caller still sees stt.Transcriber and one method.
type HTTPTranscriber struct {
url string
token string
lang string
http *http.Client
}
// NewHTTPTranscriber builds the remote client. token may be empty for a
// service on a trusted socket, but audio is the most sensitive thing that
// crosses this seam, so a LAN deployment should always set one.
func NewHTTPTranscriber(url, token, lang string, timeout time.Duration) *HTTPTranscriber {
return &HTTPTranscriber{
url: url,
token: token,
lang: lang,
http: &http.Client{Timeout: timeout},
}
}
// ErrFormat — the audio is not the one canonical shape. Refused at the seam
// rather than sent to a model that expects something else.
var ErrFormat = errors.New("stt: audio is not 16kHz mono pcm_s16le")
type httpTranscript struct {
Text string `json:"text"`
Confidence float64 `json:"confidence"`
}
// Transcribe posts the raw PCM and reads back the text.
//
// The body is the PCM bytes themselves rather than JSON. A minute of 16kHz
// mono is under 2MB raw and about 2.6MB base64, and the format is fixed by
// audio.PCM16kMono, so a header carries it more cheaply than an envelope.
func (t *HTTPTranscriber) Transcribe(ctx context.Context, a audio.Audio) (string, float64, error) {
if !a.Format.IsValid() {
return "", 0, ErrFormat
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, t.url, bytes.NewReader(a.Bytes))
if err != nil {
return "", 0, fmt.Errorf("stt: build request: %w", err)
}
req.Header.Set("Content-Type", "application/octet-stream")
req.Header.Set("X-Sample-Rate", strconv.Itoa(a.Format.SampleRate))
req.Header.Set("X-Channels", strconv.Itoa(a.Format.Channels))
req.Header.Set("X-Sample-Bits", strconv.Itoa(a.Format.SampleBits))
req.Header.Set("X-Language", t.lang)
if t.token != "" {
req.Header.Set("Authorization", "Bearer "+t.token)
}
resp, err := t.http.Do(req)
if err != nil {
return "", 0, fmt.Errorf("stt: post audio: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return "", 0, fmt.Errorf("stt: remote returned %d", resp.StatusCode)
}
var out httpTranscript
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
return "", 0, fmt.Errorf("stt: decode transcript: %w", err)
}
return out.Text, out.Confidence, nil
}
var _ Transcriber = (*HTTPTranscriber)(nil)
+89
View File
@@ -0,0 +1,89 @@
package stt
import (
"context"
"errors"
"io"
"net/http"
"net/http/httptest"
"strconv"
"testing"
"time"
"github.com/kami/maven/internal/audio"
)
func TestHTTPTranscriberSendsRawPCM(t *testing.T) {
t.Parallel()
var gotBody []byte
var gotHeader http.Header
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotBody, _ = io.ReadAll(r.Body)
gotHeader = r.Header.Clone()
w.Header().Set("Content-Type", "application/json")
_, _ = io.WriteString(w, `{"text":"привет","confidence":0.82}`)
}))
defer srv.Close()
a := audio.Audio{Format: audio.PCM16kMono, Bytes: []byte("pcm-bytes")}
tr := NewHTTPTranscriber(srv.URL, "s3cret", "ru", 2*time.Second)
text, conf, err := tr.Transcribe(context.Background(), a)
if err != nil {
t.Fatalf("Transcribe: %v", err)
}
if text != "привет" || conf != 0.82 {
t.Fatalf("got %q %v", text, conf)
}
if string(gotBody) != "pcm-bytes" {
t.Fatalf("body should be the PCM itself, got %q", gotBody)
}
if got := gotHeader.Get("X-Sample-Rate"); got != strconv.Itoa(audio.PCM16kMono.SampleRate) {
t.Fatalf("X-Sample-Rate = %q", got)
}
if got := gotHeader.Get("X-Language"); got != "ru" {
t.Fatalf("X-Language = %q", got)
}
// Audio is the most sensitive thing crossing this seam.
if got := gotHeader.Get("Authorization"); got != "Bearer s3cret" {
t.Fatalf("Authorization = %q", got)
}
}
func TestHTTPTranscriberOmitsEmptyToken(t *testing.T) {
t.Parallel()
var auth string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
auth = r.Header.Get("Authorization")
_, _ = io.WriteString(w, `{"text":"x"}`)
}))
defer srv.Close()
a := audio.Audio{Format: audio.PCM16kMono, Bytes: []byte("x")}
if _, _, err := NewHTTPTranscriber(srv.URL, "", "ru", time.Second).Transcribe(context.Background(), a); err != nil {
t.Fatalf("Transcribe: %v", err)
}
if auth != "" {
t.Fatalf("Authorization should be absent, got %q", auth)
}
}
func TestHTTPTranscriberRefusesWrongFormat(t *testing.T) {
t.Parallel()
a := audio.Audio{Format: audio.Format{SampleRate: 44100, Channels: 2, SampleBits: 16, Encoding: "pcm_s16le"}}
_, _, err := NewHTTPTranscriber("http://example.invalid", "", "ru", time.Second).Transcribe(context.Background(), a)
if !errors.Is(err, ErrFormat) {
t.Fatalf("want ErrFormat, got %v", err)
}
}
func TestHTTPTranscriberErrorsOnBadStatus(t *testing.T) {
t.Parallel()
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusUnauthorized)
}))
defer srv.Close()
a := audio.Audio{Format: audio.PCM16kMono, Bytes: []byte("x")}
_, _, err := NewHTTPTranscriber(srv.URL, "", "ru", time.Second).Transcribe(context.Background(), a)
if err == nil {
t.Fatal("a 401 must be an error, so the Pair falls back")
}
}
+157
View File
@@ -0,0 +1,157 @@
package stt
import (
"context"
"errors"
"log"
"net/http"
"sync"
"sync/atomic"
"time"
"github.com/kami/maven/internal/audio"
)
// Pair — a preferred transcriber on the workstation, with mavsttd as the floor.
//
// Same arrangement as llm.Pair and for the same reason. The microphone is at
// workpc, the card there has 16GB, and CrisperWhisper 2.0 turbo scores 10.4%
// WER in Russian against 27.5% for the ggml-small.bin homesrv loads
// (docs/evals/2026-08-09-crisperwhisper2-russian-wer.md). The workstation is
// never assumed up: it sleeps, and the card is often held by a training run.
//
// Speech-to-text has only the silent half of the degradation rule. A worse
// transcript is still a turn, and there is nothing to name a gap about, so
// Transcribe always falls back. That is the whole difference from llm.Pair,
// which also carries CompleteRemote for callers that must refuse instead.
type Pair struct {
remote Transcriber
floor Transcriber
// up — the cached admission answer, written only by the prober and read by
// every turn. A voice turn must never wait on a machine that may be asleep.
up atomic.Bool
health string
interval time.Duration
http *http.Client
stop chan struct{}
stopOnce sync.Once
}
const (
probeTimeout = 2 * time.Second
defaultProbeInterval = 15 * time.Second
)
// ErrNoFloor — a Pair was built with no local transcriber to fall back to. A
// configuration mistake: the floor is what makes the remote optional.
var ErrNoFloor = errors.New("stt: no floor transcriber")
// NewPair builds the two-transcriber arrangement. remote may be nil, which is
// the unconfigured deploy: every turn goes to the floor and nothing probes.
func NewPair(remote, floor Transcriber, health string, interval time.Duration) *Pair {
if interval <= 0 {
// The config normalises this, so a zero here is a caller that built the
// Pair directly. Panicking in a ticker is the wrong way to say so.
interval = defaultProbeInterval
}
return &Pair{
remote: remote,
floor: floor,
health: health,
interval: interval,
http: &http.Client{Timeout: probeTimeout},
stop: make(chan struct{}),
}
}
// Start begins probing. The first probe runs before the first tick, so a
// workstation that is already up serves the first utterance rather than the
// second. Safe with a nil remote.
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 and safe from two goroutines.
func (p *Pair) Stop() {
p.stopOnce.Do(func() { close(p.stop) })
}
// Available reports whether the workstation will transcribe right now.
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 transitions. A machine that
// sleeps nightly 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("stt: workstation transcriber available at %s", p.health)
} else {
log.Print("stt: workstation transcriber unavailable, falling back to mavsttd")
}
}
// Transcribe sends the audio to the workstation when it will take work, and to
// mavsttd otherwise. A remote that fails mid-request falls back too, because
// the admission answer is a cache and can be one interval out of date.
//
// Killing the remote mid-session must not drop the turn. That is the whole
// point of the floor, and it is what TestPairFallsBackWhenRemoteFails pins.
func (p *Pair) Transcribe(ctx context.Context, a audio.Audio) (string, float64, error) {
if p.floor == nil {
return "", 0, ErrNoFloor
}
if p.Available() {
text, conf, err := p.remote.Transcribe(ctx, a)
if err == nil {
log.Print("stt: transcribed on the workstation")
return text, conf, nil
}
// The cached answer was wrong. Correct it now rather than sending the
// next utterance into the same hole, then fall back.
p.set(false)
log.Printf("stt: workstation failed mid-request, falling back: %v", err)
}
return p.floor.Transcribe(ctx, a)
}
var _ Transcriber = (*Pair)(nil)
+143
View File
@@ -0,0 +1,143 @@
package stt
import (
"context"
"errors"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"time"
"github.com/kami/maven/internal/audio"
)
// scripted — a Transcriber that answers with a fixed text, or fails.
type scripted struct {
text string
err error
calls atomic.Int32
}
func (s *scripted) Transcribe(_ context.Context, _ audio.Audio) (string, float64, error) {
s.calls.Add(1)
if s.err != nil {
return "", 0, s.err
}
return s.text, 0.9, nil
}
func sample() audio.Audio {
return audio.Audio{Format: audio.PCM16kMono, Bytes: make([]byte, 3200)}
}
// up builds a Pair whose admission answer is already true, without probing.
func up(remote, floor Transcriber) *Pair {
p := NewPair(remote, floor, "", time.Minute)
p.up.Store(true)
return p
}
func TestPairPrefersTheWorkstation(t *testing.T) {
t.Parallel()
remote := &scripted{text: "с рабочей станции"}
floor := &scripted{text: "с homesrv"}
text, _, err := up(remote, floor).Transcribe(context.Background(), sample())
if err != nil {
t.Fatalf("Transcribe: %v", err)
}
if text != "с рабочей станции" {
t.Fatalf("want the remote transcript, got %q", text)
}
if floor.calls.Load() != 0 {
t.Fatalf("floor was called %d times, want 0", floor.calls.Load())
}
}
// The turn is what matters. A remote that dies mid-session must cost a worse
// transcript and nothing else. This is the V-486 bar.
func TestPairFallsBackWhenRemoteFails(t *testing.T) {
t.Parallel()
remote := &scripted{err: errors.New("connection refused")}
floor := &scripted{text: "с homesrv"}
p := up(remote, floor)
text, conf, err := p.Transcribe(context.Background(), sample())
if err != nil {
t.Fatalf("a failed remote must not fail the turn: %v", err)
}
if text != "с homesrv" {
t.Fatalf("want the floor transcript, got %q", text)
}
if conf != 0.9 {
t.Fatalf("want the floor confidence, got %v", conf)
}
if p.Available() {
t.Fatal("a failed request must correct the cached admission answer")
}
// The next utterance goes straight to the floor rather than into the
// same hole.
if _, _, err := p.Transcribe(context.Background(), sample()); err != nil {
t.Fatalf("second turn: %v", err)
}
if remote.calls.Load() != 1 {
t.Fatalf("remote called %d times, want 1", remote.calls.Load())
}
}
func TestPairWithNoRemoteIsTheFloor(t *testing.T) {
t.Parallel()
floor := &scripted{text: "с homesrv"}
p := NewPair(nil, floor, "", time.Minute)
p.Start(context.Background()) // no health url, so this is a no-op
if p.Available() {
t.Fatal("an unconfigured remote is never available")
}
text, _, err := p.Transcribe(context.Background(), sample())
if err != nil {
t.Fatalf("Transcribe: %v", err)
}
if text != "с homesrv" {
t.Fatalf("want the floor transcript, got %q", text)
}
}
func TestPairWithNoFloorRefuses(t *testing.T) {
t.Parallel()
_, _, err := NewPair(nil, nil, "", time.Minute).Transcribe(context.Background(), sample())
if !errors.Is(err, ErrNoFloor) {
t.Fatalf("want ErrNoFloor, got %v", err)
}
}
func TestPairProbeReadsHealth(t *testing.T) {
t.Parallel()
var ok atomic.Bool
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
if !ok.Load() {
w.WriteHeader(http.StatusServiceUnavailable)
return
}
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
p := NewPair(&scripted{text: "remote"}, &scripted{text: "floor"}, srv.URL, time.Minute)
p.probe(context.Background())
if p.Available() {
t.Fatal("a 503 means the card is busy, so the workstation is not available")
}
ok.Store(true)
p.probe(context.Background())
if !p.Available() {
t.Fatal("a 200 means the workstation will take work")
}
}
func TestPairStopIsIdempotent(t *testing.T) {
t.Parallel()
p := NewPair(nil, &scripted{}, "", time.Minute)
p.Stop()
p.Stop()
}