Compare commits

...

2 Commits

Author SHA1 Message Date
kami 7c7bd8ceeb Ship voice enrolment, and report recognition as blocked (#255)
Maven can now be told who someone is. She cannot yet tell who is speaking,
and this commit is careful to say so rather than pretend otherwise.

What works: profiles are enrolled from several deliberately recorded samples,
listed, and deleted. They live in the existing memory_vectors table under a
"speaker:" id prefix, so there is no migration; what that needed was a wider
interface than memory.Store, hence memory.Catalog with ByPrefix and Delete.
Delete is the load-bearing half — a voiceprint someone asked to be rid of has
to actually go, and a search-only store cannot do that. InMemoryStore.Insert
became an upsert by id to match what the persistent store already did.

What does not work, and why it is not faked: there is no speaker-embedding
model on this box. Sixteen ggufs in /mnt/hdd1/llms, all text; no ECAPA, no
x-vector, no titanet, no wespeaker, no .onnx anywhere under /mnt/hdd1. So
newSpeakerEmbedder returns nil, internal/speaker falls back to
speaker.Disabled, Identify answers ErrDisabled, and the daemon logs which
half is off at startup. The plan's "simple MFCC + GMM" floor is refused in
the package comment: MFCC cosine distance detects channel and loudness as
much as voice, and a biometric that is confidently wrong writes false claims
about named people into his memory. A bad floor is worse than none here.

Refused as well, and the reason is in enroll.go's doc comment: the plan asked
for unknown speakers to be enrolled on first interaction with a TTS "кто
это?". There is no request shape in the protocol that could express that.
Taking a biometric of whoever walks past the microphone does it to guests who
are not party to the exchange, and a synthesised question into a room is not
consent from whoever answers.

Authority: enrolment is AuthStepUp, because it is a deliberate sit-down act
that writes a biometric of a named person and never something done by voice
mid-conversation. Deletion is one rung lower at AuthWrite, deliberately
inverting the usual pattern — getting rid of a biometric must never be the
harder half. Listing is AuthRead and never returns the vectors themselves.

Off unless configured: no speaker block means the three methods answer
ErrUnknownMethod, so a default box has no wire path that takes a voiceprint.

make build and make test pass.

Vikunja #255

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TrVSBKe3RFDF4fGYKWYQnX
2026-08-01 05:23:03 +04:00
kami aa1a26532c Add meeting capture with explicit start and stop (#253)
Maven can record a meeting when she is told to, transcribe it through the
STT she already has, and write a summary note. The audio lives in the blob
store #252 introduced, under the same retention loop.

Nothing here listens. Recorder.Append is the only way audio enters and it
refuses every frame unless someone explicitly started a session, so audio
arriving at an idle core is dropped rather than buffered. The plan document
asked for a keyword trigger ("maven record" heard in the room) and that is
refused: noticing a keyword means listening to the room, which is the one
behaviour this capability must not have.

Off unless configured twice over. No media block means nowhere to keep
audio, no capture block means no recorder, and in either case the four IPC
methods answer ErrUnknownMethod. On an unconfigured box there is no wire
path that begins a recording at all.

A forgotten session ends itself at max_minutes, checked on every append,
and the audio collected before the cap is kept. Stop with discard set is
what "забудь, не записывай" maps to and it leaves nothing behind. The
verbatim transcript is not saved unless save_transcript says so; the
summary is.

Long audio against n_ctx 4096 is handled by map-reduce over 3000-rune
windows rather than by truncation, because a truncated meeting summary
reads as complete and is not. Transcription is windowed at five minutes so
the whisper worker stays responsive to the voice path.

No second STT: internal/capture takes the stt.Transcriber the voice path
already holds. Capture with voice off is refused rather than degraded,
since hours of unreadable audio of other people is worse than no recording.

The three write methods are AuthWrite, not AuthStepUp: step-up needs a
passkey gesture the voice path cannot make, which would leave "запиши
встречу" impossible by voice. capture_status is AuthRead.

make build and make test both pass.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TrVSBKe3RFDF4fGYKWYQnX
2026-08-01 05:08:08 +04:00
29 changed files with 3929 additions and 58 deletions
+263
View File
@@ -0,0 +1,263 @@
// mavend/capture.go — core's half of the meeting recorder (Vikunja #253,
// docs/plans/08-hearing.md).
//
// The split: a client that has a microphone (mavenclient, or a phone on the PWA)
// is told to start, streams frames over ipc.MethodCaptureAppend, and is told to
// stop. Core keeps the PCM, stores it as a WAV blob under the same media store
// and the same retention as images, transcribes it through the ONE STT Maven has
// (mavsttd's whisper.cpp, reused — not a second engine), and summarises the
// transcript on the resident model in windows that fit n_ctx 4096.
//
// # Off unless configured, twice over
//
// No `media` block ⇒ nowhere to keep audio ⇒ the four capture methods do not
// exist. No `capture` block with enabled ⇒ they still do not exist. On an
// unconfigured box there is no wire path that starts a recording, which is the
// only guarantee worth making about a capability like this one.
//
// # What this file refuses to do
//
// - Nothing listens. There is no VAD hook here, no wake-word branch, no
// "start when you hear a meeting". The plan document's keyword-triggered
// recorder is refused in internal/capture's package comment for the reason
// that applies here too: noticing a keyword requires listening, which is
// the behaviour this capability must not have.
// - No transcript note by default. The summary is written where he will read
// it; the verbatim record of what other people said takes a deliberate
// capture.save_transcript.
// - The transcript is never search input beyond this box, and the audio never
// leaves it at all.
package main
import (
"context"
"errors"
"fmt"
"log"
"time"
"github.com/kami/maven/internal/capture"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/llm"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/store"
)
// captureSummaryTimeout — the budget for one Stop, which is a map-reduce over
// the whole meeting: one model call per transcript window plus a reduce, each of
// which is seconds on this box. Forty windows is the configured ceiling, so the
// budget has to be minutes, not the 60s the reply path uses.
const captureSummaryTimeout = 20 * time.Minute
// llmCompleter adapts *llm.Client to capture.Completer. The pure package names
// the two strings it needs and stays free of the llm request struct; the client
// itself is the swap-aware one from llmClientFor, so a model swap re-points it.
type llmCompleter struct {
c *llm.Client
maxTokens int
}
func (l llmCompleter) Complete(ctx context.Context, system, user string) (string, error) {
return l.c.Complete(ctx, llm.Req{System: system, User: user, MaxTokens: l.maxTokens})
}
// captureWiring — the recorder plus what it needs to write the result down.
type captureWiring struct {
rec *capture.Recorder
st *store.Store
emb router.Embedder
cfg *config.CaptureConfig
now func() time.Time
}
// newCaptureWiring returns nil when the recorder should not exist: no media
// store, no capture block, capture disabled, or no STT to transcribe with.
//
// A missing llama-server is NOT a reason to return nil. Without one the
// recording is still made, stored and transcribed, and the summary is simply
// absent — the honest degradation, and much better than refusing to record a
// meeting that is happening now.
func newCaptureWiring(keeper *mediaKeeper, st *store.Store, voiceW *voiceWiring, phr phraser.Phraser, emb router.Embedder, cfg *config.Config) *captureWiring {
if keeper == nil || !cfg.Capture.Records() {
return nil
}
tr := transcriberOf(voiceW)
if tr == nil {
// Voice off ⇒ no STT client ⇒ nothing could turn the audio into words.
// Storing hours of unreadable audio of other people is worse than not
// recording, so this is a refusal, not a degradation.
log.Printf("capture: enabled but voice/stt is not wired — meeting capture disabled")
return nil
}
cc := cfg.Capture
var sum *capture.Summarizer
if lp, ok := phr.(*phraser.LLMPhraser); ok {
client := llmClientFor(lp, captureSummaryTimeout)
sum = capture.NewSummarizer(
llmCompleter{c: client, maxTokens: 512},
cc.ChunkRunes, cc.MaxChunks, contextBlockFn(cfg, time.Now),
)
} else {
log.Printf("capture: no llama-server phraser — meetings are transcribed, not summarised")
}
rec, err := capture.New(keeper.store, tr, sum, capture.Config{
MaxDuration: cc.MaxDuration(),
STTWindow: time.Duration(cc.STTWindow),
})
if err != nil {
log.Printf("capture: %v — meeting capture disabled", err)
return nil
}
log.Printf("capture: enabled, sessions capped at %s", rec.MaxDuration())
return &captureWiring{rec: rec, st: st, emb: emb, cfg: cc, now: time.Now}
}
// start handles ipc.MethodCaptureStart.
func (c *captureWiring) start(_ context.Context, req ipc.CaptureStartReq) (ipc.CaptureStartResp, error) {
s, err := c.rec.Start(req.Label)
if err != nil {
return ipc.CaptureStartResp{}, err
}
// The label is logged; nothing that was said ever is.
log.Printf("capture: started %q", s.Label)
return ipc.CaptureStartResp{
Label: s.Label,
Started: s.Started,
MaxSeconds: int(c.rec.MaxDuration().Seconds()),
}, nil
}
// append handles ipc.MethodCaptureAppend. ErrExpired is reported as a successful
// response with Expired set rather than an error: the cap firing is the designed
// behaviour, and the client needs the flag to stop sending and call stop.
func (c *captureWiring) append(_ context.Context, req ipc.CaptureAppendReq) (ipc.CaptureAppendResp, error) {
err := c.rec.Append(req.Audio)
st := c.rec.Status()
if errors.Is(err, capture.ErrExpired) {
log.Printf("capture: %q hit the %s cap — stopping", st.Label, c.rec.MaxDuration())
return ipc.CaptureAppendResp{Seconds: st.Duration.Seconds(), Expired: true}, nil
}
if err != nil {
return ipc.CaptureAppendResp{}, err
}
return ipc.CaptureAppendResp{Seconds: st.Duration.Seconds()}, nil
}
// stop handles ipc.MethodCaptureStop.
//
// The error handling here mirrors vision's, and for the same reason: the audio is
// stored first, so a transcription or summary failure returns what exists rather
// than nothing. A response can carry a blob id with no transcript (STT failed,
// re-runnable), or a transcript with no summary (the model failed, the words are
// kept) — both are degraded successes and neither is an error to the caller.
func (c *captureWiring) stop(ctx context.Context, req ipc.CaptureStopReq) (ipc.CaptureStopResp, error) {
if req.Discard {
// "забудь, не записывай" — nothing is stored, transcribed or noted.
if !c.rec.Abort() {
return ipc.CaptureStopResp{}, capture.ErrNoSession
}
log.Printf("capture: session discarded on request")
return ipc.CaptureStopResp{Discarded: true}, nil
}
res, err := c.rec.Stop(ctx)
resp := ipc.CaptureStopResp{
BlobID: res.BlobID,
Label: res.Label,
Started: res.Started,
Seconds: res.Duration.Seconds(),
Transcript: res.Transcript,
Summary: res.Summary,
Chunks: res.Chunks,
}
if err != nil {
if res.BlobID == "" && res.Transcript == "" {
// Nothing survived: no session, or an empty recording. There is
// nothing to hand back, so this is a real error.
return ipc.CaptureStopResp{}, err
}
log.Printf("capture: %q partially finished: %v", res.Label, err)
}
if id, werr := c.writeNotes(ctx, res); werr != nil {
log.Printf("capture: note write for %q failed: %v", res.Label, werr)
} else {
resp.NoteID = id
}
log.Printf("capture: finished %q — %s of audio, %d summary chunk(s)",
res.Label, res.Duration.Round(time.Second), res.Chunks)
return resp, nil
}
// writeNotes stores the summary as a note, and the transcript too when
// capture.save_transcript is set. Returns the summary note's id, or 0 when there
// was no summary to write.
//
// The note source carries the blob id, which is the only link back to the audio.
// When retention prunes the blob the note remains — words about a meeting are a
// far lighter thing to keep than a recording of it.
func (c *captureWiring) writeNotes(ctx context.Context, res capture.Result) (int64, error) {
source := "capture:meeting"
if res.BlobID != "" {
source = "capture:meeting:" + res.BlobID[:12]
}
var id int64
if text := res.Summary; text != "" {
var err error
id, err = c.writeNote(ctx, text, source)
if err != nil {
return 0, fmt.Errorf("summary note: %w", err)
}
}
if c.cfg.SaveTranscript && res.Transcript != "" {
if _, err := c.writeNote(ctx, res.Transcript, source+":transcript"); err != nil {
return id, fmt.Errorf("transcript note: %w", err)
}
}
return id, nil
}
func (c *captureWiring) writeNote(ctx context.Context, text, source string) (int64, error) {
var vec []float32
if c.emb != nil {
// EmbedPassage, not Embed: this is text being searched FOR, and the e5
// embedder is asymmetric. Backwards here makes the meeting unfindable by
// the question that should have matched it.
var err error
vec, err = router.EmbedPassage(ctx, c.emb, text)
if err != nil {
return 0, fmt.Errorf("embed: %w", err)
}
}
return c.st.WriteNote(ctx, c.now(), text, vec, source)
}
// status handles ipc.MethodCaptureStatus.
func (c *captureWiring) status(_ context.Context) (ipc.CaptureStatusResp, error) {
st := c.rec.Status()
return ipc.CaptureStatusResp{
Running: st.Running,
Label: st.Label,
Started: st.Started,
Seconds: st.Duration.Seconds(),
Bytes: st.Bytes,
}, nil
}
// wireCapture installs the four IPC hooks, or leaves them nil so every capture
// method reports ErrUnknownMethod. Takes the media keeper wireVision already
// opened: one blob store, one retention loop, images and audio side by side.
func wireCapture(srv *ipc.Server, keeper *mediaKeeper, st *store.Store, voiceW *voiceWiring, phr phraser.Phraser, cfg *config.Config) {
cw := newCaptureWiring(keeper, st, voiceW, phr, embedderOf(voiceW), cfg)
if cw == nil {
return
}
srv.CaptureStartFn = cw.start
srv.CaptureAppendFn = cw.append
srv.CaptureStopFn = cw.stop
srv.CaptureStatusFn = cw.status
}
+12
View File
@@ -29,6 +29,7 @@ import (
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/rss"
"github.com/kami/maven/internal/stt"
"github.com/kami/maven/internal/webfetch"
)
@@ -116,6 +117,17 @@ func embedderOf(w *voiceWiring) router.Embedder {
return w.embedder
}
// transcriberOf — the STT the voice path is using, or nil when voice is off.
// The meeting recorder reuses it rather than dialling mavsttd a second time:
// Maven has one speech-to-text engine and adding a second would mean two
// whisper contexts competing for the same iGPU.
func transcriberOf(w *voiceWiring) stt.Transcriber {
if w == nil {
return nil
}
return w.transcriber
}
// feedFetcher adapts webfetch to rss.Fetcher — the pure package names the two
// fields it needs and stays free of net/http.
type feedFetcher struct{ f *webfetch.Fetcher }
+15 -2
View File
@@ -335,7 +335,15 @@ func run(args []string) error {
wireModelSwap(srv, phr, cfg)
// Vision + the media blob store (Vikunja #252). Both stay dark without a
// media block; MethodDescribeImage answers ErrUnknownMethod then.
wireVision(ctx, srv, st, embedderOf(voiceW), cfg)
keeper := wireVision(ctx, srv, st, embedderOf(voiceW), cfg)
// The meeting recorder (Vikunja #253) shares that blob store and its
// retention loop. Off unless a capture block enables it, in which case
// all four capture methods answer ErrUnknownMethod.
wireCapture(srv, keeper, st, voiceW, phr, cfg)
// Voice identification (Vikunja #255). Enrolment plumbing only until a
// speaker-embedding model exists on disk; off entirely without a speaker
// block, so no wire path takes a voiceprint on a default box.
wireSpeaker(srv, st, cfg)
}
// WrapKeyFn — wraps the env key with a passkey credential public key and
@@ -479,7 +487,12 @@ func run(args []string) error {
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
wireMailIntake(srv, st, phr, cfg)
wireModelSwap(srv, phr, cfg)
wireVision(ctx, srv, st, embedderOf(voiceW), cfg)
keeper := wireVision(ctx, srv, st, embedderOf(voiceW), cfg)
wireCapture(srv, keeper, st, voiceW, phr, cfg)
// Voice identification (Vikunja #255). Enrolment plumbing only until a
// speaker-embedding model exists on disk; off entirely without a speaker
// block, so no wire path takes a voiceprint on a default box.
wireSpeaker(srv, st, cfg)
// Start voice server.
if voiceW != nil {
+146
View File
@@ -0,0 +1,146 @@
// mavend/speaker.go — core's half of voice identification (Vikunja #255,
// docs/plans/10-speaker-recognition.md).
//
// # What is actually wired here, and what is not
//
// The enrolment plumbing is real: profiles are stored, listed and deleted, and
// the wire methods exist as soon as a speaker block is configured. The
// recognising half is NOT, and cannot be on this box, because there is no
// speaker-embedding model on disk — no ECAPA, no x-vector, no titanet, no
// wespeaker, nothing in /mnt/hdd1/llms but text ggufs. Until one is downloaded,
// newSpeakerEmbedder returns nil, internal/speaker falls back to
// speaker.Disabled, and every Identify answers ErrDisabled. The daemon logs
// which half is off at startup rather than pretending.
//
// This is deliberately not papered over with a hand-rolled MFCC floor. A
// biometric that is confidently wrong writes false claims about named people
// into his memory, and that is worse than a capability that is honestly absent.
//
// # Off unless configured
//
// No speaker block, or one without enabled, ⇒ the three methods do not exist and
// answer ErrUnknownMethod. On an unconfigured box there is no wire path that
// takes a voiceprint at all.
//
// # The refused design step
//
// The plan asks for unknown speakers to be enrolled on first interaction. That
// is refused in internal/speaker/enroll.go and there is no handler for it here:
// no request shape in the protocol enrols whoever just spoke. Taking a biometric
// of a guest who walked past the microphone is not something this daemon does.
package main
import (
"context"
"errors"
"log"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/speaker"
"github.com/kami/maven/internal/store"
)
// speakerWiring holds the recognizer behind the three IPC handlers.
type speakerWiring struct {
rec *speaker.Recognizer
}
// newSpeakerEmbedder loads the speaker-embedding model named by the config.
//
// It always returns nil today. The seam exists so that wiring a real model is a
// change to this one function and nothing else: give it a loader, and Identify
// starts working with no change to the store, the protocol, the auth table or
// the handlers. See the plan document for what to download.
func newSpeakerEmbedder(cfg *config.SpeakerConfig) speaker.Embedder {
if cfg == nil || cfg.ModelPath == "" {
return nil
}
log.Printf("speaker: model_path %q is configured but no embedding backend is built yet; "+
"enrolment and deletion work, recognition does not (Vikunja #255)", cfg.ModelPath)
return nil
}
// newSpeakerWiring builds the recognizer, or nil when the capability is off.
func newSpeakerWiring(st *store.Store, cfg *config.Config) *speakerWiring {
if cfg == nil || cfg.Speaker == nil || !cfg.Speaker.Enabled {
return nil
}
if st == nil {
log.Print("speaker: enabled but there is no store to keep profiles in; staying off")
return nil
}
rec, err := speaker.New(newSpeakerEmbedder(cfg.Speaker), st.VectorMemory(), speaker.Config{
Threshold: cfg.Speaker.Threshold,
MinSeconds: cfg.Speaker.MinSeconds,
})
if err != nil {
log.Printf("speaker: %v; staying off", err)
return nil
}
if rec.Enabled() {
log.Printf("speaker: recognition on, threshold %.2f", rec.Threshold())
} else {
log.Print("speaker: enrolment on, recognition BLOCKED — no speaker-embedding model " +
"on this box (see docs/plans/10-speaker-recognition.md)")
}
return &speakerWiring{rec: rec}
}
func (w *speakerWiring) enroll(ctx context.Context, req ipc.EnrollSpeakerReq) (ipc.EnrollSpeakerResp, error) {
p, err := w.rec.Enroll(ctx, req.ID, req.Name, req.Samples)
if err != nil {
return ipc.EnrollSpeakerResp{}, speakerErr(err)
}
return ipc.EnrollSpeakerResp{Speaker: toWireSpeaker(p)}, nil
}
func (w *speakerWiring) list(ctx context.Context) (ipc.ListSpeakersResp, error) {
ps, err := w.rec.List(ctx)
if err != nil {
return ipc.ListSpeakersResp{}, speakerErr(err)
}
out := make([]ipc.Speaker, 0, len(ps))
for _, p := range ps {
out = append(out, toWireSpeaker(p))
}
return ipc.ListSpeakersResp{Speakers: out, Enabled: w.rec.Enabled()}, nil
}
func (w *speakerWiring) forget(ctx context.Context, req ipc.ForgetSpeakerReq) error {
return speakerErr(w.rec.Forget(ctx, req.ID))
}
// toWireSpeaker drops the voiceprint. A listing says who is enrolled; it does
// not hand the biometric back out over the socket.
func toWireSpeaker(p speaker.Profile) ipc.Speaker {
return ipc.Speaker{ID: p.ID, Name: p.Name, Enrolled: p.Enrolled, Samples: p.Samples}
}
// speakerErr maps the package sentinels onto the wire vocabulary so a surface
// can tell "you asked wrong" from "core broke".
func speakerErr(err error) error {
switch {
case err == nil:
return nil
case errors.Is(err, speaker.ErrNotFound):
return ipc.ErrNoFact
case errors.Is(err, speaker.ErrBadID),
errors.Is(err, speaker.ErrBadFormat),
errors.Is(err, speaker.ErrTooShort):
return errors.Join(ipc.ErrBadParams, err)
default:
return err
}
}
// wireSpeaker attaches the three handlers when the capability is configured.
func wireSpeaker(srv *ipc.Server, st *store.Store, cfg *config.Config) {
w := newSpeakerWiring(st, cfg)
if w == nil {
return
}
srv.EnrollSpeakerFn = w.enroll
srv.ListSpeakersFn = w.list
srv.ForgetSpeakerFn = w.forget
}
+9 -3
View File
@@ -236,16 +236,22 @@ func sourceOrDefault(s string) string {
// hook nil so ipc.MethodDescribeImage reports ErrUnknownMethod. Called on both
// startup paths (unlocked boot and passkey unlock) so vision behaves the same
// either way.
func wireVision(ctx context.Context, srv *ipc.Server, st *store.Store, emb router.Embedder, cfg *config.Config) {
//
// Returns the media keeper so the meeting recorder can share it: one blob store
// with one retention loop holds both the images and the audio, which is the
// whole point of internal/media being a shared package. nil ⇒ no media block,
// and neither capability exists.
func wireVision(ctx context.Context, srv *ipc.Server, st *store.Store, emb router.Embedder, cfg *config.Config) *mediaKeeper {
keeper := openMediaStore(cfg)
if keeper == nil {
return
return nil
}
go keeper.runPrune(ctx)
vi := newVisionIntake(keeper, st, emb, cfg)
if vi == nil {
return
return keeper
}
srv.DescribeImageFn = vi.describe
return keeper
}
+6
View File
@@ -40,6 +40,11 @@ type voiceWiring struct {
// mavsttd / mavttsd don't keep a stale conn into a restarting daemon.
sttClient *worker.Client
ttsClient *worker.Client
// transcriber — the STT in use, exposed so the meeting recorder
// (cmd/mavend/capture.go) can reuse it. Maven has exactly one STT and does
// not grow a second one for capture: this is the same whisper.cpp worker the
// voice path talks to.
transcriber stt.Transcriber
// mcp — the MCP client, nil unless the `mcp` block configures an enabled
// server (Vikunja #251). Its tools land in the same allowlist as every
// other act, so nothing else here has to know about it.
@@ -92,6 +97,7 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
} else {
transcriber = stt.NewStub()
}
w.transcriber = transcriber
// ----- tts (Stub in-process OR Remote) -----
var synthesizer tts.Synthesizer
+117 -24
View File
@@ -1,28 +1,121 @@
# Plan: Hearing — Audio Stream Monitoring & Meeting Summarization
# Plan: Hearing — Meeting Capture & Summarisation
**Goal:** Maven can "hear" ambient audio from workpc — microphone input during meetings, system audio — and on demand (or on trigger) produce transcripts, summaries, or extract action items. A typical use case: "Maven, запиши встречу" starts capture, "хватит" stops it, and Maven writes a summary note.
**Goal:** "Maven, запиши встречу" starts a recording, "хватит" stops it, and she writes a
summary note. The audio stays on the box, is pruned by retention, and nothing is recorded that
nobody asked for.
**Done when:**
- `internal/audio/capture.go` — remote microphone capture client (receives PCM stream from workpc over WebSocket or the existing voice TCP protocol)
- `internal/stt/` — streaming transcription (uses existing `stt.Transcriber` interface, extended with streaming support)
- Meeting capture triggered by voice command (IntentCapture) or configurable keyword ("maven record")
- Raw audio is either streamed to STT in real-time or saved to a WAV file and transcribed after capture ends
- Transcription + LLM summary is written as a note (`source:capture:meeting`) through `ipc.CoreAPI`
- New `mavheary` module (`cmd/mavheard/`) — the workpc-side agent that captures mic/speaker audio and streams it to mavend
**Status (2026-08-01):** the recorder, the storage, the chunked transcription, the map-reduce
summariser, the config seam and the four IPC methods are shipped and tested. What is not
shipped is the workpc-side microphone agent and the router intent — see "Still open".
**Scope:**
- New `cmd/mavheard/` — workpc-side agent: captures microphone (PortAudio or ALSA `arecord`), streams over WebSocket to mavend
- `internal/audio/` extended with capture types: `MicCapture`, `SystemCapture`, `FileCapture`
- `internal/stt/stt.go` extended with `StreamingTranscriber` interface (or reuse existing with chunked input)
- Router: new `IntentCapture` intent for start/stop commands
- Reuses `internal/llm.Client` for summarization
- Reuses `internal/voice/server.go` TCP protocol for streaming audio
## What shipped
**Steps:**
1. Create `cmd/mavheard/main.go` — workpc-side daemon: captures microphone via `arecord` pipe or PortAudio, opens WebSocket or TCP connection to mavend, streams PCM frames
2. Create `internal/audio/capture.go``Capture` interface: `Start()`, `Stop()`, `AudioCh <-chan Audio`; implement `MicCapture` (reads from `mavheard` stream) and `FileCapture` (reads WAV)
3. Extend `internal/stt/stt.go` — add `TranscribeStream(ctx, audio <-chan Audio) (string, error)` to `Transcriber` interface; `Stub` returns empty; `Remote` forwards chunks to worker socket
4. Add `IntentCapture` to `internal/router/intent.go` — slots: `Action` ("start"/"stop"/"status"), `Duration`
5. Wire capture handler in `cmd/mavend/voice.go:reactiveHandler` — start = spawn goroutine receiving audio, stream to STT; stop = finalize, send to LLM for summarization, write note via `WriteNote`
6. Add capture config to `voice` block in `config.Config` `{capture_enabled, capture_timeout}`
7. Test with a recorded WAV file — simulate a meeting, verify transcription + summary note is created
| Piece | Where |
|---|---|
| Session state machine: start / append / stop / abort / status | `internal/capture/capture.go` |
| Map-reduce summarisation against `n_ctx` 4096 | `internal/capture/summarize.go` |
| Audio blobs in the shared store, pruned by `media.retention` | `internal/media` (from #252) |
| Config block `capture`, off by default | `internal/config/config.go` |
| IPC `capture_start` / `capture_append` / `capture_stop` / `capture_status` | `internal/ipc/{wire,api,client,server}.go` |
| Authority: the three write methods `AuthWrite`, status `AuthRead` | `internal/auth/policy.go` |
| Daemon wiring, note write, STT reuse | `cmd/mavend/capture.go` |
The audio lands in the same content-addressed blob store as images, under the same retention
loop, because #252 and #253 have the same intake problem and solving it twice would mean two
directories to remember to prune.
## The refusals, and why
**Nothing listens.** The original step 8 called for capture "triggered by voice command
(IntentCapture) **or configurable keyword ('maven record')**". The keyword half is refused.
Noticing a keyword requires listening to the room continuously, which is precisely the
behaviour this capability must not have, and the refusal is in the code rather than in a
comment: `Recorder.Append` is the only way audio enters, and it returns `ErrNoSession` unless
someone explicitly started a session. Audio arriving at an idle core is dropped, not buffered
"just in case".
**Off unless configured, twice over.** No `media` block ⇒ nowhere to keep audio ⇒ the four
methods do not exist. No `capture` block with `enabled: true` ⇒ they still do not exist. On an
unconfigured box there is no wire path at all that begins a recording. That is the only
guarantee worth making here, and it is the reason the hooks use the nil-hook ⇒
`ErrUnknownMethod` pattern rather than an in-handler check.
**A forgotten session ends itself.** `max_minutes` defaults to 120 and is checked on every
append, not on a timer that could be missed. Past the cap `Append` returns `ErrExpired`
permanently, so a client that ignores the error cannot grow the recording; the audio collected
before the cap is kept and `Stop` still works.
**"Забудь, не записывай" leaves nothing behind.** `capture_stop` with `discard: true` throws
the session away without storing, transcribing or summarising anything — not a blob with a note
saying it was abandoned. Nothing.
**The transcript is not saved by default.** The summary is written where he will read it; the
verbatim record of what other people said in a room is a heavier thing to keep and takes a
deliberate `save_transcript: true`. The audio blob is pruned by `media.retention` either way.
**No second STT.** Step 3 of the original plan extended the `Transcriber` interface with
streaming. Not needed and not done: whisper.cpp already runs as `mavsttd`, and `internal/capture`
takes the ordinary `stt.Transcriber` the voice path already holds (exposed as
`voiceWiring.transcriber`). Long recordings are handed over in five-minute windows —
`chunkAudio`, cut on sample boundaries — for the same reason whisper itself works in 30-second
windows: an hour of PCM in one call either times out or blocks the voice path for minutes.
Capture with voice off is refused rather than degraded, because storing hours of unreadable
audio of other people is worse than not recording.
**Not `AuthStepUp`.** Recording people is invasive enough to argue for the top rung, and it is
still wrong: step-up needs a passkey gesture, which the voice path cannot make, so
"запиши встречу" could never work by voice — the only way he will actually use this. `AuthWrite`
plus the off-unless-configured gate is the honest combination.
## Long audio against a 4096-token context
The resident model is a Thinking variant at `n_ctx` 4096, so an hour of transcript does not fit
in one prompt and never will. `summarize.go` does map-reduce and nothing cleverer: split the
transcript on sentence boundaries into 3000-rune windows (about 1100 Qwen tokens of Russian,
leaving room for the persona block, the reasoning and the answer), summarise each, then
summarise the summaries. A transcript that fits in one window skips the reduce step.
Truncation was the alternative and is rejected: a truncated meeting summary reads as complete
and is not, and he would act on it. Past `max_chunks` (40, roughly the two-hour cap) the
transcript *is* cut, and the summary says so in the note.
Two degradations are deliberate and both are reported rather than hidden:
- No llama-server ⇒ transcript, no summary. The words exist.
- The reduce call fails ⇒ the per-chunk summaries are returned joined. Real work, not thrown
away over the last call.
The map and reduce prompts contain no first person at all, so the persona's feminine-form rules
have nothing to get wrong in them; the reply she actually gives him is phrased by the ordinary
replier, which does carry the persona.
## Config
```json
"media": { "dir": "media", "retention": "168h" },
"capture": {
"enabled": true,
"max_minutes": 120,
"stt_window": "5m",
"chunk_runes": 3000,
"max_chunks": 40,
"save_transcript": false
}
```
Both absent by default. `capture` alone does nothing without `media`.
## Still open
- **`cmd/mavheard`** — the workpc-side microphone agent. Deferred, not refused: the core half
is the part with the invariants in it, and a mic client is straightforward once there is a
stable wire to stream at. It should be an explicit-start process, not a resident one, for the
same reason the recorder has no keyword trigger. The four IPC methods are the wire it will
use; `mavenclient` already has the mic plumbing to borrow.
- **Router intent.** "запиши встречу" / "хватит" does not route anywhere yet. It needs the
`system` intent plus slots, and it needs care: "хватит" is also how someone tells her to stop
talking, so the recorder's stop and the speech barge-in must not collide.
- **A `/dash` panel** showing a running session, so a recording is visible on a surface and not
only in a log line.
- **Speaker attribution** — who said what — is #255 and is blocked on a model; see
`docs/plans/10-speaker-recognition.md`.
+120 -22
View File
@@ -1,27 +1,125 @@
# Plan: Speaker Recognition
**Goal:** Maven can distinguish between different speakers on the voice channel — recognize known voices (the user, family members) and tag facts/notes/transcripts with a speaker identity.
**Goal:** Maven can tell who is speaking on the voice channel, and tag what she writes with
who said it.
**Done when:**
- Speaker embedding extractor (e.g., ECAPA-TDNN or a simple MFCC + GMM) runs on incoming voice PCM before STT
- Embedding is compared against enrolled speaker profiles (stored as vectors in the `memory_vectors` table alongside semantic memory)
- Unknown speakers are enrolled on first interaction (prompt: "кто это?")
- All voice fact/note writes are tagged with `speaker:<id>` in the value/source metadata
- Speaker identity is available as context to the router, phraser, and replier ("ok, <name>")
**Status (2026-08-01, Vikunja #255):** the enrolment half is shipped. The recognising half is
**BLOCKED on a model download** — there is no speaker-embedding model on this box, and one
was not invented to fill the gap. See "Blocked, and on what" below.
**Scope:**
- New `internal/speaker/` package — enrollment, recognition, embedding extraction
- Reuses `internal/store.MemoryStore` for speaker vector storage (same `memory_vectors` table, different `source` prefix)
- Reuses `internal/audio` for PCM preprocessing
- Integration point: `cmd/mavend/voice.go:HandlePushToTalk` — speaker ID extracted before STT, passed through context
## What shipped
**Steps:**
1. Research speaker embedding approaches — simplest floor: MFCC + cosine similarity via `github.com/mjibson/go-dsp` or a pre-trained ONNX model (SpeechBrain ECAPA)
2. Create `internal/speaker/recognizer.go``Recognizer` interface: `Identify(pcm []float32) (SpeakerID, confidence)`, `Enroll(id, pcm)`
3. Create `internal/speaker/store.go` — speaker profile CRUD via `store.MemoryStore`: `Insert("speaker:<id>", embedding, meta)`, `Search(embedding, k)`
4. Create `internal/speaker/enroll.go` — enrollment flow: capture N seconds of audio, extract embedding, prompt for name via TTS + STT round-trip
5. Wire into `cmd/mavend/voice.go:HandlePushToTalk` — run speaker ID on the PCM before STT; pass speaker ID through `context.Context` to `applyAction`
6. Tag all voice-written facts/notes with speaker ID — `Source` becomes `tap:voice:speaker:<id>` or metadata field
7. Add IPC methods `MethodEnrollSpeaker`, `MethodListSpeakers`, `MethodRemoveSpeaker`
8. Add speaker config block to `voice` in `config.Config``{speaker_recognition: true, model_path}`
9. Test with 2+ recorded voice samples — verify correct identification and rejection of unknown speakers
| Piece | Where | State |
|---|---|---|
| `Recognizer` — identify, list, get, forget | `internal/speaker/recognizer.go` | done; `Identify` answers `ErrDisabled` until a model exists |
| Enrolment — several samples, averaged, re-normalised | `internal/speaker/enroll.go` | done |
| Profile shape, id validation, cosine similarity | `internal/speaker/speaker.go` | done |
| Profile storage as `speaker:<id>` vectors | `internal/memory` `Catalog` + `internal/store/memory.go` | done, no schema migration |
| Config block, off by default | `internal/config` `SpeakerConfig` | done |
| `enroll_speaker` / `list_speakers` / `forget_speaker` | `internal/ipc` | done, absent unless configured |
| Authority rows | `internal/auth/policy.go` | done — enrol step-up, forget write, list read |
| Daemon wiring + honest startup log | `cmd/mavend/speaker.go` | done |
| Embedding backend | `newSpeakerEmbedder` | **BLOCKED** — returns nil, seam only |
| Tagging voice writes with the speaker | `cmd/mavend/voice.go` | not wired; nothing to tag with yet |
## Blocked, and on what
A voiceprint needs a speaker-embedding model. The box was searched: `/mnt/hdd1/llms` holds
sixteen ggufs across seven families and every one of them is a text model. There is no ECAPA,
no x-vector, no titanet, no wespeaker, and no `.onnx` under `/mnt/hdd1` at all. There are also
no enrolment samples, because nothing has ever recorded any.
To unblock, two things are needed and neither can be done from inside the repo:
1. **A model.** SpeechBrain ECAPA-TDNN exported to ONNX (`speechbrain/spkrec-ecapa-voxceleb`,
192-dim) is the usual choice and runs on CPU in well under a second for a few seconds of
audio. Download it per the recipe in `AGENTS.md`, put it beside the other models so the
bind mount picks it up, and point `speaker.model_path` at it.
2. **An implementation of one function.** `newSpeakerEmbedder` in `cmd/mavend/speaker.go` is
the entire seam: give it an ONNX session that turns `audio.Audio` into a `[]float32` and
`Identify` starts working. Nothing else changes — not the store, not the protocol, not the
authority table, not the handlers. `internal/onnx` already loads the e5 embedder, so the
runtime wiring exists to copy.
3. **Enrolment samples**, three or more per person, recorded deliberately.
### Why there is no fallback
The original plan offered "a simple MFCC + GMM" as the floor. That is refused. MFCC cosine
distance is a channel and loudness detector as much as a voice detector: it will happily match
two different people who sit at the same distance from the same microphone, and it drifts when
the room changes. A general classifier that is sometimes wrong is a nuisance; a **biometric**
that is confidently wrong writes false claims about named people into his memory, and then
those claims get recalled as fact. For this capability a bad floor is worse than none, so the
shipped state is honest absence: `speaker.Disabled`, `ErrDisabled`, and a startup line saying
so.
## The refusals, and why
- **Unknown speakers are NOT enrolled on first interaction.** The plan's fourth "done when"
bullet asked for exactly that, with a TTS "кто это?" prompt. It is refused in
`enroll.go`'s doc comment and there is no request shape in the protocol that could express
it. Enrolling a voice is taking a biometric of a person; doing it automatically to whoever
walks past the microphone does it to guests who are not party to the exchange, and a
synthesised question into a room is not consent from whoever happens to answer. Enrolment is
an explicit act: an id, a name, and samples recorded for the purpose.
- **One sample is not enough.** Three separate utterances and nine seconds minimum. A profile
built from one sentence encodes that sentence as much as the person, and the threshold then
behaves unpredictably against everything else.
- **An unknown voice stays unknown.** Below threshold, `Identify` returns `ErrUnknown` naming
the closest profile in the error text for diagnosis, never as an answer. Guessing who is in
the room is how false memories about people get written.
- **Deletion is one authority rung below enrolment.** Everywhere else in `policy.go` the
destructive direction is gated at least as hard as the constructive one. Here that would be
backwards: getting rid of a biometric must never be the harder half.
- **The voiceprint never crosses the socket.** `ListSpeakersResp` carries ids, names, dates
and sample counts. The vector stays in core.
- **Off unless configured.** No `speaker` block ⇒ the three methods answer
`ErrUnknownMethod`. There is no wire path on a default box that takes a voiceprint.
## Storage
Profiles live in the existing `memory_vectors` table under the `speaker:` id prefix, as the
plan intended, so there is no migration. What that needed was a wider interface than
`memory.Store`: `memory.Catalog` adds `ByPrefix` and `Delete`. `Delete` is the load-bearing
one — a voiceprint someone asked to be rid of has to actually go, and a search-only store
cannot do that. `InMemoryStore.Insert` also became an upsert by id, matching what the
persistent store already did, so re-enrolling replaces a profile instead of stacking a second
one behind the first.
Profiles do not collide with note or fact vectors: they are only ever read through
`ByPrefix("speaker:")`, and a note search never returns one because the prefix is not in its
query path.
## Config
```json
"speaker": {
"enabled": true,
"model_path": "/opt/maven/models/spk/ecapa-voxceleb.onnx",
"lib_path": "/opt/maven/lib",
"threshold": 0.7,
"min_seconds": 2.0
}
```
`Recognizes()` requires both `enabled` and a `model_path`, so a half-filled block reads as off
rather than as a capability that fails every turn. With `enabled` and no model the daemon still
attaches the three methods — profiles can be created, listed and deleted — and logs that
recognition is blocked.
## Still open
- The embedding backend (above). Everything below waits on it.
- **Tagging voice writes.** `Profile.Source("tap:voice")` already produces
`tap:voice:speaker:kami`, which is the shape step 6 asked for, but nothing calls it yet:
with no recogniser there is no id to tag with. When the model lands, the hook is in the
voice path before STT.
- **Speaker as router/phraser context.** Same dependency. Note the persona constraint when it
arrives: Maven addresses the owner informally and speaks to him, so "ok, <name>" needs care
for anyone who is not him.
- **An enrolment surface.** The three IPC methods exist; no page drives them. Enrolment is
step-up, so it belongs on `/dash` behind a passkey, with a per-profile forget button next to
each row — that button is the reason `list_speakers` exists.
- **A speaker column on the meeting recorder** (#253). Attributing lines in a transcript is
the obvious pairing, and it is the place where getting attribution wrong is most damaging,
so it waits for a real model too.
+56
View File
@@ -416,3 +416,59 @@ func TestRequirement_SwapModel(t *testing.T) {
t.Errorf("SwapModel with no asserted step-up = %v; want ErrForbidden", err)
}
}
// TestRequirement_Capture — recording other people is a write, not a read: it
// puts audio of them on disk. The read side, "что ты записываешь?", is not.
//
// It is deliberately NOT AuthStepUp. Step-up needs a passkey gesture, which the
// voice path cannot make, so putting it there would mean "запиши встречу" could
// never work by voice. The real gate on this capability is that the methods do
// not exist at all unless the operator enabled a capture block.
func TestRequirement_Capture(t *testing.T) {
for _, m := range []ipc.Method{
ipc.MethodCaptureStart, ipc.MethodCaptureAppend, ipc.MethodCaptureStop,
} {
if got := Requirement(m); got != AuthWrite {
t.Errorf("%s authority = %v; want AuthWrite", m, got)
}
}
if got := Requirement(ipc.MethodCaptureStatus); got != AuthRead {
t.Errorf("CaptureStatus authority = %v; want AuthRead", got)
}
// Voice can start one: it is the surface he will actually use to say
// "запиши встречу", and it carries AuthWrite.
voice := Scope{Surface: SurfaceVoice, Module: "voice", SourceScope: []string{"*"}}
if err := Can(ipc.MethodCaptureStart, voice, nil); err != nil {
t.Errorf("voice starting a capture = %v; want allowed", err)
}
}
// TestRequirement_Speaker — a voiceprint is a biometric of a named person, so
// taking one is step-up: a deliberate act from a surface that can carry a
// passkey gesture, never something the voice path does mid-conversation.
//
// Deletion is one rung lower, and that asymmetry is the point. Everywhere else
// in the table the destructive direction is gated at least as hard as the
// constructive one; for a biometric that would be backwards, because getting
// rid of it must never be the harder half.
func TestRequirement_Speaker(t *testing.T) {
if got := Requirement(ipc.MethodEnrollSpeaker); got != AuthStepUp {
t.Errorf("EnrollSpeaker authority = %v; want AuthStepUp", got)
}
if got := Requirement(ipc.MethodForgetSpeaker); got != AuthWrite {
t.Errorf("ForgetSpeaker authority = %v; want AuthWrite", got)
}
if got := Requirement(ipc.MethodListSpeakers); got != AuthRead {
t.Errorf("ListSpeakers authority = %v; want AuthRead", got)
}
// Voice cannot enrol anybody, however the utterance is phrased.
voice := Scope{Surface: SurfaceVoice, Module: "voice", SourceScope: []string{"*"}}
if err := Can(ipc.MethodEnrollSpeaker, voice, nil); err == nil {
t.Error("voice enrolling a speaker was allowed; want refused")
}
// But it can read the roster, which is what answering "кого ты знаешь?"
// needs.
if err := Can(ipc.MethodListSpeakers, voice, nil); err != nil {
t.Errorf("voice listing speakers = %v; want allowed", err)
}
}
+38
View File
@@ -60,6 +60,37 @@ func Requirement(m ipc.Method) Authority {
// and for the same reason: nothing Maven says or does may reach it.
// MethodModelStatus is only the read side, so it stays at AuthRead.
return AuthStepUp
case ipc.MethodCaptureStart, ipc.MethodCaptureAppend, ipc.MethodCaptureStop:
// Recording a meeting (Vikunja #253). AuthWrite, not AuthRead: it puts
// audio of other people on disk, which is a heavier thing than reading a
// fact, and it is not something a read-only surface should be able to
// begin. Append and Stop sit on the same rung as Start deliberately —
// a surface that may not start a recording has no business feeding or
// harvesting one either.
//
// Not AuthStepUp, and this is the interesting line: step-up needs a
// passkey gesture, which the voice path cannot make. Putting it here
// would mean "запиши встречу" could never work by voice, and the real
// gate on this capability is elsewhere and stronger — the methods do not
// exist at all unless the operator enabled a capture block, and no
// recording can begin without someone saying so.
return AuthWrite
case ipc.MethodEnrollSpeaker:
// Taking a voiceprint (Vikunja #255). AuthStepUp, and unlike recording a
// meeting there is no reason to soften it: enrolment is not a thing anyone
// does by voice mid-conversation. It is a deliberate sit-down with a
// surface that can carry a passkey gesture, and it writes a biometric of a
// named person. If the gesture is inconvenient, that is the correct amount
// of friction for this particular write.
return AuthStepUp
case ipc.MethodForgetSpeaker:
// Deleting a voiceprint. One rung BELOW enrolment on purpose. Everywhere
// else in this table the destructive direction is gated at least as hard
// as the constructive one, and here that would be wrong: getting rid of a
// biometric must never be the harder half. The worst a caller at this rung
// can do is make Maven stop recognising someone, which is the state the
// box ships in anyway.
return AuthWrite
case ipc.MethodWriteFact:
return AuthWrite
case ipc.MethodAssertStepUp:
@@ -95,6 +126,13 @@ func Requirement(m ipc.Method) Authority {
// the box, which internal/vision enforces by refusing a non-private
// endpoint.
ipc.MethodDescribeImage,
// "что ты записываешь?" — the read side of the recorder. It reports a
// label, a start time and a byte count, begins nothing and keeps nothing.
ipc.MethodCaptureStatus,
// Who is enrolled. Returns ids, names and enrolment dates — never the
// voiceprints themselves, which stay in core. Listing the people Maven can
// recognise is exactly the read a surface needs to offer a "forget" button.
ipc.MethodListSpeakers,
// The read side of the model swap: which model is resident, which ones are
// allowlisted. It loads nothing and changes nothing.
ipc.MethodModelStatus:
+404
View File
@@ -0,0 +1,404 @@
// Package capture is Maven's meeting recorder (Vikunja #253,
// docs/plans/08-hearing.md).
//
// One session at a time, with an explicit start and an explicit stop:
//
// Start("встреча") → audio frames appended → Stop() → transcript → summary
//
// # Nothing here listens
//
// This is the most invasive capability in the backlog and the design is
// constrained accordingly. The constraints are the code, not a preamble:
//
// - There is no ambient path. `Session.Append` is the only way audio enters,
// and it only accepts frames while a session someone started is running.
// A keyword-triggered recorder ("maven record" heard in the room) was in the
// plan document and is refused: it requires listening in order to notice the
// keyword, which is the exact behaviour this capability must not have.
// - A session that is not stopped stops itself. MaxDuration is a hard cap
// checked on every Append, not a suggestion; a forgotten recording is a
// recording that ends, not one that runs until the disk is full.
// - Audio is stored under internal/media, which means retention prunes it and
// it never leaves the box. Both the audio blob and the transcript stay
// local; only the summary is written where he will read it.
// - The transcript is never search input for anything outside this box. It is
// text about a conversation with other people in it.
//
// # Long audio against a 4096-token context
//
// The resident model is a Thinking variant at n_ctx 4096, so an hour of meeting
// transcript does not fit in one prompt and never will. summarize.go does the
// obvious map-reduce: split the transcript on sentence boundaries into windows
// that fit, summarise each, then summarise the summaries. That is handled
// explicitly rather than by truncation, because a truncated meeting summary is
// worse than none — it looks complete and is not.
//
// # Transcription
//
// There is exactly one STT in Maven and this package does not add a second: it
// takes an stt.Transcriber, which in deploy is the whisper.cpp worker behind
// cmd/mavsttd. Long audio is transcribed in windows too (see chunkAudio), for
// the same reason whisper itself works in 30s windows — handing a worker an hour
// of PCM in one call is a request that either times out or blocks everything
// else for minutes.
package capture
import (
"context"
"errors"
"fmt"
"strings"
"sync"
"time"
"github.com/kami/maven/internal/audio"
"github.com/kami/maven/internal/media"
"github.com/kami/maven/internal/stt"
)
// DefaultMaxDuration — how long one capture may run before it stops itself.
// Two hours covers a long meeting and bounds the damage of a forgotten session:
// at 16 kHz mono that is about 230 MB of PCM, which is over media's default
// per-blob cap, so a session at the limit is stored truncated rather than
// refused. That trade is deliberate — a partial recording of a meeting he asked
// for beats an error after two hours.
const DefaultMaxDuration = 2 * time.Hour
// DefaultSTTWindow — how much audio goes to the transcriber in one call. Five
// minutes of 16 kHz mono is under 10 MB, transcribes in well under whisper's
// own timeout on this box, and keeps the worker responsive to the voice path
// between windows.
const DefaultSTTWindow = 5 * time.Minute
// Errors callers distinguish.
var (
// ErrDisabled — capture is not configured. A capability is off unless
// configured, and a recorder most of all.
ErrDisabled = errors.New("capture: not configured")
// ErrBusy — a session is already running. One at a time: two concurrent
// recordings would make "хватит" ambiguous.
ErrBusy = errors.New("capture: a session is already running")
// ErrNoSession — stop or append with nothing running.
ErrNoSession = errors.New("capture: nothing is being recorded")
// ErrBadFormat — a frame is not the canonical 16 kHz mono PCM shape.
ErrBadFormat = errors.New("capture: audio format not supported")
// ErrEmptyCapture — the session ended with no audio in it.
ErrEmptyCapture = errors.New("capture: nothing was recorded")
// ErrExpired — the session hit MaxDuration and was closed. Returned from
// Append so the caller stops sending; the audio collected so far is kept.
ErrExpired = errors.New("capture: session reached its time limit")
)
// Session — one recording in progress. Not created directly; Recorder.Start
// makes it. Guarded by a mutex because frames arrive from a network goroutine
// while a status call may read from another.
type Session struct {
Label string
Started time.Time
mu sync.Mutex
pcm []byte
format audio.Format
expired bool
}
// Duration is how much audio has been collected, from the bytes rather than the
// wall clock: a stream that dropped frames should report the audio that exists,
// not the time that passed.
func (s *Session) Duration() time.Duration {
s.mu.Lock()
defer s.mu.Unlock()
return s.duration()
}
func (s *Session) duration() time.Duration {
a := audio.Audio{Format: s.format, Bytes: s.pcm}
return time.Duration(a.Duration() * float64(time.Second))
}
// Bytes is how much PCM has been collected. For a status line.
func (s *Session) Bytes() int {
s.mu.Lock()
defer s.mu.Unlock()
return len(s.pcm)
}
// Status — what a "что записываешь?" answer needs, and what /dash shows. It is
// the read side of a running session and is safe to ask for at any time.
type Status struct {
Running bool `json:"running"`
Label string `json:"label,omitempty"`
Started time.Time `json:"started,omitempty"`
Duration time.Duration `json:"duration,omitempty"`
Bytes int `json:"bytes,omitempty"`
}
// Recorder owns the single session slot, the blob store and the two models a
// finished capture needs. Build it with New; a zero Recorder is not usable.
type Recorder struct {
blobs *media.Store
tr stt.Transcriber
sum *Summarizer
maxDuration time.Duration
sttWindow time.Duration
now func() time.Time
mu sync.Mutex
current *Session
}
// Config — the recorder's knobs, built from config.CaptureConfig by the daemon.
type Config struct {
// MaxDuration — hard cap on one session. 0 ⇒ DefaultMaxDuration.
MaxDuration time.Duration
// STTWindow — audio per transcription call. 0 ⇒ DefaultSTTWindow.
STTWindow time.Duration
}
// New builds a Recorder. blobs and tr are required — a recorder with nowhere to
// put the audio, or nothing to transcribe it with, is not a recorder. sum may be
// nil: the transcript is still produced and stored, and the summary is simply
// absent, which is the honest degradation when there is no llama-server.
func New(blobs *media.Store, tr stt.Transcriber, sum *Summarizer, cfg Config) (*Recorder, error) {
if blobs == nil {
return nil, errors.New("capture: no blob store")
}
if tr == nil {
return nil, errors.New("capture: no transcriber")
}
maxDur := cfg.MaxDuration
if maxDur <= 0 {
maxDur = DefaultMaxDuration
}
window := cfg.STTWindow
if window <= 0 {
window = DefaultSTTWindow
}
return &Recorder{
blobs: blobs,
tr: tr,
sum: sum,
maxDuration: maxDur,
sttWindow: window,
now: time.Now,
}, nil
}
// MaxDuration is the configured hard cap. For the reply that tells him how long
// she will keep going if he forgets to say "хватит".
func (r *Recorder) MaxDuration() time.Duration { return r.maxDuration }
// Start opens a session. label is what the meeting is called ("встреча с
// подрядчиком"); it ends up in the summary note so the note is findable.
// ErrBusy if one is already running — the caller says so rather than silently
// discarding the first recording.
func (r *Recorder) Start(label string) (*Session, error) {
r.mu.Lock()
defer r.mu.Unlock()
if r.current != nil {
return nil, fmt.Errorf("%w: %q since %s", ErrBusy, r.current.Label,
r.current.Started.Format(time.Kitchen))
}
s := &Session{
Label: strings.TrimSpace(label),
Started: r.now().UTC(),
format: audio.PCM16kMono,
}
r.current = s
return s, nil
}
// Append adds one frame to the running session. ErrNoSession when nothing is
// running, which is the guard that makes an ambient path impossible: a stream
// arriving at a Recorder nobody started is refused frame by frame.
//
// ErrExpired once the session is at MaxDuration. The audio collected so far is
// kept and Stop still works — the cap ends the recording, it does not throw it
// away.
func (r *Recorder) Append(a audio.Audio) error {
if !a.Format.IsValid() {
return fmt.Errorf("%w: %+v", ErrBadFormat, a.Format)
}
r.mu.Lock()
s := r.current
r.mu.Unlock()
if s == nil {
return ErrNoSession
}
s.mu.Lock()
defer s.mu.Unlock()
if s.expired {
return ErrExpired
}
s.pcm = append(s.pcm, a.Bytes...)
if s.duration() >= r.maxDuration {
s.expired = true
return ErrExpired
}
return nil
}
// Status reports the running session, or Running=false.
func (r *Recorder) Status() Status {
r.mu.Lock()
s := r.current
r.mu.Unlock()
if s == nil {
return Status{}
}
return Status{
Running: true,
Label: s.Label,
Started: s.Started,
Duration: s.Duration(),
Bytes: s.Bytes(),
}
}
// Result — a finished capture.
type Result struct {
// BlobID — the stored audio, content-addressed. Empty only if storing failed.
BlobID string
// Label / Started / Duration — what was recorded and when.
Label string
Started time.Time
Duration time.Duration
// Transcript — the full text, joined across STT windows.
Transcript string
// Summary — the map-reduced summary, or empty when no summarizer was wired
// or the model failed. Empty summary with a non-empty transcript is a
// degraded success, not a failure: the words are there.
Summary string
// Chunks — how many windows the transcript was summarised in. 1 means it fit
// in one prompt. Reported so a suspiciously vague summary can be explained.
Chunks int
}
// Stop ends the session and produces the result: store the audio, transcribe it
// in windows, summarise it in windows. The session slot is freed before any of
// the slow work starts, so a stuck model cannot block the next recording.
//
// The order matters and is the same as vision's: the audio is stored FIRST. If
// transcription or summarisation fails, the recording is still on disk and can
// be run again; a meeting that happened once must not be lost to a model error.
func (r *Recorder) Stop(ctx context.Context) (Result, error) {
r.mu.Lock()
s := r.current
r.current = nil
r.mu.Unlock()
if s == nil {
return Result{}, ErrNoSession
}
s.mu.Lock()
pcm := s.pcm
format := s.format
s.mu.Unlock()
res := Result{Label: s.Label, Started: s.Started}
if len(pcm) == 0 {
return res, ErrEmptyCapture
}
full := audio.Audio{Format: format, Bytes: pcm}
res.Duration = time.Duration(full.Duration() * float64(time.Second))
// Stored as WAV, not headerless PCM: a blob on disk that `aplay` and whisper
// can both open without being told the format is worth 44 bytes.
wav, err := audio.WAVFromPCM(format, pcm)
if err != nil {
return res, fmt.Errorf("capture: wav: %w", err)
}
blob, err := r.blobs.Put(media.KindAudio, "audio/wav", "capture:meeting", wav)
if err != nil {
// Over the per-blob cap is the expected case for a very long meeting.
// Report it and keep going: a transcript without the audio still beats
// nothing, and the words are what he will read.
return res, fmt.Errorf("capture: store audio: %w", err)
}
res.BlobID = blob.ID
text, err := r.transcribe(ctx, full)
if err != nil {
return res, fmt.Errorf("capture: transcribe: %w", err)
}
res.Transcript = text
if strings.TrimSpace(text) == "" {
return res, ErrEmptyCapture
}
if r.sum == nil {
return res, nil
}
summary, chunks, err := r.sum.Summarize(ctx, s.Label, text)
res.Chunks = chunks
if err != nil {
// Degraded success: the transcript is real and stored, only the summary
// is missing. The caller writes the transcript note and says so.
return res, fmt.Errorf("capture: summarize: %w", err)
}
res.Summary = summary
return res, nil
}
// Abort throws the running session away without transcribing or storing it.
// This is what "забудь, не записывай" must map to: a recording someone changed
// their mind about leaves nothing behind, not a blob with a note saying it was
// abandoned. Returns whether anything was running.
func (r *Recorder) Abort() bool {
r.mu.Lock()
defer r.mu.Unlock()
if r.current == nil {
return false
}
r.current = nil
return true
}
// transcribe runs the transcriber over the audio in windows and joins the text.
// A window that fails is fatal: a summary of a meeting with a silent hole in the
// middle is a summary that misleads.
func (r *Recorder) transcribe(ctx context.Context, a audio.Audio) (string, error) {
windows := chunkAudio(a, r.sttWindow)
parts := make([]string, 0, len(windows))
for i, w := range windows {
text, _, err := r.tr.Transcribe(ctx, w)
if err != nil {
return "", fmt.Errorf("window %d/%d: %w", i+1, len(windows), err)
}
if t := strings.TrimSpace(text); t != "" {
parts = append(parts, t)
}
}
return strings.Join(parts, " "), nil
}
// chunkAudio splits audio into windows of at most window duration, cut on
// sample boundaries. A window shorter than one sample is impossible; audio
// shorter than one window comes back as a single element, so the caller never
// special-cases the short case.
func chunkAudio(a audio.Audio, window time.Duration) []audio.Audio {
bytesPerSample := a.Format.SampleBits / 8 * a.Format.Channels
if bytesPerSample <= 0 || a.Format.SampleRate <= 0 || window <= 0 {
return []audio.Audio{a}
}
per := int(window.Seconds()) * a.Format.SampleRate * bytesPerSample
if per <= 0 || len(a.Bytes) <= per {
return []audio.Audio{a}
}
var out []audio.Audio
for off := 0; off < len(a.Bytes); off += per {
end := off + per
if end > len(a.Bytes) {
end = len(a.Bytes)
}
// Never cut mid-sample: a split inside an int16 shifts every following
// sample by a byte and turns the tail of the window into noise.
end -= (end - off) % bytesPerSample
if end <= off {
break
}
out = append(out, audio.Audio{Format: a.Format, Bytes: a.Bytes[off:end]})
}
return out
}
+363
View File
@@ -0,0 +1,363 @@
package capture
import (
"context"
"errors"
"fmt"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/audio"
"github.com/kami/maven/internal/media"
)
// fakeTranscriber returns a fixed phrase per call so a windowed transcription is
// visible in the joined output.
type fakeTranscriber struct {
calls int
err error
phrase string
}
func (f *fakeTranscriber) Transcribe(_ context.Context, a audio.Audio) (string, float64, error) {
f.calls++
if f.err != nil {
return "", 0, f.err
}
p := f.phrase
if p == "" {
p = "окно"
}
return fmt.Sprintf("%s%d", p, f.calls), 1.0, nil
}
// fakeCompleter records prompts and replies from a script.
type fakeCompleter struct {
replies []string
systems []string
users []string
err error
}
func (f *fakeCompleter) Complete(_ context.Context, system, user string) (string, error) {
f.systems = append(f.systems, system)
f.users = append(f.users, user)
if f.err != nil {
return "", f.err
}
if len(f.replies) == 0 {
return "итог", nil
}
r := f.replies[0]
f.replies = f.replies[1:]
return r, nil
}
// frame builds n seconds of silence in the canonical format.
func frame(seconds float64) audio.Audio {
n := int(seconds*16000) * 2
return audio.Audio{Format: audio.PCM16kMono, Bytes: make([]byte, n)}
}
func testRecorder(t *testing.T, tr *fakeTranscriber, sum *Summarizer, cfg Config) (*Recorder, *media.Store) {
t.Helper()
blobs, err := media.Open(t.TempDir(), 0, 0)
if err != nil {
t.Fatal(err)
}
r, err := New(blobs, tr, sum, cfg)
if err != nil {
t.Fatal(err)
}
return r, blobs
}
func TestNewRequiresStoreAndTranscriber(t *testing.T) {
blobs, err := media.Open(t.TempDir(), 0, 0)
if err != nil {
t.Fatal(err)
}
if _, err := New(nil, &fakeTranscriber{}, nil, Config{}); err == nil {
t.Error("recorder built with no blob store")
}
if _, err := New(blobs, nil, nil, Config{}); err == nil {
t.Error("recorder built with no transcriber")
}
}
// The invariant that matters most: audio arriving at a recorder nobody started
// is refused. There is no ambient path in.
func TestAppendWithoutStartIsRefused(t *testing.T) {
r, _ := testRecorder(t, &fakeTranscriber{}, nil, Config{})
if err := r.Append(frame(1)); !errors.Is(err, ErrNoSession) {
t.Fatalf("got %v, want ErrNoSession", err)
}
if r.Status().Running {
t.Error("a refused frame started a session")
}
}
func TestStopWithoutStartIsRefused(t *testing.T) {
r, _ := testRecorder(t, &fakeTranscriber{}, nil, Config{})
if _, err := r.Stop(context.Background()); !errors.Is(err, ErrNoSession) {
t.Fatalf("got %v, want ErrNoSession", err)
}
}
func TestOneSessionAtATime(t *testing.T) {
r, _ := testRecorder(t, &fakeTranscriber{}, nil, Config{})
if _, err := r.Start("встреча"); err != nil {
t.Fatal(err)
}
if _, err := r.Start("вторая"); !errors.Is(err, ErrBusy) {
t.Fatalf("got %v, want ErrBusy", err)
}
if _, err := r.Stop(context.Background()); !errors.Is(err, ErrEmptyCapture) {
t.Fatalf("empty stop: %v", err)
}
// The slot is free again after a stop, even a failed one.
if _, err := r.Start("третья"); err != nil {
t.Errorf("slot not released: %v", err)
}
}
func TestRoundTripStoresAudioTranscriptAndSummary(t *testing.T) {
tr := &fakeTranscriber{phrase: "совещание"}
sum := NewSummarizer(&fakeCompleter{replies: []string{"— решили купить насос"}}, 0, 0, nil)
r, blobs := testRecorder(t, tr, sum, Config{})
if _, err := r.Start("встреча с подрядчиком"); err != nil {
t.Fatal(err)
}
for i := 0; i < 3; i++ {
if err := r.Append(frame(2)); err != nil {
t.Fatal(err)
}
}
res, err := r.Stop(context.Background())
if err != nil {
t.Fatalf("stop: %v", err)
}
if res.BlobID == "" {
t.Error("no audio blob stored")
}
blob, data, err := blobs.Read(res.BlobID)
if err != nil {
t.Fatalf("blob unreadable: %v", err)
}
if blob.Kind != media.KindAudio || blob.Source != "capture:meeting" {
t.Errorf("blob metadata = %+v", blob)
}
if string(data[:4]) != "RIFF" {
t.Error("audio was not stored as a playable WAV")
}
if res.Transcript == "" {
t.Error("no transcript")
}
if !strings.Contains(res.Summary, "насос") {
t.Errorf("summary = %q", res.Summary)
}
if !strings.Contains(res.Summary, "встреча с подрядчиком") {
t.Errorf("label missing from summary: %q", res.Summary)
}
if res.Duration != 6*time.Second {
t.Errorf("duration = %v, want 6s", res.Duration)
}
}
// A forgotten session stops itself, and the audio collected before the cap is
// kept rather than thrown away.
func TestMaxDurationEndsTheSessionAndKeepsAudio(t *testing.T) {
tr := &fakeTranscriber{}
r, _ := testRecorder(t, tr, nil, Config{MaxDuration: 4 * time.Second})
if _, err := r.Start("длинная"); err != nil {
t.Fatal(err)
}
if err := r.Append(frame(3)); err != nil {
t.Fatalf("first frame: %v", err)
}
if err := r.Append(frame(3)); !errors.Is(err, ErrExpired) {
t.Fatalf("got %v, want ErrExpired", err)
}
// Further frames keep being refused, so a client that ignores the error
// cannot grow the recording past the cap.
if err := r.Append(frame(3)); !errors.Is(err, ErrExpired) {
t.Fatalf("post-expiry frame: %v", err)
}
res, err := r.Stop(context.Background())
if err != nil {
t.Fatalf("stop after expiry: %v", err)
}
if res.Duration != 6*time.Second {
t.Errorf("duration = %v, want the 6s collected before the cap", res.Duration)
}
}
func TestAppendRejectsWrongFormat(t *testing.T) {
r, _ := testRecorder(t, &fakeTranscriber{}, nil, Config{})
if _, err := r.Start("x"); err != nil {
t.Fatal(err)
}
bad := audio.Audio{Format: audio.Format{SampleRate: 44100, Channels: 2, SampleBits: 16, Encoding: "pcm_s16le"}, Bytes: make([]byte, 100)}
if err := r.Append(bad); !errors.Is(err, ErrBadFormat) {
t.Fatalf("got %v, want ErrBadFormat", err)
}
}
// "забудь, не записывай" must leave nothing behind — no blob, no transcript.
func TestAbortLeavesNothing(t *testing.T) {
tr := &fakeTranscriber{}
r, blobs := testRecorder(t, tr, nil, Config{})
if _, err := r.Start("зря начали"); err != nil {
t.Fatal(err)
}
if err := r.Append(frame(5)); err != nil {
t.Fatal(err)
}
if !r.Abort() {
t.Fatal("Abort reported nothing running")
}
if r.Status().Running {
t.Error("session survived Abort")
}
list, err := blobs.List(media.KindAudio)
if err != nil {
t.Fatal(err)
}
if len(list) != 0 {
t.Errorf("Abort stored %d blob(s)", len(list))
}
if tr.calls != 0 {
t.Errorf("Abort transcribed anyway (%d calls)", tr.calls)
}
if r.Abort() {
t.Error("second Abort reported a session")
}
}
func TestStatusReportsTheRunningSession(t *testing.T) {
r, _ := testRecorder(t, &fakeTranscriber{}, nil, Config{})
if got := r.Status(); got.Running {
t.Error("idle recorder reports running")
}
if _, err := r.Start("планёрка"); err != nil {
t.Fatal(err)
}
if err := r.Append(frame(10)); err != nil {
t.Fatal(err)
}
st := r.Status()
if !st.Running || st.Label != "планёрка" {
t.Fatalf("status = %+v", st)
}
if st.Duration != 10*time.Second {
t.Errorf("duration = %v", st.Duration)
}
if st.Bytes != 10*16000*2 {
t.Errorf("bytes = %d", st.Bytes)
}
}
// Long audio goes to the transcriber in windows: handing a whisper worker an
// hour of PCM in one call blocks the voice path for minutes.
func TestLongAudioIsTranscribedInWindows(t *testing.T) {
tr := &fakeTranscriber{}
r, _ := testRecorder(t, tr, nil, Config{STTWindow: 2 * time.Second})
if _, err := r.Start("длинная"); err != nil {
t.Fatal(err)
}
if err := r.Append(frame(9)); err != nil {
t.Fatal(err)
}
res, err := r.Stop(context.Background())
if err != nil {
t.Fatalf("stop: %v", err)
}
if tr.calls != 5 { // 2+2+2+2+1
t.Errorf("transcriber called %d times, want 5", tr.calls)
}
if !strings.Contains(res.Transcript, "окно5") {
t.Errorf("last window missing from transcript: %q", res.Transcript)
}
}
// A hole in the middle of a meeting summary would mislead, so a failed window is
// fatal — but the audio is already stored and re-runnable.
func TestTranscriptionFailureKeepsTheAudio(t *testing.T) {
tr := &fakeTranscriber{err: errors.New("whisper is down")}
r, blobs := testRecorder(t, tr, nil, Config{})
if _, err := r.Start("встреча"); err != nil {
t.Fatal(err)
}
if err := r.Append(frame(2)); err != nil {
t.Fatal(err)
}
res, err := r.Stop(context.Background())
if err == nil {
t.Fatal("transcription failure was not reported")
}
if res.BlobID == "" {
t.Fatal("no blob id to retry with")
}
if _, _, err := blobs.Read(res.BlobID); err != nil {
t.Errorf("audio was not kept: %v", err)
}
}
// No llama-server ⇒ transcript only. That is the honest degradation, not an
// error.
func TestNoSummarizerStillProducesATranscript(t *testing.T) {
r, _ := testRecorder(t, &fakeTranscriber{}, nil, Config{})
if _, err := r.Start("встреча"); err != nil {
t.Fatal(err)
}
if err := r.Append(frame(1)); err != nil {
t.Fatal(err)
}
res, err := r.Stop(context.Background())
if err != nil {
t.Fatalf("stop: %v", err)
}
if res.Transcript == "" {
t.Error("no transcript")
}
if res.Summary != "" {
t.Errorf("summary appeared from nowhere: %q", res.Summary)
}
}
// A summariser failure is a degraded success: the words exist and are returned.
func TestSummaryFailureStillReturnsTheTranscript(t *testing.T) {
sum := NewSummarizer(&fakeCompleter{err: errors.New("llama is down")}, 0, 0, nil)
r, _ := testRecorder(t, &fakeTranscriber{}, sum, Config{})
if _, err := r.Start("встреча"); err != nil {
t.Fatal(err)
}
if err := r.Append(frame(1)); err != nil {
t.Fatal(err)
}
res, err := r.Stop(context.Background())
if err == nil {
t.Fatal("summary failure was not reported")
}
if res.Transcript == "" {
t.Error("transcript lost to a summary failure")
}
}
func TestChunkAudioNeverCutsMidSample(t *testing.T) {
a := audio.Audio{Format: audio.PCM16kMono, Bytes: make([]byte, 16000*2*5+1)}
for _, w := range chunkAudio(a, 2*time.Second) {
if len(w.Bytes)%2 != 0 {
t.Fatalf("window of %d bytes cuts an int16 in half", len(w.Bytes))
}
}
}
func TestChunkAudioShortInputIsOneWindow(t *testing.T) {
a := frame(1)
if got := chunkAudio(a, time.Minute); len(got) != 1 {
t.Errorf("got %d windows, want 1", len(got))
}
}
+261
View File
@@ -0,0 +1,261 @@
package capture
import (
"context"
"errors"
"fmt"
"strings"
"unicode"
)
// DefaultChunkRunes — how much transcript goes into one summarisation prompt.
//
// The resident model runs at n_ctx 4096 and is a Thinking variant, so reasoning
// tokens need room too. Russian runs roughly 2.53 characters per token on a
// Qwen tokenizer, so 3000 runes is about 1100 tokens of transcript, leaving the
// prompt, the persona block, the reasoning and the answer comfortable space.
// This is the same reasoning internal/crawl used to land on 4000 runes, tightened
// because a meeting transcript is denser in named entities than a web page and
// the reduce step has to fit several summaries at once.
const DefaultChunkRunes = 3000
// DefaultMaxChunks — how many windows one meeting may be summarised in. Forty
// chunks at 3000 runes is roughly a two-hour meeting, which is MaxDuration; past
// that the transcript is truncated and the summary says so, because forty-one
// sequential model calls on this box is half an hour of work nobody is waiting
// through.
const DefaultMaxChunks = 40
// ErrNoSummary — the model returned nothing usable for every chunk.
var ErrNoSummary = errors.New("capture: model produced no summary")
// Completer is the one thing the summarizer needs from a model: text in, text
// out. It is an interface rather than an *llm.Client so this package stays pure
// and testable, and so the daemon can pass whatever it already has.
type Completer interface {
Complete(ctx context.Context, system, user string) (string, error)
}
// Summarizer turns a transcript into something worth reading. It is map-reduce
// and nothing cleverer: summarise each window, then summarise the summaries.
//
// Truncation was the alternative and is rejected. A truncated meeting summary
// reads as complete and is not, which is worse than no summary at all — he would
// act on it.
type Summarizer struct {
llm Completer
chunkRunes int
maxChunks int
// context is the persona/context block the daemon prepends to every prompt,
// or empty. Passed in rather than built here so this package does not import
// internal/persona and the feminine self-reference rules stay in one place.
context func() string
}
// NewSummarizer wires a summarizer. llm nil ⇒ nil Summarizer, which Recorder
// treats as "transcript only", the honest degradation with no llama-server.
// chunkRunes ≤ 0 ⇒ DefaultChunkRunes; maxChunks ≤ 0 ⇒ DefaultMaxChunks.
func NewSummarizer(llm Completer, chunkRunes, maxChunks int, contextBlock func() string) *Summarizer {
if llm == nil {
return nil
}
if chunkRunes <= 0 {
chunkRunes = DefaultChunkRunes
}
if maxChunks <= 0 {
maxChunks = DefaultMaxChunks
}
if contextBlock == nil {
contextBlock = func() string { return "" }
}
return &Summarizer{llm: llm, chunkRunes: chunkRunes, maxChunks: maxChunks, context: contextBlock}
}
// chunkPrompt — the map step. Deliberately plain: this is not Maven speaking to
// him, it is a model condensing text, so there is no first person in it at all
// and therefore nothing for the persona's gender rules to get wrong. The reply
// she gives him afterwards is phrased by the ordinary replier, which does carry
// the persona.
const chunkPrompt = `Ты обрабатываешь фрагмент расшифровки разговора.
Сожми его до 2-4 пунктов: о чём говорили, какие решения приняли, какие задачи назвали.
Без вступлений и выводов. Только по тексту — не придумывай того, чего в нём нет.
Если во фрагменте нет ничего содержательного, ответь одним словом: пусто.`
// reducePrompt — the reduce step. Same rules, over the chunk summaries.
const reducePrompt = `Ниже — конспекты фрагментов одной встречи, по порядку.
Собери из них один короткий итог: о чём была встреча, какие решения приняли, что кому делать.
Не повторяйся, не придумывай, не добавляй вступлений.`
// emptyMarker — what the map step answers for a chunk with nothing in it. Such
// chunks are dropped before the reduce step rather than padding it with noise.
const emptyMarker = "пусто"
// Summarize returns the summary and the number of chunks the transcript was
// split into. One chunk means it fit in a single prompt and the reduce step was
// skipped, which is the common case for a short meeting and saves a model call.
func (s *Summarizer) Summarize(ctx context.Context, label, transcript string) (string, int, error) {
if s == nil {
return "", 0, ErrDisabled
}
chunks := ChunkText(transcript, s.chunkRunes)
if len(chunks) == 0 {
return "", 0, ErrEmptyCapture
}
truncated := false
if len(chunks) > s.maxChunks {
chunks = chunks[:s.maxChunks]
truncated = true
}
system := s.context() + chunkPrompt
parts := make([]string, 0, len(chunks))
for i, c := range chunks {
out, err := s.llm.Complete(ctx, system, c)
if err != nil {
return "", len(chunks), fmt.Errorf("chunk %d/%d: %w", i+1, len(chunks), err)
}
out = strings.TrimSpace(out)
if out == "" || strings.EqualFold(out, emptyMarker) {
continue
}
parts = append(parts, out)
}
if len(parts) == 0 {
return "", len(chunks), ErrNoSummary
}
summary := parts[0]
if len(parts) > 1 {
joined := strings.Join(parts, "\n\n")
reduced, err := s.llm.Complete(ctx, s.context()+reducePrompt, joined)
if err != nil {
// The per-chunk summaries are real work; hand them over rather than
// losing them to a failure in the last step.
return joined, len(chunks), fmt.Errorf("reduce: %w", err)
}
if r := strings.TrimSpace(reduced); r != "" {
summary = r
} else {
summary = joined
}
}
if label != "" {
summary = label + "\n\n" + summary
}
if truncated {
// Said in the note, not swallowed: a summary that silently covers the
// first hour of a three-hour meeting is the failure mode this guards.
summary += fmt.Sprintf("\n\n(расшифровка обрезана: обработано %d фрагментов из большего числа)", s.maxChunks)
}
return summary, len(chunks), nil
}
// ChunkText splits text into windows of at most maxRunes runes, cutting on
// sentence boundaries where it can and on a word boundary otherwise. Exported
// because it is the part worth testing on its own and the part a future
// transcript viewer will want.
//
// A sentence longer than maxRunes (a transcript with no punctuation at all,
// which whisper does produce) is cut on whitespace rather than dropped or run
// past the limit.
func ChunkText(text string, maxRunes int) []string {
text = strings.TrimSpace(text)
if text == "" {
return nil
}
if maxRunes <= 0 {
maxRunes = DefaultChunkRunes
}
if len([]rune(text)) <= maxRunes {
return []string{text}
}
var out []string
var cur []rune
flush := func() {
if s := strings.TrimSpace(string(cur)); s != "" {
out = append(out, s)
}
cur = cur[:0]
}
for _, sent := range splitSentences(text) {
sr := []rune(sent)
if len(sr) > maxRunes {
// Oversized sentence: emit what is buffered, then cut this one on
// word boundaries.
flush()
for _, piece := range splitWords(sr, maxRunes) {
out = append(out, piece)
}
continue
}
if len(cur)+len(sr) > maxRunes {
flush()
}
cur = append(cur, sr...)
}
flush()
return out
}
// splitSentences cuts on sentence-ending punctuation followed by a space,
// keeping the punctuation with the sentence it ends. Good enough for a
// transcript: whisper emits periods and question marks, and being wrong about an
// abbreviation costs a slightly uneven chunk, nothing more.
func splitSentences(text string) []string {
runes := []rune(text)
var out []string
start := 0
for i := 0; i < len(runes); i++ {
if runes[i] != '.' && runes[i] != '!' && runes[i] != '?' && runes[i] != '\n' {
continue
}
// Consume a run of punctuation ("?!", "...") so it stays together.
j := i
for j+1 < len(runes) && isSentenceEnd(runes[j+1]) {
j++
}
if j+1 < len(runes) && !unicode.IsSpace(runes[j+1]) {
i = j
continue
}
end := j + 1
for end < len(runes) && unicode.IsSpace(runes[end]) {
end++
}
out = append(out, string(runes[start:end]))
start = end
i = end - 1
}
if start < len(runes) {
out = append(out, string(runes[start:]))
}
return out
}
func isSentenceEnd(r rune) bool {
return r == '.' || r == '!' || r == '?'
}
// splitWords cuts an oversized run on whitespace, falling back to a hard cut
// when a single "word" is itself longer than the limit.
func splitWords(runes []rune, maxRunes int) []string {
var out []string
for len(runes) > maxRunes {
cut := maxRunes
for cut > 0 && !unicode.IsSpace(runes[cut]) {
cut--
}
if cut == 0 {
cut = maxRunes
}
if s := strings.TrimSpace(string(runes[:cut])); s != "" {
out = append(out, s)
}
runes = runes[cut:]
}
if s := strings.TrimSpace(string(runes)); s != "" {
out = append(out, s)
}
return out
}
+244
View File
@@ -0,0 +1,244 @@
package capture
import (
"context"
"errors"
"fmt"
"strings"
"testing"
)
func TestNilSummarizerWithoutAModel(t *testing.T) {
if s := NewSummarizer(nil, 0, 0, nil); s != nil {
t.Fatal("a summarizer with no model is not nil")
}
var s *Summarizer
if _, _, err := s.Summarize(context.Background(), "x", "текст"); !errors.Is(err, ErrDisabled) {
t.Fatalf("got %v, want ErrDisabled", err)
}
}
// The common case: a short meeting fits in one prompt, so there is exactly one
// model call and no reduce step.
func TestShortTranscriptSkipsTheReduceStep(t *testing.T) {
f := &fakeCompleter{replies: []string{"— договорились о смете"}}
s := NewSummarizer(f, 0, 0, nil)
out, chunks, err := s.Summarize(context.Background(), "смета", "Обсудили смету. Решили подписать.")
if err != nil {
t.Fatal(err)
}
if chunks != 1 {
t.Errorf("chunks = %d, want 1", chunks)
}
if len(f.users) != 1 {
t.Fatalf("%d model calls, want 1", len(f.users))
}
if !strings.Contains(out, "смете") || !strings.HasPrefix(out, "смета") {
t.Errorf("summary = %q", out)
}
}
func TestLongTranscriptIsMappedThenReduced(t *testing.T) {
f := &roleCompleter{mapReply: "часть", reduceReply: "общий итог"}
s := NewSummarizer(f, 40, 0, nil)
long := strings.Repeat("Говорили про насос и трубы. ", 12)
out, chunks, err := s.Summarize(context.Background(), "", long)
if err != nil {
t.Fatal(err)
}
if chunks < 2 {
t.Fatalf("chunks = %d, want the transcript split", chunks)
}
// One map call per chunk, then exactly one reduce.
if f.maps != chunks {
t.Errorf("%d map calls for %d chunks", f.maps, chunks)
}
if f.reduces != 1 {
t.Errorf("%d reduce calls, want 1", f.reduces)
}
if out != "общий итог" {
t.Errorf("summary = %q, want the reduced text", out)
}
}
// Losing every per-chunk summary because the last call failed would throw away
// most of the work.
func TestReduceFailureReturnsTheJoinedParts(t *testing.T) {
f := &roleCompleter{mapReply: "часть", reduceFails: true}
s := NewSummarizer(f, 40, 0, nil)
long := strings.Repeat("Говорили про насос и трубы. ", 12)
out, _, err := s.Summarize(context.Background(), "", long)
if err == nil {
t.Fatal("reduce failure was not reported")
}
if !strings.Contains(out, "часть1") || !strings.Contains(out, "часть2") {
t.Errorf("per-chunk work was lost: %q", out)
}
}
func TestChunkFailureIsReported(t *testing.T) {
f := &fakeCompleter{err: errors.New("llama is down")}
s := NewSummarizer(f, 0, 0, nil)
if _, _, err := s.Summarize(context.Background(), "", "текст"); err == nil {
t.Fatal("chunk failure was not reported")
}
}
// "пусто" chunks are noise; they must not pad the reduce prompt, and a
// transcript that is entirely empty chunks is an honest ErrNoSummary rather than
// an invented summary.
func TestEmptyChunksAreDropped(t *testing.T) {
f := &roleCompleter{mapReply: "пусто", literalMap: true, reduceReply: "не должно вызываться"}
s := NewSummarizer(f, 40, 0, nil)
long := strings.Repeat("Тишина в комнате. ", 12)
if _, _, err := s.Summarize(context.Background(), "", long); !errors.Is(err, ErrNoSummary) {
t.Fatalf("got %v, want ErrNoSummary", err)
}
}
func TestEmptyTranscriptIsRefused(t *testing.T) {
s := NewSummarizer(&fakeCompleter{}, 0, 0, nil)
if _, _, err := s.Summarize(context.Background(), "", " \n "); !errors.Is(err, ErrEmptyCapture) {
t.Fatalf("got %v, want ErrEmptyCapture", err)
}
}
// A summary that silently covers the first fraction of a long meeting is the
// failure mode; it has to say so.
func TestTruncationIsStatedInTheSummary(t *testing.T) {
f := &fakeCompleter{replies: []string{"a", "b", "итог"}}
s := NewSummarizer(f, 30, 2, nil)
long := strings.Repeat("Говорили про насос и про трубы. ", 20)
out, chunks, err := s.Summarize(context.Background(), "", long)
if err != nil {
t.Fatal(err)
}
if chunks != 2 {
t.Errorf("chunks = %d, want the cap of 2", chunks)
}
if !strings.Contains(out, "обрезана") {
t.Errorf("truncation not stated: %q", out)
}
}
// The persona block belongs to the daemon, not this package, and must reach the
// model when it is supplied.
func TestContextBlockIsPrependedToEveryPrompt(t *testing.T) {
f := &fakeCompleter{replies: []string{"итог"}}
s := NewSummarizer(f, 0, 0, func() string { return "ПЕРСОНА\n\n" })
if _, _, err := s.Summarize(context.Background(), "", "Обсудили смету."); err != nil {
t.Fatal(err)
}
for i, sys := range f.systems {
if !strings.HasPrefix(sys, "ПЕРСОНА") {
t.Errorf("call %d lost the context block: %q", i, sys)
}
}
}
// The map/reduce prompts must contain no first person at all: the persona's
// feminine forms live in the replier, and a first-person instruction here is a
// place for the model to write "я рад".
func TestPromptsHaveNoFirstPerson(t *testing.T) {
for name, p := range map[string]string{"chunk": chunkPrompt, "reduce": reducePrompt} {
for _, bad := range []string{" я ", "рад", "поняла", "мне ", "вы ", "ваш"} {
if strings.Contains(strings.ToLower(" "+p+" "), bad) {
t.Errorf("%s prompt contains %q", name, bad)
}
}
}
}
func TestChunkTextSplitsOnSentenceBoundaries(t *testing.T) {
text := "Раз два три. Четыре пять шесть. Семь восемь девять."
got := ChunkText(text, 20)
if len(got) != 3 {
t.Fatalf("got %d chunks: %q", len(got), got)
}
for _, c := range got {
if !strings.HasSuffix(c, ".") {
t.Errorf("chunk does not end on a sentence: %q", c)
}
}
}
func TestChunkTextPacksSentencesUpToTheLimit(t *testing.T) {
text := "Раз. Два. Три. Четыре."
got := ChunkText(text, 12)
if len(got) < 2 {
t.Fatalf("nothing was split: %q", got)
}
for _, c := range got {
if n := len([]rune(c)); n > 12 {
t.Errorf("chunk of %d runes exceeds the limit: %q", n, c)
}
}
}
// whisper does emit long unpunctuated runs; those must be cut on whitespace, not
// dropped and not run past the context limit.
func TestChunkTextCutsUnpunctuatedRuns(t *testing.T) {
text := strings.TrimSpace(strings.Repeat("слово ", 50))
got := ChunkText(text, 30)
if len(got) < 2 {
t.Fatalf("unpunctuated run was not split: %d chunks", len(got))
}
total := 0
for _, c := range got {
if n := len([]rune(c)); n > 30 {
t.Errorf("chunk of %d runes exceeds the limit", n)
}
total += strings.Count(c, "слово")
}
if total != 50 {
t.Errorf("%d of 50 words survived chunking", total)
}
}
// A single token longer than the window must still come out, hard-cut.
func TestChunkTextHandlesOneOversizedWord(t *testing.T) {
text := strings.Repeat("я", 70)
got := ChunkText(text, 20)
if len(got) != 4 {
t.Fatalf("got %d chunks, want 4", len(got))
}
if joined := strings.Join(got, ""); len([]rune(joined)) != 70 {
t.Errorf("%d runes survived, want 70", len([]rune(joined)))
}
}
func TestChunkTextShortInputAndEmpty(t *testing.T) {
if got := ChunkText("коротко", 100); len(got) != 1 || got[0] != "коротко" {
t.Errorf("got %q", got)
}
if got := ChunkText(" ", 100); got != nil {
t.Errorf("blank text produced %q", got)
}
}
// roleCompleter answers by which prompt it was handed, so a test does not have
// to predict how many chunks the text splits into. Map replies are numbered
// ("часть1", "часть2", …) unless literalMap is set.
type roleCompleter struct {
mapReply string
literalMap bool
reduceReply string
reduceFails bool
maps int
reduces int
}
func (f *roleCompleter) Complete(_ context.Context, system, _ string) (string, error) {
if strings.Contains(system, "конспекты фрагментов") {
f.reduces++
if f.reduceFails {
return "", errors.New("llama fell over")
}
return f.reduceReply, nil
}
f.maps++
if f.literalMap {
return f.mapReply, nil
}
return fmt.Sprintf("%s%d", f.mapReply, f.maps), nil
}
+105
View File
@@ -205,6 +205,18 @@ type Config struct {
// in this repo holds a recording only in memory. See MediaConfig.
Media *MediaConfig `json:"media,omitempty"`
// Capture — meeting recording and summarisation (Vikunja #253). nil /
// absent ⇒ the recorder does not exist: the start/stop methods are not
// served at all, so nothing on this box can begin a recording. This is the
// most invasive capability Maven has and it is the one most firmly off by
// default. See CaptureConfig.
Capture *CaptureConfig `json:"capture,omitempty"`
// Speaker — voice identification (Vikunja #255). nil / absent ⇒ no
// voiceprint is ever computed and nobody can be enrolled. Enabling it needs
// a speaker-embedding model, which is not on this box. See SpeakerConfig.
Speaker *SpeakerConfig `json:"speaker,omitempty"`
// MCP — Model Context Protocol servers Maven connects OUT to (Vikunja
// #251). nil / absent / no enabled server ⇒ no connection is made and no
// tool is discovered, like every other capability that reaches outside the
@@ -558,6 +570,99 @@ func (v *VisionConfig) LooksAtImages() bool {
return v != nil && v.Enabled && strings.TrimSpace(v.Endpoint) != ""
}
// CaptureConfig — the meeting recorder (internal/capture,
// docs/plans/08-hearing.md).
//
// Absent, or enabled=false, ⇒ the recorder is not wired and the capture methods
// return "unknown method", so no client can start a recording however it asks.
// A media block is required too: audio is never held only in memory.
//
// There is deliberately no "auto", no keyword trigger and no duration default
// long enough to be forgotten about. Recording other people is an explicit act
// with a start, a stop, and a cap.
type CaptureConfig struct {
// Enabled — may she record a meeting when asked. Default false.
Enabled bool `json:"enabled,omitempty"`
// MaxMinutes — hard cap on one session; it stops itself there. 0 ⇒
// capture.DefaultMaxDuration (120 minutes).
MaxMinutes int `json:"max_minutes,omitempty"`
// STTWindow — audio handed to whisper per call. 0 ⇒
// capture.DefaultSTTWindow (5m). Larger windows transcribe slightly better
// and block the STT worker for longer.
STTWindow Duration `json:"stt_window,omitempty"`
// ChunkRunes — transcript runes per summarisation prompt. 0 ⇒
// capture.DefaultChunkRunes (3000), sized for the resident model's n_ctx of
// 4096. Raise this only if the resident model's context grows.
ChunkRunes int `json:"chunk_runes,omitempty"`
// MaxChunks — how many windows one meeting may be summarised in before the
// transcript is truncated and the summary says so. 0 ⇒
// capture.DefaultMaxChunks (40).
MaxChunks int `json:"max_chunks,omitempty"`
// SaveTranscript — write the full transcript as a note alongside the
// summary. Default false: a verbatim record of what other people said in a
// room is a heavier thing to keep than a four-line summary, so it takes a
// deliberate yes. The audio blob is pruned by media.retention either way.
SaveTranscript bool `json:"save_transcript,omitempty"`
}
// Records reports whether the recorder should be wired. Safe on a nil receiver.
func (c *CaptureConfig) Records() bool {
return c != nil && c.Enabled
}
// MaxDuration is the configured session cap as a duration, or 0 for the
// package default. Safe on a nil receiver.
func (c *CaptureConfig) MaxDuration() time.Duration {
if c == nil || c.MaxMinutes <= 0 {
return 0
}
return time.Duration(c.MaxMinutes) * time.Minute
}
// SpeakerConfig — voice identification (internal/speaker,
// docs/plans/10-speaker-recognition.md).
//
// Absent, or enabled=false, ⇒ no voiceprint is computed for any turn, the
// enrolment methods do not exist, and nobody can be enrolled. A voiceprint is
// biometric data about a person, so this one is off until someone typed a model
// path on purpose.
//
// It cannot currently be turned on: there is no speaker-embedding model on this
// box. See the plan document for what to download.
type SpeakerConfig struct {
// Enabled — may she work out who is speaking. Default false.
Enabled bool `json:"enabled,omitempty"`
// ModelPath — an ECAPA-TDNN (or equivalent) speaker-embedding ONNX model.
// Required; without it the recognizer runs disabled and says so once.
ModelPath string `json:"model_path,omitempty"`
// LibPath — onnxruntime shared library, as for the text embedder. Empty ⇒
// the same default the embedder block uses.
LibPath string `json:"lib_path,omitempty"`
// Threshold — cosine similarity a match must beat. 0 ⇒
// speaker.DefaultThreshold (0.7). Lower it and she starts calling guests by
// his name, which is the expensive direction of this error.
Threshold float64 `json:"threshold,omitempty"`
// MinSeconds — least speech an identification will look at. 0 ⇒
// speaker.DefaultMinSeconds (2s).
MinSeconds float64 `json:"min_seconds,omitempty"`
}
// Recognizes reports whether voice identification should be wired. Safe on a
// nil receiver, and false without a model path — enabled with nothing to embed
// with is a misconfiguration, not a capability.
func (s *SpeakerConfig) Recognizes() bool {
return s != nil && s.Enabled && strings.TrimSpace(s.ModelPath) != ""
}
// WeatherConfig configures the weather provider for voice queries.
type WeatherConfig struct {
Provider string `json:"provider,omitempty"` // "open-meteo" or "" → stub
+127
View File
@@ -16,6 +16,73 @@ func TestSensesOffByDefault(t *testing.T) {
if cfg.Vision.LooksAtImages() {
t.Error("vision is on with no vision block")
}
if cfg.Capture.Records() {
t.Error("the recorder is on with no capture block")
}
if cfg.Capture.MaxDuration() != 0 {
t.Error("a nil capture block invented a duration")
}
if cfg.Speaker.Recognizes() {
t.Error("speaker recognition is on with no speaker block")
}
}
// The recorder is the capability that most needs its default to be off, so it
// gets its own test rather than a line in the one above.
func TestCaptureIsOffUntilExplicitlyEnabled(t *testing.T) {
cases := []struct {
name string
c *CaptureConfig
want bool
}{
{"absent", nil, false},
{"present but not enabled", &CaptureConfig{MaxMinutes: 60}, false},
{"enabled", &CaptureConfig{Enabled: true}, true},
}
for _, c := range cases {
if got := c.c.Records(); got != c.want {
t.Errorf("%s: Records() = %v, want %v", c.name, got, c.want)
}
}
}
func TestCaptureBlockParsesFromJSON(t *testing.T) {
raw := `{"capture":{"enabled":true,"max_minutes":45,"stt_window":"2m",
"chunk_runes":2000,"max_chunks":10,"save_transcript":true}}`
var cfg Config
if err := json.Unmarshal([]byte(raw), &cfg); err != nil {
t.Fatalf("unmarshal: %v", err)
}
if !cfg.Capture.Records() {
t.Fatal("capture did not parse as enabled")
}
if cfg.Capture.MaxDuration() != 45*time.Minute {
t.Errorf("max duration = %v", cfg.Capture.MaxDuration())
}
if time.Duration(cfg.Capture.STTWindow) != 2*time.Minute {
t.Errorf("stt window = %v", time.Duration(cfg.Capture.STTWindow))
}
if cfg.Capture.ChunkRunes != 2000 || cfg.Capture.MaxChunks != 10 {
t.Errorf("summariser limits = %+v", cfg.Capture)
}
if !cfg.Capture.SaveTranscript {
t.Error("save_transcript did not parse")
}
}
// Keeping the verbatim record of what other people said is the heavier act, so
// it is separately opt-in from recording at all.
func TestTranscriptIsNotSavedByDefault(t *testing.T) {
var cfg Config
if err := json.Unmarshal([]byte(`{"capture":{"enabled":true}}`), &cfg); err != nil {
t.Fatal(err)
}
if cfg.Capture.SaveTranscript {
t.Error("transcripts are saved without anyone asking")
}
if cfg.Capture.MaxDuration() != 0 {
t.Error("max_minutes defaulted in config instead of in the package")
}
}
// enabled with nothing to talk to is a misconfiguration, not a capability.
@@ -94,3 +161,63 @@ func TestMediaWithoutVisionIsValid(t *testing.T) {
t.Error("vision came on by itself")
}
}
// A voiceprint is a biometric of a named person. Nothing about it turns on by
// itself: no speaker block means no recognition, and no enrolment either.
func TestSpeakerIsOffUntilExplicitlyEnabled(t *testing.T) {
var cfg Config
if err := json.Unmarshal([]byte(`{}`), &cfg); err != nil {
t.Fatal(err)
}
if cfg.Speaker.Recognizes() {
t.Error("speaker recognition came on with no config at all")
}
var empty Config
if err := json.Unmarshal([]byte(`{"speaker":{}}`), &empty); err != nil {
t.Fatal(err)
}
if empty.Speaker.Recognizes() {
t.Error("an empty speaker block enabled recognition")
}
}
// Enabled alone is not enough: recognition needs a model, and on this box there
// is none. Recognizes() must stay false so the daemon reports the honest state
// instead of claiming a capability it cannot perform.
func TestSpeakerNeedsBothEnabledAndAModel(t *testing.T) {
var cfg Config
if err := json.Unmarshal([]byte(`{"speaker":{"enabled":true}}`), &cfg); err != nil {
t.Fatal(err)
}
if cfg.Speaker.Recognizes() {
t.Error("enabled with no model_path claimed to recognise")
}
var only Config
if err := json.Unmarshal([]byte(`{"speaker":{"model_path":"/opt/x.onnx"}}`), &only); err != nil {
t.Fatal(err)
}
if only.Speaker.Recognizes() {
t.Error("a model_path alone enabled recognition")
}
}
func TestSpeakerBlockParsesFromJSON(t *testing.T) {
const raw = `{"speaker":{"enabled":true,"model_path":"/opt/maven/models/spk/ecapa.onnx",` +
`"lib_path":"/opt/maven/lib","threshold":0.62,"min_seconds":1.5}}`
var cfg Config
if err := json.Unmarshal([]byte(raw), &cfg); err != nil {
t.Fatalf("unmarshal: %v", err)
}
if !cfg.Speaker.Recognizes() {
t.Fatal("speaker did not parse as enabled")
}
if cfg.Speaker.ModelPath != "/opt/maven/models/spk/ecapa.onnx" {
t.Errorf("model_path = %q", cfg.Speaker.ModelPath)
}
if cfg.Speaker.LibPath != "/opt/maven/lib" {
t.Errorf("lib_path = %q", cfg.Speaker.LibPath)
}
if cfg.Speaker.Threshold != 0.62 || cfg.Speaker.MinSeconds != 1.5 {
t.Errorf("thresholds = %+v", cfg.Speaker)
}
}
+125
View File
@@ -4,6 +4,8 @@ import (
"context"
"errors"
"time"
"github.com/kami/maven/internal/audio"
)
// DTOs — wire-level data. Decoupled from internal/store so the protocol is
@@ -221,6 +223,129 @@ type DescribeImageResp struct {
NoteID int64 `json:"note_id,omitempty"`
}
// CaptureStartReq — begin recording a meeting (Vikunja #253).
//
// Label is what the meeting is called ("встреча с подрядчиком"); it goes into
// the summary note so the note is findable later. Empty is allowed.
//
// There is no "auto", no keyword and no schedule in this request, and there will
// not be: the only way audio enters the recorder is a client that was told to
// start, appending frames it was told to append. All four capture methods answer
// ErrUnknownMethod unless the operator enabled a capture block, so a surface
// cannot start a recording by asking nicely.
type CaptureStartReq struct {
Label string `json:"label,omitempty"`
}
// CaptureStartResp — the session that opened. MaxSeconds is the hard cap after
// which it stops itself; the caller tells him, so a forgotten recording is his
// own informed choice rather than a surprise.
type CaptureStartResp struct {
Label string `json:"label,omitempty"`
Started time.Time `json:"started"`
MaxSeconds int `json:"max_seconds"`
}
// CaptureAppendReq — one chunk of audio for the running session. Refused with
// "nothing is being recorded" when no session is open, which is the guard that
// makes an ambient path impossible: audio arriving at an idle core is dropped on
// the floor, not buffered "just in case".
type CaptureAppendReq struct {
Audio audio.Audio `json:"audio"`
}
// CaptureAppendResp — how much has been collected, so a client can show a timer
// and notice the cap coming. Expired means the session hit its limit and closed;
// stop sending and call capture_stop, the audio so far is kept.
type CaptureAppendResp struct {
Seconds float64 `json:"seconds"`
Expired bool `json:"expired,omitempty"`
}
// CaptureStopReq — end the running session.
//
// Discard throws the recording away without transcribing, storing or
// summarising anything. This is what "забудь, не записывай" maps to, and it is a
// flag rather than a separate method so the client that says "stop" and the
// client that says "stop and forget" take the same path to the same session.
type CaptureStopReq struct {
Discard bool `json:"discard,omitempty"`
}
// CaptureStopResp — the finished capture. BlobID is the stored WAV, kept under
// media.retention like any other blob and pruned with it.
//
// A response with a Transcript and an empty Summary is a degraded success: the
// words exist, only the model failed. A response with a BlobID and neither is
// the audio surviving a transcription failure — the same id can be run again.
// Discarded is true when nothing was kept.
type CaptureStopResp struct {
BlobID string `json:"blob_id,omitempty"`
Label string `json:"label,omitempty"`
Started time.Time `json:"started,omitempty"`
Seconds float64 `json:"seconds,omitempty"`
Transcript string `json:"transcript,omitempty"`
Summary string `json:"summary,omitempty"`
Chunks int `json:"chunks,omitempty"`
NoteID int64 `json:"note_id,omitempty"`
Discarded bool `json:"discarded,omitempty"`
}
// CaptureStatusResp — what "что ты записываешь?" needs, and what /dash shows.
// Running=false with everything else empty is the normal state.
type CaptureStatusResp struct {
Running bool `json:"running"`
Label string `json:"label,omitempty"`
Started time.Time `json:"started,omitempty"`
Seconds float64 `json:"seconds,omitempty"`
Bytes int `json:"bytes,omitempty"`
}
// EnrollSpeakerReq — register a voice (Vikunja #255).
//
// Samples are separate utterances recorded deliberately for this purpose, not
// audio harvested from ordinary turns. internal/speaker requires several of
// them totalling enough seconds, and refuses one long clip: a profile built
// from a single sentence encodes that sentence as much as the person.
//
// There is no "enrol whoever just spoke" request shape, and that omission is
// the point. Taking a biometric of a guest because they walked past the
// microphone is not something a wire protocol should make easy.
type EnrollSpeakerReq struct {
ID string `json:"id"`
Name string `json:"name,omitempty"`
Samples []audio.Audio `json:"samples"`
}
// Speaker — one enrolled voice as a surface sees it. The voiceprint itself is
// never sent: a listing says who is enrolled, it does not hand out the
// biometric.
type Speaker struct {
ID string `json:"id"`
Name string `json:"name"`
Enrolled time.Time `json:"enrolled"`
Samples int `json:"samples"`
}
// EnrollSpeakerResp — the profile that was written.
type EnrollSpeakerResp struct {
Speaker Speaker `json:"speaker"`
}
// ListSpeakersResp — who is enrolled, sorted by id. Enabled is false when no
// embedding model is wired, which is this box's state: the profiles can be
// listed and deleted, nothing can be recognised.
type ListSpeakersResp struct {
Speakers []Speaker `json:"speakers"`
Enabled bool `json:"enabled"`
}
// ForgetSpeakerReq — delete one voiceprint. This is the request that must
// always work; a biometric someone asked to be rid of has to actually go.
type ForgetSpeakerReq struct {
ID string `json:"id"`
}
// SwapModelReq — load another resident model without restarting the daemon
// (Vikunja #250). ModelPath must be one of the paths in phraser.swap_models;
// anything else is ErrForbidden, and an unconfigured allowlist makes the whole
+123
View File
@@ -0,0 +1,123 @@
package ipc
import (
"context"
"errors"
"testing"
"time"
"github.com/kami/maven/internal/audio"
)
// The load-bearing default for the most invasive capability Maven has: on a core
// that was never configured to record, there is no wire path that starts a
// recording, feeds one, or harvests one. Every one of the four methods refuses.
func TestCapture_OffUnlessConfigured(t *testing.T) {
_, _, cli, _ := newServerWithStore(t)
ctx := context.Background()
if _, err := cli.CaptureStart(ctx, CaptureStartReq{Label: "встреча"}); !errors.Is(err, ErrUnknownMethod) {
t.Errorf("CaptureStart error = %v, want ErrUnknownMethod", err)
}
if _, err := cli.CaptureAppend(ctx, CaptureAppendReq{}); !errors.Is(err, ErrUnknownMethod) {
t.Errorf("CaptureAppend error = %v, want ErrUnknownMethod", err)
}
if _, err := cli.CaptureStop(ctx, CaptureStopReq{}); !errors.Is(err, ErrUnknownMethod) {
t.Errorf("CaptureStop error = %v, want ErrUnknownMethod", err)
}
if _, err := cli.CaptureStatus(ctx); !errors.Is(err, ErrUnknownMethod) {
t.Errorf("CaptureStatus error = %v, want ErrUnknownMethod", err)
}
}
// With the hooks wired, a whole session crosses the boundary intact: the label
// out, the audio in, the summary back.
func TestCapture_RoundTrip(t *testing.T) {
_, srv, cli, _ := newServerWithStore(t)
ctx := context.Background()
started := time.Now().UTC().Truncate(time.Second)
var gotLabel string
var gotBytes int
var gotDiscard bool
srv.CaptureStartFn = func(_ context.Context, req CaptureStartReq) (CaptureStartResp, error) {
gotLabel = req.Label
return CaptureStartResp{Label: req.Label, Started: started, MaxSeconds: 7200}, nil
}
srv.CaptureAppendFn = func(_ context.Context, req CaptureAppendReq) (CaptureAppendResp, error) {
gotBytes = len(req.Audio.Bytes)
return CaptureAppendResp{Seconds: 1.5}, nil
}
srv.CaptureStopFn = func(_ context.Context, req CaptureStopReq) (CaptureStopResp, error) {
gotDiscard = req.Discard
return CaptureStopResp{BlobID: "abc", Summary: "— решили купить насос", Chunks: 1}, nil
}
srv.CaptureStatusFn = func(context.Context) (CaptureStatusResp, error) {
return CaptureStatusResp{Running: true, Label: "встреча", Seconds: 1.5}, nil
}
start, err := cli.CaptureStart(ctx, CaptureStartReq{Label: "встреча с подрядчиком"})
if err != nil {
t.Fatalf("CaptureStart: %v", err)
}
if gotLabel != "встреча с подрядчиком" || start.MaxSeconds != 7200 {
t.Errorf("start = %+v (label seen: %q)", start, gotLabel)
}
if !start.Started.Equal(started) {
t.Errorf("started = %v, want %v", start.Started, started)
}
// Audio must survive the JSON round trip byte for byte — a base64 mistake
// here would be silence in the transcript, not a visible error.
pcm := []byte{1, 2, 3, 4, 5, 6, 7, 8}
ap, err := cli.CaptureAppend(ctx, CaptureAppendReq{
Audio: audio.Audio{Format: audio.PCM16kMono, Bytes: pcm},
})
if err != nil {
t.Fatalf("CaptureAppend: %v", err)
}
if gotBytes != len(pcm) {
t.Errorf("%d bytes arrived, sent %d", gotBytes, len(pcm))
}
if ap.Seconds != 1.5 || ap.Expired {
t.Errorf("append resp = %+v", ap)
}
st, err := cli.CaptureStatus(ctx)
if err != nil {
t.Fatalf("CaptureStatus: %v", err)
}
if !st.Running || st.Label != "встреча" {
t.Errorf("status = %+v", st)
}
stop, err := cli.CaptureStop(ctx, CaptureStopReq{})
if err != nil {
t.Fatalf("CaptureStop: %v", err)
}
if gotDiscard {
t.Error("a plain stop arrived as a discard")
}
if stop.BlobID != "abc" || stop.Summary == "" {
t.Errorf("stop = %+v", stop)
}
}
// "забудь, не записывай" has to reach core as a discard, not as an ordinary
// stop that quietly keeps everything.
func TestCapture_DiscardCrossesTheWire(t *testing.T) {
_, srv, cli, _ := newServerWithStore(t)
var gotDiscard bool
srv.CaptureStopFn = func(_ context.Context, req CaptureStopReq) (CaptureStopResp, error) {
gotDiscard = req.Discard
return CaptureStopResp{Discarded: req.Discard}, nil
}
resp, err := cli.CaptureStop(context.Background(), CaptureStopReq{Discard: true})
if err != nil {
t.Fatalf("CaptureStop: %v", err)
}
if !gotDiscard || !resp.Discarded {
t.Errorf("discard lost: sent true, core saw %v, resp %+v", gotDiscard, resp)
}
}
+70
View File
@@ -473,6 +473,76 @@ func (c *Client) DescribeImage(ctx context.Context, req DescribeImageReq) (Descr
return r, nil
}
// CaptureStart begins recording a meeting (Vikunja #253). ErrUnknownMethod
// means the operator has not enabled capture — the caller should say so and stop
// asking, not retry.
func (c *Client) CaptureStart(ctx context.Context, req CaptureStartReq) (CaptureStartResp, error) {
var r CaptureStartResp
if err := c.call(ctx, MethodCaptureStart, req, &r); err != nil {
return CaptureStartResp{}, err
}
return r, nil
}
// CaptureAppend hands one chunk of audio to the running session. An error means
// the frame was not kept: either nothing is being recorded, or the session hit
// its time limit. Either way the client stops sending.
func (c *Client) CaptureAppend(ctx context.Context, req CaptureAppendReq) (CaptureAppendResp, error) {
var r CaptureAppendResp
if err := c.call(ctx, MethodCaptureAppend, req, &r); err != nil {
return CaptureAppendResp{}, err
}
return r, nil
}
// CaptureStop ends the session. Slow — it transcribes and summarises the whole
// recording — so pass a context with room. Set Discard to throw the recording
// away instead.
func (c *Client) CaptureStop(ctx context.Context, req CaptureStopReq) (CaptureStopResp, error) {
var r CaptureStopResp
if err := c.call(ctx, MethodCaptureStop, req, &r); err != nil {
return CaptureStopResp{}, err
}
return r, nil
}
// CaptureStatus reports the running session, if any.
func (c *Client) CaptureStatus(ctx context.Context) (CaptureStatusResp, error) {
var r CaptureStatusResp
if err := c.call(ctx, MethodCaptureStatus, nil, &r); err != nil {
return CaptureStatusResp{}, err
}
return r, nil
}
// EnrollSpeaker registers a voice from several deliberately recorded samples
// (Vikunja #255). ErrUnknownMethod means no speaker block is configured, which
// is the default: on an unconfigured box there is no way to take a voiceprint.
func (c *Client) EnrollSpeaker(ctx context.Context, req EnrollSpeakerReq) (EnrollSpeakerResp, error) {
var r EnrollSpeakerResp
if err := c.call(ctx, MethodEnrollSpeaker, req, &r); err != nil {
return EnrollSpeakerResp{}, err
}
return r, nil
}
// ListSpeakers reports who is enrolled. The voiceprints themselves stay in
// core. Enabled is false when profiles exist but no embedding model is wired,
// so a surface can say "enrolled, not recognising" rather than implying Maven
// knows who is talking.
func (c *Client) ListSpeakers(ctx context.Context) (ListSpeakersResp, error) {
var r ListSpeakersResp
if err := c.call(ctx, MethodListSpeakers, nil, &r); err != nil {
return ListSpeakersResp{}, err
}
return r, nil
}
// ForgetSpeaker deletes one voiceprint.
func (c *Client) ForgetSpeaker(ctx context.Context, id string) error {
return c.call(ctx, MethodForgetSpeaker, ForgetSpeakerReq{ID: id}, nil)
}
// SwapModel asks core to load another resident model (Vikunja #250).
// ErrUnknownMethod means core has no phraser.swap_models allowlist configured;
// ErrForbidden means the path is not on it, or step-up was not asserted. A
+128 -4
View File
@@ -459,6 +459,27 @@ type Server struct {
// other CoreAPI implementation should have to carry it.
DescribeImageFn DescribeImageFunc
// Capture* — the meeting recorder (Vikunja #253). Set by the daemon only
// when a media store is configured AND capture.enabled is true; nil ⇒ all
// four methods answer ErrUnknownMethod. That is the load-bearing default for
// this capability: on an unconfigured box there is no wire path that begins a
// recording, so nothing can be recorded by accident, by a bug in a surface,
// or by a model deciding it would be helpful.
//
// They bypass CoreAPI because a recorder needs a blob store, an STT worker
// and a llama-server, none of which is a store operation.
CaptureStartFn CaptureStartFunc
CaptureAppendFn CaptureAppendFunc
CaptureStopFn CaptureStopFunc
CaptureStatusFn CaptureStatusFunc
// Speaker* — voice identification (Vikunja #255). Set by the daemon only
// when a speaker block is configured; nil ⇒ all three methods answer
// ErrUnknownMethod, so on an unconfigured box no wire path enrols a voice.
EnrollSpeakerFn EnrollSpeakerFunc
ListSpeakersFn ListSpeakersFunc
ForgetSpeakerFn ForgetSpeakerFunc
// UnlockFn — unwraps the store encryption key from the wrapped blob using
// the passkey credential public key, opens the encrypted store, and wires
// the rest of the daemon (voice, loop, delivery). Set by the daemon when
@@ -489,6 +510,19 @@ type IngestMailFunc func(ctx context.Context, req IngestMailReq) (IngestMailResp
// DescribeImageFunc — core-side image intake + description.
type DescribeImageFunc func(ctx context.Context, req DescribeImageReq) (DescribeImageResp, error)
// CaptureStartFunc / CaptureAppendFunc / CaptureStopFunc / CaptureStatusFunc —
// the four core-side halves of the meeting recorder.
type CaptureStartFunc func(ctx context.Context, req CaptureStartReq) (CaptureStartResp, error)
type CaptureAppendFunc func(ctx context.Context, req CaptureAppendReq) (CaptureAppendResp, error)
type CaptureStopFunc func(ctx context.Context, req CaptureStopReq) (CaptureStopResp, error)
type CaptureStatusFunc func(ctx context.Context) (CaptureStatusResp, error)
// EnrollSpeakerFunc / ListSpeakersFunc / ForgetSpeakerFunc — the core-side
// halves of voice enrolment.
type EnrollSpeakerFunc func(ctx context.Context, req EnrollSpeakerReq) (EnrollSpeakerResp, error)
type ListSpeakersFunc func(ctx context.Context) (ListSpeakersResp, error)
type ForgetSpeakerFunc func(ctx context.Context, req ForgetSpeakerReq) error
// CheckFunc — the auth hook signature. Wired by the daemon (auth.Gate.Check
// satisfies this); dispatch calls it once per request after param-unmarshal
// independence (it gets the raw params, may unmarshal what it needs — ipc
@@ -657,10 +691,11 @@ func withoutParams[R any](fn func(ctx context.Context, api CoreAPI) (R, error))
// is still honored on the very next request with no extra plumbing here.
//
// MethodAssertStepUp, MethodStoreEncryptionKey, MethodUnlock,
// MethodIngestMail, MethodSwapModel, MethodModelStatus and
// MethodDescribeImage are NOT in this table: they bypass CoreAPI entirely
// (s.StepUp / s.WrapKeyFn / s.UnlockFn / s.IngestMailFn / s.DescribeImageFn),
// so dispatch special-cases them before consulting the table.
// MethodIngestMail, MethodSwapModel, MethodModelStatus,
// MethodDescribeImage and the four MethodCapture* methods are NOT in this
// table: they bypass CoreAPI entirely (s.StepUp / s.WrapKeyFn / s.UnlockFn /
// s.IngestMailFn / s.DescribeImageFn / s.Capture*Fn), so dispatch
// special-cases them before consulting the table.
var methodTable = map[Method]handlerFunc{
MethodWriteFact: withParams(func(ctx context.Context, api CoreAPI, p WriteFactReq) (idResp, error) {
id, err := api.WriteFact(ctx, p)
@@ -950,6 +985,95 @@ func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, er
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodCaptureStart:
if s.CaptureStartFn != nil {
var p CaptureStartReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.CaptureStartFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodCaptureAppend:
if s.CaptureAppendFn != nil {
var p CaptureAppendReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.CaptureAppendFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodCaptureStop:
if s.CaptureStopFn != nil {
var p CaptureStopReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.CaptureStopFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodCaptureStatus:
if s.CaptureStatusFn != nil {
resp, err := s.CaptureStatusFn(ctx)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodEnrollSpeaker:
if s.EnrollSpeakerFn != nil {
var p EnrollSpeakerReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.EnrollSpeakerFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodListSpeakers:
if s.ListSpeakersFn != nil {
resp, err := s.ListSpeakersFn(ctx)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodForgetSpeaker:
if s.ForgetSpeakerFn != nil {
var p ForgetSpeakerReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
if err := s.ForgetSpeakerFn(ctx, p); err != nil {
return nil, err
}
return marshalResult(nil), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodModelStatus:
if s.ModelStatusFn != nil {
resp, err := s.ModelStatusFn(ctx)
+123
View File
@@ -0,0 +1,123 @@
package ipc
import (
"context"
"errors"
"testing"
"time"
"github.com/kami/maven/internal/audio"
)
// The default that matters most for a biometric: on a core that was never
// configured with a speaker block, there is no wire path that takes a
// voiceprint, and none that lists the ones that might exist.
func TestSpeaker_OffUnlessConfigured(t *testing.T) {
_, _, cli, _ := newServerWithStore(t)
ctx := context.Background()
if _, err := cli.EnrollSpeaker(ctx, EnrollSpeakerReq{ID: "kami"}); !errors.Is(err, ErrUnknownMethod) {
t.Errorf("EnrollSpeaker error = %v, want ErrUnknownMethod", err)
}
if _, err := cli.ListSpeakers(ctx); !errors.Is(err, ErrUnknownMethod) {
t.Errorf("ListSpeakers error = %v, want ErrUnknownMethod", err)
}
if err := cli.ForgetSpeaker(ctx, "kami"); !errors.Is(err, ErrUnknownMethod) {
t.Errorf("ForgetSpeaker error = %v, want ErrUnknownMethod", err)
}
}
// Enrolment carries several samples across the boundary byte for byte — a
// profile averaged over the wrong bytes is a profile of nobody.
func TestSpeaker_EnrollCrossesTheWire(t *testing.T) {
_, srv, cli, _ := newServerWithStore(t)
ctx := context.Background()
enrolled := time.Now().UTC().Truncate(time.Second)
var gotID, gotName string
var gotSamples [][]byte
srv.EnrollSpeakerFn = func(_ context.Context, req EnrollSpeakerReq) (EnrollSpeakerResp, error) {
gotID, gotName = req.ID, req.Name
for _, s := range req.Samples {
gotSamples = append(gotSamples, s.Bytes)
}
return EnrollSpeakerResp{Speaker: Speaker{
ID: req.ID, Name: req.Name, Enrolled: enrolled, Samples: len(req.Samples),
}}, nil
}
mk := func(b byte, n int) audio.Audio {
buf := make([]byte, n)
for i := range buf {
buf[i] = b
}
return audio.Audio{Format: audio.PCM16kMono, Bytes: buf}
}
samples := []audio.Audio{mk(1, 64), mk(2, 96), mk(3, 128)}
resp, err := cli.EnrollSpeaker(ctx, EnrollSpeakerReq{ID: "kami", Name: "Ками", Samples: samples})
if err != nil {
t.Fatalf("EnrollSpeaker: %v", err)
}
if gotID != "kami" || gotName != "Ками" {
t.Errorf("server saw id=%q name=%q", gotID, gotName)
}
if len(gotSamples) != 3 {
t.Fatalf("server saw %d samples, want 3", len(gotSamples))
}
for i, want := range samples {
if string(gotSamples[i]) != string(want.Bytes) {
t.Errorf("sample %d altered in transit", i)
}
}
if resp.Speaker.Samples != 3 || !resp.Speaker.Enrolled.Equal(enrolled) {
t.Errorf("profile came back wrong: %+v", resp.Speaker)
}
}
// A listing says who is enrolled and whether recognition actually works. On
// this box the honest answer is "enrolled, not recognising", and the response
// has to be able to say so — otherwise a surface implies Maven knows who is
// talking when nothing on disk can tell.
func TestSpeaker_ListReportsDisabledRecognition(t *testing.T) {
_, srv, cli, _ := newServerWithStore(t)
ctx := context.Background()
srv.ListSpeakersFn = func(context.Context) (ListSpeakersResp, error) {
return ListSpeakersResp{
Speakers: []Speaker{{ID: "kami", Name: "Ками", Samples: 3}},
Enabled: false,
}, nil
}
resp, err := cli.ListSpeakers(ctx)
if err != nil {
t.Fatalf("ListSpeakers: %v", err)
}
if len(resp.Speakers) != 1 || resp.Speakers[0].ID != "kami" {
t.Fatalf("speakers = %+v", resp.Speakers)
}
if resp.Enabled {
t.Error("Enabled = true; the seam must be able to report that nothing recognises")
}
}
// Deletion reaches core with the id intact and reports success. This is the
// request that must always work.
func TestSpeaker_ForgetReachesCore(t *testing.T) {
_, srv, cli, _ := newServerWithStore(t)
ctx := context.Background()
var forgot string
srv.ForgetSpeakerFn = func(_ context.Context, req ForgetSpeakerReq) error {
forgot = req.ID
return nil
}
if err := cli.ForgetSpeaker(ctx, "гость"); err != nil {
t.Fatalf("ForgetSpeaker: %v", err)
}
if forgot != "гость" {
t.Errorf("core forgot %q, want %q", forgot, "гость")
}
}
+7
View File
@@ -55,6 +55,13 @@ const (
MethodSwapModel Method = "swap_model"
MethodModelStatus Method = "model_status"
MethodDescribeImage Method = "describe_image"
MethodCaptureStart Method = "capture_start"
MethodCaptureAppend Method = "capture_append"
MethodCaptureStop Method = "capture_stop"
MethodCaptureStatus Method = "capture_status"
MethodEnrollSpeaker Method = "enroll_speaker"
MethodListSpeakers Method = "list_speakers"
MethodForgetSpeaker Method = "forget_speaker"
)
// Request — one frame from module to core. Params is the JSON-encoded argument
+69 -1
View File
@@ -3,6 +3,7 @@ package memory
import (
"context"
"sort"
"strings"
"sync"
)
@@ -19,6 +20,33 @@ type Store interface {
Search(ctx context.Context, vec []float32, topK int) ([]Result, error)
}
// Record is a stored vector read back whole — id, vector and metadata — as
// opposed to Result, which is a search hit and carries a score instead of the
// vector.
type Record struct {
ID string
Vec []float32
Meta map[string]string
}
// Catalog is a Store that can also be enumerated by id prefix and deleted from.
//
// Search is not enough for every user of the vector table. Speaker profiles
// (internal/speaker) need to list exactly their own rows without scoring
// anything, because listing enrolled voices is not a similarity question, and
// they need Delete because a voiceprint is data about a person and "forget this
// voice" has to actually remove it. Note and fact recall use plain Store and are
// unaffected.
type Catalog interface {
Store
// ByPrefix returns every row whose id starts with prefix, in no particular
// order. An empty prefix returns everything.
ByPrefix(ctx context.Context, prefix string) ([]Record, error)
// Delete removes one row by id. Deleting a row that is not there is not an
// error: the caller asked for it to be gone and it is gone.
Delete(ctx context.Context, id string) error
}
// item is a single stored vector with metadata.
type item struct {
id string
@@ -32,14 +60,54 @@ type InMemoryStore struct {
items []item
}
// compile-time check: InMemoryStore satisfies Catalog.
var _ Catalog = (*InMemoryStore)(nil)
func NewInMemoryStore() *InMemoryStore {
return &InMemoryStore{}
}
// Insert upserts by id, matching the persistent store.MemoryStore: a repeated
// id replaces the prior row rather than accumulating a second copy. Re-indexing
// a note is an update, and re-enrolling a voice must replace the old voiceprint
// rather than leave it searchable.
func (s *InMemoryStore) Insert(_ context.Context, id string, vec []float32, meta map[string]string) error {
s.mu.Lock()
defer s.mu.Unlock()
for i := range s.items {
if s.items[i].id == id {
s.items[i] = item{id: id, vec: vec, meta: meta}
return nil
}
}
s.items = append(s.items, item{id: id, vec: vec, meta: meta})
s.mu.Unlock()
return nil
}
// ByPrefix implements Catalog.
func (s *InMemoryStore) ByPrefix(_ context.Context, prefix string) ([]Record, error) {
s.mu.RLock()
defer s.mu.RUnlock()
var out []Record
for _, it := range s.items {
if !strings.HasPrefix(it.id, prefix) {
continue
}
out = append(out, Record{ID: it.id, Vec: append([]float32(nil), it.vec...), Meta: it.meta})
}
return out, nil
}
// Delete implements Catalog.
func (s *InMemoryStore) Delete(_ context.Context, id string) error {
s.mu.Lock()
defer s.mu.Unlock()
for i := range s.items {
if s.items[i].id == id {
s.items = append(s.items[:i], s.items[i+1:]...)
return nil
}
}
return nil
}
+3 -1
View File
@@ -2,6 +2,7 @@ package memory
import (
"context"
"fmt"
"math"
"testing"
)
@@ -40,8 +41,9 @@ func TestTopKTruncation(t *testing.T) {
s := NewInMemoryStore()
ctx := context.Background()
// Distinct ids: Insert upserts by id, so ten rows need ten ids.
for i := 0; i < 10; i++ {
s.Insert(ctx, "", []float32{float32(i) / 10, 0, 0}, nil)
s.Insert(ctx, fmt.Sprintf("n%d", i), []float32{float32(i) / 10, 0, 0}, nil)
}
results, err := s.Search(ctx, []float32{1, 0, 0}, 3)
+109
View File
@@ -0,0 +1,109 @@
package speaker
import (
"context"
"fmt"
"strconv"
"strings"
"time"
"github.com/kami/maven/internal/audio"
)
// Enroll registers a voice under an id and a spoken name.
//
// Several separate samples are required (MinEnrollSamples, MinEnrollSeconds
// total): a profile built from one sentence encodes that sentence as much as the
// person, and the resulting threshold behaviour is unpredictable. The samples
// are embedded individually and the voiceprints averaged, then re-normalised.
//
// Re-enrolling an existing id REPLACES the profile. That is the intended way to
// improve a weak one, and it is why the store upserts by id.
//
// # The refused step
//
// The plan document's fourth bullet reads "unknown speakers are enrolled on
// first interaction (prompt: 'кто это?')". That is refused. Enrolling a voice is
// taking a biometric of a person; doing it automatically the first time someone
// walks past the microphone is doing it to guests, without them being part of
// the exchange, and a TTS question into a room is not consent from whoever
// happens to answer. Enrolment here is an explicit act: an id, a name, and
// samples deliberately recorded for the purpose. An unknown voice stays unknown,
// which the rest of the system is built to cope with.
func (r *Recognizer) Enroll(ctx context.Context, id, name string, samples []audio.Audio) (Profile, error) {
id = NormalizeID(id)
if !ValidID(id) {
return Profile{}, fmt.Errorf("%w: %q", ErrBadID, id)
}
name = strings.TrimSpace(name)
if name == "" {
name = id
}
if len(samples) < MinEnrollSamples {
return Profile{}, fmt.Errorf("%w: %d sample(s), need %d separate ones",
ErrTooShort, len(samples), MinEnrollSamples)
}
var total float64
for i, s := range samples {
if !s.Format.IsValid() {
return Profile{}, fmt.Errorf("%w: sample %d: %+v", ErrBadFormat, i+1, s.Format)
}
total += seconds(s)
}
if total < MinEnrollSeconds {
return Profile{}, fmt.Errorf("%w: %.1fs total, need %.1fs",
ErrTooShort, total, MinEnrollSeconds)
}
// Embed first, store second. A model failure halfway through must not leave
// a half-built profile that would then be matched against.
var (
sum []float32
dim int
)
for i, s := range samples {
vec, err := r.embed(ctx, s)
if err != nil {
return Profile{}, fmt.Errorf("speaker: enroll %q sample %d: %w", id, i+1, err)
}
if sum == nil {
sum = make([]float32, len(vec))
dim = len(vec)
} else if len(vec) != dim {
// One model, one width. A mixed-width average would be nonsense.
return Profile{}, fmt.Errorf("%w: sample %d is %d wide, expected %d",
ErrBadVector, i+1, len(vec), dim)
}
for j, f := range vec {
sum[j] += f
}
}
mean, err := Normalize(sum)
if err != nil {
// Samples that cancel each other out to zero are not one voice.
return Profile{}, fmt.Errorf("speaker: enroll %q: %w", id, err)
}
p := Profile{
ID: id,
Name: name,
Enrolled: r.now().UTC(),
Samples: len(samples),
Dim: dim,
Vec: mean,
}
meta := map[string]string{
"name": p.Name,
"samples": strconv.Itoa(p.Samples),
"enrolled": p.Enrolled.Format(time.RFC3339),
// kind marks the row for anything walking the vector table, so a future
// export or debug page can tell a voiceprint from a note embedding
// without parsing the id.
"kind": "speaker",
}
if err := r.cat.Insert(ctx, Prefix+id, mean, meta); err != nil {
return Profile{}, fmt.Errorf("speaker: enroll %q: %w", id, err)
}
return p, nil
}
+212
View File
@@ -0,0 +1,212 @@
package speaker
import (
"context"
"fmt"
"sort"
"time"
"github.com/kami/maven/internal/audio"
"github.com/kami/maven/internal/memory"
)
// Recognizer holds the embedder and the enrolled profiles.
//
// The profiles live in the shared vector table under the "speaker:" id prefix,
// which is what the plan asked for and what keeps them inside the encrypted
// store rather than in a sidecar file. They are read through memory.Catalog
// (ByPrefix / Delete) rather than Search, because "who is enrolled" is not a
// similarity question and note recall must never rank a voiceprint.
type Recognizer struct {
emb Embedder
cat memory.Catalog
threshold float64
minSec float64
now func() time.Time
}
// Config — the recognizer's knobs, built from config.SpeakerConfig.
type Config struct {
// Threshold — cosine similarity a match must beat. 0 ⇒ DefaultThreshold.
Threshold float64
// MinSeconds — least speech an identification will look at. 0 ⇒
// DefaultMinSeconds.
MinSeconds float64
}
// New builds a Recognizer. emb nil ⇒ Disabled, which is this box's state and
// makes every Identify answer ErrDisabled while enrolment and listing still
// behave sensibly (they refuse for the same reason, with the same error).
func New(emb Embedder, cat memory.Catalog, cfg Config) (*Recognizer, error) {
if cat == nil {
return nil, fmt.Errorf("speaker: no profile store")
}
if emb == nil {
emb = Disabled{}
}
th := cfg.Threshold
if th <= 0 {
th = DefaultThreshold
}
min := cfg.MinSeconds
if min <= 0 {
min = DefaultMinSeconds
}
return &Recognizer{emb: emb, cat: cat, threshold: th, minSec: min, now: time.Now}, nil
}
// Enabled reports whether an embedding model is actually wired. Surfaces use it
// to say "recognition is off" once instead of failing every turn.
func (r *Recognizer) Enabled() bool {
_, disabled := r.emb.(Disabled)
return !disabled
}
// Threshold is the configured match floor, for a status line.
func (r *Recognizer) Threshold() float64 { return r.threshold }
// Identify names the voice in a. ErrUnknown when nothing is close enough, which
// is a normal answer and not a failure: a guest is a guest, and the caller
// carries on with no speaker attached rather than guessing.
//
// Identification never decides whether Maven listens. It annotates the turn.
func (r *Recognizer) Identify(ctx context.Context, a audio.Audio) (Match, error) {
if !a.Format.IsValid() {
return Match{}, fmt.Errorf("%w: %+v", ErrBadFormat, a.Format)
}
if seconds(a) < r.minSec {
return Match{}, fmt.Errorf("%w: %.1fs, need %.1fs", ErrTooShort, seconds(a), r.minSec)
}
vec, err := r.embed(ctx, a)
if err != nil {
return Match{}, err
}
profiles, err := r.List(ctx)
if err != nil {
return Match{}, err
}
if len(profiles) == 0 {
return Match{}, ErrNoProfiles
}
best := Match{Score: -2}
for _, p := range profiles {
if s := Similarity(vec, p.Vec); s > best.Score {
best = Match{Profile: p, Score: s}
}
}
if best.Score < r.threshold {
// The closest profile is reported in the error for a log line, because
// "не узнала, ближе всего Ками на 0.61" is what makes a threshold
// tunable. The caller must not use it as an identification.
return Match{}, fmt.Errorf("%w (closest %s at %.2f, need %.2f)",
ErrUnknown, best.Profile.ID, best.Score, r.threshold)
}
return best, nil
}
// List returns every enrolled profile, sorted by id so a listing is stable.
func (r *Recognizer) List(ctx context.Context) ([]Profile, error) {
recs, err := r.cat.ByPrefix(ctx, Prefix)
if err != nil {
return nil, fmt.Errorf("speaker: list: %w", err)
}
out := make([]Profile, 0, len(recs))
for _, rec := range recs {
out = append(out, profileFromRecord(rec))
}
sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
return out, nil
}
// Get returns one profile by id.
func (r *Recognizer) Get(ctx context.Context, id string) (Profile, error) {
id = NormalizeID(id)
if !ValidID(id) {
return Profile{}, fmt.Errorf("%w: %q", ErrBadID, id)
}
recs, err := r.cat.ByPrefix(ctx, Prefix+id)
if err != nil {
return Profile{}, fmt.Errorf("speaker: get: %w", err)
}
for _, rec := range recs {
if rec.ID == Prefix+id {
return profileFromRecord(rec), nil
}
}
return Profile{}, fmt.Errorf("%w: %q", ErrNotFound, id)
}
// Forget deletes a profile. This is the one operation that must always work:
// a voiceprint is data about a person, and "перестань узнавать её" has to
// actually remove it, not mark it inactive.
func (r *Recognizer) Forget(ctx context.Context, id string) error {
id = NormalizeID(id)
if !ValidID(id) {
return fmt.Errorf("%w: %q", ErrBadID, id)
}
if _, err := r.Get(ctx, id); err != nil {
return err
}
if err := r.cat.Delete(ctx, Prefix+id); err != nil {
return fmt.Errorf("speaker: forget %q: %w", id, err)
}
return nil
}
// embed runs the model and normalises the result.
func (r *Recognizer) embed(ctx context.Context, a audio.Audio) ([]float32, error) {
raw, err := r.emb.Embed(ctx, a)
if err != nil {
return nil, err
}
vec, err := Normalize(raw)
if err != nil {
return nil, err
}
return vec, nil
}
// profileFromRecord reads a stored row back into a Profile. A row with
// unreadable metadata still yields a usable voiceprint — the vector is the part
// that matters, and losing a name should not lose the enrolment.
func profileFromRecord(rec memory.Record) Profile {
p := Profile{
ID: trimPrefix(rec.ID),
Vec: rec.Vec,
Dim: len(rec.Vec),
Name: rec.Meta["name"],
}
if s := rec.Meta["samples"]; s != "" {
p.Samples = atoi(s)
}
if ts := rec.Meta["enrolled"]; ts != "" {
if t, err := time.Parse(time.RFC3339, ts); err == nil {
p.Enrolled = t
}
}
if p.Name == "" {
p.Name = p.ID
}
return p
}
func trimPrefix(id string) string {
if len(id) > len(Prefix) && id[:len(Prefix)] == Prefix {
return id[len(Prefix):]
}
return id
}
// atoi is a tolerant small-integer parse: metadata that is not a number reads
// as 0 rather than failing the whole listing.
func atoi(s string) int {
n := 0
for _, r := range s {
if r < '0' || r > '9' {
return 0
}
n = n*10 + int(r-'0')
}
return n
}
+223
View File
@@ -0,0 +1,223 @@
// Package speaker is voice identification (Vikunja #255,
// docs/plans/10-speaker-recognition.md).
//
// The shape is the same seam internal/vision uses: an Embedder turns audio into
// a voiceprint, a Recognizer compares one against the enrolled profiles, and a
// Disabled floor refuses politely when nothing is wired. On this box nothing is
// wired, and that is the honest state — see "Blocked" below.
//
// # A voiceprint is not like the other vectors
//
// Everything else in the vector table is something he wrote or said. A speaker
// profile is biometric data about a person, quite possibly a person who never
// asked for Maven to exist. The rules that follow from that are in the code:
//
// - Enrolment is explicit and named. There is no "enrol the unknown voice
// automatically" path; see the refusal in enroll.go.
// - A profile is deletable, individually, and Forget really removes the row.
// - Below the threshold the answer is "I do not know", never the closest
// guess. A misattributed fact is worse than an unattributed one.
// - Nothing here gates whether Maven listens or answers. Identification
// annotates a turn; it never authorises one, and an unrecognised voice is
// not turned away.
// - Voiceprints never leave the box. They live in the encrypted store with
// everything else and are never search input to anything external.
//
// # Blocked
//
// There is no speaker-embedding model on this box: no ECAPA-TDNN, no x-vector,
// no wespeaker or titanet ONNX anywhere under /mnt/hdd1 or models/ (checked
// 2026-08-01; the only ONNX files are the e5 text embedder and the piper voice).
// There are also no enrolment samples. So Recognizer runs against Disabled and
// every Identify answers ErrDisabled until a model lands.
//
// The MFCC + GMM "simplest floor" in the plan document is refused rather than
// deferred. A hand-rolled spectral distance would identify people confidently
// and wrongly, and its output would be written into facts as "Ками said this".
// For a biometric, a bad floor is worse than none: no answer is honest, and a
// wrong answer is a false memory about a person.
package speaker
import (
"context"
"errors"
"math"
"strings"
"time"
"github.com/kami/maven/internal/audio"
)
// Prefix — the id prefix speaker profiles carry in the shared vector table.
// It is what ByPrefix enumerates and what keeps voiceprints out of note recall.
const Prefix = "speaker:"
// DefaultThreshold — cosine similarity a match must beat to be a match.
//
// 0.7 is the usual operating point for ECAPA-style embeddings on clean speech
// and it is deliberately on the strict side here. The two error directions are
// not symmetric: refusing to name a voice costs a "не узнала", while naming the
// wrong person writes his wife's remark into a fact attributed to him.
const DefaultThreshold = 0.7
// DefaultMinSeconds — how much speech an identification needs. Under about two
// seconds a voiceprint is mostly noise and the similarity score is not worth
// reading.
const DefaultMinSeconds = 2.0
// MinEnrollSamples / MinEnrollSeconds — what enrolment requires. Several
// separate utterances, not one long one: a profile built from a single sentence
// encodes that sentence's prosody as much as the voice.
const (
MinEnrollSamples = 3
MinEnrollSeconds = 9.0
)
// Errors callers distinguish.
var (
// ErrDisabled — no embedding model is wired. The state of this box.
ErrDisabled = errors.New("speaker: recognition is not configured")
// ErrTooShort — not enough speech to say anything about.
ErrTooShort = errors.New("speaker: not enough audio")
// ErrUnknown — audio embedded fine, but no enrolled profile is close
// enough. Not an error in the sense of something being broken: it is the
// correct answer for a guest, and the caller should carry on without a
// speaker rather than treat the turn as failed.
ErrUnknown = errors.New("speaker: voice not recognised")
// ErrNoProfiles — nobody is enrolled yet.
ErrNoProfiles = errors.New("speaker: nobody is enrolled")
// ErrNotFound — no profile with that id.
ErrNotFound = errors.New("speaker: no such profile")
// ErrBadID — an id that is empty or carries characters an id should not.
ErrBadID = errors.New("speaker: invalid profile id")
// ErrBadFormat — audio that is not the canonical 16 kHz mono PCM shape.
ErrBadFormat = errors.New("speaker: audio format not supported")
// ErrBadVector — an embedder returned something unusable (empty, or all
// zeroes, which normalises to nothing and would match everything equally).
ErrBadVector = errors.New("speaker: embedder returned an unusable vector")
)
// Embedder turns speech into a voiceprint. Implementations are expected to
// return an L2-normalised vector, because the whole store compares by dot
// product; Normalize is applied anyway rather than trusted.
//
// This is the seam a downloaded ECAPA-TDNN ONNX model plugs into. It is an
// interface rather than a concrete ONNX type so the package is testable with no
// model on disk, which is the only way it could be tested here at all.
type Embedder interface {
Embed(ctx context.Context, a audio.Audio) ([]float32, error)
// Dim is the vector width, used to reject a profile recorded with a
// different model rather than silently scoring it as zero.
Dim() int
}
// Disabled is the floor: no model, no answers, no guesses.
type Disabled struct{}
// Embed always fails with ErrDisabled.
func (Disabled) Embed(context.Context, audio.Audio) ([]float32, error) { return nil, ErrDisabled }
// Dim is 0 for the disabled embedder.
func (Disabled) Dim() int { return 0 }
// Profile — one enrolled voice.
//
// Name is what she calls the person out loud ("Ками"). ID is the stable handle
// used in sources and metadata. Samples records how many utterances the
// voiceprint was averaged from, so a profile enrolled from the bare minimum is
// visibly weaker than one built from ten.
type Profile struct {
ID string `json:"id"`
Name string `json:"name"`
Enrolled time.Time `json:"enrolled"`
Samples int `json:"samples"`
Dim int `json:"dim"`
// Vec is the voiceprint. Not serialised to any surface: a listing tells him
// who is enrolled, it does not hand out the biometric itself.
Vec []float32 `json:"-"`
}
// Source is what a fact or note written during this speaker's turn is tagged
// with, e.g. "tap:voice:speaker:kami". Attribution belongs in the source rather
// than in the text, so it can be corrected or dropped later without rewriting
// what was said.
func (p Profile) Source(base string) string {
if p.ID == "" {
return base
}
return base + ":" + Prefix + p.ID
}
// Match — an identification result. Score is cosine similarity in [-1, 1].
type Match struct {
Profile Profile
Score float64
}
// ValidID reports whether an id is usable as a profile handle. Deliberately
// narrow: lowercase letters, digits, dash and underscore. Ids end up in note
// sources and in vector-table keys, so a permissive id would be a way to write
// into a neighbouring key space.
func ValidID(id string) bool {
if id == "" || len(id) > 64 {
return false
}
for _, r := range id {
switch {
case r >= 'a' && r <= 'z', r >= '0' && r <= '9', r == '-', r == '_':
default:
return false
}
}
return true
}
// NormalizeID lowercases and trims a proposed id before validating it, so
// "Ками" typed as "Kami " does not fail for a reason nobody can see.
func NormalizeID(id string) string {
return strings.ToLower(strings.TrimSpace(id))
}
// Normalize returns an L2-normalised copy of v, or ErrBadVector when there is
// nothing to normalise. A zero vector is refused rather than passed on: it
// scores 0 against everything, which reads as "no match" but for the wrong
// reason and would hide a broken embedder.
func Normalize(v []float32) ([]float32, error) {
if len(v) == 0 {
return nil, ErrBadVector
}
var sum float64
for _, f := range v {
if math.IsNaN(float64(f)) || math.IsInf(float64(f), 0) {
return nil, ErrBadVector
}
sum += float64(f) * float64(f)
}
norm := math.Sqrt(sum)
if norm == 0 {
return nil, ErrBadVector
}
out := make([]float32, len(v))
for i, f := range v {
out[i] = float32(float64(f) / norm)
}
return out, nil
}
// Similarity is the cosine similarity of two L2-normalised vectors. Different
// widths score 0: a profile enrolled with another model must not accidentally
// match, and 0 is below every sane threshold.
func Similarity(a, b []float32) float64 {
if len(a) != len(b) || len(a) == 0 {
return 0
}
var sum float64
for i := range a {
sum += float64(a[i]) * float64(b[i])
}
return sum
}
// seconds is the playback length of a frame, for the minimum-audio checks.
func seconds(a audio.Audio) float64 { return a.Duration() }
+397
View File
@@ -0,0 +1,397 @@
package speaker
import (
"context"
"errors"
"math"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/audio"
"github.com/kami/maven/internal/memory"
)
// fakeEmbedder returns a fixed vector per "voice", so a test can enrol one
// person and present another without a model. Wobble adds a small perturbation
// so repeated samples of one voice are close but not identical, which is what a
// real embedder produces.
type fakeEmbedder struct {
vec []float32
err error
calls int
wobble float32
}
func (f *fakeEmbedder) Embed(_ context.Context, _ audio.Audio) ([]float32, error) {
f.calls++
if f.err != nil {
return nil, f.err
}
out := append([]float32(nil), f.vec...)
if f.wobble != 0 && len(out) > 1 {
out[0] += f.wobble * float32(f.calls)
out[1] -= f.wobble * float32(f.calls)
}
return out, nil
}
func (f *fakeEmbedder) Dim() int { return len(f.vec) }
// speech builds n seconds of the canonical audio shape.
func speech(sec float64) audio.Audio {
return audio.Audio{Format: audio.PCM16kMono, Bytes: make([]byte, int(sec*16000)*2)}
}
func newRec(t *testing.T, emb Embedder) (*Recognizer, memory.Catalog) {
t.Helper()
cat := memory.NewInMemoryStore()
r, err := New(emb, cat, Config{})
if err != nil {
t.Fatal(err)
}
return r, cat
}
func enrolSamples(n int, sec float64) []audio.Audio {
out := make([]audio.Audio, n)
for i := range out {
out[i] = speech(sec)
}
return out
}
// The state of this box: no model on disk. Every identification refuses rather
// than guessing, and it says why.
func TestDisabledRefusesEverything(t *testing.T) {
r, _ := newRec(t, nil)
if r.Enabled() {
t.Error("a recognizer with no model reports itself enabled")
}
if _, err := r.Identify(context.Background(), speech(5)); !errors.Is(err, ErrDisabled) {
t.Errorf("Identify = %v, want ErrDisabled", err)
}
if _, err := r.Enroll(context.Background(), "kami", "Ками", enrolSamples(3, 4)); !errors.Is(err, ErrDisabled) {
t.Errorf("Enroll = %v, want ErrDisabled", err)
}
// Listing still works: knowing that nobody is enrolled needs no model.
got, err := r.List(context.Background())
if err != nil || len(got) != 0 {
t.Errorf("List = %v, %v", got, err)
}
}
func TestNewRequiresAProfileStore(t *testing.T) {
if _, err := New(nil, nil, Config{}); err == nil {
t.Error("built a recognizer with nowhere to keep profiles")
}
}
func TestEnrollThenIdentify(t *testing.T) {
emb := &fakeEmbedder{vec: []float32{1, 0, 0, 0}, wobble: 0.01}
r, _ := newRec(t, emb)
ctx := context.Background()
p, err := r.Enroll(ctx, "Kami ", "Ками", enrolSamples(3, 4))
if err != nil {
t.Fatalf("enroll: %v", err)
}
if p.ID != "kami" {
t.Errorf("id = %q, want the normalised %q", p.ID, "kami")
}
if p.Name != "Ками" || p.Samples != 3 || p.Dim != 4 {
t.Errorf("profile = %+v", p)
}
m, err := r.Identify(ctx, speech(5))
if err != nil {
t.Fatalf("identify: %v", err)
}
if m.Profile.ID != "kami" || m.Profile.Name != "Ками" {
t.Errorf("match = %+v", m)
}
if m.Score < r.Threshold() {
t.Errorf("score %.3f is below the threshold it supposedly passed", m.Score)
}
}
// The error direction that matters. Naming the wrong person writes a false
// memory about them, so a voice that is not close enough gets no name at all.
func TestUnfamiliarVoiceIsNotGuessed(t *testing.T) {
emb := &fakeEmbedder{vec: []float32{1, 0, 0, 0}}
r, _ := newRec(t, emb)
ctx := context.Background()
if _, err := r.Enroll(ctx, "kami", "Ками", enrolSamples(3, 4)); err != nil {
t.Fatal(err)
}
// A different voice: orthogonal voiceprint, similarity 0.
emb.vec = []float32{0, 1, 0, 0}
m, err := r.Identify(ctx, speech(5))
if !errors.Is(err, ErrUnknown) {
t.Fatalf("Identify = %v, want ErrUnknown", err)
}
if m.Profile.ID != "" {
t.Errorf("a refused identification still handed back %q", m.Profile.ID)
}
// The log line needs the near miss to make the threshold tunable.
if !contains(err.Error(), "kami") {
t.Errorf("error does not name the closest profile: %v", err)
}
}
// Just under the threshold is still unknown. A boundary this important gets its
// own test rather than being implied.
func TestThresholdIsAFloorNotASuggestion(t *testing.T) {
cat := memory.NewInMemoryStore()
emb := &fakeEmbedder{vec: []float32{1, 0}}
r, err := New(emb, cat, Config{Threshold: 0.9})
if err != nil {
t.Fatal(err)
}
ctx := context.Background()
if _, err := r.Enroll(ctx, "kami", "Ками", enrolSamples(3, 4)); err != nil {
t.Fatal(err)
}
// cos ≈ 0.866, comfortably similar and still not similar enough.
emb.vec = []float32{0.866, 0.5}
if _, err := r.Identify(ctx, speech(5)); !errors.Is(err, ErrUnknown) {
t.Fatalf("0.866 against a 0.9 threshold = %v, want ErrUnknown", err)
}
}
func TestShortAudioIsRefusedBeforeTheModelRuns(t *testing.T) {
emb := &fakeEmbedder{vec: []float32{1, 0}}
r, _ := newRec(t, emb)
if _, err := r.Identify(context.Background(), speech(0.5)); !errors.Is(err, ErrTooShort) {
t.Fatalf("got %v, want ErrTooShort", err)
}
if emb.calls != 0 {
t.Error("a half-second of audio was sent to the model anyway")
}
}
func TestWrongAudioFormatIsRefused(t *testing.T) {
r, _ := newRec(t, &fakeEmbedder{vec: []float32{1, 0}})
bad := audio.Audio{
Format: audio.Format{SampleRate: 44100, Channels: 2, SampleBits: 16, Encoding: "pcm_s16le"},
Bytes: make([]byte, 44100*4*5),
}
if _, err := r.Identify(context.Background(), bad); !errors.Is(err, ErrBadFormat) {
t.Fatalf("got %v, want ErrBadFormat", err)
}
}
func TestIdentifyWithNobodyEnrolled(t *testing.T) {
r, _ := newRec(t, &fakeEmbedder{vec: []float32{1, 0}})
if _, err := r.Identify(context.Background(), speech(5)); !errors.Is(err, ErrNoProfiles) {
t.Fatalf("got %v, want ErrNoProfiles", err)
}
}
// Enrolment is an explicit act with real samples behind it, not a byproduct of
// someone speaking once.
func TestEnrollmentRequiresSeveralRealSamples(t *testing.T) {
r, _ := newRec(t, &fakeEmbedder{vec: []float32{1, 0}})
ctx := context.Background()
cases := []struct {
name string
samples []audio.Audio
}{
{"one long sample", enrolSamples(1, 30)},
{"two samples", enrolSamples(2, 10)},
{"three samples but seconds of audio", enrolSamples(3, 1)},
{"none at all", nil},
}
for _, c := range cases {
if _, err := r.Enroll(ctx, "kami", "Ками", c.samples); !errors.Is(err, ErrTooShort) {
t.Errorf("%s: %v, want ErrTooShort", c.name, err)
}
}
}
func TestEnrollRejectsBadIDs(t *testing.T) {
r, _ := newRec(t, &fakeEmbedder{vec: []float32{1, 0}})
for _, id := range []string{"", " ", "../etc/passwd", "speaker:kami", "имя", "a/b", "x y"} {
if _, err := r.Enroll(context.Background(), id, "n", enrolSamples(3, 4)); !errors.Is(err, ErrBadID) {
t.Errorf("id %q accepted or wrong error: %v", id, err)
}
}
}
// Re-enrolling replaces the voiceprint. Leaving the old one searchable would
// mean a person's rejected profile keeps matching them.
func TestReEnrollReplaces(t *testing.T) {
emb := &fakeEmbedder{vec: []float32{1, 0, 0}}
r, _ := newRec(t, emb)
ctx := context.Background()
if _, err := r.Enroll(ctx, "kami", "Ками", enrolSamples(3, 4)); err != nil {
t.Fatal(err)
}
emb.vec = []float32{0, 1, 0}
if _, err := r.Enroll(ctx, "kami", "Ками", enrolSamples(4, 4)); err != nil {
t.Fatal(err)
}
list, err := r.List(ctx)
if err != nil {
t.Fatal(err)
}
if len(list) != 1 {
t.Fatalf("%d profiles after re-enrolling one person", len(list))
}
if list[0].Samples != 4 {
t.Errorf("sample count = %d, want the new 4", list[0].Samples)
}
// The new voiceprint is the one that matches.
if m, err := r.Identify(ctx, speech(5)); err != nil || m.Score < 0.99 {
t.Errorf("identify after re-enrol: %v (score %.3f)", err, m.Score)
}
}
// "Перестань узнавать её" has to actually delete the biometric.
func TestForgetRemovesTheVoiceprint(t *testing.T) {
emb := &fakeEmbedder{vec: []float32{1, 0}}
r, cat := newRec(t, emb)
ctx := context.Background()
if _, err := r.Enroll(ctx, "guest", "Гостья", enrolSamples(3, 4)); err != nil {
t.Fatal(err)
}
if err := r.Forget(ctx, "Guest "); err != nil {
t.Fatalf("forget: %v", err)
}
recs, err := cat.ByPrefix(ctx, Prefix)
if err != nil {
t.Fatal(err)
}
if len(recs) != 0 {
t.Errorf("%d row(s) survived Forget", len(recs))
}
if _, err := r.Get(ctx, "guest"); !errors.Is(err, ErrNotFound) {
t.Errorf("Get after Forget = %v, want ErrNotFound", err)
}
if err := r.Forget(ctx, "guest"); !errors.Is(err, ErrNotFound) {
t.Errorf("second Forget = %v, want ErrNotFound", err)
}
}
// Voiceprints share the vector table with note and fact embeddings, so the
// prefix has to actually partition it.
func TestProfilesDoNotCollideWithNoteVectors(t *testing.T) {
emb := &fakeEmbedder{vec: []float32{1, 0}}
r, cat := newRec(t, emb)
ctx := context.Background()
if err := cat.Insert(ctx, "note:1", []float32{1, 0}, map[string]string{"text": "заметка"}); err != nil {
t.Fatal(err)
}
if _, err := r.Enroll(ctx, "kami", "Ками", enrolSamples(3, 4)); err != nil {
t.Fatal(err)
}
list, err := r.List(ctx)
if err != nil {
t.Fatal(err)
}
if len(list) != 1 || list[0].ID != "kami" {
t.Errorf("listing picked up a non-speaker row: %+v", list)
}
// And an identical note vector is never returned as a match.
m, err := r.Identify(ctx, speech(5))
if err != nil {
t.Fatal(err)
}
if m.Profile.ID != "kami" {
t.Errorf("matched %q", m.Profile.ID)
}
}
func TestEmbedderFailurePropagates(t *testing.T) {
emb := &fakeEmbedder{vec: []float32{1, 0}, err: errors.New("onnx fell over")}
r, _ := newRec(t, emb)
if _, err := r.Identify(context.Background(), speech(5)); err == nil {
t.Error("a model failure was reported as a successful identification")
}
if _, err := r.Enroll(context.Background(), "kami", "К", enrolSamples(3, 4)); err == nil {
t.Error("a model failure produced a profile")
}
}
// A zero vector scores 0 against everything, which reads as "no match" for the
// wrong reason and would hide a broken model.
func TestUnusableVectorsAreRefused(t *testing.T) {
r, _ := newRec(t, &fakeEmbedder{vec: []float32{0, 0, 0}})
if _, err := r.Enroll(context.Background(), "kami", "К", enrolSamples(3, 4)); !errors.Is(err, ErrBadVector) {
t.Errorf("zero vector: %v, want ErrBadVector", err)
}
if _, err := Normalize(nil); !errors.Is(err, ErrBadVector) {
t.Errorf("empty: %v", err)
}
if _, err := Normalize([]float32{float32(nan())}); !errors.Is(err, ErrBadVector) {
t.Errorf("NaN: %v", err)
}
}
func TestNormalizeProducesAUnitVector(t *testing.T) {
v, err := Normalize([]float32{3, 4})
if err != nil {
t.Fatal(err)
}
if got := Similarity(v, v); got < 0.999 || got > 1.001 {
t.Errorf("self-similarity = %f, want 1", got)
}
}
// A profile enrolled with another model must not accidentally match.
func TestDifferentWidthsScoreZero(t *testing.T) {
if got := Similarity([]float32{1, 0}, []float32{1, 0, 0}); got != 0 {
t.Errorf("mismatched widths scored %f", got)
}
}
// Attribution belongs in the source, so it can be corrected without rewriting
// what was said.
func TestProfileSource(t *testing.T) {
p := Profile{ID: "kami"}
if got := p.Source("tap:voice"); got != "tap:voice:speaker:kami" {
t.Errorf("source = %q", got)
}
var anon Profile
if got := anon.Source("tap:voice"); got != "tap:voice" {
t.Errorf("unattributed source = %q, want the base unchanged", got)
}
}
func TestValidID(t *testing.T) {
for _, ok := range []string{"kami", "guest-2", "a_b", "x"} {
if !ValidID(ok) {
t.Errorf("%q rejected", ok)
}
}
for _, bad := range []string{"", "Kami", "имя", "a b", "a/b", "a:b", "..", strings.Repeat("a", 65)} {
if ValidID(bad) {
t.Errorf("%q accepted", bad)
}
}
}
func TestProfileMetadataSurvivesARoundTrip(t *testing.T) {
emb := &fakeEmbedder{vec: []float32{1, 0}}
r, _ := newRec(t, emb)
r.now = func() time.Time { return time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC) }
ctx := context.Background()
if _, err := r.Enroll(ctx, "kami", "Ками", enrolSamples(3, 4)); err != nil {
t.Fatal(err)
}
got, err := r.Get(ctx, "kami")
if err != nil {
t.Fatal(err)
}
if got.Name != "Ками" || got.Samples != 3 {
t.Errorf("profile = %+v", got)
}
if !got.Enrolled.Equal(time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC)) {
t.Errorf("enrolled = %v", got.Enrolled)
}
}
func contains(s, sub string) bool { return strings.Contains(s, sub) }
func nan() float64 { return math.NaN() }
+54 -1
View File
@@ -8,6 +8,7 @@ import (
"fmt"
"math"
"sort"
"strings"
"time"
"github.com/kami/maven/internal/memory"
@@ -37,8 +38,10 @@ func (s *Store) VectorMemory() *MemoryStore {
return &MemoryStore{db: s.db}
}
// compile-time check: MemoryStore satisfies the memory.Store interface.
// compile-time check: MemoryStore satisfies the memory.Store interface, and the
// wider Catalog that speaker profiles need (enumerate by prefix, delete by id).
var _ memory.Store = (*MemoryStore)(nil)
var _ memory.Catalog = (*MemoryStore)(nil)
// Insert upserts a vector by id: a repeated id replaces the prior row rather
// than accumulating duplicates (the note/fact ids are stable and unique, so a
@@ -95,6 +98,56 @@ func (m *MemoryStore) Search(ctx context.Context, vec []float32, topK int) ([]me
return out, nil
}
// ByPrefix returns every row whose id starts with prefix, vectors included.
//
// This is not a similarity query and deliberately does not score anything:
// listing the enrolled voices is a question about which rows exist, and asking
// it through Search would mean inventing a query vector to rank them by. The
// prefix is matched with LIKE against an escaped pattern, so a profile id
// containing % or _ cannot widen the match.
func (m *MemoryStore) ByPrefix(ctx context.Context, prefix string) ([]memory.Record, error) {
pattern := escapeLike(prefix) + "%"
rows, err := m.db.QueryContext(ctx,
`SELECT id, vec, meta FROM memory_vectors WHERE id LIKE ? ESCAPE '\'`, pattern)
if err != nil {
return nil, fmt.Errorf("memory: by prefix %q: %w", prefix, err)
}
defer rows.Close()
var out []memory.Record
for rows.Next() {
var id, metaJSON string
var blob []byte
if err := rows.Scan(&id, &blob, &metaJSON); err != nil {
return nil, fmt.Errorf("memory: row: %w", err)
}
meta := map[string]string{}
if err := json.Unmarshal([]byte(metaJSON), &meta); err != nil {
return nil, fmt.Errorf("memory: unmarshal meta for %q: %w", id, err)
}
out = append(out, memory.Record{ID: id, Vec: decodeVec(blob), Meta: meta})
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("memory: rows: %w", err)
}
return out, nil
}
// Delete removes one vector by id. A row that is not there is not an error —
// "forget this voice" is satisfied either way.
func (m *MemoryStore) Delete(ctx context.Context, id string) error {
if _, err := m.db.ExecContext(ctx, `DELETE FROM memory_vectors WHERE id = ?`, id); err != nil {
return fmt.Errorf("memory: delete %q: %w", id, err)
}
return nil
}
// escapeLike neutralises the LIKE wildcards in a literal prefix.
func escapeLike(s string) string {
r := strings.NewReplacer(`\`, `\\`, `%`, `\%`, `_`, `\_`)
return r.Replace(s)
}
// encodeVec serializes a float32 slice as little-endian IEEE-754 bytes (4 bytes
// per element) for the BLOB column.
func encodeVec(v []float32) []byte {