diff --git a/cmd/mavend/capture.go b/cmd/mavend/capture.go new file mode 100644 index 0000000..de53a70 --- /dev/null +++ b/cmd/mavend/capture.go @@ -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 +} diff --git a/cmd/mavend/feeds.go b/cmd/mavend/feeds.go index 7797697..99d9d8f 100644 --- a/cmd/mavend/feeds.go +++ b/cmd/mavend/feeds.go @@ -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 } diff --git a/cmd/mavend/main.go b/cmd/mavend/main.go index f020711..16c4921 100644 --- a/cmd/mavend/main.go +++ b/cmd/mavend/main.go @@ -335,7 +335,11 @@ 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) } // WrapKeyFn — wraps the env key with a passkey credential public key and @@ -479,7 +483,8 @@ 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) // Start voice server. if voiceW != nil { diff --git a/cmd/mavend/vision.go b/cmd/mavend/vision.go index 5516e7f..93b0c2d 100644 --- a/cmd/mavend/vision.go +++ b/cmd/mavend/vision.go @@ -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 } diff --git a/cmd/mavend/voicewire.go b/cmd/mavend/voicewire.go index 77e787e..5b4a60a 100644 --- a/cmd/mavend/voicewire.go +++ b/cmd/mavend/voicewire.go @@ -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 diff --git a/docs/plans/08-hearing.md b/docs/plans/08-hearing.md index a7cb69f..41c30a0 100644 --- a/docs/plans/08-hearing.md +++ b/docs/plans/08-hearing.md @@ -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`. diff --git a/internal/auth/auth_test.go b/internal/auth/auth_test.go index 67907af..cf13d1a 100644 --- a/internal/auth/auth_test.go +++ b/internal/auth/auth_test.go @@ -416,3 +416,29 @@ 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) + } +} diff --git a/internal/auth/policy.go b/internal/auth/policy.go index a22ecd4..809bf04 100644 --- a/internal/auth/policy.go +++ b/internal/auth/policy.go @@ -60,6 +60,21 @@ 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.MethodWriteFact: return AuthWrite case ipc.MethodAssertStepUp: @@ -95,6 +110,9 @@ 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, // The read side of the model swap: which model is resident, which ones are // allowlisted. It loads nothing and changes nothing. ipc.MethodModelStatus: diff --git a/internal/capture/capture.go b/internal/capture/capture.go new file mode 100644 index 0000000..c9cb852 --- /dev/null +++ b/internal/capture/capture.go @@ -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 +} diff --git a/internal/capture/capture_test.go b/internal/capture/capture_test.go new file mode 100644 index 0000000..e1c4268 --- /dev/null +++ b/internal/capture/capture_test.go @@ -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)) + } +} diff --git a/internal/capture/summarize.go b/internal/capture/summarize.go new file mode 100644 index 0000000..d539472 --- /dev/null +++ b/internal/capture/summarize.go @@ -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.5–3 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 +} diff --git a/internal/capture/summarize_test.go b/internal/capture/summarize_test.go new file mode 100644 index 0000000..e5e633a --- /dev/null +++ b/internal/capture/summarize_test.go @@ -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 +} diff --git a/internal/config/config.go b/internal/config/config.go index 4381c6a..5940110 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -205,6 +205,13 @@ 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"` + // 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 +565,60 @@ 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 +} + // WeatherConfig configures the weather provider for voice queries. type WeatherConfig struct { Provider string `json:"provider,omitempty"` // "open-meteo" or "" → stub diff --git a/internal/config/senses_test.go b/internal/config/senses_test.go index b0cc5aa..3997cc0 100644 --- a/internal/config/senses_test.go +++ b/internal/config/senses_test.go @@ -16,6 +16,70 @@ 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") + } +} + +// 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. diff --git a/internal/ipc/api.go b/internal/ipc/api.go index ae6f69a..1e961b7 100644 --- a/internal/ipc/api.go +++ b/internal/ipc/api.go @@ -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,84 @@ 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"` +} + // 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 diff --git a/internal/ipc/capture_test.go b/internal/ipc/capture_test.go new file mode 100644 index 0000000..11ce16f --- /dev/null +++ b/internal/ipc/capture_test.go @@ -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) + } +} diff --git a/internal/ipc/client.go b/internal/ipc/client.go index 9f0ca88..6e8bd0f 100644 --- a/internal/ipc/client.go +++ b/internal/ipc/client.go @@ -473,6 +473,48 @@ 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 +} + // 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 diff --git a/internal/ipc/server.go b/internal/ipc/server.go index ebce760..15b380a 100644 --- a/internal/ipc/server.go +++ b/internal/ipc/server.go @@ -459,6 +459,20 @@ 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 + // 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 +503,13 @@ 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) + // 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 +678,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 +972,58 @@ 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 MethodModelStatus: if s.ModelStatusFn != nil { resp, err := s.ModelStatusFn(ctx) diff --git a/internal/ipc/wire.go b/internal/ipc/wire.go index 7a351ff..da56b80 100644 --- a/internal/ipc/wire.go +++ b/internal/ipc/wire.go @@ -55,6 +55,10 @@ 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" ) // Request — one frame from module to core. Params is the JSON-encoded argument