Compare commits

..

4 Commits

Author SHA1 Message Date
kami 45b5e16eff Normalize every intake path into one event envelope (#283)
Things arrive at Maven from eight directions — a relayed Android
notification on POST /api/ambient, mail candidates from mavmaild, RSS
items, changed pages from the crawler, zenmoney and wg reads from
mavpoll, CalDAV events, presence probes, meeting transcripts and image
descriptions. Each grew its own shape and its own log line, and nothing
could answer "what came in today, from where".

internal/event is that answer: a flat source-agnostic envelope (Source,
Kind, EntityIDs, Title, Body, Priority, OccurredAt, Payload) plus a
bounded in-memory journal. Both are pure — Publish and Normalize take
`now` as a parameter, so no clock read sits on a path a replay would
drive.

Adopting it did not touch eight callers, because every intake path
already converges on three ipc.CoreAPI methods: WriteFact, WriteNote and
CaptureTask. cmd/mavend/intake.go decorates that ONE interface, so
mavweb, mavcaldav, mavpoll, mavmaild and the in-core feed/crawl/capture/
vision workers publish envelopes without knowing events exist. The lone
exception is cmd/mavend/mail.go, which captures through the store
directly and now publishes explicitly.

Nothing dispatches on an event. It is a report that something arrived,
never an instruction to speak — "a feed item appeared" becoming a
notification is the nag this repo refuses. Digestion may read the
journal later; it will still go through internal/loop's rules and the
severity/presence routing table.

Read surface: ipc.MethodRecentEvents (AuthRead, daemon-cached like
TickTrace — a bare store cannot serve a ring) and a read-only /events
page in mavweb.

Production is unchanged when nobody is watching: a nil *event.Bus makes
Publish a no-op and newIntakeAPI returns the wrapped API untouched, so
config.intake_journal < 0 leaves no decorator on the call path at all.
The default is 512 entries; the "off unless configured" rule is for
capabilities that reach out, and a bounded in-memory log of writes core
already performed reaches nowhere.

Verified: make build, make test (go test -race) both clean. New tests
cover the envelope and ring (internal/event, 95.7%), the decorator's
invariants — a failed write publishes nothing, a deduped capture
publishes nothing, OccurredAt is the fact's Ts and not notice time — and
the /events page including escaping of feed-supplied titles.
2026-08-01 06:05:00 +04:00
kami 4eca20bd94 Derive the cold-start unlock key from the passkey PRF, not the public key (#14)
Cold-start unlock wrapped the database key under the credential *public* key.
A public key is public: mavweb writes it verbatim to passkeys.json, normally in
the same state dir as db_key.wrapped, so anyone holding both files recovered the
database key offline with no authenticator involved. The wrapped blob was a
plaintext key with extra steps.

The secret is now the WebAuthn PRF extension output — 32 bytes the authenticator
computes over a fixed salt and never stores anywhere. The blob gains a version:

  v2:  "MVNKW2\x00" || salt || nonce || AES-256-GCM(key), magic as AAD
  v1:  salt || nonce || AES-256-GCM(key)                  (read-only)

v1 still opens so an existing deployment is not bricked, and reports itself so
the daemon can log a SECURITY line telling him to re-enroll. Nothing writes v1.
The magic is authenticated, so a v2 blob cannot be stripped and re-read as v1.

Four other defects on the same path:

  - The locked-boot store was opened on an IPC goroutine inside UnlockFn and
    never closed. Close is what re-encrypts the tmpfs working copy back over
    the ciphertext, so every write of a cold-started session was lost silently
    on the next boot. daemonLock now owns the store and seals it at shutdown.
  - MethodUnlock was reachable by anything on the box; the socket is same-uid
    and cannot authenticate its caller. It now requires a passkey assertion
    that mavweb verified first.
  - Concurrent unlocks would each open a store and wire a daemon. One at a
    time, and never a second one.
  - The hand-rolled HKDF keyed the expand step with the salt instead of the
    PRK. Replaced with crypto/hkdf.

Key wrapping moves from enrolment to the first assertion, because create() does
not produce a PRF result on most authenticators — only a support flag. An
authenticator without PRF now writes no wrapped file at all rather than one
that looks protected and is not, and the page says so.

Verified: make build, make test. New tests cover the v2 round trip, a wrong
secret, every single-bit tamper, truncation, the v1 downgrade attempt, legacy
v1 reads, non-32-byte and all-zero secrets, the ipc wire field, locked-mode
default-deny, a forged assertion never reaching the unlock path, seal-on-
shutdown after a cold start, and that nothing in the state dir contains the
plaintext key. The PRF round trip against real hardware is a QA step.

Vikunja #14
2026-08-01 05:49:27 +04:00
kami fed33a4e16 Stop mavwaked from hearing itself, and add barge-in (#287)
Playback was `go playAudio(reply)` — fire and forget, nobody holding the
process handle. Two audible consequences fell out of that.

She answered herself. The capture loop kept feeding the VAD while the
speaker was running, so her own reply came back in through the mic,
tripped the VAD, and was shipped to the daemon as a fresh command. There
is no acoustic echo canceller in this pipeline, so the fix is
half-duplex: while she is speaking, the capture side is muted. That part
is unconditional — it repairs a defect, it is not a new capability.

And talking over her did nothing, because there was no handle to cancel.
-barge-in now cuts playback when sustained energy clears a room-tuned
threshold (-barge-in-rms, default 0.12 normalised, over -barge-in-frames
consecutive frames, default 5). It is off by default: without an echo
canceller the only way to tell "he is talking over her" from "the mic is
hearing her" is that he is much louder, and how much louder depends on
where the mic sits.

The frame decision moved out of main.go into session.feed, behind a
player and an utteranceSender interface, so all of it is testable with
no mic, no speaker and no daemon. Nine tests cover the self-hearing
case, the off-by-default case, the consecutive-frame requirement,
speaker-leak-level audio not triggering, capturing the interrupting
utterance after a cut, and failed round-trips not starting playback.

The other seven items on #287 (partial STT, per-segment retry, mic
profiles, noise-floor calibration, short-response-while-speaking) are
untouched and stay on the task.
2026-08-01 05:36:13 +04:00
kami 62cc072f8c Add golden-audio STT tests against real whisper.cpp (#288)
Four committed WAV fixtures go through the real whisper.cpp binding in
cmd/mavsttd, so a wrong model, a wrong language hint, a broken resample
or a regressed silence gate fails `make test` instead of surfacing as
Maven mishearing him.

The fixtures are piper-synthesised, not recorded: scripts/gen-stt-fixtures.sh
drives the vendored piper with the ru_RU-irina voice Maven already speaks
with, so nothing of the owner's voice is committed and every fixture is
reproducible. 360K total for three Russian clips and one English.

Matching is tolerant on purpose. Golden transcripts move with the model,
so each case asserts intent-carrying keywords (prefix match, so Russian
inflection does not fail it) plus a word error rate ceiling, not an exact
string. The matcher is unit-tested on its own and needs no model.

TestGoldenAudioTranscription skips when models/stt/ggml-small.bin is
absent, so `make test` still passes on a box without models.
TestGoldenFixturesAreCanonical runs everywhere and checks the WAVs are
16k mono s16le and would clear mavsttd's own silence gate.
2026-08-01 05:32:10 +04:00
42 changed files with 3638 additions and 279 deletions
+12 -1
View File
@@ -16,7 +16,7 @@ PIPER_BIN := $(shell pwd)/deps/piper/piper
PIPER_MODEL := $(shell pwd)/models/tts/ru_RU-irina-medium.onnx
PIPER_ESPEAK := $(shell pwd)/deps/piper/espeak-ng-data
.PHONY: all build build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav clean test fmt-check vet run-stt run-tts run-web download-embedder deps-go eval-router eval-recall eval-phrasing eval-models
.PHONY: stt-fixtures test-stt-golden all build build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav clean test fmt-check vet run-stt run-tts run-web download-embedder deps-go eval-router eval-recall eval-phrasing eval-models
all: build
@@ -139,6 +139,17 @@ eval-models:
MAVEN_LLM_URL="$(MAVEN_LLM_URL)" $(GO) test -v -count=1 -timeout 60m \
-run TestLLMRouterBaseline ./internal/router/eval/
# stt-fixtures — regenerate the golden STT audio in cmd/mavsttd/testdata from
# the piper voices (#288). The committed WAVs are synthesised, never recorded,
# so this is the only way they should ever change. TestGoldenAudioTranscription
# then scores them against ggml-small; it self-skips when the model is absent.
stt-fixtures:
./scripts/gen-stt-fixtures.sh
test-stt-golden:
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
$(GO) test -v -count=1 -run TestGolden ./cmd/mavsttd/
run-stt: build-stt
LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
./mavsttd -socket /tmp/maven/stt.sock -model $(WHISPER_MODEL)
+190
View File
@@ -0,0 +1,190 @@
package main
import (
"bytes"
"context"
"crypto/rand"
"io"
"os"
"path/filepath"
"testing"
"time"
"github.com/kami/maven/internal/store"
"github.com/kami/maven/internal/webauthn"
)
func randBytes(t *testing.T, n int) []byte {
t.Helper()
b := make([]byte, n)
if _, err := io.ReadFull(rand.Reader, b); err != nil {
t.Fatalf("rand: %v", err)
}
b[0] |= 1
return b
}
func TestDaemonLockStartsLockedAndFlips(t *testing.T) {
dl := newDaemonLock(true)
if !dl.isLocked() {
t.Fatal("newDaemonLock(true) is not locked")
}
dl.unlock(nil)
if dl.isLocked() {
t.Fatal("still locked after unlock")
}
if newDaemonLock(false).isLocked() {
t.Fatal("newDaemonLock(false) reports locked")
}
}
// closeStore must be safe on a daemon that never unlocked and safe twice —
// shutdown runs it unconditionally.
func TestDaemonLockCloseStoreIsSafeWhenNeverUnlocked(t *testing.T) {
dl := newDaemonLock(true)
if err := dl.closeStore(); err != nil {
t.Fatalf("closeStore with no store: %v", err)
}
if err := dl.closeStore(); err != nil {
t.Fatalf("second closeStore: %v", err)
}
}
// The data-loss bug: in locked mode the store is opened on an IPC goroutine
// inside UnlockFn, and shutdown runs on main. Without the handoff nothing
// calls Close, and Close is what re-encrypts the tmpfs working copy back over
// the ciphertext file — so every write of a cold-started session vanished.
func TestDaemonLockSealsTheStoreOpenedAfterUnlock(t *testing.T) {
dir := t.TempDir()
dbPath := filepath.Join(dir, "maven.db")
tmpfs := filepath.Join(dir, "work")
key := randBytes(t, 32)
// Store.Close zeroes the key slice it was handed (encState.key is the
// caller's backing array), so the next boot needs its own copy — exactly
// as mavend keeps envKeyBytes separate from the config's key.
nextBoot := bytes.Clone(key)
ctx := context.Background()
// Cold start: locked, no store.
dl := newDaemonLock(true)
// ... unlock arrives, opens the store and hands it over.
st, err := store.OpenEncrypted(ctx, dbPath, tmpfs, key)
if err != nil {
t.Fatalf("OpenEncrypted: %v", err)
}
dl.unlock(st)
if _, err := st.WriteNote(ctx, time.Now(), "заметка после холодного старта", nil, "test"); err != nil {
t.Fatalf("WriteNote: %v", err)
}
// Shutdown.
if err := dl.closeStore(); err != nil {
t.Fatalf("closeStore: %v", err)
}
if err := dl.closeStore(); err != nil {
t.Fatalf("second closeStore after a real store: %v", err)
}
// Next boot with the same key must see the write.
st2, err := store.OpenEncrypted(ctx, dbPath, tmpfs, nextBoot)
if err != nil {
t.Fatalf("reopen: %v", err)
}
defer st2.Close()
notes, err := st2.RecentNotes(ctx, 10)
if err != nil {
t.Fatalf("RecentNotes: %v", err)
}
if len(notes) != 1 {
t.Fatalf("got %d notes after a cold-started session, want 1 — the session was lost", len(notes))
}
}
// The whole point of the wrapped blob: what sits in the state dir must not let
// anyone open the database. Nothing written there may contain the key, and the
// ciphertext must not be readable with a wrong one.
func TestColdStartLeavesNoPlaintextKeyOnDisk(t *testing.T) {
dir := t.TempDir()
dbPath := filepath.Join(dir, "maven.db")
tmpfs := filepath.Join(dir, "work")
wrappedPath := filepath.Join(dir, "db_key.wrapped")
key := randBytes(t, 32)
secret := randBytes(t, 32)
ctx := context.Background()
blob, err := webauthn.WrapKey(key, secret)
if err != nil {
t.Fatalf("WrapKey: %v", err)
}
if err := os.WriteFile(wrappedPath, blob, 0o600); err != nil {
t.Fatalf("write wrapped key: %v", err)
}
st, err := store.OpenEncrypted(ctx, dbPath, tmpfs, key)
if err != nil {
t.Fatalf("OpenEncrypted: %v", err)
}
if _, err := st.WriteNote(ctx, time.Now(), "секрет", nil, "test"); err != nil {
t.Fatalf("WriteNote: %v", err)
}
if err := st.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
// Walk everything in the state dir; none of it may contain the key.
err = filepath.Walk(dir, func(p string, info os.FileInfo, err error) error {
if err != nil || info.IsDir() {
return err
}
b, rerr := os.ReadFile(p)
if rerr != nil {
return nil // unreadable is not a leak
}
if bytes.Contains(b, key) {
t.Errorf("%s contains the plaintext encryption key", p)
}
return nil
})
if err != nil {
t.Fatalf("walk: %v", err)
}
// The wrapped file must have owner-only permissions.
fi, err := os.Stat(wrappedPath)
if err != nil {
t.Fatalf("stat: %v", err)
}
if perm := fi.Mode().Perm(); perm != 0o600 {
t.Errorf("wrapped key file mode = %o, want 600", perm)
}
// A wrong passkey must not open the store.
if _, _, err := webauthn.UnwrapKey(blob, randBytes(t, 32)); err == nil {
t.Fatal("a wrong PRF secret unwrapped the key")
}
if _, err := store.OpenEncrypted(ctx, dbPath, filepath.Join(dir, "work2"), randBytes(t, 32)); err == nil {
t.Fatal("the encrypted store opened under a wrong key")
}
// And the right one round-trips back to a readable database.
got, version, err := webauthn.UnwrapKey(blob, secret)
if err != nil {
t.Fatalf("UnwrapKey: %v", err)
}
if version != webauthn.BlobV2 {
t.Errorf("blob version = %v, want v2", version)
}
st2, err := store.OpenEncrypted(ctx, dbPath, tmpfs, got)
if err != nil {
t.Fatalf("reopen with the unwrapped key: %v", err)
}
defer st2.Close()
notes, err := st2.RecentNotes(ctx, 10)
if err != nil {
t.Fatalf("RecentNotes: %v", err)
}
if len(notes) != 1 {
t.Fatalf("got %d notes, want 1", len(notes))
}
}
+235
View File
@@ -0,0 +1,235 @@
// mavend/intake.go — the unified event intake envelope, wired (Vikunja #283).
//
// internal/event defines the envelope and the bounded in-memory journal. This
// file is the one place that FILLS it, and the reason it is one place is worth
// stating, because the alternative was eight patches:
//
// Every intake path in Maven already converges on three writes, and all three
// are ipc.CoreAPI methods —
//
// WriteFact ← POST /api/ambient, mavcaldav, mavpoll's zenmoney + wg reads,
// /api/signal presence probes, the RSS/crawl watermarks
// WriteNote ← the RSS poller, the page crawler, meeting transcripts,
// image descriptions
// CaptureTask ← the voice path, the web form, and the mail reader
//
// — so decorating that ONE interface with a publish covers the lot without a
// caller knowing about events at all. cmd/mavmaild, cmd/mavcaldav, cmd/mavpoll,
// cmd/mavweb and the in-core feed/crawl/capture/vision workers are unchanged:
// they call the same interface they always called, and it now also narrates.
//
// The exception is cmd/mavend/mail.go, which reaches past the interface to
// st.CaptureTask directly. It publishes explicitly; see mailIntake.ingest.
//
// # Production behaviour when nobody is watching
//
// A nil *event.Bus makes Publish a no-op, and newIntakeAPI with a nil bus
// returns the wrapped API unchanged, so there is not even a decorator on the
// call path. The journal is memory-only and is never consulted by the tick
// loop, the router, or delivery — nothing Maven says depends on it. It is a
// read surface (`/events`, `recent_events`) and an observation seam for the
// simulator.
//
// # What is deliberately NOT here
//
// No dispatch. An event is a report that something arrived, never an
// instruction to speak: "a feed item appeared" becoming a notification is the
// nag this repo refuses. Digestion may one day read the journal; it will still
// go through internal/loop's rules and the severity/presence routing table.
package main
import (
"context"
"log"
"strings"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/store"
)
// newEventBus builds the journal, or returns nil when the operator turned it
// off (a negative config.intake_journal). nil is the "behave exactly as before"
// value all the way down: no decorator, no ring, no /events rows.
func newEventBus(cfg *config.Config) *event.Bus {
if cfg == nil || cfg.IntakeJournal < 0 {
log.Printf("intake journal: off (intake_journal < 0)")
return nil
}
n := cfg.IntakeJournal
if n == 0 {
n = config.DefaultIntakeJournal
}
log.Printf("intake journal: keeping the last %d intake events in memory", n)
return event.NewBus(n)
}
// intakeEventsFn is the daemonAPI.getEvents closure: the bus's ring rendered as
// the wire type. Returns nil for a nil bus, which the daemonAPI reports as an
// empty journal rather than an error.
func intakeEventsFn(bus *event.Bus) func(n int) []ipc.IntakeEvent {
if bus == nil {
return nil
}
return func(n int) []ipc.IntakeEvent {
evs := bus.Recent(n)
out := make([]ipc.IntakeEvent, 0, len(evs))
for _, e := range evs {
out = append(out, ipc.IntakeEvent{
Source: e.Source,
Kind: e.Kind,
EntityIDs: e.EntityIDs,
Title: e.Title,
Body: e.Body,
Priority: e.Priority,
OccurredAt: e.OccurredAt,
})
}
return out
}
}
// intakeAPI decorates a CoreAPI, publishing one envelope per successful
// intake write. Embedding the interface means every other method passes
// through untouched, and a new CoreAPI method is inherited rather than
// silently dropped.
type intakeAPI struct {
ipc.CoreAPI
bus *event.Bus
now func() time.Time
}
// newIntakeAPI wraps api so its intake writes are journalled. A nil bus
// returns api itself — no decorator, no allocation, no behaviour change.
func newIntakeAPI(api ipc.CoreAPI, bus *event.Bus, now func() time.Time) ipc.CoreAPI {
if bus == nil || api == nil {
return api
}
if now == nil {
now = time.Now
}
return &intakeAPI{CoreAPI: api, bus: bus, now: now}
}
// WriteFact journals the fact after it lands. Order matters: an event is a
// report of something that HAPPENED, so a failed write publishes nothing.
func (a *intakeAPI) WriteFact(ctx context.Context, req ipc.WriteFactReq) (int64, error) {
id, err := a.CoreAPI.WriteFact(ctx, req)
if err != nil {
return id, err
}
// OccurredAt is req.Ts, not now: mavpoll's wg read carries the handshake
// instant and the ambient path carries the meeting's start. Flattening
// those to notice-time would make the journal lie about when things
// happened, which is the one thing it is for.
a.bus.Publish(event.Event{
Source: req.Source,
Kind: event.SourceKind(req.Source, event.KindFact),
Title: req.Key,
Body: req.Value,
Priority: factPriority(req),
OccurredAt: req.Ts,
EntityIDs: entityIDs(req.Subject),
}, a.now())
return id, nil
}
// WriteNote journals a note. This is the RSS and crawler path, and also the
// meeting transcript and image description paths, which write their derived
// text as ordinary notes.
func (a *intakeAPI) WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error) {
id, err := a.CoreAPI.WriteNote(ctx, ts, text, embedding, source)
if err != nil {
return id, err
}
title, body := splitFirstLine(text)
a.bus.Publish(event.Event{
Source: source,
Kind: event.SourceKind(source, event.KindNote),
Title: title,
Body: body,
Priority: event.PriorityLow,
OccurredAt: ts,
}, a.now())
return id, nil
}
// CaptureTask journals a captured task, but only when a row was actually
// created. CaptureTask dedupes on normalised text among live rows, so a
// mailbox re-read after a restart must not refill the journal with tasks that
// were already there.
func (a *intakeAPI) CaptureTask(ctx context.Context, req ipc.CaptureTaskReq) (ipc.CaptureTaskResp, error) {
resp, err := a.CoreAPI.CaptureTask(ctx, req)
if err != nil || !resp.Created {
return resp, err
}
a.bus.Publish(publishableTask(store.Task{
CreatedTs: req.Ts,
Text: req.Text,
Source: req.Source,
Evidence: req.Evidence,
Status: req.Status,
Due: req.Due,
}, a.now()), a.now())
return resp, nil
}
// publishableTask is the task→envelope shape, shared with mail.go, which
// captures through the store directly rather than through the interface.
//
// Priority is high for a candidate with a due date and normal otherwise. That
// is the only place this file makes a judgement, and it is a display hint on a
// review page — nothing routes on it.
func publishableTask(t store.Task, now time.Time) event.Event {
occurred := t.CreatedTs
if occurred.IsZero() {
occurred = now
}
prio := event.PriorityNormal
if t.Due != nil {
prio = event.PriorityHigh
}
return event.Event{
Source: t.Source,
Kind: event.KindTask,
Title: t.Text,
Body: t.Evidence,
Priority: prio,
OccurredAt: occurred,
}
}
// factPriority is the attention hint for a fact write. Deliberately crude:
// a low-confidence inference (the ambient notification path writes below 1.0)
// is worth less attention than a read he or a credentialled poller made, and
// nothing else is distinguishable from here.
func factPriority(req ipc.WriteFactReq) string {
if req.Confidence > 0 && req.Confidence < 1.0 {
return event.PriorityLow
}
return event.PriorityNormal
}
// entityIDs turns a fact's free-text Subject into the EntityIDs slot when it
// already looks resolved. Intake runs BEFORE the fact enrichment worker
// resolves a subject against Nexus, so this is almost always empty — the slot
// exists for the paths that do know (the ecosystem acts), not for guessing.
func entityIDs(subject string) []string {
subject = strings.TrimSpace(subject)
if subject == "" || !strings.HasPrefix(subject, "entity:") {
return nil
}
return []string{strings.TrimPrefix(subject, "entity:")}
}
// splitFirstLine renders a note as title + body. Feed and crawl notes are
// written "headline\nsummary\nlink", so the first line is already the title.
func splitFirstLine(text string) (title, body string) {
text = strings.TrimSpace(text)
if i := strings.IndexByte(text, '\n'); i >= 0 {
return strings.TrimSpace(text[:i]), strings.TrimSpace(text[i+1:])
}
return text, ""
}
+178
View File
@@ -0,0 +1,178 @@
package main
import (
"context"
"errors"
"testing"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
)
var intakeNow = time.Date(2026, 8, 1, 10, 0, 0, 0, time.UTC)
func intakeClock() time.Time { return intakeNow }
// failingAPI wraps the store adapter, failing the three intake writes on
// demand, so the "a failed write publishes nothing" invariant is testable.
type failingAPI struct {
ipc.CoreAPI
fail bool
}
func (f *failingAPI) WriteFact(ctx context.Context, req ipc.WriteFactReq) (int64, error) {
if f.fail {
return 0, errors.New("injected")
}
return f.CoreAPI.WriteFact(ctx, req)
}
func newIntakeTestAPI(t *testing.T) (ipc.CoreAPI, *event.Bus) {
t.Helper()
st := newTestStore(t)
bus := event.NewBus(32)
return newIntakeAPI(ipc.NewStoreAPI(st), bus, intakeClock), bus
}
func TestIntakeAPIWithoutBusIsTheBareAPI(t *testing.T) {
// The adoption invariant: with the journal off there is not even a
// decorator on the intake path, so production behaves exactly as before.
st := newTestStore(t)
bare := ipc.NewStoreAPI(st)
if got := newIntakeAPI(bare, nil, intakeClock); got != ipc.CoreAPI(bare) {
t.Errorf("newIntakeAPI with a nil bus returned a wrapper, want the bare API")
}
}
func TestNewEventBusOffWhenNegative(t *testing.T) {
if b := newEventBus(&config.Config{IntakeJournal: -1}); b != nil {
t.Error("intake_journal = -1 still built a bus")
}
if b := newEventBus(&config.Config{IntakeJournal: 4}); b == nil {
t.Error("intake_journal = 4 built no bus")
}
}
func TestIntakeJournalsAFactWrite(t *testing.T) {
api, bus := newIntakeTestAPI(t)
ctx := context.Background()
// The ambient path's shape: an env fact below full confidence, timestamped
// at the meeting's start rather than at notice time.
start := intakeNow.Add(2 * time.Hour)
if _, err := api.WriteFact(ctx, ipc.WriteFactReq{
Ts: start, Kind: "env", Key: "calendar_event_20260801_планёрка",
Value: "10:00-11:00 планёрка", Source: "ambient:notif", Confidence: 0.6,
}); err != nil {
t.Fatalf("WriteFact: %v", err)
}
got := bus.Recent(0)
if len(got) != 1 {
t.Fatalf("journal has %d entries, want 1", len(got))
}
e := got[0]
if e.Source != "ambient:notif" || e.Kind != event.KindFact {
t.Errorf("source/kind = %q/%q", e.Source, e.Kind)
}
if e.Title != "calendar_event_20260801_планёрка" {
t.Errorf("title = %q, want the fact key", e.Title)
}
if !e.OccurredAt.Equal(start) {
t.Errorf("occurred_at = %v, want the fact's Ts %v — the journal must not flatten intake to notice time", e.OccurredAt, start)
}
if e.Priority != event.PriorityLow {
t.Errorf("priority = %q, want %q for a sub-1.0 confidence read", e.Priority, event.PriorityLow)
}
}
func TestIntakeDoesNotJournalAFailedWrite(t *testing.T) {
st := newTestStore(t)
bus := event.NewBus(8)
api := newIntakeAPI(&failingAPI{CoreAPI: ipc.NewStoreAPI(st), fail: true}, bus, intakeClock)
if _, err := api.WriteFact(context.Background(), ipc.WriteFactReq{
Ts: intakeNow, Kind: "env", Key: "k", Value: "v", Source: "poll:zenmoney", Confidence: 1,
}); err == nil {
t.Fatal("expected the injected error")
}
if bus.Len() != 0 {
t.Errorf("journal has %d entries after a failed write, want 0 — an event reports something that happened", bus.Len())
}
}
func TestIntakeJournalsANoteAsTitlePlusBody(t *testing.T) {
api, bus := newIntakeTestAPI(t)
// The RSS shape: "headline\nsummary\nlink".
if _, err := api.WriteNote(context.Background(), intakeNow,
"Вышло ядро 6.19\nкраткое содержание\nhttps://example.org/a", nil, "rss:tech"); err != nil {
t.Fatalf("WriteNote: %v", err)
}
got := bus.Recent(1)
if len(got) != 1 {
t.Fatalf("journal has %d entries, want 1", len(got))
}
if got[0].Title != "Вышло ядро 6.19" {
t.Errorf("title = %q, want the headline", got[0].Title)
}
if got[0].Kind != event.KindNote {
t.Errorf("kind = %q, want %q", got[0].Kind, event.KindNote)
}
if got[0].Body == "" {
t.Error("body is empty, want the rest of the note")
}
}
func TestIntakeJournalsOnlyCreatedTasks(t *testing.T) {
api, bus := newIntakeTestAPI(t)
ctx := context.Background()
req := ipc.CaptureTaskReq{Text: "оплатить интернет", Source: "email:inbox", Status: "candidate", Ts: intakeNow}
if _, err := api.CaptureTask(ctx, req); err != nil {
t.Fatalf("CaptureTask: %v", err)
}
// Same text again: CaptureTask dedupes among live rows, and a re-read of a
// mailbox must not refill the journal.
resp, err := api.CaptureTask(ctx, req)
if err != nil {
t.Fatalf("CaptureTask (repeat): %v", err)
}
if resp.Created {
t.Fatal("store did not dedupe; the test cannot check what it means to")
}
if bus.Len() != 1 {
t.Errorf("journal has %d entries, want 1 — a deduped capture must not publish", bus.Len())
}
if got := bus.Recent(1)[0]; got.Kind != event.KindTask || got.Title != "оплатить интернет" {
t.Errorf("entry = %+v, want the captured task", got)
}
}
func TestIntakeEventsFnRendersNewestFirst(t *testing.T) {
api, bus := newIntakeTestAPI(t)
ctx := context.Background()
for _, key := range []string{"a", "b", "c"} {
if _, err := api.WriteFact(ctx, ipc.WriteFactReq{
Ts: intakeNow, Kind: "env", Key: key, Value: "1", Source: "poll:zenmoney", Confidence: 1,
}); err != nil {
t.Fatalf("WriteFact %s: %v", key, err)
}
}
fn := intakeEventsFn(bus)
got := fn(2)
if len(got) != 2 || got[0].Title != "c" || got[1].Title != "b" {
t.Errorf("intakeEventsFn(2) = %+v, want the two newest, newest first", got)
}
if intakeEventsFn(nil) != nil {
t.Error("intakeEventsFn(nil) returned a closure, want nil so daemonAPI reports an empty journal")
}
}
func TestDaemonAPIRecentEventsEmptyWithoutABus(t *testing.T) {
d := &daemonAPI{CoreAPI: ipc.UnimplementedCoreAPI{}}
got, err := d.RecentEvents(context.Background(), 10)
if err != nil {
t.Fatalf("RecentEvents with no journal errored: %v", err)
}
if len(got) != 0 {
t.Errorf("got %d events, want none", len(got))
}
}
+14 -4
View File
@@ -26,6 +26,7 @@ import (
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/email"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/store"
@@ -42,6 +43,11 @@ type mailIntake struct {
ex *email.Extractor
timeout time.Duration
now func() time.Time
// bus — the unified intake journal (Vikunja #283). This path captures
// through the store directly rather than through ipc.CoreAPI, so the
// decorator in intake.go does not see it and the publish is explicit here.
// nil is a working no-op.
bus *event.Bus
}
// newMailIntake returns nil when mail ingestion must not be available, which is
@@ -52,7 +58,7 @@ type mailIntake struct {
// no keyword fallback: "the subject line became a task" is not extraction,
// it is a mailbox rendered as a to-do list, and it would fill the review
// page faster than he could clear it.
func newMailIntake(st *store.Store, phr phraser.Phraser, cfg *config.Config) *mailIntake {
func newMailIntake(st *store.Store, phr phraser.Phraser, cfg *config.Config, bus *event.Bus) *mailIntake {
if cfg.Email == nil {
return nil
}
@@ -67,7 +73,7 @@ func newMailIntake(st *store.Store, phr phraser.Phraser, cfg *config.Config) *ma
}
ex := email.NewExtractor(llmClientFor(lp, timeout), cfg.Email.MaxTasks, contextBlockFn(cfg, time.Now))
log.Printf("mail intake: enabled (max %d candidates per message, timeout %s)", cfg.Email.MaxTasks, timeout)
return &mailIntake{st: st, ex: ex, timeout: timeout, now: time.Now}
return &mailIntake{st: st, ex: ex, timeout: timeout, now: time.Now, bus: bus}
}
// ingest handles one ipc.MethodIngestMail call.
@@ -128,6 +134,10 @@ func (m *mailIntake) ingest(ctx context.Context, req ipc.IngestMailReq) (ipc.Ing
resp.TaskIDs = append(resp.TaskIDs, id)
if created {
resp.Created++
// Only a row that was actually created. CaptureTask dedupes on
// normalised text among live rows, so a mailbox re-read after a
// restart must not refill the journal with tasks already in it.
m.bus.Publish(publishableTask(t, now), now)
}
}
// Counts only: the log line names the mailbox and the UID, never the subject,
@@ -139,8 +149,8 @@ func (m *mailIntake) ingest(ctx context.Context, req ipc.IngestMailReq) (ipc.Ing
// wireMailIntake installs the IPC hook, or leaves it nil so the method reports
// ErrUnknownMethod. Called on both startup paths (unlocked boot and passkey
// unlock) so mail behaves the same either way.
func wireMailIntake(srv *ipc.Server, st *store.Store, phr phraser.Phraser, cfg *config.Config) {
mi := newMailIntake(st, phr, cfg)
func wireMailIntake(srv *ipc.Server, st *store.Store, phr phraser.Phraser, cfg *config.Config, bus *event.Bus) {
mi := newMailIntake(st, phr, cfg, bus)
if mi == nil {
return
}
+2 -2
View File
@@ -171,12 +171,12 @@ func TestIngestTruncatesEvidence(t *testing.T) {
// exist at all.
func TestNewMailIntakeOffWithoutConfig(t *testing.T) {
st := newTestStore(t)
if mi := newMailIntake(st, nil, &config.Config{}); mi != nil {
if mi := newMailIntake(st, nil, &config.Config{}, nil); mi != nil {
t.Error("no email block must mean no mail intake")
}
// Configured but with a non-LLM phraser: still off — there is no fallback
// extraction, by design.
if mi := newMailIntake(st, nil, &config.Config{Email: &config.EmailConfig{}}); mi != nil {
if mi := newMailIntake(st, nil, &config.Config{Email: &config.EmailConfig{}}, nil); mi != nil {
t.Error("without a llama-server phraser there is nothing to extract with")
}
}
+92 -24
View File
@@ -66,12 +66,19 @@ import (
var errLocked = errors.New("mavend: daemon locked — complete passkey assertion first")
// daemonLock tracks whether the daemon is in locked (pre-unlock) mode.
// In locked mode, all CoreAPI methods return errLocked. The unlock path
// replaces the CoreAPI with the real store adapter and flips the flag.
// daemonLock tracks whether the daemon is in locked (pre-unlock) mode, and
// owns the store handle the unlock path creates.
//
// The store matters here because of who runs when. In locked mode there is no
// store at boot; one is opened inside UnlockFn, on an IPC goroutine, minutes
// or days later. Shutdown runs on the main goroutine. Without a handoff the
// main goroutine has nothing to close, and store.Close is what re-encrypts
// the tmpfs working copy back over the ciphertext file — so a daemon that
// cold-started lost every write of that session, silently, on the next boot.
type daemonLock struct {
mu sync.Mutex
locked bool
st *store.Store
}
func newDaemonLock(locked bool) *daemonLock {
@@ -84,10 +91,25 @@ func (l *daemonLock) isLocked() bool {
return l.locked
}
func (l *daemonLock) unlock() {
// unlock flips the flag and takes ownership of the store opened by UnlockFn.
func (l *daemonLock) unlock(st *store.Store) {
l.mu.Lock()
defer l.mu.Unlock()
l.locked = false
l.st = st
}
// closeStore seals the store the unlock path opened, if any. Safe to call
// when the daemon never unlocked, and safe to call twice.
func (l *daemonLock) closeStore() error {
l.mu.Lock()
st := l.st
l.st = nil
l.mu.Unlock()
if st == nil {
return nil
}
return st.Close()
}
func main() {
@@ -155,6 +177,14 @@ func run(args []string) error {
return fmt.Errorf("open store: %w", err)
}
defer st.Close()
} else {
// Locked boot: the store does not exist yet. Seal whatever UnlockFn
// opened, at shutdown, on this goroutine.
defer func() {
if err := dl.closeStore(); err != nil {
log.Printf("mavend: seal store on shutdown: %v", err)
}
}()
}
// ----- daemon components (only wired when unlocked) -----
@@ -174,6 +204,16 @@ func run(args []string) error {
crawlWkr *crawlWorker // nil ⇒ no page is watched (the default)
)
// The unified intake journal (Vikunja #283). Built before anything else
// that holds a CoreAPI, because intakeAPI wraps that one interface and
// every intake path in the daemon reaches its sink through it. nil (the
// operator set intake_journal negative) means no decorator at all.
evBus := newEventBus(cfg)
// coreFor is what every in-process holder of a CoreAPI now takes, instead
// of a bare ipc.NewStoreAPI(st). Identical behaviour plus one published
// envelope per successful intake write.
coreFor := func() ipc.CoreAPI { return newIntakeAPI(ipc.NewStoreAPI(st), evBus, time.Now) }
if !locked {
rules = loop.DefaultRules()
gatherer = loop.NewGatherer(st, rules)
@@ -217,7 +257,7 @@ func run(args []string) error {
eco = wireEcosystem(cfg)
// voice
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st, eco)
voiceW, err = wireVoice(cfg, coreFor(), phr, st.VectorMemory(), st, eco)
if err != nil {
return fmt.Errorf("wire voice: %w", err)
}
@@ -267,14 +307,15 @@ func run(args []string) error {
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines), cfg.PatternProposals)
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
evalWorker = newMemoryEvalWorker(st, phr, cfg)
feedWkr = newFeedWorker(ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
crawlWkr = newCrawlWorker(newCrawler(cfg), ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
feedWkr = newFeedWorker(coreFor(), embedderOf(voiceW), cfg)
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
coreAPI = &daemonAPI{
CoreAPI: ipc.NewStoreAPI(st),
CoreAPI: coreFor(),
getTrace: tl.trace,
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
getDayPlan: func(ctx context.Context) ipc.DayPlan { return tl.dayPlan(ctx, time.Now()) },
getEvents: intakeEventsFn(evBus),
}
if voiceW != nil && voiceW.handler != nil {
api := coreAPI.(*daemonAPI)
@@ -331,7 +372,7 @@ func run(args []string) error {
// configured and there is a llama-server to extract with, in which case
// ipc.MethodIngestMail reports ErrUnknownMethod.
if !locked {
wireMailIntake(srv, st, phr, cfg)
wireMailIntake(srv, st, phr, cfg, evBus)
wireModelSwap(srv, phr, cfg)
// Vision + the media blob store (Vikunja #252). Both stay dark without a
// media block; MethodDescribeImage answers ErrUnknownMethod then.
@@ -346,12 +387,16 @@ func run(args []string) error {
wireSpeaker(srv, st, cfg)
}
// WrapKeyFn — wraps the env key with a passkey credential public key and
// persists the wrapped blob. Only wired when the daemon has the key in
// memory (env key mode). Called by mavweb after passkey enrollment.
// WrapKeyFn — wraps the env key under the passkey PRF secret and persists
// the wrapped blob. Only wired when the daemon has the key in memory (env
// key mode). Called by mavweb after passkey enrollment.
//
// webauthn.WrapKey refuses anything that is not a 32-byte PRF output, so
// an authenticator without PRF support produces no wrapped file at all
// rather than a file that looks protected and is not.
if envKeyBytes != nil {
srv.WrapKeyFn = func(ctx context.Context, publicKey []byte) error {
blob, err := webauthn.WrapKey(envKeyBytes, publicKey)
srv.WrapKeyFn = func(ctx context.Context, secret []byte) error {
blob, err := webauthn.WrapKey(envKeyBytes, secret)
if err != nil {
return fmt.Errorf("wrap encryption key: %w", err)
}
@@ -367,20 +412,42 @@ func run(args []string) error {
}
}
// UnlockFn — cold-start unlock: unwraps the encryption key from the wrapped
// blob using the passkey credential public key, opens the store, wires all
// UnlockFn — cold-start unlock: unwraps the encryption key from the
// wrapped blob using the passkey PRF secret, opens the store, wires all
// daemon components, and replaces the locked API.
if locked {
srv.UnlockFn = func(ctx context.Context, publicKey []byte) error {
var unlockMu sync.Mutex
srv.UnlockFn = func(ctx context.Context, secret []byte) error {
// One unlock at a time, and never a second one. Without this a
// concurrent pair of Unlock calls would each open a store and
// wire a full daemon, and the loser's goroutines would run
// against a store nobody closes.
unlockMu.Lock()
defer unlockMu.Unlock()
if !dl.isLocked() {
return nil // already unlocked; the caller does not need to know
}
// The wire cannot authenticate its caller — the socket is
// same-uid — so the unlock path requires a passkey assertion
// that mavweb verified cryptographically first. Without this,
// MethodUnlock is reachable by anything on the box.
if !passkeySess.IsStepUp() {
return errors.New("unlock: no verified passkey assertion (assert first)")
}
wp := *wrappedKeyPath
blob, err := os.ReadFile(wp)
if err != nil {
return fmt.Errorf("read wrapped key: %w", err)
}
key, err := webauthn.UnwrapKey(blob, publicKey)
key, version, err := webauthn.UnwrapKey(blob, secret)
if err != nil {
return fmt.Errorf("unwrap key: %w", err)
}
if version == webauthn.BlobV1 {
log.Printf("SECURITY: %s was unwrapped from a %s blob. The wrapping key is derived from the credential PUBLIC key, which mavweb also writes to its passkeys.json — anyone holding both files can recover the database key with no authenticator. Re-enroll the passkey on an authenticator that supports the PRF extension to rewrite it as v2.", wp, version)
}
// Open the store with the unwrapped key.
st, err = store.OpenEncrypted(ctx, cfg.DBPath, cfg.DBTmpfs, key)
if err != nil {
@@ -426,7 +493,7 @@ func run(args []string) error {
eco = wireEcosystem(cfg)
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st, eco)
voiceW, err = wireVoice(cfg, coreFor(), phr, st.VectorMemory(), st, eco)
if err != nil {
return fmt.Errorf("wire voice: %w", err)
}
@@ -470,22 +537,23 @@ func run(args []string) error {
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines), cfg.PatternProposals)
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
evalWorker = newMemoryEvalWorker(st, phr, cfg)
feedWkr = newFeedWorker(ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
crawlWkr = newCrawlWorker(newCrawler(cfg), ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
feedWkr = newFeedWorker(coreFor(), embedderOf(voiceW), cfg)
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
// Swap the CoreAPI from the locked placeholder to the real store adapter.
newAPI := &daemonAPI{
CoreAPI: ipc.NewStoreAPI(st),
CoreAPI: coreFor(),
getTrace: tl.trace,
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
getDayPlan: func(ctx context.Context) ipc.DayPlan { return tl.dayPlan(ctx, time.Now()) },
getEvents: intakeEventsFn(evBus),
}
if voiceW != nil && voiceW.handler != nil {
newAPI.chatFn = voiceW.handler.handleText
}
srv.SetAPI(newAPI)
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
wireMailIntake(srv, st, phr, cfg)
wireMailIntake(srv, st, phr, cfg, evBus)
wireModelSwap(srv, phr, cfg)
keeper := wireVision(ctx, srv, st, embedderOf(voiceW), cfg)
wireCapture(srv, keeper, st, voiceW, phr, cfg)
@@ -543,7 +611,7 @@ func run(args []string) error {
go voiceW.mcp.run(ctx)
}
dl.unlock()
dl.unlock(st)
log.Printf("mavend: unlocked via passkey assertion")
return nil
}
+12
View File
@@ -942,6 +942,18 @@ type daemonAPI struct {
getDayPlan func(ctx context.Context) ipc.DayPlan
chatFn func(ctx context.Context, text string) string
getMCPServers func() []ipc.MCPServerStatus
getEvents func(n int) []ipc.IntakeEvent
}
// RecentEvents — the unified intake journal (Vikunja #283). Empty, not an
// error, when no bus was wired: "nothing has arrived" and "the journal is off"
// look the same to a reader on purpose, because neither is a fault and the
// page renders both as an empty table.
func (d *daemonAPI) RecentEvents(ctx context.Context, n int) ([]ipc.IntakeEvent, error) {
if d.getEvents == nil {
return nil, nil
}
return d.getEvents(n), nil
}
func (d *daemonAPI) Chat(ctx context.Context, text string) (string, error) {
+323
View File
@@ -0,0 +1,323 @@
package main
// Golden-audio STT tests (Vikunja #288).
//
// These push real audio through the real whisper.cpp binding, so a bad model
// path, a wrong language hint, a broken resample or a regressed silence gate
// is caught by `make test` rather than by the owner talking to a daemon that
// mishears him.
//
// The fixtures are piper-synthesised, not recorded — see
// scripts/gen-stt-fixtures.sh. Nothing of the owner's voice is committed, and
// any fixture can be rebuilt from the script plus a voice model.
//
// Matching is deliberately tolerant. Golden transcripts are model-dependent:
// swapping ggml-small for a different whisper build moves punctuation, casing
// and the odd word ending, and an exact-string assertion would turn every
// model swap into a fixture rewrite. Each case therefore asserts two things —
// the words that carry the intent are present, and the word error rate
// against the reference stays under a per-case ceiling.
import (
"context"
"encoding/json"
"os"
"path/filepath"
"strings"
"testing"
"unicode"
"github.com/kami/maven/internal/audio"
"github.com/kami/maven/internal/worker"
)
// goldenModelPath — the whisper model the golden tests run against. Same file
// the Makefile's run-stt target uses. Overridable so a box that keeps its
// models elsewhere can still run these.
func goldenModelPath() string {
if p := os.Getenv("MAVEN_WHISPER_MODEL"); p != "" {
return p
}
return filepath.Join("..", "..", "models", "stt", "ggml-small.bin")
}
type goldenCase struct {
Name string `json:"name"`
WAV string `json:"wav"`
Lang string `json:"lang"`
Text string `json:"text"`
Keywords []string `json:"keywords"`
MaxWER float64 `json:"max_wer"`
}
type goldenManifest struct {
Cases []goldenCase `json:"cases"`
}
func loadGoldenManifest(t *testing.T) goldenManifest {
t.Helper()
raw, err := os.ReadFile(filepath.Join("testdata", "golden_v1.json"))
if err != nil {
t.Fatalf("read golden manifest: %v", err)
}
var m goldenManifest
if err := json.Unmarshal(raw, &m); err != nil {
t.Fatalf("parse golden manifest: %v", err)
}
if len(m.Cases) == 0 {
t.Fatal("golden manifest has no cases")
}
return m
}
// normalizeTranscript lowercases, drops punctuation, folds the Russian ё onto
// е (whisper is inconsistent about it and the router does not care), and
// collapses whitespace. Everything the comparison does happens on this form.
func normalizeTranscript(s string) []string {
var b strings.Builder
for _, r := range strings.ToLower(s) {
switch {
case r == 'ё':
b.WriteRune('е')
case unicode.IsLetter(r) || unicode.IsDigit(r):
b.WriteRune(r)
default:
b.WriteRune(' ')
}
}
return strings.Fields(b.String())
}
// wordErrorRate is the Levenshtein distance between two word sequences,
// divided by the length of the reference. 0 means identical; it can exceed 1
// when the hypothesis is much longer than the reference.
func wordErrorRate(ref, hyp []string) float64 {
if len(ref) == 0 {
if len(hyp) == 0 {
return 0
}
return 1
}
prev := make([]int, len(hyp)+1)
cur := make([]int, len(hyp)+1)
for j := range prev {
prev[j] = j
}
for i := 1; i <= len(ref); i++ {
cur[0] = i
for j := 1; j <= len(hyp); j++ {
cost := 1
if ref[i-1] == hyp[j-1] {
cost = 0
}
cur[j] = min(prev[j]+1, min(cur[j-1]+1, prev[j-1]+cost))
}
prev, cur = cur, prev
}
return float64(prev[len(hyp)]) / float64(len(ref))
}
// missingKeywords returns the keywords absent from the hypothesis. A keyword
// matches on prefix, so a different case ending ("воды" vs "воду") does not
// fail the assertion — the router's stage-0 grammar is stem-shaped too.
func missingKeywords(keywords []string, hyp []string) []string {
var missing []string
for _, kw := range keywords {
want := normalizeTranscript(kw)
if len(want) == 0 {
continue
}
if !containsSeq(hyp, want) {
missing = append(missing, kw)
}
}
return missing
}
func containsSeq(hyp, want []string) bool {
for i := 0; i+len(want) <= len(hyp); i++ {
ok := true
for j, w := range want {
// Prefix match, so inflection differences pass but
// distinct words do not.
if !looseWordMatch(hyp[i+j], w) {
ok = false
break
}
}
if ok {
return true
}
}
return false
}
func looseWordMatch(got, want string) bool {
if got == want {
return true
}
g, w := []rune(got), []rune(want)
n := len(w) - 1
if len(w) > 6 {
n = len(w) - 2
}
// Words of three runes or fewer have no room for a safe prefix: require
// an exact match rather than letting "час" pass for "часть".
if n < 3 || len(g) < n {
return false
}
return string(g[:n]) == string(w[:n])
}
// --- the model-backed test -------------------------------------------------
func TestGoldenAudioTranscription(t *testing.T) {
m := loadGoldenManifest(t)
model := goldenModelPath()
if _, err := os.Stat(model); err != nil {
t.Skipf("whisper model %s absent (%v) — set MAVEN_WHISPER_MODEL or see AGENTS.md", model, err)
}
// Same gate thresholds as mavsttd's defaults, so a regression in the
// silence gate shows up here as an empty transcript.
h, err := newWhisperHandler(model, 300, 0.01)
if err != nil {
t.Fatalf("load whisper model %s: %v", model, err)
}
defer h.Close()
for _, c := range m.Cases {
t.Run(c.Name, func(t *testing.T) {
path := filepath.Join("testdata", c.WAV)
raw, err := os.ReadFile(path)
if err != nil {
t.Skipf("fixture %s absent (%v) — run scripts/gen-stt-fixtures.sh", path, err)
}
format, pcm, err := audio.PCMFromWAV(raw)
if err != nil {
t.Fatalf("%s is not canonical 16k mono PCM: %v", path, err)
}
resp, err := h.Transcribe(context.Background(), worker.TranscribeReq{
Audio: audio.Audio{Format: format, Bytes: pcm},
Lang: c.Lang,
})
if err != nil {
t.Fatalf("transcribe %s: %v", c.WAV, err)
}
t.Logf("%s → %q (confidence %.3f)", c.WAV, resp.Text, resp.Confidence)
if strings.TrimSpace(resp.Text) == "" {
t.Fatalf("%s transcribed to empty text — the silence gate ate real speech", c.WAV)
}
if resp.Confidence <= 0 {
t.Errorf("%s: confidence %v, want > 0", c.WAV, resp.Confidence)
}
hyp := normalizeTranscript(resp.Text)
ref := normalizeTranscript(c.Text)
if missing := missingKeywords(c.Keywords, hyp); len(missing) > 0 {
t.Errorf("%s: missing keywords %v in %q", c.WAV, missing, resp.Text)
}
if wer := wordErrorRate(ref, hyp); wer > c.MaxWER {
t.Errorf("%s: WER %.2f > %.2f\n want: %q\n got: %q", c.WAV, wer, c.MaxWER, c.Text, resp.Text)
}
})
}
}
// TestGoldenFixturesAreCanonical checks the committed audio without needing a
// model, so a fixture regenerated at the wrong sample rate fails on every box.
func TestGoldenFixturesAreCanonical(t *testing.T) {
m := loadGoldenManifest(t)
for _, c := range m.Cases {
path := filepath.Join("testdata", c.WAV)
raw, err := os.ReadFile(path)
if err != nil {
t.Errorf("fixture %s missing: %v", path, err)
continue
}
format, pcm, err := audio.PCMFromWAV(raw)
if err != nil {
t.Errorf("%s: %v", path, err)
continue
}
if !format.IsValid() {
t.Errorf("%s: format %+v is not canonical", path, format)
}
a := audio.Audio{Format: format, Bytes: pcm}
if d := a.Duration(); d < 0.5 || d > 10 {
t.Errorf("%s: duration %.2fs outside the sane 0.510s fixture range", path, d)
}
// The fixture must clear mavsttd's own silence gate, otherwise the
// model test below would be asserting on a gated empty string.
if reason := gateReason(pcmToF32(pcm), whisperSampleRate, 300, 0.01); reason != "" {
t.Errorf("%s: would be gated as %s", path, reason)
}
if len(c.Keywords) == 0 {
t.Errorf("%s: manifest case has no keywords", c.Name)
}
if c.MaxWER <= 0 || c.MaxWER > 1 {
t.Errorf("%s: max_wer %v outside (0,1]", c.Name, c.MaxWER)
}
}
}
func pcmToF32(b []byte) []float32 {
out := make([]float32, len(b)/2)
for i := range out {
s := int16(b[i*2]) | int16(b[i*2+1])<<8
out[i] = float32(s) / 32768.0
}
return out
}
// --- matcher unit tests (no model, no fixtures) ----------------------------
func TestNormalizeTranscript(t *testing.T) {
got := normalizeTranscript(" Ещё, Раз... ")
want := []string{"еще", "раз"}
if len(got) != len(want) || got[0] != want[0] || got[1] != want[1] {
t.Fatalf("normalizeTranscript = %v, want %v", got, want)
}
}
func TestWordErrorRate(t *testing.T) {
cases := []struct {
name string
ref, hyp string
want float64
}{
{"identical", "напомни мне через час", "Напомни мне через час.", 0},
{"one substitution", "напомни мне через час", "напомни мне через день", 0.25},
{"one deletion", "напомни мне через час", "напомни мне час", 0.25},
{"empty hypothesis", "напомни мне", "", 1},
{"both empty", "", "", 0},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
got := wordErrorRate(normalizeTranscript(c.ref), normalizeTranscript(c.hyp))
if got != c.want {
t.Fatalf("WER = %v, want %v", got, c.want)
}
})
}
}
func TestMissingKeywords(t *testing.T) {
hyp := normalizeTranscript("Отметь, что я выпил воду.")
if got := missingKeywords([]string{"воды", "отметь"}, hyp); len(got) != 0 {
t.Fatalf("missingKeywords = %v, want none (inflection must not fail the match)", got)
}
if got := missingKeywords([]string{"календарю"}, hyp); len(got) != 1 {
t.Fatalf("missingKeywords = %v, want the absent keyword reported", got)
}
// A short word must match exactly — no 4-rune prefix shortcut that would
// let "час" pass for "часть".
hyp2 := normalizeTranscript("через час")
if got := missingKeywords([]string{"часть"}, hyp2); len(got) != 1 {
t.Fatalf("missingKeywords = %v, want %q reported missing", got, "часть")
}
}
BIN
View File
Binary file not shown.
+37
View File
@@ -0,0 +1,37 @@
{
"note": "Golden STT fixtures. Audio is piper-synthesised, not recorded — see scripts/gen-stt-fixtures.sh. Regenerate with that script; do not hand-edit `wav`.",
"cases": [
{
"name": "ru_reminder",
"wav": "ru_reminder.wav",
"lang": "ru",
"text": "напомни мне через час позвонить маме",
"keywords": ["напомни", "час", "позвонить"],
"max_wer": 0.34
},
{
"name": "ru_fact",
"wav": "ru_fact.wav",
"lang": "ru",
"text": "отметь что я выпил воды",
"keywords": ["отметь", "воды"],
"max_wer": 0.34
},
{
"name": "ru_query",
"wav": "ru_query.wav",
"lang": "ru",
"text": "что у меня сегодня по календарю",
"keywords": ["сегодня", "календарю"],
"max_wer": 0.34
},
{
"name": "en_act",
"wav": "en_act.wav",
"lang": "en",
"text": "restart the web server and check the disk space",
"keywords": ["restart", "server", "disk"],
"max_wer": 0.34
}
]
}
Binary file not shown.
Binary file not shown.
Binary file not shown.
+33 -92
View File
@@ -12,8 +12,15 @@
// (30ms frames, 16kHz PCM) matches silero-vad's input interface exactly, so
// swapping energy-threshold for ONNX-inference is a local change in vad.go.
//
// While a reply is playing the capture side is muted (half-duplex): without
// it, Maven's own voice comes back in through the mic and she answers
// herself. -barge-in punches one hole in that gate — sustained energy well
// above the speaker's leak level cuts playback so he can talk over her. It is
// off by default because the threshold is room-specific; see playback.go.
//
// usage:
// mavwaked # default ALSA device, 127.0.0.1:9100
// mavwaked -barge-in # let him interrupt her mid-reply
// mavwaked -device hw:1,0 -addr 10.42.0.1:9100
// mavwaked -test file.wav # read from file, no arecord
package main
@@ -60,6 +67,9 @@ func run(args []string) error {
silenceMs := flag.Int("silence-ms", defaultSilenceMs, "silence ms to end utterance")
maxMs := flag.Int("max-ms", defaultMaxMs, "max utterance ms")
testFile := flag.String("test", "", "read PCM from file instead of arecord (testing only)")
bargeIn := flag.Bool("barge-in", false, "cut Maven off when he talks over her (needs a room-tuned -barge-in-rms)")
bargeRMS := flag.Int("barge-in-rms", defaultBargeRMS, "RMS x10000 a frame must clear to count as barge-in")
bargeFrames := flag.Int("barge-in-frames", defaultBargeFrames, "consecutive frames over -barge-in-rms before playback is cut")
flag.CommandLine.Parse(args)
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM, syscall.SIGHUP)
@@ -117,14 +127,21 @@ func run(args []string) error {
defer src.Close()
return captureLoop(ctx, src, vad, vc, *lang)
var barge bargeInConfig
if *bargeIn {
barge = bargeInConfig{RMS: float64(*bargeRMS) / 10000.0, Frames: *bargeFrames}
log.Printf("mavwaked: barge-in on (rms %.4f x %d frames)", barge.RMS, barge.Frames)
}
sess := newSession(vad, newAplayPlayer(), &voiceSender{vc: vc}, *lang, barge)
return captureLoop(ctx, src, sess)
}
// captureLoop reads PCM from src, runs VAD, and sends complete utterances to
// the voice server. Returns when ctx is done or src is exhausted.
func captureLoop(ctx context.Context, src io.Reader, vad *VAD, vc *voice.Client, lang string) error {
// captureLoop reads PCM from src and hands whole frames to the session.
// Returns when ctx is done or src is exhausted.
func captureLoop(ctx context.Context, src io.Reader, sess *session) error {
br := bufio.NewReaderSize(src, defaultReadSize)
frameBytes := vad.FrameSamples() * 2 // 480 samples × 2 bytes = 960 bytes per 30ms
frameBytes := sess.vad.FrameSamples() * 2 // 480 samples × 2 bytes = 960 bytes per 30ms
log.Printf("mavwaked: capture loop starting (frame=%d bytes, %dms)",
frameBytes, defaultFrameMs)
@@ -147,7 +164,7 @@ func captureLoop(ctx context.Context, src io.Reader, vad *VAD, vc *voice.Client,
// Flush partial frame.
partial = append(partial, buf[:n]...)
if len(partial) >= frameBytes {
if err := processFrame(partial[:frameBytes], vad, vc, lang); err != nil {
if err := sess.feed(ctx, partial[:frameBytes]); err != nil {
log.Printf("mavwaked: process frame: %v", err)
}
partial = partial[frameBytes:]
@@ -165,107 +182,31 @@ func captureLoop(ctx context.Context, src io.Reader, vad *VAD, vc *voice.Client,
partial = nil
}
if err := processFrame(full, vad, vc, lang); err != nil {
if err := sess.feed(ctx, full); err != nil {
log.Printf("mavwaked: process frame: %v", err)
}
}
}
// processFrame feeds one 30ms PCM frame to the VAD and sends any completed
// utterance to the voice server.
func processFrame(frame []byte, vad *VAD, vc *voice.Client, lang string) error {
samples := PCMToI16(frame)
utt, state := vad.Feed(samples)
// voiceSender is the production utteranceSender: one PushToTalk round-trip
// over the voice wire. SurfaceVoice (not the default SurfacePCClient that
// c.PushToTalk uses) caps everything at L0, which is what makes an accidental
// VAD trigger safe.
type voiceSender struct{ vc *voice.Client }
if state == StateSpeech {
// Speech is in progress; nothing to send yet.
return nil
}
if utt.Bytes == nil {
// Still in silence, or short speech that didn't trigger.
return nil
}
// We have a complete utterance — send it to the voice server.
return sendUtterance(context.Background(), utt, vc, lang)
}
// sendUtterance sends audio to the voice server and plays the reply.
func sendUtterance(ctx context.Context, utt audio.Audio, vc *voice.Client, lang string) error {
dur := utt.Duration()
log.Printf("mavwaked: utterance complete (%.2fs, %d bytes), sending...",
dur, len(utt.Bytes))
// Use SendRequest directly so we can set SurfaceVoice instead of the
// default SurfacePCClient that c.PushToTalk uses.
func (s *voiceSender) Send(ctx context.Context, utt audio.Audio, lang string) (audio.Audio, error) {
var resp voice.PushToTalkResp
err := vc.SendRequest(ctx, voice.MethodPushToTalk, voice.PushToTalkReq{
err := s.vc.SendRequest(ctx, voice.MethodPushToTalk, voice.PushToTalkReq{
Audio: utt,
Lang: lang,
Surface: voice.SurfaceVoice,
}, &resp)
if err != nil {
return fmt.Errorf("push-to-talk: %w", err)
return audio.Audio{}, fmt.Errorf("push-to-talk: %w", err)
}
log.Printf("mavwaked: reply: %q (%.2fs audio)", resp.ReplyText, resp.ReplyAudio.Duration())
// Play the reply audio.
if len(resp.ReplyAudio.Bytes) > 0 {
go playAudio(resp.ReplyAudio)
} else {
log.Printf("mavwaked: empty reply audio (text only)")
}
if len(resp.RoutedChannels) > 0 {
log.Printf("mavwaked: also routed to: %v", resp.RoutedChannels)
}
return nil
}
// playAudio pipes PCM audio to aplay(1) for playback. Runs in a goroutine.
func playAudio(a audio.Audio) {
// Build WAV header for aplay (or pipe raw PCM with the right format flags).
cmd := exec.Command("aplay",
"-f", "S16_LE",
"-r", fmt.Sprintf("%d", a.Format.SampleRate),
"-c", fmt.Sprintf("%d", a.Format.Channels),
"-t", "raw",
)
stdin, err := cmd.StdinPipe()
if err != nil {
log.Printf("mavwaked: aplay stdin pipe: %v", err)
return
}
if err := cmd.Start(); err != nil {
log.Printf("mavwaked: start aplay: %v", err)
return
}
// Write audio to aplay's stdin.
if _, err := stdin.Write(a.Bytes); err != nil {
log.Printf("mavwaked: write to aplay: %v", err)
}
_ = stdin.Close()
// Wait for playback to finish (with a timeout).
done := make(chan error, 1)
go func() {
done <- cmd.Wait()
}()
select {
case err := <-done:
if err != nil {
log.Printf("mavwaked: aplay: %v", err)
}
case <-time.After(30 * time.Second):
log.Printf("mavwaked: aplay timeout, killing")
_ = cmd.Process.Kill()
<-done
}
return resp.ReplyAudio, nil
}
+139
View File
@@ -0,0 +1,139 @@
package main
// Reply playback, and the half-duplex gate around it (Vikunja #287).
//
// Before this, playback was `go playAudio(reply)` — fire and forget, with no
// handle on the running aplay. Two things fell out of that, and both are
// audible:
//
// 1. Self-trigger. The capture loop keeps feeding the VAD while the speaker
// is playing, so Maven's own reply comes back in through the mic, trips
// the VAD, and is sent to the daemon as a fresh utterance. She answers
// herself. There is no acoustic echo canceller in this pipeline, so the
// only correct fix is half-duplex: while she is speaking, the capture
// side is muted.
//
// 2. No barge-in. Talking over her did nothing — there was nothing to
// cancel, because nobody held the process handle.
//
// The two are the same mechanism seen from opposite sides, so they live
// together here. Echo suppression is unconditional (it fixes a bug). Barge-in
// is off unless -barge-in is passed, because it needs a room-specific energy
// threshold: with no echo canceller, the only way to tell "he is talking over
// her" from "the mic is hearing her" is that he is louder, and how much
// louder depends on where the mic sits relative to the speaker.
import (
"log"
"os/exec"
"strconv"
"sync"
"time"
"github.com/kami/maven/internal/audio"
)
// player plays one reply at a time and can be cut off mid-utterance.
type player interface {
// Play starts playback of a, replacing anything already playing, and
// returns immediately.
Play(a audio.Audio)
// Stop ends playback now. A no-op when nothing is playing.
Stop()
// Playing reports whether audio is currently going out of the speaker.
Playing() bool
}
// aplayPlayer pipes raw PCM to aplay(1). Stop kills the child, which is what
// makes barge-in instant rather than "instant at the end of the sentence".
type aplayPlayer struct {
mu sync.Mutex
cmd *exec.Cmd
playing bool
// gen rises on every Play/Stop so a finishing playback cannot clear the
// playing flag of the one that replaced it.
gen uint64
}
func newAplayPlayer() *aplayPlayer { return &aplayPlayer{} }
func (p *aplayPlayer) Play(a audio.Audio) {
if len(a.Bytes) == 0 {
return
}
p.Stop()
cmd := exec.Command("aplay",
"-f", "S16_LE",
"-r", strconv.Itoa(a.Format.SampleRate),
"-c", strconv.Itoa(a.Format.Channels),
"-t", "raw",
)
stdin, err := cmd.StdinPipe()
if err != nil {
log.Printf("mavwaked: aplay stdin pipe: %v", err)
return
}
if err := cmd.Start(); err != nil {
log.Printf("mavwaked: start aplay: %v", err)
_ = stdin.Close()
return
}
p.mu.Lock()
p.gen++
gen := p.gen
p.cmd = cmd
p.playing = true
p.mu.Unlock()
go func() {
if _, err := stdin.Write(a.Bytes); err != nil {
// Broken pipe is the expected outcome of Stop().
log.Printf("mavwaked: write to aplay: %v", err)
}
_ = stdin.Close()
done := make(chan error, 1)
go func() { done <- cmd.Wait() }()
select {
case err := <-done:
if err != nil {
log.Printf("mavwaked: aplay: %v", err)
}
case <-time.After(30 * time.Second):
log.Printf("mavwaked: aplay timeout, killing")
if pr := cmd.Process; pr != nil {
_ = pr.Kill()
}
<-done
}
p.mu.Lock()
if p.gen == gen {
p.playing = false
p.cmd = nil
}
p.mu.Unlock()
}()
}
func (p *aplayPlayer) Stop() {
p.mu.Lock()
cmd := p.cmd
if cmd != nil {
p.gen++
p.playing = false
p.cmd = nil
}
p.mu.Unlock()
if cmd != nil && cmd.Process != nil {
_ = cmd.Process.Kill()
}
}
func (p *aplayPlayer) Playing() bool {
p.mu.Lock()
defer p.mu.Unlock()
return p.playing
}
+33
View File
@@ -0,0 +1,33 @@
package main
import (
"testing"
"github.com/kami/maven/internal/audio"
)
// The real player must be safe to poke when nothing is playing — the capture
// loop calls Playing() on every 30ms frame, and Stop() lands on an idle
// player whenever a barge-in races the end of a reply. Neither may need
// aplay(1) to be installed.
func TestAplayPlayerIdleIsSafe(t *testing.T) {
p := newAplayPlayer()
if p.Playing() {
t.Fatal("a fresh player reports playing")
}
p.Stop()
p.Stop()
if p.Playing() {
t.Fatal("playing after Stop on an idle player")
}
// Empty audio is a text-only turn: nothing to play, no process to spawn.
p.Play(audio.Audio{Format: audio.PCM16kMono})
if p.Playing() {
t.Fatal("empty audio started playback")
}
}
func TestAplayPlayerSatisfiesPlayer(t *testing.T) {
var _ player = newAplayPlayer()
var _ player = &fakePlayer{}
}
+123
View File
@@ -0,0 +1,123 @@
package main
// The capture session: what happens to one 30ms frame, given whether Maven is
// currently speaking. Split out of main.go's processFrame so the decision is
// testable without a mic, a speaker, or a daemon (Vikunja #287).
import (
"context"
"log"
"github.com/kami/maven/internal/audio"
)
// utteranceSender ships one complete utterance to the voice server and
// returns the reply audio to play. The real one round-trips over the voice
// wire; tests substitute a recorder.
type utteranceSender interface {
Send(ctx context.Context, utt audio.Audio, lang string) (audio.Audio, error)
}
// bargeInConfig holds the two numbers barge-in needs. Zero Frames disables
// barge-in entirely — the half-duplex gate still runs.
type bargeInConfig struct {
// RMS is the normalised energy a frame must exceed to count as him
// talking over her rather than the mic hearing her. It is deliberately
// far above the VAD's own floor: the speaker leaks into the mic at
// roughly ambient level, a person talking at the mic does not.
RMS float64
// Frames is how many consecutive frames must clear RMS before playback
// is cut. One loud frame is a door closing; five in a row is a voice.
Frames int
}
// Enabled reports whether barge-in should be attempted at all.
func (c bargeInConfig) Enabled() bool { return c.Frames > 0 && c.RMS > 0 }
// session is the per-client capture state machine.
type session struct {
vad *VAD
player player
sender utteranceSender
lang string
barge bargeInConfig
// loudFrames counts consecutive over-threshold frames seen while she is
// speaking. Reset whenever a frame falls back under the threshold, and
// whenever playback ends.
loudFrames int
// counters, read by tests and logged on the way out.
suppressed int // frames dropped because she was speaking
bargeIns int // times playback was cut because he spoke over her
sent int // utterances shipped to the daemon
}
func newSession(vad *VAD, p player, s utteranceSender, lang string, barge bargeInConfig) *session {
return &session{vad: vad, player: p, sender: s, lang: lang, barge: barge}
}
// feed processes one 30ms PCM frame.
//
// While the player is running the capture side is muted: the VAD is not fed
// and no utterance can be produced, so Maven's own reply cannot come back in
// as a new command. The one thing that gets through is barge-in — sustained
// energy well above the speaker's leak level cuts playback, and capture
// resumes on the very next frame with a clean VAD.
func (s *session) feed(ctx context.Context, frame []byte) error {
if s.player.Playing() {
s.suppressed++
if !s.barge.Enabled() {
return nil
}
if frameRMS(PCMToI16(frame)) < s.barge.RMS {
s.loudFrames = 0
return nil
}
s.loudFrames++
if s.loudFrames < s.barge.Frames {
return nil
}
// He is talking over her. Cut her off, drop the VAD state that
// accumulated from the echo, and start listening for real.
s.player.Stop()
s.bargeIns++
s.loudFrames = 0
s.vad.Reset()
log.Printf("mavwaked: barge-in — stopped playback")
return nil
}
// Not speaking. If we just stopped, make sure no echo-era state leaks
// into the next utterance.
if s.loudFrames != 0 {
s.loudFrames = 0
s.vad.Reset()
}
utt, state := s.vad.Feed(PCMToI16(frame))
if state == StateSpeech || utt.Bytes == nil {
return nil
}
return s.dispatch(ctx, utt)
}
// dispatch ships a complete utterance and plays whatever comes back.
func (s *session) dispatch(ctx context.Context, utt audio.Audio) error {
log.Printf("mavwaked: utterance complete (%.2fs, %d bytes), sending...", utt.Duration(), len(utt.Bytes))
reply, err := s.sender.Send(ctx, utt, s.lang)
s.sent++
if err != nil {
return err
}
if len(reply.Bytes) == 0 {
log.Printf("mavwaked: empty reply audio (text only)")
return nil
}
// The VAD has been accumulating from the buffered mic stream while the
// round-trip blocked. None of it is a command — reset before the
// speaker opens, so the first post-reply frame starts clean.
s.vad.Reset()
s.player.Play(reply)
return nil
}
+282
View File
@@ -0,0 +1,282 @@
package main
import (
"context"
"errors"
"math"
"testing"
"github.com/kami/maven/internal/audio"
)
// fakePlayer records Play/Stop instead of shelling out to aplay.
type fakePlayer struct {
playing bool
plays int
stops int
last audio.Audio
}
func (p *fakePlayer) Play(a audio.Audio) { p.playing = true; p.plays++; p.last = a }
func (p *fakePlayer) Stop() { p.playing = false; p.stops++ }
func (p *fakePlayer) Playing() bool { return p.playing }
// fakeSender records what was shipped and hands back a canned reply.
type fakeSender struct {
sent []audio.Audio
reply audio.Audio
err error
}
func (s *fakeSender) Send(_ context.Context, utt audio.Audio, _ string) (audio.Audio, error) {
s.sent = append(s.sent, utt)
return s.reply, s.err
}
func replyAudio() audio.Audio {
return audio.Audio{Format: audio.PCM16kMono, Bytes: make([]byte, 16000)}
}
// frameAt returns a 30ms frame whose RMS is approximately rms.
func frameAt(rms float64) []byte {
amp := rms * math.Sqrt2 * 32768
f := make([]int16, frameSamples)
for i := range f {
f[i] = int16(amp * math.Sin(2*math.Pi*440*float64(i)/16000))
}
return pcmBytes(f)
}
func silentBytes() []byte { return make([]byte, frameSamples*2) }
// newTestSession wires a session with fakes and a default VAD.
func newTestSession(barge bargeInConfig) (*session, *fakePlayer, *fakeSender) {
p := &fakePlayer{}
s := &fakeSender{reply: replyAudio()}
return newSession(NewVAD(0, 0, 0, 0), p, s, "ru", barge), p, s
}
// speakThenPause drives a full utterance through the session: enough loud
// frames to trigger, then enough silence to end it.
func speakThenPause(t *testing.T, sess *session) {
t.Helper()
speechFrames := (defaultSpeechMs + defaultFrameMs - 1) / defaultFrameMs
silenceFrames := (defaultSilenceMs+defaultFrameMs-1)/defaultFrameMs + 2
loud := frameAt(0.35)
for i := 0; i < speechFrames+5; i++ {
if err := sess.feed(context.Background(), loud); err != nil {
t.Fatalf("feed loud frame %d: %v", i, err)
}
}
for i := 0; i < silenceFrames; i++ {
if err := sess.feed(context.Background(), silentBytes()); err != nil {
t.Fatalf("feed silent frame %d: %v", i, err)
}
}
}
func TestSessionSendsUtteranceAndPlaysReply(t *testing.T) {
sess, p, snd := newTestSession(bargeInConfig{})
speakThenPause(t, sess)
if len(snd.sent) != 1 {
t.Fatalf("sent %d utterances, want 1", len(snd.sent))
}
if snd.sent[0].Format != audio.PCM16kMono {
t.Errorf("utterance format = %+v, want canonical", snd.sent[0].Format)
}
if p.plays != 1 {
t.Errorf("plays = %d, want 1", p.plays)
}
}
// The bug this whole file exists for: while the speaker is running, the mic
// hears Maven and the old code shipped that back as a fresh command.
func TestSessionDoesNotHearItselfWhilePlaying(t *testing.T) {
sess, p, snd := newTestSession(bargeInConfig{})
speakThenPause(t, sess)
if !p.Playing() {
t.Fatal("expected playback to be running after the reply")
}
// Feed a long stretch of loud audio — Maven's own voice coming back in.
base := sess.suppressed
loud := frameAt(0.35)
for i := 0; i < 200; i++ {
if err := sess.feed(context.Background(), loud); err != nil {
t.Fatalf("feed echo frame %d: %v", i, err)
}
}
if len(snd.sent) != 1 {
t.Fatalf("sent %d utterances, want 1 — her own reply was captured as a command", len(snd.sent))
}
if got := sess.suppressed - base; got != 200 {
t.Errorf("suppressed %d of the 200 echo frames, want all of them", got)
}
if p.stops != 0 {
t.Errorf("stops = %d, want 0 — barge-in is off, nothing should cut her off", p.stops)
}
}
// With barge-in off, no amount of noise stops playback.
func TestSessionBargeInDisabledByDefault(t *testing.T) {
sess, p, _ := newTestSession(bargeInConfig{})
if sess.barge.Enabled() {
t.Fatal("zero bargeInConfig must be disabled")
}
speakThenPause(t, sess)
veryLoud := frameAt(0.6)
for i := 0; i < 50; i++ {
_ = sess.feed(context.Background(), veryLoud)
}
if p.stops != 0 || sess.bargeIns != 0 {
t.Fatalf("stops = %d, bargeIns = %d, want 0 with barge-in off", p.stops, sess.bargeIns)
}
}
func TestSessionBargeInCutsPlayback(t *testing.T) {
barge := bargeInConfig{RMS: 0.12, Frames: 5}
sess, p, _ := newTestSession(barge)
speakThenPause(t, sess)
if !p.Playing() {
t.Fatal("expected playback after the reply")
}
// Four loud frames must not be enough — a door closing is not a voice.
veryLoud := frameAt(0.35)
for i := 0; i < 4; i++ {
_ = sess.feed(context.Background(), veryLoud)
}
if p.stops != 0 {
t.Fatalf("playback cut after 4 frames, want it to hold until %d", barge.Frames)
}
// The fifth cuts her off.
_ = sess.feed(context.Background(), veryLoud)
if p.stops != 1 || sess.bargeIns != 1 {
t.Fatalf("stops = %d, bargeIns = %d, want 1 and 1", p.stops, sess.bargeIns)
}
if p.Playing() {
t.Fatal("still playing after barge-in")
}
}
// A burst that falls back under the threshold resets the counter, so noise
// spread over a whole reply never accumulates into a false barge-in.
func TestSessionBargeInNeedsConsecutiveFrames(t *testing.T) {
sess, p, _ := newTestSession(bargeInConfig{RMS: 0.12, Frames: 5})
speakThenPause(t, sess)
veryLoud := frameAt(0.35)
quiet := frameAt(0.02)
for i := 0; i < 20; i++ {
_ = sess.feed(context.Background(), veryLoud)
_ = sess.feed(context.Background(), veryLoud)
_ = sess.feed(context.Background(), quiet)
}
if p.stops != 0 || sess.bargeIns != 0 {
t.Fatalf("stops = %d, bargeIns = %d, want 0 — two-frame bursts must not accumulate", p.stops, sess.bargeIns)
}
}
// Speaker leak sits near the room floor; it must never reach the barge-in bar.
func TestSessionEchoLevelAudioNeverBargesIn(t *testing.T) {
sess, p, _ := newTestSession(bargeInConfig{RMS: 0.12, Frames: 5})
speakThenPause(t, sess)
base := sess.suppressed
leak := frameAt(0.05) // loud enough for the VAD, far under the barge bar
for i := 0; i < 300; i++ {
_ = sess.feed(context.Background(), leak)
}
if p.stops != 0 {
t.Fatalf("stops = %d, want 0 — speaker leak must not read as barge-in", p.stops)
}
if got := sess.suppressed - base; got != 300 {
t.Errorf("suppressed %d of the 300 leak frames, want all of them", got)
}
}
// After barge-in the VAD must start clean, so the interrupting speech is
// captured as a whole utterance rather than joined onto echo state.
func TestSessionCapturesTheInterruptingUtterance(t *testing.T) {
sess, p, snd := newTestSession(bargeInConfig{RMS: 0.12, Frames: 5})
speakThenPause(t, sess)
veryLoud := frameAt(0.35)
for i := 0; i < 5; i++ {
_ = sess.feed(context.Background(), veryLoud)
}
if p.stops != 1 {
t.Fatalf("expected barge-in, stops = %d", p.stops)
}
// He keeps talking; that is a new command.
speakThenPause(t, sess)
if len(snd.sent) != 2 {
t.Fatalf("sent %d utterances, want 2 — the interruption itself must be heard", len(snd.sent))
}
if p.plays != 2 {
t.Errorf("plays = %d, want 2", p.plays)
}
}
// A failed round-trip must surface as an error and must not start playback.
func TestSessionSendErrorDoesNotPlay(t *testing.T) {
p := &fakePlayer{}
snd := &fakeSender{err: errors.New("boom")}
sess := newSession(NewVAD(0, 0, 0, 0), p, snd, "ru", bargeInConfig{})
speechFrames := (defaultSpeechMs + defaultFrameMs - 1) / defaultFrameMs
silenceFrames := (defaultSilenceMs+defaultFrameMs-1)/defaultFrameMs + 2
loud := frameAt(0.35)
var lastErr error
for i := 0; i < speechFrames+5; i++ {
_ = sess.feed(context.Background(), loud)
}
for i := 0; i < silenceFrames; i++ {
if err := sess.feed(context.Background(), silentBytes()); err != nil {
lastErr = err
}
}
if lastErr == nil {
t.Fatal("send error was swallowed")
}
if p.plays != 0 || p.Playing() {
t.Fatalf("plays = %d, playing = %v, want no playback on a failed round-trip", p.plays, p.Playing())
}
}
// An empty reply (text-only turn) must leave the capture side open.
func TestSessionEmptyReplyLeavesCaptureOpen(t *testing.T) {
p := &fakePlayer{}
snd := &fakeSender{reply: audio.Audio{Format: audio.PCM16kMono}}
sess := newSession(NewVAD(0, 0, 0, 0), p, snd, "ru", bargeInConfig{})
speakThenPause(t, sess)
if p.plays != 0 {
t.Fatalf("plays = %d, want 0 for an empty reply", p.plays)
}
speakThenPause(t, sess)
if len(snd.sent) != 2 {
t.Fatalf("sent %d, want 2 — capture must stay open when there is no audio reply", len(snd.sent))
}
}
func TestBargeInConfigEnabled(t *testing.T) {
cases := []struct {
c bargeInConfig
want bool
}{
{bargeInConfig{}, false},
{bargeInConfig{RMS: 0.12}, false},
{bargeInConfig{Frames: 5}, false},
{bargeInConfig{RMS: 0.12, Frames: 5}, true},
}
for _, tc := range cases {
if got := tc.c.Enabled(); got != tc.want {
t.Errorf("%+v.Enabled() = %v, want %v", tc.c, got, tc.want)
}
}
}
+9
View File
@@ -23,6 +23,15 @@ const (
defaultSilenceMs = 800 // silence hold before declaring end-of-utterance
defaultMaxMs = 10000 // cap single utterance at 10s
defaultMinRMS = 0.01 // RMS floor (same as mavsttd)
// Barge-in thresholds. Only used when -barge-in is passed. The RMS is
// x10000 like -min-rms, and sits an order of magnitude above the VAD's
// own floor on purpose: with no acoustic echo canceller, a frame only
// counts as "he is talking over her" if it is far louder than what the
// speaker leaks back into the mic. 5 frames is 150ms — long enough that
// a door or a cough does not cut her off mid-sentence.
defaultBargeRMS = 1200 // 0.12 normalised RMS
defaultBargeFrames = 5
)
// frameSamples — samples per 30ms frame at 16kHz.
+24
View File
@@ -0,0 +1,24 @@
{{template "shellTop" "events"}}
<h1>Intake</h1>
<div class=hint>Everything that arrived, newest first — a relayed notification, a mail candidate, a feed
item, a changed page, a spend, a presence probe. One envelope per write; the durable row is still the
fact, note or task itself. Held in memory only, so a restart empties this.</div>
{{if .Err}}<div class=hint>journal unavailable: {{.Err}}</div>{{end}}
{{if and (not .Events) (not .Err)}}
<div class=hint>nothing has arrived yet</div>
{{end}}
{{if .Events}}
<div class=scroll><table class=mono>
<tr><th>when<th>source<th>kind<th>pri<th>what<th>detail</tr>
{{range .Events}}<tr>
<td>{{.OccurredAt.Format "02.01 15:04:05"}}</td>
<td class=gray>{{.Source}}</td>
<td class=gray>{{.Kind}}</td>
<td class=gray>{{.Priority}}</td>
<td>{{.Title}}</td>
<td class=gray>{{.Body}}</td>
</tr>{{end}}
</table></div>
{{end}}
{{template "shellBottom"}}
</html>
+107
View File
@@ -0,0 +1,107 @@
package main
import (
"context"
"errors"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/ipc"
)
// eventsCore serves a canned intake journal. Embedding
// ipc.UnimplementedCoreAPI means any other call fails loudly.
type eventsCore struct {
ipc.UnimplementedCoreAPI
events []ipc.IntakeEvent
err error
gotN int
}
func (c *eventsCore) RecentEvents(_ context.Context, n int) ([]ipc.IntakeEvent, error) {
c.gotN = n
return c.events, c.err
}
func getEvents(t *testing.T, core ipc.CoreAPI) *httptest.ResponseRecorder {
t.Helper()
w := httptest.NewRecorder()
handleEvents(w, httptest.NewRequest(http.MethodGet, "/events", nil), core)
return w
}
func TestEventsPageRendersTheJournal(t *testing.T) {
core := &eventsCore{events: []ipc.IntakeEvent{
{Source: "rss:tech", Kind: "note", Title: "Вышло ядро 6.19", Priority: "low",
OccurredAt: time.Date(2026, 8, 1, 7, 15, 0, 0, time.UTC)},
{Source: "ambient:notif", Kind: "fact", Title: "calendar_event_20260801_планёрка",
Body: "10:00-11:00 планёрка", Priority: "low",
OccurredAt: time.Date(2026, 8, 1, 10, 0, 0, 0, time.UTC)},
}}
w := getEvents(t, core)
if w.Code != http.StatusOK {
t.Fatalf("status = %d, want 200", w.Code)
}
body := w.Body.String()
for _, want := range []string{"rss:tech", "Вышло ядро 6.19", "ambient:notif", "10:00-11:00 планёрка", "01.08 10:00:00"} {
if !strings.Contains(body, want) {
t.Errorf("page does not mention %q", want)
}
}
if core.gotN != eventsPageLimit {
t.Errorf("asked core for %d events, want %d", core.gotN, eventsPageLimit)
}
}
func TestEventsPageSaysNothingArrived(t *testing.T) {
w := getEvents(t, &eventsCore{})
if w.Code != http.StatusOK {
t.Fatalf("status = %d, want 200", w.Code)
}
if !strings.Contains(w.Body.String(), "nothing has arrived yet") {
t.Error("empty journal did not render the empty-state line")
}
}
func TestEventsPageReportsAReadFailure(t *testing.T) {
// An unreachable journal must say so rather than render an empty table,
// which would imply nothing arrived.
w := getEvents(t, &eventsCore{err: errors.New("core is down")})
if w.Code != http.StatusOK {
t.Fatalf("status = %d, want 200 with the error rendered", w.Code)
}
body := w.Body.String()
if !strings.Contains(body, "journal unavailable") || !strings.Contains(body, "core is down") {
t.Errorf("page did not report the read failure: %s", body)
}
if strings.Contains(body, "nothing has arrived yet") {
t.Error("a failed read rendered as an empty journal")
}
}
func TestEventsPageWithoutCore(t *testing.T) {
w := getEvents(t, nil)
if w.Code != http.StatusServiceUnavailable {
t.Errorf("status = %d, want 503", w.Code)
}
}
func TestEventsPageEscapesIntakeText(t *testing.T) {
// Titles come from outside — a feed headline, a notification. They are shown
// on a page and must never be able to inject markup into it.
core := &eventsCore{events: []ipc.IntakeEvent{{
Source: "rss:x", Kind: "note", Priority: "low",
Title: `<script>alert(1)</script>`,
OccurredAt: time.Date(2026, 8, 1, 7, 0, 0, 0, time.UTC),
}}}
body := getEvents(t, core).Body.String()
if strings.Contains(body, "<script>alert(1)</script>") {
t.Error("intake title was not escaped")
}
if !strings.Contains(body, "&lt;script&gt;") {
t.Error("intake title is missing from the page entirely")
}
}
+42
View File
@@ -71,6 +71,9 @@ var ecosystemHTML string
//go:embed morning.html
var morningHTML string
//go:embed events.html
var eventsHTML string
// ── Ethos Workstation Shell ──
//
// Two template pieces that wrap every page:
@@ -105,6 +108,7 @@ var sidebarSections = []struct {
{Label: "Reminders", URL: "/reminders", Key: "reminders"},
{Label: "Routines", URL: "/routines", Key: "routines"},
{Label: "Morning", URL: "/morning", Key: "morning"},
{Label: "Intake", URL: "/events", Key: "events"},
},
},
{
@@ -318,6 +322,10 @@ var ecosystemTmpl = template.Must(template.New("ecosystem").Funcs(shellFuncs()).
// morning routine (internal/morning). Same shape as trace.html: a plain
// server-rendered page, refreshed on reload — no live-update loop, since
// checklist state changes on the scale of minutes, not seconds.
// eventsTmpl — the unified intake journal (Vikunja #283), read-only. Same
// shape as trace.html and morning.html: server-rendered, refreshed on reload.
var eventsTmpl = template.Must(template.New("events").Funcs(shellFuncs()).Parse(shellTopHTML + eventsHTML + shellBottomHTML))
var morningTmpl = template.Must(template.New("morning").Funcs(shellFuncs()).Parse(shellTopHTML + morningHTML + shellBottomHTML))
func noCache(h http.Handler) http.Handler {
@@ -430,6 +438,9 @@ func main() {
mux.HandleFunc("/morning", func(w http.ResponseWriter, r *http.Request) {
handleMorning(w, r, core)
})
mux.HandleFunc("/events", func(w http.ResponseWriter, r *http.Request) {
handleEvents(w, r, core)
})
ecoURLsCfg := ecoURLs{nexus: *nexusURL, praxis: *praxisURL, hexis: *hexisURL}
mux.HandleFunc("/ecosystem", func(w http.ResponseWriter, r *http.Request) {
handleEcosystem(w, r, ecoURLsCfg)
@@ -1206,6 +1217,37 @@ type morningView struct {
Routines []ipc.MorningRoutineStatus
}
// eventsView — what /events renders. Err is set instead of Events when the
// core could not serve the journal, so the page says why rather than showing an
// empty intake and implying nothing arrived.
type eventsView struct {
Events []ipc.IntakeEvent
Err string
}
// eventsPageLimit — how many envelopes the page shows. The ring holds more; a
// page is for scanning what just happened, not for archaeology.
const eventsPageLimit = 200
func handleEvents(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI) {
if core == nil {
http.Error(w, "intake journal disabled (no -core)", http.StatusServiceUnavailable)
return
}
var view eventsView
evs, err := core.RecentEvents(r.Context(), eventsPageLimit)
if err != nil {
log.Printf("events: %v", err)
view.Err = err.Error()
} else {
view.Events = evs
}
w.Header().Set("Content-Type", "text/html; charset=utf-8")
if err := eventsTmpl.Execute(w, view); err != nil {
log.Printf("events render: %v", err)
}
}
func handleVoice(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/html; charset=utf-8")
if err := voiceTmpl.Execute(w, nil); err != nil {
+293
View File
@@ -0,0 +1,293 @@
package main
import (
"bytes"
"context"
"crypto/ecdsa"
"crypto/elliptic"
"crypto/rand"
"crypto/sha256"
"encoding/base64"
"encoding/binary"
"encoding/json"
"errors"
"net/http"
"net/http/httptest"
"path/filepath"
"strings"
"testing"
"github.com/kami/maven/internal/webauthn"
)
const prfTestOrigin = "https://maven.test"
const prfTestRPID = "maven.test"
// fakeKeyIPC stands in for the mavend socket and records exactly what secret
// each call received — the point of the whole test file is that it is the PRF
// output and never the credential public key.
type fakeKeyIPC struct {
unlockSecret []byte
wrapSecret []byte
unlockCalls int
wrapCalls int
unlockErr error
}
func (f *fakeKeyIPC) Unlock(_ context.Context, secret []byte) error {
f.unlockCalls++
f.unlockSecret = bytes.Clone(secret)
return f.unlockErr
}
func (f *fakeKeyIPC) StoreEncryptionKey(_ context.Context, secret []byte) error {
f.wrapCalls++
f.wrapSecret = bytes.Clone(secret)
return nil
}
func b64u(b []byte) string { return base64.RawURLEncoding.EncodeToString(b) }
// prfAuthenticator is a minimal software authenticator: a P-256 key plus the
// COSE encoding of its public half.
type prfAuthenticator struct {
key *ecdsa.PrivateKey
credID []byte
cose []byte
}
func newPRFAuthenticator(t *testing.T) *prfAuthenticator {
t.Helper()
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
t.Fatalf("generate key: %v", err)
}
x := key.PublicKey.X.FillBytes(make([]byte, 32))
y := key.PublicKey.Y.FillBytes(make([]byte, 32))
// COSE_Key: {1: 2 (EC2), 3: -7 (ES256), -1: 1 (P-256), -2: x, -3: y}
var c []byte
c = append(c, 0xa5) // map(5)
c = append(c, 0x01, 0x02) // 1: 2
c = append(c, 0x03, 0x26) // 3: -7
c = append(c, 0x20, 0x01) // -1: 1
c = append(c, 0x21, 0x58, 0x20) // -2: bytes(32)
c = append(c, x...)
c = append(c, 0x22, 0x58, 0x20) // -3: bytes(32)
c = append(c, y...)
return &prfAuthenticator{key: key, credID: []byte("prf-cred"), cose: c}
}
func (a *prfAuthenticator) authData(flags byte, counter uint32, attested bool) []byte {
h := sha256.Sum256([]byte(prfTestRPID))
d := append([]byte{}, h[:]...)
d = append(d, flags)
cb := make([]byte, 4)
binary.BigEndian.PutUint32(cb, counter)
d = append(d, cb...)
if attested {
d = append(d, make([]byte, 16)...) // aaguid
l := make([]byte, 2)
binary.BigEndian.PutUint16(l, uint16(len(a.credID)))
d = append(d, l...)
d = append(d, a.credID...)
d = append(d, a.cose...)
}
return d
}
func clientDataJSON(typ, challenge string) []byte {
b, _ := json.Marshal(map[string]string{"type": typ, "challenge": challenge, "origin": prfTestOrigin})
return b
}
// register drives POST /register/finish with a valid attestation.
func (a *prfAuthenticator) register(t *testing.T, h *PasskeyHandle) {
t.Helper()
_, chal, err := h.rp.CreationOptions([]byte("u"), "user")
if err != nil {
t.Fatalf("CreationOptions: %v", err)
}
// {"fmt":"none","attStmt":{},"authData":<bytes>}
att := []byte{0xa3}
att = append(att, 0x63, 'f', 'm', 't', 0x64, 'n', 'o', 'n', 'e')
att = append(att, 0x67, 'a', 't', 't', 'S', 't', 'm', 't', 0xa0)
ad := a.authData(1<<6|0x05, 0, true)
att = append(att, 0x68, 'a', 'u', 't', 'h', 'D', 'a', 't', 'a')
att = append(att, 0x59, byte(len(ad)>>8), byte(len(ad)))
att = append(att, ad...)
body, _ := json.Marshal(map[string]any{
"challenge": chal,
"credential": map[string]any{
"id": b64u(a.credID),
"type": "public-key",
"response": map[string]any{
"clientDataJSON": b64u(clientDataJSON("webauthn.create", chal)),
"attestationObject": b64u(att),
},
},
})
w := httptest.NewRecorder()
h.RegisterFinish(w, httptest.NewRequest(http.MethodPost, "/auth/webauthn/register/finish", bytes.NewReader(body)))
if w.Code != http.StatusOK {
t.Fatalf("RegisterFinish: %d %s", w.Code, w.Body.String())
}
}
// assert drives POST /assert/finish with a valid assertion and the given
// base64url PRF result.
func (a *prfAuthenticator) assert(t *testing.T, h *PasskeyHandle, prf string) *httptest.ResponseRecorder {
t.Helper()
_, chal, err := h.rp.AssertionOptions()
if err != nil {
t.Fatalf("AssertionOptions: %v", err)
}
ad := a.authData(0x05, 7, false)
cdj := clientDataJSON("webauthn.get", chal)
hash := sha256.Sum256(cdj)
sig, err := ecdsa.SignASN1(rand.Reader, a.key, append(append([]byte{}, ad...), hash[:]...))
if err != nil {
t.Fatalf("sign: %v", err)
}
body, _ := json.Marshal(map[string]any{
"challenge": chal,
"prf": prf,
"credential": map[string]any{
"id": b64u(a.credID),
"type": "public-key",
"response": map[string]any{
"clientDataJSON": b64u(cdj),
"authenticatorData": b64u(ad),
"signature": b64u(sig),
},
},
})
w := httptest.NewRecorder()
h.AssertFinish(w, httptest.NewRequest(http.MethodPost, "/auth/webauthn/assert/finish", bytes.NewReader(body)))
return w
}
func newPRFHandle(t *testing.T, key *fakeKeyIPC) *PasskeyHandle {
t.Helper()
store, err := newCredentialStore(filepath.Join(t.TempDir(), "passkeys.json"))
if err != nil {
t.Fatalf("credential store: %v", err)
}
return &PasskeyHandle{
rp: webauthn.NewRP(webauthn.Config{Origin: prfTestOrigin, RPID: prfTestRPID, RPName: "maven"}),
encryptFn: key,
store: store,
session: webauthn.NewPasskeySession(0),
}
}
// The fix for Vikunja #14: what goes over IPC is the PRF secret from the
// authenticator, not the credential public key sitting in passkeys.json.
func TestAssertSendsPRFSecretNotPublicKey(t *testing.T) {
key := &fakeKeyIPC{}
h := newPRFHandle(t, key)
auth := newPRFAuthenticator(t)
auth.register(t, h)
// Enrolment must not wrap anything: create() yields no PRF result.
if key.wrapCalls != 0 || key.unlockCalls != 0 {
t.Fatalf("registration touched the key IPC (wrap=%d unlock=%d)", key.wrapCalls, key.unlockCalls)
}
secret := make([]byte, 32)
for i := range secret {
secret[i] = byte(i + 1)
}
if w := auth.assert(t, h, b64u(secret)); w.Code != http.StatusOK {
t.Fatalf("AssertFinish: %d %s", w.Code, w.Body.String())
}
if key.unlockCalls != 1 || key.wrapCalls != 1 {
t.Fatalf("unlock=%d wrap=%d, want 1 and 1", key.unlockCalls, key.wrapCalls)
}
if !bytes.Equal(key.unlockSecret, secret) {
t.Errorf("Unlock got %x, want the PRF secret %x", key.unlockSecret, secret)
}
if !bytes.Equal(key.wrapSecret, secret) {
t.Errorf("StoreEncryptionKey got %x, want the PRF secret %x", key.wrapSecret, secret)
}
// And explicitly: not the credential public key.
pub, _, err := h.store.Lookup(b64u(auth.credID))
if err != nil {
t.Fatalf("lookup: %v", err)
}
if bytes.Equal(key.unlockSecret, pub) {
t.Fatal("the credential public key was sent as the unlock secret")
}
}
// An authenticator without PRF must produce no unlock attempt at all — the
// assertion still succeeds (step-up works), but cold-start unlock stays off
// rather than falling back to something weaker.
func TestAssertWithoutPRFDoesNotUnlock(t *testing.T) {
for _, prf := range []string{"", "!!!not-base64!!!", b64u(make([]byte, 32)), b64u(make([]byte, 16))} {
key := &fakeKeyIPC{}
h := newPRFHandle(t, key)
auth := newPRFAuthenticator(t)
auth.register(t, h)
w := auth.assert(t, h, prf)
if w.Code != http.StatusOK {
t.Fatalf("prf=%q: AssertFinish %d %s", prf, w.Code, w.Body.String())
}
if key.unlockCalls != 0 || key.wrapCalls != 0 {
t.Errorf("prf=%q: unlock=%d wrap=%d, want no key IPC at all", prf, key.unlockCalls, key.wrapCalls)
}
}
}
// A failed unlock must not fail the assertion: step-up is independently valid,
// and a locked daemon degrades rather than breaking the login.
func TestAssertSucceedsWhenUnlockFails(t *testing.T) {
key := &fakeKeyIPC{unlockErr: errors.New("wrong credential")}
h := newPRFHandle(t, key)
auth := newPRFAuthenticator(t)
auth.register(t, h)
secret := bytes.Repeat([]byte{3}, 32)
if w := auth.assert(t, h, b64u(secret)); w.Code != http.StatusOK {
t.Fatalf("AssertFinish: %d %s", w.Code, w.Body.String())
}
if key.unlockCalls != 1 {
t.Errorf("unlock attempted %d times, want 1", key.unlockCalls)
}
}
// A forged assertion must never reach the unlock path.
func TestForgedAssertionNeverUnlocks(t *testing.T) {
key := &fakeKeyIPC{}
h := newPRFHandle(t, key)
auth := newPRFAuthenticator(t)
auth.register(t, h)
// A different key signing over the same credential id.
attacker := newPRFAuthenticator(t)
attacker.credID = auth.credID
w := attacker.assert(t, h, b64u(bytes.Repeat([]byte{4}, 32)))
if w.Code == http.StatusOK {
t.Fatal("an assertion signed by the wrong key was accepted")
}
if key.unlockCalls != 0 || key.wrapCalls != 0 {
t.Fatalf("a forged assertion reached the key IPC (unlock=%d wrap=%d)", key.unlockCalls, key.wrapCalls)
}
}
// The browser side is the only place the PRF result exists. If the page stops
// asking for it or stops reading it back, cold-start unlock silently dies with
// nothing failing, so the page source is asserted directly.
func TestPasskeyPageRequestsAndPostsPRF(t *testing.T) {
for _, want := range []string{
"getClientExtensionResults",
"ext.prf.results.first",
"body:JSON.stringify({challenge,prf,",
} {
if !strings.Contains(passkeyPageHTML, want) {
t.Errorf("the passkey page no longer contains %q", want)
}
}
}
+52 -33
View File
@@ -23,8 +23,8 @@ type assertIPC interface {
// is *ipc.Client; in-process CoreAPI adapters do not implement it. When nil,
// StoreEncryptionKey and Unlock are silently skipped.
type keyIPC interface {
StoreEncryptionKey(ctx context.Context, publicKey []byte) error
Unlock(ctx context.Context, publicKey []byte) error
StoreEncryptionKey(ctx context.Context, secret []byte) error
Unlock(ctx context.Context, secret []byte) error
}
// PasskeyHandle holds the WebAuthn relying party, a local in-memory credential
@@ -102,17 +102,31 @@ async function enroll(){try{
const r=await fetch('/auth/webauthn/register/finish',{method:'POST',headers:{'content-type':'application/json'},
body:JSON.stringify({challenge,credential:{id:c.id,type:c.type,response:{
clientDataJSON:b64u(c.response.clientDataJSON),attestationObject:b64u(c.response.attestationObject)}}})});
say(r.ok?'enrolled ✓':'enroll failed: '+await r.text(),r.ok);
if(!r.ok){say('enroll failed: '+await r.text(),false);return;}
// The wrapped key can only be written from an assertion: PRF results are
// not produced at create() time on most authenticators. Enrolment reports
// whether PRF is available at all so he is not told cold-start works when
// it cannot.
const ext=c.getClientExtensionResults?c.getClientExtensionResults():{};
const prfOK=!!(ext.prf&&ext.prf.enabled);
say(prfOK?'enrolled ✓ — now assert once to write the cold-start key':
'enrolled ✓ — but this authenticator has no PRF: cold-start unlock unavailable',true);
}catch(e){say('enroll error: '+e,false);}}
async function assert(){try{
const {challenge,options}=await (await fetch('/auth/webauthn/assert/begin')).json();
options.challenge=ub64(options.challenge);
const c=await navigator.credentials.get({publicKey:options});
// The PRF result is the cold-start secret. It never touches localStorage
// and is posted once, over the same request as the assertion.
const ext=c.getClientExtensionResults?c.getClientExtensionResults():{};
const prf=ext.prf&&ext.prf.results&&ext.prf.results.first?b64u(ext.prf.results.first):'';
const r=await fetch('/auth/webauthn/assert/finish',{method:'POST',headers:{'content-type':'application/json'},
body:JSON.stringify({challenge,credential:{id:c.id,type:c.type,response:{
body:JSON.stringify({challenge,prf,credential:{id:c.id,type:c.type,response:{
clientDataJSON:b64u(c.response.clientDataJSON),authenticatorData:b64u(c.response.authenticatorData),
signature:b64u(c.response.signature)}}})});
say(r.ok?'stepped up ✓ — enable tools now':'assert failed: '+await r.text(),r.ok);
if(!r.ok){say('assert failed: '+await r.text(),false);return;}
say(prf?'stepped up ✓ — enable tools now':
'stepped up ✓ — no PRF from this authenticator, so cold-start unlock stayed unavailable',true);
}catch(e){say('assert error: '+e,false);}}
</script>`
@@ -140,9 +154,7 @@ func (h *PasskeyHandle) RegisterFinish(w http.ResponseWriter, r *http.Request) {
http.Error(w, "bad request: "+err.Error(), http.StatusBadRequest)
return
}
var enrolledPublicKey []byte
save := func(id string, publicKey []byte, _ []byte, _ string) error {
enrolledPublicKey = publicKey
return h.store.Save(id, publicKey)
}
credID, err := h.rp.FinishRegistration(save, body.Challenge, body.Credential)
@@ -153,19 +165,15 @@ func (h *PasskeyHandle) RegisterFinish(w http.ResponseWriter, r *http.Request) {
}
log.Printf("webauthn: registered credential %s", credID)
// If mavend is reachable and supports key wrapping, store the encryption
// key wrapped with this credential's public key — enables cold-start unlock.
if h.encryptFn != nil && enrolledPublicKey != nil {
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
defer cancel()
if err := h.encryptFn.StoreEncryptionKey(ctx, enrolledPublicKey); err != nil {
log.Printf("webauthn: store encryption key: %v", err)
// Non-fatal: enrollment still succeeded, the wrapped key can be
// created later via the same endpoint.
} else {
log.Printf("webauthn: encryption key wrapped with credential %s", credID)
}
}
// Note what does NOT happen here: the encryption key is not wrapped at
// enrolment. Wrapping needs the authenticator's PRF output, and create()
// does not produce one on most authenticators — it only reports whether
// the extension is supported. The wrapped key is written on the first
// assertion instead (see AssertFinish).
//
// This used to wrap the key under the credential *public* key, which is
// written to passkeys.json next to the wrapped blob. See the header of
// internal/webauthn/keywrap.go.
json.NewEncoder(w).Encode(map[string]string{"credential_id": credID})
}
@@ -189,6 +197,11 @@ func (h *PasskeyHandle) AssertFinish(w http.ResponseWriter, r *http.Request) {
var body struct {
Challenge string `json:"challenge"`
Credential map[string]any `json:"credential"`
// PRF is the base64url WebAuthn PRF output the browser read out of
// getClientExtensionResults(). Empty when the authenticator has no
// PRF extension: cold-start unlock is then unavailable and we say so
// rather than falling back to something weaker.
PRF string `json:"prf"`
}
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
http.Error(w, "bad request: "+err.Error(), http.StatusBadRequest)
@@ -222,26 +235,32 @@ func (h *PasskeyHandle) AssertFinish(w http.ResponseWriter, r *http.Request) {
}
}
// If the daemon is locked (cold-start), send the credential's public key
// over IPC so mavend can unwrap its encryption key and open the store.
// The public key comes from the local credential store (it was stored
// during enrollment). Non-fatal: if IPC doesn't support Unlock or the
// daemon is already unlocked, the call is a no-op on the server side.
// Cold-start unlock and key wrapping, both keyed on the PRF secret that
// this assertion just produced. The secret is used here and dropped; it is
// never stored on this side.
//
// Order matters: unlock first (if the daemon is locked there is nothing to
// wrap yet), then re-wrap, which writes the blob on the first assertion
// after enrolment and is a harmless rewrite afterwards. Both are
// best-effort — the assertion itself is valid either way.
if h.encryptFn != nil {
publicKey, _, err := h.store.Lookup(credID)
if err == nil && publicKey != nil {
secret, err := webauthn.DecodePRFResult(body.PRF)
switch {
case err != nil:
log.Printf("webauthn: no usable PRF secret from credential %s: %v", credID, err)
default:
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
defer cancel()
if err := h.encryptFn.Unlock(ctx, publicKey); err != nil {
if err := h.encryptFn.Unlock(ctx, secret); err != nil {
log.Printf("webauthn: unlock via credential %s: %v", credID, err)
// Non-fatal: assertion succeeded; if the daemon stays locked
// the user will see errors on subsequent pages, but the
// assertion itself is valid.
} else {
log.Printf("webauthn: daemon unlocked via credential %s", credID)
}
} else if err != nil {
log.Printf("webauthn: lookup credential %s for unlock: %v", credID, err)
if err := h.encryptFn.StoreEncryptionKey(ctx, secret); err != nil {
log.Printf("webauthn: wrap encryption key: %v", err)
} else {
log.Printf("webauthn: encryption key wrapped for credential %s", credID)
}
}
}
+7 -1
View File
@@ -135,7 +135,13 @@ func Requirement(m ipc.Method) Authority {
ipc.MethodListSpeakers,
// The read side of the model swap: which model is resident, which ones are
// allowlisted. It loads nothing and changes nothing.
ipc.MethodModelStatus:
ipc.MethodModelStatus,
// The unified intake journal (Vikunja #283). AuthRead, and listed
// explicitly rather than inherited so the reasoning is on the record: it
// reports what already arrived — sources, keys, note headlines — which is
// the same material RecentFacts and RecentNotes already return at this
// rung. It writes nothing, and it holds nothing a fact read does not.
ipc.MethodRecentEvents:
return AuthRead
}
// Unknown method ⇒ AuthRead, but ipc.dispatch returns ErrUnknownMethod
+22 -1
View File
@@ -169,6 +169,20 @@ type Config struct {
// live in the reader (cmd/mavmaild), never here.
Email *EmailConfig `json:"email,omitempty"`
// IntakeJournal — how many entries the unified intake journal keeps
// (Vikunja #283): one envelope per thing that arrived, whatever direction it
// came from. Absent ⇒ DefaultIntakeJournal. A NEGATIVE value turns the
// journal off entirely, and then there is no decorator on the intake path at
// all.
//
// Not gated behind an "off unless configured" block like feeds or telegram,
// and the distinction is the one CLAUDE.md draws: that rule exists for
// capabilities that reach OUT — a fetch, a send, a third party. This reaches
// nowhere. It is a bounded in-memory log of writes core already performed,
// it is read only by /events and the simulator, and nothing Maven says
// depends on it.
IntakeJournal int `json:"intake_journal,omitempty"`
// Feeds — RSS/Atom feed reading (Vikunja #258). nil / absent ⇒ no feed is
// ever fetched: reading the outside world is off unless configured, like
// the weather and telegram. See FeedsConfig.
@@ -968,7 +982,11 @@ const (
DefaultRepeatInterval = 5 * time.Minute
DefaultAutotuneInterval = 10 * time.Minute
DefaultRouterThreshold = 0.55
DefaultQueryMinScore = 0.55
// DefaultIntakeJournal — entries kept in the unified intake journal
// (Vikunja #283). A busy day is a few hundred intake writes, so this is
// roughly "today and yesterday" at a few hundred KB of memory.
DefaultIntakeJournal = 512
DefaultQueryMinScore = 0.55
// Read off the margin sweep in internal/memory/recalleval on the e5
// embedder: 0.008 answers 68% of real questions (down from 72%) and cuts
// false recall from 5/5 to 1/5. Every larger delta costs real recall
@@ -1018,6 +1036,9 @@ func Load(path string) (*Config, error) {
}
func (c *Config) applyDefaults() {
if c.IntakeJournal == 0 {
c.IntakeJournal = DefaultIntakeJournal
}
if c.TickInterval == 0 {
c.TickInterval = Duration(DefaultTickInterval)
}
+129
View File
@@ -0,0 +1,129 @@
package event
import (
"sync"
"time"
)
// Bus — the in-memory intake journal: a bounded ring of recent Events plus
// zero or more subscribers.
//
// Two properties are load-bearing, both about not changing production
// behaviour when nobody is watching:
//
// - A nil *Bus is a working no-op. Publish on nil returns immediately, so
// an intake path can call b.Publish(...) unconditionally and a daemon that
// never built a bus behaves exactly as it did before. This is what let
// eight callers adopt the envelope without a config flag each.
// - Publish never blocks on a subscriber and never propagates a panic from
// one. Intake is on the request path of POST /api/ambient and of every
// fact write; a slow or broken observer must not be able to stall or kill
// a write that already succeeded.
//
// The ring is bounded because it is memory that nothing prunes otherwise. Its
// contents are a window, not a record: the durable consequence of an event is
// the fact, note or task the intake path wrote.
type Bus struct {
mu sync.Mutex
ring []Event // len == cap once full; oldest at (next % cap)
next int
n int
subs []func(Event)
}
// DefaultCapacity — how many recent events a bus keeps. A busy day is a few
// hundred intake events (a feed poll is one per new item), so this is roughly
// "today and yesterday" at a few hundred KB.
const DefaultCapacity = 512
// NewBus returns a bus keeping the last capacity events. capacity <= 0 uses
// DefaultCapacity.
func NewBus(capacity int) *Bus {
if capacity <= 0 {
capacity = DefaultCapacity
}
return &Bus{ring: make([]Event, capacity)}
}
// Publish normalizes e, drops it if it is not Valid, appends it to the ring and
// hands it to every subscriber. Safe on a nil receiver and safe from any
// goroutine.
//
// now is passed in rather than read from the clock: the whole point of #284's
// replay is that no time.Now() sits inside a path a scenario drives.
func (b *Bus) Publish(e Event, now time.Time) {
if b == nil {
return
}
e = e.Normalize(now)
if !e.Valid() {
return
}
b.mu.Lock()
b.ring[b.next] = e
b.next = (b.next + 1) % len(b.ring)
if b.n < len(b.ring) {
b.n++
}
subs := make([]func(Event), len(b.subs))
copy(subs, b.subs)
b.mu.Unlock()
for _, fn := range subs {
notify(fn, e)
}
}
// notify calls one subscriber, swallowing a panic. A test double or a page
// renderer must not be able to take down a daemon from the intake path.
func notify(fn func(Event), e Event) {
defer func() { _ = recover() }()
fn(e)
}
// Subscribe registers fn to be called for every subsequent event, in publish
// order. There is no unsubscribe: subscribers are wired at startup and live as
// long as the daemon. Safe on a nil receiver (the subscription is dropped,
// which is the honest outcome when there is no bus to subscribe to).
func (b *Bus) Subscribe(fn func(Event)) {
if b == nil || fn == nil {
return
}
b.mu.Lock()
defer b.mu.Unlock()
b.subs = append(b.subs, fn)
}
// Recent returns up to limit events, newest first. limit <= 0 returns
// everything held. Safe on a nil receiver (returns nil).
func (b *Bus) Recent(limit int) []Event {
if b == nil {
return nil
}
b.mu.Lock()
defer b.mu.Unlock()
if b.n == 0 {
return nil
}
if limit <= 0 || limit > b.n {
limit = b.n
}
out := make([]Event, 0, limit)
// next points one past the newest; walk backwards.
for i := 0; i < limit; i++ {
idx := (b.next - 1 - i + len(b.ring)*2) % len(b.ring)
out = append(out, b.ring[idx])
}
return out
}
// Len reports how many events the ring currently holds. Safe on nil.
func (b *Bus) Len() int {
if b == nil {
return 0
}
b.mu.Lock()
defer b.mu.Unlock()
return b.n
}
+198
View File
@@ -0,0 +1,198 @@
// Package event is the unified intake envelope (Vikunja #283,
// 20-07-2026-BACKLOG.md item 1).
//
// # The problem it solves
//
// Things arrive at Maven from a lot of directions: a relayed Android
// notification (POST /api/ambient), a mail the reader extracted candidates
// from (ingest_mail), an RSS item, a changed page the crawler noticed, a
// zenmoney spend, a CalDAV event, a wg handshake that means he is home, a
// photo he sent, a meeting she was asked to record. Each of those grew its own
// shape, its own storage decision and its own log line. Nothing could answer
// "what came in today, from where" without reading eight packages.
//
// An Event is that answer: one flat, source-agnostic description of "something
// arrived". It is deliberately NOT a new storage layer and NOT a new transport.
// Every intake path keeps writing exactly what it wrote before — a fact, a
// note, a candidate task — and additionally describes what it did as an Event.
// The envelope is a VIEW over intake, not a replacement for it, which is why
// adopting it did not require touching eight callers.
//
// # What it is not
//
// - Not a command. An Event is a report of something that happened; nothing
// in Maven executes one. Digestion may read them; it may not be driven by
// an event alone, because "a thing arrived" is not "a thing must be said".
// - Not durable. The bus is a bounded in-memory ring. An event's durable
// consequence is the fact/note/task the intake path already wrote; the
// envelope is the recent-history window on top. A restart losing the ring
// loses nothing that mattered.
// - Not a secret store. Body carries what the intake path was already willing
// to log or store. Nothing puts a mail body, an IMAP password or a
// voiceprint in here, and callers must keep it that way.
package event
import (
"encoding/json"
"strings"
"time"
)
// Event — one thing that arrived, normalized.
//
// The field set is the one recorded in the backlog, and it is intentionally
// small: anything source-specific goes in Payload, so adding a source never
// widens the struct and never breaks a reader.
type Event struct {
// Source — provenance, in the facts vocabulary already used across the
// repo: "ambient:notif", "caldav:personal", "poll:zenmoney", "rss:<feed>",
// "crawl:<watch>", "email:<mailbox>", "infer:wg", "tap:voice". Same string
// the fact or note was written under, so an event and its row can be
// matched up by eye.
Source string `json:"source"`
// Kind — what sort of thing arrived, from the closed set below. This is the
// field digestion switches on; Source is for provenance and display.
Kind string `json:"kind"`
// EntityIDs — Nexus entity ids this event is about, when the intake path
// knew any. Usually empty: most intake happens before enrichment resolves a
// subject to an entity.
EntityIDs []string `json:"entity_ids,omitempty"`
// Title — one short line, safe to show on a page. For a fact it is the key,
// for a note the first line, for a task the task text.
Title string `json:"title"`
// Body — optional detail, already truncated by the caller.
Body string `json:"body,omitempty"`
// Priority — one of PriorityLow / PriorityNormal / PriorityHigh. It is a
// hint about attention, not a delivery instruction: nothing here decides
// whether Maven speaks. That stays with internal/loop and internal/delivery,
// where the severity/presence routing table lives.
Priority string `json:"priority"`
// OccurredAt — when the thing happened, NOT when Maven noticed it. A wg
// handshake carries the handshake instant; an RSS item carries its publish
// time. Intake paths already make this distinction when writing facts, and
// the envelope must not flatten it.
OccurredAt time.Time `json:"occurred_at"`
// Payload — source-specific extra, opaque here. Optional.
Payload json.RawMessage `json:"payload,omitempty"`
}
// Kinds. Closed set: a reader may switch on these exhaustively. A new intake
// path picks the closest existing kind before it adds one — the point of the
// envelope is that digestion has a small stable input.
const (
// KindFact — something was written to the facts table: a calendar read, a
// zenmoney window, a presence probe, a crawler watermark.
KindFact = "fact"
// KindNote — something was written to the notes table: an RSS item, a
// changed page, a meeting transcript, an image description.
KindNote = "note"
// KindTask — a candidate task was captured: the mail reader, the web form,
// the voice path.
KindTask = "task"
// KindMessage — an inbound message on a reach channel. Nothing produces
// this yet (telegram is send-only today); the kind exists so the bridge,
// when it lands, is a constructor and not a schema change.
KindMessage = "message"
// KindHealth — a service or probe reported its own state.
KindHealth = "health"
)
// Priorities.
const (
PriorityLow = "low"
PriorityNormal = "normal"
PriorityHigh = "high"
)
// TitleMaxRunes / BodyMaxRunes bound what an envelope carries. The ring is
// in memory and served to a web page; a 40 KB crawled article has no business
// in either. Cut on a rune boundary — most of this text is Russian and half a
// cyrillic letter is a broken line.
const (
TitleMaxRunes = 120
BodyMaxRunes = 400
)
// Normalize returns e with its fields put in range: whitespace collapsed out
// of Title, Title and Body truncated, an unknown or empty Priority forced to
// PriorityNormal, and a zero OccurredAt filled from now.
//
// It takes now as a parameter rather than reading the clock, so the whole
// package stays pure and the simulator (Vikunja #284) can replay intake against
// a scripted clock.
func (e Event) Normalize(now time.Time) Event {
e.Title = truncateRunes(strings.Join(strings.Fields(e.Title), " "), TitleMaxRunes)
e.Body = truncateRunes(strings.TrimSpace(e.Body), BodyMaxRunes)
if !validPriority(e.Priority) {
e.Priority = PriorityNormal
}
if e.Kind == "" {
e.Kind = KindFact
}
if e.OccurredAt.IsZero() {
e.OccurredAt = now
}
return e
}
// Valid reports whether e carries the minimum a reader can rely on: a source,
// a known kind, a title and a time. The bus drops anything that fails — an
// envelope with no provenance is worse than no envelope, because it looks like
// evidence.
func (e Event) Valid() bool {
return e.Source != "" && validKind(e.Kind) && e.Title != "" && !e.OccurredAt.IsZero()
}
func validKind(k string) bool {
switch k {
case KindFact, KindNote, KindTask, KindMessage, KindHealth:
return true
}
return false
}
func validPriority(p string) bool {
switch p {
case PriorityLow, PriorityNormal, PriorityHigh:
return true
}
return false
}
// truncateRunes cuts s to n runes, marking the cut.
func truncateRunes(s string, n int) string {
r := []rune(s)
if len(r) <= n {
return s
}
return string(r[:n]) + "…"
}
// SourceKind guesses the Kind for a source string when the caller has not said
// otherwise. It exists so the one intake decorator in cmd/mavend does not need
// a switch per writer: the source prefix already tells you what arrived.
//
// Unknown prefixes get fallback, which is what the caller was going to write
// anyway (a WriteFact call knows it is a fact).
func SourceKind(source, fallback string) string {
switch {
case strings.HasPrefix(source, "rss:"), strings.HasPrefix(source, "crawl:"):
return KindNote
case strings.HasPrefix(source, "email:"):
return KindTask
case strings.HasPrefix(source, "probe:"), strings.HasPrefix(source, "health:"):
return KindHealth
}
return fallback
}
+185
View File
@@ -0,0 +1,185 @@
package event
import (
"strings"
"sync"
"testing"
"time"
)
var testNow = time.Date(2026, 8, 1, 9, 30, 0, 0, time.UTC)
func TestNormalizeFillsDefaults(t *testing.T) {
got := Event{Source: "poll:zenmoney", Title: " spent today "}.Normalize(testNow)
if got.Title != "spent today" {
t.Errorf("title = %q, want collapsed whitespace", got.Title)
}
if got.Priority != PriorityNormal {
t.Errorf("priority = %q, want %q", got.Priority, PriorityNormal)
}
if got.Kind != KindFact {
t.Errorf("kind = %q, want %q", got.Kind, KindFact)
}
if !got.OccurredAt.Equal(testNow) {
t.Errorf("occurred_at = %v, want %v", got.OccurredAt, testNow)
}
}
func TestNormalizeKeepsRealOccurredAt(t *testing.T) {
// A wg handshake carries the handshake instant, not "now". Flattening that
// would make every intake look like it happened at notice time.
real := testNow.Add(-3 * time.Hour)
got := Event{Source: "infer:wg", Title: "wg_handshake", OccurredAt: real}.Normalize(testNow)
if !got.OccurredAt.Equal(real) {
t.Errorf("occurred_at = %v, want the supplied %v", got.OccurredAt, real)
}
}
func TestNormalizeTruncatesOnRuneBoundary(t *testing.T) {
long := strings.Repeat("я", TitleMaxRunes+50)
got := Event{Source: "rss:x", Title: long}.Normalize(testNow)
r := []rune(got.Title)
if len(r) != TitleMaxRunes+1 { // +1 for the ellipsis marker
t.Fatalf("title runes = %d, want %d", len(r), TitleMaxRunes+1)
}
if r[len(r)-1] != '…' {
t.Errorf("truncated title does not mark the cut: %q", string(r[len(r)-3:]))
}
for _, c := range r[:TitleMaxRunes] {
if c != 'я' {
t.Fatalf("truncation broke a rune: got %q", c)
}
}
}
func TestNormalizeRejectsUnknownPriority(t *testing.T) {
got := Event{Source: "s", Title: "t", Priority: "URGENT!!"}.Normalize(testNow)
if got.Priority != PriorityNormal {
t.Errorf("priority = %q, want %q", got.Priority, PriorityNormal)
}
}
func TestValid(t *testing.T) {
base := Event{Source: "rss:tech", Kind: KindNote, Title: "заголовок", OccurredAt: testNow}
if !base.Valid() {
t.Fatal("well-formed event reported invalid")
}
for name, mut := range map[string]func(Event) Event{
"no source": func(e Event) Event { e.Source = ""; return e },
"no title": func(e Event) Event { e.Title = ""; return e },
"no time": func(e Event) Event { e.OccurredAt = time.Time{}; return e },
"bad kind": func(e Event) Event { e.Kind = "whatever"; return e },
} {
if mut(base).Valid() {
t.Errorf("%s: reported valid", name)
}
}
}
func TestSourceKind(t *testing.T) {
cases := map[string]string{
"rss:tech": KindNote,
"crawl:kernel": KindNote,
"email:inbox": KindTask,
"probe:netdata": KindHealth,
"ambient:notif": KindFact,
"tap:voice": KindFact,
}
for src, want := range cases {
if got := SourceKind(src, KindFact); got != want {
t.Errorf("SourceKind(%q) = %q, want %q", src, got, want)
}
}
}
func TestBusNilIsANoOp(t *testing.T) {
// The whole adoption story depends on this: an intake path calls Publish
// unconditionally, and a daemon with no bus behaves as it did before.
var b *Bus
b.Publish(Event{Source: "s", Kind: KindFact, Title: "t"}, testNow)
b.Subscribe(func(Event) { t.Error("nil bus delivered to a subscriber") })
if got := b.Recent(10); got != nil {
t.Errorf("Recent on nil bus = %v, want nil", got)
}
if got := b.Len(); got != 0 {
t.Errorf("Len on nil bus = %d, want 0", got)
}
}
func TestBusRecentIsNewestFirst(t *testing.T) {
b := NewBus(8)
for _, title := range []string{"one", "two", "three"} {
b.Publish(Event{Source: "rss:t", Kind: KindNote, Title: title}, testNow)
}
got := b.Recent(0)
if len(got) != 3 {
t.Fatalf("len = %d, want 3", len(got))
}
want := []string{"three", "two", "one"}
for i, w := range want {
if got[i].Title != w {
t.Errorf("Recent()[%d] = %q, want %q", i, got[i].Title, w)
}
}
if lim := b.Recent(2); len(lim) != 2 || lim[0].Title != "three" {
t.Errorf("Recent(2) = %v, want the two newest", lim)
}
}
func TestBusRingEvicts(t *testing.T) {
b := NewBus(3)
for _, title := range []string{"a", "b", "c", "d", "e"} {
b.Publish(Event{Source: "s", Kind: KindFact, Title: title}, testNow)
}
if b.Len() != 3 {
t.Fatalf("Len = %d, want the capacity 3", b.Len())
}
got := b.Recent(0)
want := []string{"e", "d", "c"}
for i, w := range want {
if got[i].Title != w {
t.Errorf("Recent()[%d] = %q, want %q", i, got[i].Title, w)
}
}
}
func TestBusDropsInvalid(t *testing.T) {
b := NewBus(4)
b.Publish(Event{Kind: KindFact, Title: "no source"}, testNow)
b.Publish(Event{Source: "s", Kind: KindFact}, testNow)
if b.Len() != 0 {
t.Errorf("Len = %d, want 0 — an envelope with no provenance must not be kept", b.Len())
}
}
func TestBusSubscriberPanicDoesNotBreakIntake(t *testing.T) {
b := NewBus(4)
var seen int
b.Subscribe(func(Event) { panic("observer is broken") })
b.Subscribe(func(Event) { seen++ })
b.Publish(Event{Source: "s", Kind: KindFact, Title: "t"}, testNow)
if seen != 1 {
t.Errorf("healthy subscriber called %d times, want 1", seen)
}
if b.Len() != 1 {
t.Errorf("event not recorded despite a panicking subscriber")
}
}
func TestBusConcurrentPublish(t *testing.T) {
b := NewBus(256)
var wg sync.WaitGroup
for i := 0; i < 16; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for j := 0; j < 10; j++ {
b.Publish(Event{Source: "s", Kind: KindFact, Title: "t"}, testNow)
}
}()
}
wg.Wait()
if b.Len() != 160 {
t.Errorf("Len = %d, want 160", b.Len())
}
}
+36 -6
View File
@@ -660,6 +660,31 @@ type CoreAPI interface {
// (router → dialogue → action → replier) and returns the reply text.
// No audio or stt/tts — for text channels (mavweb, telegram).
Chat(ctx context.Context, text string) (string, error)
// RecentEvents returns the daemon's unified intake journal, newest first
// (Vikunja #283) — one envelope per thing that arrived, whatever direction
// it came from: a relayed notification, a mail candidate, a feed item, a
// changed page, a spend, a presence probe.
//
// Read-only and daemon-cached, the same shape as TickTrace and DayPlan:
// the store adapter returns an error, because the journal is a bounded
// in-memory ring and not a table. Its contents are a window over intake,
// never the durable record — that is still the fact, note or task the
// intake path wrote.
RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error)
}
// IntakeEvent — one entry of the unified intake journal on the wire. Mirrors
// event.Event field for field; the ipc package does not import internal/event
// so the wire shape stays independent of the in-process type.
type IntakeEvent struct {
Source string `json:"source"`
Kind string `json:"kind"`
EntityIDs []string `json:"entity_ids,omitempty"`
Title string `json:"title"`
Body string `json:"body,omitempty"`
Priority string `json:"priority"`
OccurredAt time.Time `json:"occurred_at"`
}
// --- Rule trace / explanation DTOs ---
@@ -729,17 +754,22 @@ type DayPlan struct {
Spoken string `json:"spoken"`
}
// storeEncryptionKeyReq — passkey credential public key for wrapping the store
// storeEncryptionKeyReq — the passkey-derived secret used to wrap the store
// encryption key at enrollment time. Called by mavweb after RegisterFinish.
//
// Secret is the 32-byte WebAuthn PRF output, NOT the credential public key.
// The field used to carry the public key and that was the bug: a public key
// sits in passkeys.json next to the wrapped blob, so the blob protected
// nothing. See internal/webauthn/keywrap.go.
type storeEncryptionKeyReq struct {
PublicKey []byte `json:"public_key"`
Secret []byte `json:"secret"`
}
// unlockReq — passkey credential public key for unwrapping the store
// encryption key at cold-start. mavend reads the wrapped blob from its own
// configured path; the public key is the other half needed for unwrapping.
// unlockReq — the passkey-derived secret for unwrapping the store encryption
// key at cold-start. mavend reads the wrapped blob from its own configured
// path; this is the other half. Same PRF-output contract as above.
type unlockReq struct {
PublicKey []byte `json:"public_key"`
Secret []byte `json:"secret"`
}
// ErrToolNotFound — no tool row with this name (re-exported store sentinel for
+17 -4
View File
@@ -74,6 +74,7 @@ var readOnlyMethods = map[Method]bool{
MethodMorningStatus: true,
MethodMCPServers: true,
MethodDayPlan: true,
MethodRecentEvents: true,
}
// Dial connects to a core socket at path and returns a Client. The module
@@ -393,12 +394,16 @@ func (c *Client) AssertStepUp(ctx context.Context) error {
return c.call(ctx, MethodAssertStepUp, nil, nil)
}
func (c *Client) StoreEncryptionKey(ctx context.Context, publicKey []byte) error {
return c.call(ctx, MethodStoreEncryptionKey, storeEncryptionKeyReq{PublicKey: publicKey}, nil)
// StoreEncryptionKey wraps the daemon's at-rest key under secret, the 32-byte
// WebAuthn PRF output for the freshly enrolled credential.
func (c *Client) StoreEncryptionKey(ctx context.Context, secret []byte) error {
return c.call(ctx, MethodStoreEncryptionKey, storeEncryptionKeyReq{Secret: secret}, nil)
}
func (c *Client) Unlock(ctx context.Context, publicKey []byte) error {
return c.call(ctx, MethodUnlock, unlockReq{PublicKey: publicKey}, nil)
// Unlock hands the daemon the PRF secret so it can unwrap its at-rest key and
// open the store. Refused unless a passkey assertion was verified first.
func (c *Client) Unlock(ctx context.Context, secret []byte) error {
return c.call(ctx, MethodUnlock, unlockReq{Secret: secret}, nil)
}
func (c *Client) LookupTool(ctx context.Context, name string) (Tool, error) {
@@ -589,6 +594,14 @@ func (c *Client) TickTrace(ctx context.Context) (TickTrace, error) {
return t, nil
}
func (c *Client) RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error) {
var e []IntakeEvent
if err := c.call(ctx, MethodRecentEvents, nReq{N: n}, &e); err != nil {
return nil, err
}
return e, nil
}
func (c *Client) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
var s []MCPServerStatus
if err := c.call(ctx, MethodMCPServers, nil, &s); err != nil {
+27 -11
View File
@@ -211,6 +211,12 @@ func (a *storeAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, e
return nil, errors.New("store: morning status not available via direct store API")
}
// RecentEvents — same shape as TickTrace: the intake journal is a bounded ring
// in the daemon's memory, not a table, so a bare store cannot serve it.
func (a *storeAPI) RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error) {
return nil, errors.New("store: intake events not available via direct store API")
}
func (a *storeAPI) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
return nil, nil // no manager behind a bare store: nothing configured
}
@@ -420,8 +426,8 @@ type Server struct {
// MethodAssertStepUp returns ErrUnknownMethod (same as pre-stepup floor).
StepUp StepUpFunc
// WrapKeyFn — wraps the in-memory store encryption key with a passkey
// credential public key (HKDF-AESGCM) and writes the wrapped blob to disk.
// WrapKeyFn — wraps the in-memory store encryption key under the passkey
// PRF secret (HKDF-AESGCM) and writes the wrapped blob to disk.
// Set by the daemon; nil ⇒ MethodStoreEncryptionKey returns ErrUnknownMethod.
WrapKeyFn WrapKeyFunc
@@ -481,7 +487,7 @@ type Server struct {
ForgetSpeakerFn ForgetSpeakerFunc
// UnlockFn — unwraps the store encryption key from the wrapped blob using
// the passkey credential public key, opens the encrypted store, and wires
// the passkey PRF secret, opens the encrypted store, and wires
// the rest of the daemon (voice, loop, delivery). Set by the daemon when
// in locked mode; nil ⇒ MethodUnlock returns ErrUnknownMethod.
UnlockFn UnlockFunc
@@ -490,13 +496,13 @@ type Server struct {
// absolute ts supplied by callers, so this isn't load-bearing for live ops.
}
// WrapKeyFunc — wraps the store encryption key with the given credential
// public key and persists the wrapped blob.
type WrapKeyFunc func(ctx context.Context, publicKey []byte) error
// WrapKeyFunc — wraps the store encryption key under the passkey-derived
// secret (a 32-byte WebAuthn PRF output) and persists the wrapped blob.
type WrapKeyFunc func(ctx context.Context, secret []byte) error
// UnlockFunc — unwraps the store encryption key using the given credential
// public key and completes daemon initialization.
type UnlockFunc func(ctx context.Context, publicKey []byte) error
// UnlockFunc — unwraps the store encryption key using the passkey-derived
// secret and completes daemon initialization.
type UnlockFunc func(ctx context.Context, secret []byte) error
// SwapModelFunc — loads another resident model in place of the live one.
type SwapModelFunc func(ctx context.Context, req SwapModelReq) (SwapModelResp, error)
@@ -884,6 +890,16 @@ var methodTable = map[Method]handlerFunc{
MethodMorningStatus: withoutParams(func(ctx context.Context, api CoreAPI) ([]MorningRoutineStatus, error) {
return api.MorningStatus(ctx)
}),
MethodRecentEvents: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]IntakeEvent, error) {
out, err := api.RecentEvents(ctx, p.N)
if err != nil {
return nil, err
}
if out == nil {
out = []IntakeEvent{}
}
return out, nil
}),
MethodMCPServers: withoutParams(func(ctx context.Context, api CoreAPI) ([]MCPServerStatus, error) {
out, err := api.MCPServers(ctx)
if err != nil {
@@ -929,7 +945,7 @@ func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, er
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
return marshalResult(nil), s.WrapKeyFn(ctx, p.PublicKey)
return marshalResult(nil), s.WrapKeyFn(ctx, p.Secret)
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
@@ -939,7 +955,7 @@ func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, er
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
return marshalResult(nil), s.UnlockFn(ctx, p.PublicKey)
return marshalResult(nil), s.UnlockFn(ctx, p.Secret)
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
+3
View File
@@ -122,6 +122,9 @@ func (UnimplementedCoreAPI) TickTrace(ctx context.Context) (TickTrace, error) {
func (UnimplementedCoreAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error) {
return nil, ErrNotImplemented
}
func (UnimplementedCoreAPI) RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error) {
return nil, ErrNotImplemented
}
func (UnimplementedCoreAPI) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
return nil, ErrNotImplemented
}
+124
View File
@@ -0,0 +1,124 @@
package ipc
import (
"bytes"
"context"
"encoding/json"
"errors"
"testing"
)
// The wire must carry the PRF secret, not the credential public key. This is
// the field rename that fixes Vikunja #14: a v1 deployment sent "public_key",
// and the value it sent was in passkeys.json next to the wrapped blob.
func TestUnlockWireCarriesSecret(t *testing.T) {
secret := bytes.Repeat([]byte{7}, 32)
for _, p := range []any{unlockReq{Secret: secret}, storeEncryptionKeyReq{Secret: secret}} {
b, err := json.Marshal(p)
if err != nil {
t.Fatalf("marshal %T: %v", p, err)
}
var m map[string]any
if err := json.Unmarshal(b, &m); err != nil {
t.Fatalf("unmarshal %T: %v", p, err)
}
if _, ok := m["secret"]; !ok {
t.Errorf("%T has no \"secret\" field: %s", p, b)
}
if _, ok := m["public_key"]; ok {
t.Errorf("%T still sends \"public_key\": %s", p, b)
}
}
}
// The secret must reach the daemon hook byte-for-byte through the socket.
func TestUnlockDeliversSecretToHook(t *testing.T) {
_, srv, cli, _ := newServerWithStore(t)
secret := make([]byte, 32)
for i := range secret {
secret[i] = byte(i + 1)
}
var gotUnlock, gotWrap []byte
srv.UnlockFn = func(_ context.Context, s []byte) error { gotUnlock = bytes.Clone(s); return nil }
srv.WrapKeyFn = func(_ context.Context, s []byte) error { gotWrap = bytes.Clone(s); return nil }
ctx := context.Background()
if err := cli.Unlock(ctx, secret); err != nil {
t.Fatalf("Unlock: %v", err)
}
if !bytes.Equal(gotUnlock, secret) {
t.Errorf("UnlockFn got %x, want %x", gotUnlock, secret)
}
if err := cli.StoreEncryptionKey(ctx, secret); err != nil {
t.Fatalf("StoreEncryptionKey: %v", err)
}
if !bytes.Equal(gotWrap, secret) {
t.Errorf("WrapKeyFn got %x, want %x", gotWrap, secret)
}
}
// A refusal from the daemon hook — a wrong passkey, or no prior assertion —
// must surface to the caller as an error, never be swallowed into success.
func TestUnlockPropagatesRefusal(t *testing.T) {
_, srv, cli, _ := newServerWithStore(t)
srv.UnlockFn = func(context.Context, []byte) error {
return errors.New("unlock: no verified passkey assertion (assert first)")
}
if err := cli.Unlock(context.Background(), bytes.Repeat([]byte{9}, 32)); err == nil {
t.Fatal("a refused unlock reported success")
}
}
// Without the hooks wired — the normal, unencrypted deployment — both methods
// answer ErrUnknownMethod rather than pretending to have done something.
func TestUnlockUnwiredIsUnknownMethod(t *testing.T) {
_, _, cli, _ := newServerWithStore(t)
ctx := context.Background()
if err := cli.Unlock(ctx, bytes.Repeat([]byte{1}, 32)); err == nil {
t.Error("Unlock succeeded with no UnlockFn wired")
}
if err := cli.StoreEncryptionKey(ctx, bytes.Repeat([]byte{1}, 32)); err == nil {
t.Error("StoreEncryptionKey succeeded with no WrapKeyFn wired")
}
}
// Locked mode: Server.Check is the whole authorization surface, and it must
// default-deny everything except the two methods the unlock flow needs.
func TestLockedCheckDefaultDenies(t *testing.T) {
_, srv, cli, _ := newServerWithStore(t)
locked := errors.New("locked")
srv.Check = func(_ context.Context, m Method, _ json.RawMessage) error {
switch m {
case MethodAssertStepUp, MethodUnlock:
return nil
default:
return locked
}
}
unlocked := false
srv.UnlockFn = func(context.Context, []byte) error { unlocked = true; return nil }
srv.StepUp = func(context.Context) error { return nil }
srv.WrapKeyFn = func(context.Context, []byte) error { return nil }
ctx := context.Background()
// A store method must be refused while locked.
if _, err := cli.RecentNotes(ctx, 5); err == nil {
t.Error("a store read went through while locked")
}
// Key wrapping is NOT on the allowlist: a locked daemon has no key to wrap.
if err := cli.StoreEncryptionKey(ctx, bytes.Repeat([]byte{2}, 32)); err == nil {
t.Error("StoreEncryptionKey was allowed while locked")
}
// The unlock flow itself must still work.
if err := cli.AssertStepUp(ctx); err != nil {
t.Errorf("AssertStepUp refused while locked: %v", err)
}
if err := cli.Unlock(ctx, bytes.Repeat([]byte{3}, 32)); err != nil {
t.Errorf("Unlock refused while locked: %v", err)
}
if !unlocked {
t.Error("UnlockFn never ran")
}
}
+1
View File
@@ -62,6 +62,7 @@ const (
MethodEnrollSpeaker Method = "enroll_speaker"
MethodListSpeakers Method = "list_speakers"
MethodForgetSpeaker Method = "forget_speaker"
MethodRecentEvents Method = "recent_events"
)
// Request — one frame from module to core. Params is the JSON-encoded argument
+159 -100
View File
@@ -1,24 +1,48 @@
// Key wrapping for cold-start unlock.
// Key wrapping for cold-start unlock (Vikunja #14).
//
// The at-rest AES-256 key is wrapped with a key derived from the passkey
// credential public key (stable across assertions) via HKDF-SHA256, then
// AES-256-GCM. The wrapped blob is stored on disk; at cold-start the passkey
// assertion provides the credential public key to unwrap it.
// The at-rest AES-256 key is never on disk in the clear. It is wrapped with a
// key derived from a secret only the authenticator can produce, so a cold boot
// needs the physical passkey and nothing else opens the store.
//
// The passkey credential is a P-256 ECDSA public key. Its raw uncompressed
// bytes (65 bytes, 0x04 || X || Y) are the HKDF input — high-entropy, stable.
// # What the secret must be
//
// Blob format: salt (16) || nonce (12) || AES-256-GCM ciphertext.
// No file magic — the caller (mavend) owns the file path.
// The WebAuthn PRF extension. On assertion, the authenticator evaluates a
// keyed pseudo-random function over a fixed salt and hands back 32 bytes that
// are stable for the credential, unpredictable to everyone else, and never
// leave the device except as that output. That is the only thing in WebAuthn
// that yields a *secret* rather than a signature, and it is what makes the
// wrapped blob worth wrapping.
//
// # What it must NOT be, and used to be
//
// v1 of this file derived the wrapping key from the credential *public* key,
// on the reasoning that it is high-entropy and stable across assertions. Both
// are true and neither matters: a public key is public. mavweb writes it
// verbatim to passkeys.json, normally in the same state dir as the wrapped
// blob, so anyone holding both files recovered the database key offline with
// no authenticator involved. A v1 blob is a plaintext key with extra steps.
//
// v1 blobs are still readable, so an existing deployment opens and can be
// re-wrapped, and UnwrapKey reports which format it read so the caller can
// say so out loud. Nothing writes v1 any more.
//
// # Blob format
//
// v2: "MVNKW2\x00" (7) || salt (16) || nonce (12) || AES-256-GCM ciphertext
// v1: salt (16) || nonce (12) || AES-256-GCM ciphertext (legacy, read-only)
//
// The magic doubles as the version discriminator: v1 had none, so anything
// that does not start with it is v1 by elimination. A random 16-byte v1 salt
// colliding with the magic is a 2^-56 event, and the GCM tag catches it.
package webauthn
import (
"crypto/aes"
"crypto/cipher"
"crypto/hmac"
"crypto/hkdf"
"crypto/rand"
"crypto/sha256"
"encoding/binary"
"crypto/subtle"
"errors"
"fmt"
"io"
@@ -31,143 +55,178 @@ const (
nonceLen = 12
// keyLen — AES-256 key length.
keyLen = 32
// wrapInfo — HKDF info string for domain separation.
wrapInfo = "maven-passkey-keywrap-v1"
// secretLen — required length of the PRF output used as key material.
// WebAuthn PRF results are 32 bytes. Requiring exactly that is not
// pedantry: it is the structural guard that stops a COSE credential
// public key (77+ bytes) being passed here again by accident.
secretLen = 32
// wrapInfoV2 — HKDF info string. Carries the version so a v1 and a v2
// derivation can never collide even given the same input.
wrapInfoV2 = "maven-passkey-keywrap-v2"
// wrapInfoV1 — the legacy info string, kept only to read old blobs.
wrapInfoV1 = "maven-passkey-keywrap-v1"
)
// blobMagicV2 prefixes every v2 blob.
var blobMagicV2 = []byte("MVNKW2\x00")
var (
ErrKeyWrap = errors.New("webauthn: key wrap failed")
ErrKeyUnwrap = errors.New("webauthn: key unwrap failed (wrong credential?)")
ErrBlobTooLong = errors.New("webauthn: wrapped blob too long")
// ErrSecretLen is returned when the caller passes something that is not a
// 32-byte PRF output — most likely a credential public key.
ErrSecretLen = errors.New("webauthn: wrapping secret must be a 32-byte PRF output")
)
// WrapKey derives a wrapping key from credPublicKey via HKDF-SHA256 and
// AES-GCM-wraps plaintextKey. Returns the blob: salt || nonce || ciphertext.
// plaintextKey must be exactly 32 bytes (AES-256).
func WrapKey(plaintextKey, credPublicKey []byte) ([]byte, error) {
// BlobVersion identifies which format a blob was read as.
type BlobVersion int
const (
// BlobV1 is the legacy public-key-derived format. Readable, never written.
BlobV1 BlobVersion = 1
// BlobV2 is the PRF-derived format.
BlobV2 BlobVersion = 2
)
func (v BlobVersion) String() string {
switch v {
case BlobV1:
return "v1 (legacy, public-key derived — NOT SECRET)"
case BlobV2:
return "v2 (PRF derived)"
}
return "unknown"
}
// maxBlobLen — sanity limit; a real blob is 67 bytes.
const maxBlobLen = 1 << 20
// WrapKey wraps plaintextKey (32 bytes, AES-256) under a key derived from
// secret via HKDF-SHA256, and returns a v2 blob.
//
// secret must be the 32-byte WebAuthn PRF output for the enrolled credential.
// Anything else is refused — see the file header for why passing a credential
// public key here is the bug this replaces.
func WrapKey(plaintextKey, secret []byte) ([]byte, error) {
if len(plaintextKey) != keyLen {
return nil, fmt.Errorf("%w: plaintext key must be %d bytes", ErrKeyWrap, keyLen)
}
if len(credPublicKey) == 0 {
return nil, fmt.Errorf("%w: empty credential public key", ErrKeyWrap)
if err := checkSecret(secret); err != nil {
return nil, fmt.Errorf("%w: %v", ErrKeyWrap, err)
}
salt := make([]byte, saltLen)
if _, err := io.ReadFull(rand.Reader, salt); err != nil {
return nil, fmt.Errorf("%w: salt: %v", ErrKeyWrap, err)
}
wrapKey := hkdfSHA256(credPublicKey, salt, []byte(wrapInfo), keyLen)
nonce := make([]byte, nonceLen)
if _, err := io.ReadFull(rand.Reader, nonce); err != nil {
return nil, fmt.Errorf("%w: nonce: %v", ErrKeyWrap, err)
}
block, err := aes.NewCipher(wrapKey)
gcm, err := gcmFor(secret, salt, wrapInfoV2)
if err != nil {
return nil, fmt.Errorf("%w: aes: %v", ErrKeyWrap, err)
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, fmt.Errorf("%w: gcm: %v", ErrKeyWrap, err)
return nil, fmt.Errorf("%w: %v", ErrKeyWrap, err)
}
// Seal appends ciphertext+tag to nonce (which becomes nonce||ct).
ct := gcm.Seal(nil, nonce, plaintextKey, nil)
// The magic is authenticated as additional data, so a v2 blob cannot be
// stripped of its header and re-read as a v1 blob.
ct := gcm.Seal(nil, nonce, plaintextKey, blobMagicV2)
out := make([]byte, 0, saltLen+nonceLen+len(ct))
out := make([]byte, 0, len(blobMagicV2)+saltLen+nonceLen+len(ct))
out = append(out, blobMagicV2...)
out = append(out, salt...)
out = append(out, nonce...)
out = append(out, ct...)
return out, nil
}
// UnwrapKey extracts the salt from blob, re-derives the wrapping key from
// credPublicKey, and AES-GCM-unwraps. Returns the plaintext 32-byte AES key.
func UnwrapKey(blob, credPublicKey []byte) ([]byte, error) {
if len(blob) < saltLen+nonceLen+1 {
return nil, fmt.Errorf("%w: blob too short (%d)", ErrKeyUnwrap, len(blob))
// UnwrapKey recovers the plaintext AES-256 key from blob.
//
// It reads both formats and reports which one it got, so the caller can warn
// that a v1 blob offers no real protection. For a v2 blob, secret must be the
// 32-byte PRF output; for a v1 blob it is the credential public key, whatever
// length that happens to be.
func UnwrapKey(blob, secret []byte) ([]byte, BlobVersion, error) {
if len(blob) > maxBlobLen {
return nil, 0, ErrBlobTooLong
}
if len(blob) > 1<<20 { // 1MB sanity limit
return nil, ErrBlobTooLong
}
if len(credPublicKey) == 0 {
return nil, fmt.Errorf("%w: empty credential public key", ErrKeyUnwrap)
if len(secret) == 0 {
return nil, 0, fmt.Errorf("%w: empty secret", ErrKeyUnwrap)
}
salt := blob[:saltLen]
nonce := blob[saltLen : saltLen+nonceLen]
ct := blob[saltLen+nonceLen:]
if len(blob) >= len(blobMagicV2) && subtle.ConstantTimeCompare(blob[:len(blobMagicV2)], blobMagicV2) == 1 {
key, err := unwrap(blob[len(blobMagicV2):], secret, wrapInfoV2, blobMagicV2, secretLen)
return key, BlobV2, err
}
key, err := unwrap(blob, secret, wrapInfoV1, nil, 0)
return key, BlobV1, err
}
wrapKey := hkdfSHA256(credPublicKey, salt, []byte(wrapInfo), keyLen)
// unwrap does the shared salt||nonce||ct work. wantSecretLen of 0 means any
// non-empty secret is accepted (the v1 case, where it is a public key).
func unwrap(body, secret []byte, info string, aad []byte, wantSecretLen int) ([]byte, error) {
if len(body) < saltLen+nonceLen+1 {
return nil, fmt.Errorf("%w: blob too short (%d)", ErrKeyUnwrap, len(body))
}
if wantSecretLen > 0 && len(secret) != wantSecretLen {
return nil, fmt.Errorf("%w: %v", ErrKeyUnwrap, ErrSecretLen)
}
block, err := aes.NewCipher(wrapKey)
salt := body[:saltLen]
nonce := body[saltLen : saltLen+nonceLen]
ct := body[saltLen+nonceLen:]
gcm, err := gcmFor(secret, salt, info)
if err != nil {
return nil, fmt.Errorf("%w: aes: %v", ErrKeyUnwrap, err)
return nil, fmt.Errorf("%w: %v", ErrKeyUnwrap, err)
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, fmt.Errorf("%w: gcm: %v", ErrKeyUnwrap, err)
}
plain, err := gcm.Open(nil, nonce, ct, nil)
plain, err := gcm.Open(nil, nonce, ct, aad)
if err != nil {
return nil, fmt.Errorf("%w: decrypt failed (wrong credential?)", ErrKeyUnwrap)
}
if len(plain) != keyLen {
return nil, fmt.Errorf("%w: unwrapped key is %d bytes, want %d", ErrKeyUnwrap, len(plain), keyLen)
}
return plain, nil
}
// hkdfSHA256 implements HKDF-SHA256 (RFC 5869) using only stdlib.
// gcmFor derives the wrapping key with HKDF-SHA256 and returns a GCM AEAD.
//
// Input:
// - secret: the input key material (credential public key bytes)
// - salt: random salt (16 bytes)
// - info: optional context string for domain separation
// - length: desired output length in bytes
//
// Output: length bytes of derived key material.
//
// HKDF is extract-then-expand. We use HMAC-SHA256 for both steps. This avoids
// importing golang.org/x/crypto/hkdf — a ~30-line function vs a new dep. The
// tradeoff is no constant-time guarantees on the extract step beyond HMAC's;
// acceptable here because the input is already high-entropy key material (a
// P-256 public key), not a low-entropy passphrase.
func hkdfSHA256(secret, salt, info []byte, length int) []byte {
// Step 1: Extract — PRK = HMAC-SHA256(salt, secret)
// If salt is nil/empty, use a zero-filled block (RFC 5869 §2.2).
if salt == nil {
salt = make([]byte, sha256.Size)
// This uses the standard library's crypto/hkdf rather than the hand-rolled
// HKDF this file used to carry. That implementation keyed the expand step with
// the salt instead of the PRK — self-consistent, so wrap and unwrap agreed,
// but not RFC 5869 and not the domain separation it claimed to provide.
func gcmFor(secret, salt []byte, info string) (cipher.AEAD, error) {
wrapKey, err := hkdf.Key(sha256.New, secret, salt, info, keyLen)
if err != nil {
return nil, fmt.Errorf("hkdf: %v", err)
}
mac := hmac.New(sha256.New, salt)
mac.Write(secret)
prk := mac.Sum(nil)
// Step 2: Expand — produce length bytes via T(i) = HMAC-SHA256(PRK, T(i-1) || info || i)
// Where T(0) = empty, i is a byte counter starting at 1.
out := make([]byte, 0, length)
block := make([]byte, 0, sha256.Size+len(info)+1)
var t []byte // T(i-1)
for counter := byte(1); len(out) < length; counter++ {
block = block[:0]
block = append(block, t...)
block = append(block, info...)
block = append(block, counter)
mac.Reset()
mac.Write(block)
t = mac.Sum(prk[:0]) // reuse prk buffer — mac.Sum appends to its arg
// t now starts with prk[:0] (empty) followed by the HMAC result.
// Since we need just the HMAC result (sha256.Size bytes), re-slice.
t = t[len(t)-sha256.Size:]
out = append(out, t...)
block, err := aes.NewCipher(wrapKey)
if err != nil {
return nil, fmt.Errorf("aes: %v", err)
}
return out[:length]
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, fmt.Errorf("gcm: %v", err)
}
return gcm, nil
}
// encodeUint32 — big-endian uint32 for the blob format header, if needed.
func encodeUint32(v uint32) []byte {
var b [4]byte
binary.BigEndian.PutUint32(b[:], v)
return b[:]
func checkSecret(secret []byte) error {
if len(secret) != secretLen {
return fmt.Errorf("%w (got %d bytes)", ErrSecretLen, len(secret))
}
// An all-zero PRF result means the authenticator returned nothing useful;
// wrapping under it would produce a blob anyone can open.
var acc byte
for _, b := range secret {
acc |= b
}
if acc == 0 {
return fmt.Errorf("%w (all zero)", ErrSecretLen)
}
return nil
}
+231
View File
@@ -0,0 +1,231 @@
package webauthn
import (
"bytes"
"crypto/rand"
"errors"
"io"
"testing"
)
func testSecret(t *testing.T) []byte {
t.Helper()
s := make([]byte, secretLen)
if _, err := io.ReadFull(rand.Reader, s); err != nil {
t.Fatalf("rand: %v", err)
}
s[0] |= 1 // never all-zero
return s
}
func testKey(t *testing.T) []byte {
t.Helper()
k := make([]byte, keyLen)
if _, err := io.ReadFull(rand.Reader, k); err != nil {
t.Fatalf("rand: %v", err)
}
return k
}
func TestWrapUnwrapRoundTrip(t *testing.T) {
key, secret := testKey(t), testSecret(t)
blob, err := WrapKey(key, secret)
if err != nil {
t.Fatalf("WrapKey: %v", err)
}
if !bytes.HasPrefix(blob, blobMagicV2) {
t.Fatalf("blob does not start with the v2 magic: %x", blob[:8])
}
// The plaintext key must not be recoverable by reading the file.
if bytes.Contains(blob, key) {
t.Fatal("the wrapped blob contains the plaintext key verbatim")
}
got, version, err := UnwrapKey(blob, secret)
if err != nil {
t.Fatalf("UnwrapKey: %v", err)
}
if version != BlobV2 {
t.Errorf("version = %v, want v2", version)
}
if !bytes.Equal(got, key) {
t.Errorf("unwrapped key differs from the wrapped one")
}
}
// Fresh salt and nonce per wrap: two blobs of the same key under the same
// secret must not be byte-identical, or the file leaks that nothing changed.
func TestWrapKeyIsNotDeterministic(t *testing.T) {
key, secret := testKey(t), testSecret(t)
a, err := WrapKey(key, secret)
if err != nil {
t.Fatalf("WrapKey: %v", err)
}
b, err := WrapKey(key, secret)
if err != nil {
t.Fatalf("WrapKey: %v", err)
}
if bytes.Equal(a, b) {
t.Fatal("two wraps of the same key produced identical blobs")
}
}
// The failure mode that matters most: a wrong passkey must not unlock.
func TestUnwrapWithWrongSecretFails(t *testing.T) {
key := testKey(t)
blob, err := WrapKey(key, testSecret(t))
if err != nil {
t.Fatalf("WrapKey: %v", err)
}
got, _, err := UnwrapKey(blob, testSecret(t))
if err == nil {
t.Fatal("a different secret unwrapped the blob")
}
if !errors.Is(err, ErrKeyUnwrap) {
t.Errorf("err = %v, want ErrKeyUnwrap", err)
}
if got != nil {
t.Error("key material returned alongside an error")
}
}
// One flipped bit anywhere must fail the GCM tag, including in the salt and
// nonce — those are not authenticated by the tag but they change the
// derivation, so the tag fails anyway.
func TestUnwrapRejectsTamperedBlob(t *testing.T) {
key, secret := testKey(t), testSecret(t)
blob, err := WrapKey(key, secret)
if err != nil {
t.Fatalf("WrapKey: %v", err)
}
for i := range blob {
bad := bytes.Clone(blob)
bad[i] ^= 0x01
if _, _, err := UnwrapKey(bad, secret); err == nil {
t.Fatalf("byte %d of %d could be flipped and the blob still opened", i, len(blob))
}
}
}
func TestUnwrapRejectsTruncatedBlob(t *testing.T) {
key, secret := testKey(t), testSecret(t)
blob, err := WrapKey(key, secret)
if err != nil {
t.Fatalf("WrapKey: %v", err)
}
for _, n := range []int{0, 1, len(blobMagicV2), len(blobMagicV2) + saltLen, len(blob) - 1} {
if _, _, err := UnwrapKey(blob[:n], secret); err == nil {
t.Errorf("a %d-byte blob unwrapped", n)
}
}
}
// A v2 blob must not be downgradeable to v1 by stripping its header: the magic
// is GCM additional data, so the tag fails once it is gone.
func TestV2BlobCannotBeStrippedToV1(t *testing.T) {
key, secret := testKey(t), testSecret(t)
blob, err := WrapKey(key, secret)
if err != nil {
t.Fatalf("WrapKey: %v", err)
}
if _, _, err := UnwrapKey(blob[len(blobMagicV2):], secret); err == nil {
t.Fatal("a header-stripped v2 blob was accepted as v1")
}
}
// v1 blobs still open, and report themselves as v1 so the daemon can warn.
// wrapV1 reproduces the legacy writer this file no longer has.
func wrapV1(t *testing.T, key, secret []byte) []byte {
t.Helper()
salt := make([]byte, saltLen)
nonce := make([]byte, nonceLen)
if _, err := io.ReadFull(rand.Reader, salt); err != nil {
t.Fatalf("rand: %v", err)
}
if _, err := io.ReadFull(rand.Reader, nonce); err != nil {
t.Fatalf("rand: %v", err)
}
gcm, err := gcmFor(secret, salt, wrapInfoV1)
if err != nil {
t.Fatalf("gcmFor: %v", err)
}
out := append([]byte{}, salt...)
out = append(out, nonce...)
return append(out, gcm.Seal(nil, nonce, key, nil)...)
}
func TestUnwrapReadsLegacyV1(t *testing.T) {
key := testKey(t)
// v1 was keyed on the credential public key: not 32 bytes, and that is
// deliberately still accepted on the read path.
pub := make([]byte, 77)
if _, err := io.ReadFull(rand.Reader, pub); err != nil {
t.Fatalf("rand: %v", err)
}
blob := wrapV1(t, key, pub)
got, version, err := UnwrapKey(blob, pub)
if err != nil {
t.Fatalf("UnwrapKey(v1): %v", err)
}
if version != BlobV1 {
t.Errorf("version = %v, want v1", version)
}
if !bytes.Equal(got, key) {
t.Error("v1 round-trip lost the key")
}
if _, _, err := UnwrapKey(blob, pub[:76]); err == nil {
t.Error("a truncated public key opened the v1 blob")
}
}
// The structural guard against the bug this replaces: a COSE public key is not
// 32 bytes, so it can never be used to write a new blob.
func TestWrapKeyRefusesNonPRFSecret(t *testing.T) {
key := testKey(t)
cases := map[string][]byte{
"nil": nil,
"empty": {},
"short": make([]byte, 16),
"cose public key": make([]byte, 77),
"all-zero 32 byte": make([]byte, 32),
}
for name, secret := range cases {
t.Run(name, func(t *testing.T) {
if _, err := WrapKey(key, secret); err == nil {
t.Fatalf("WrapKey accepted a %s secret", name)
}
})
}
}
func TestWrapKeyRefusesWrongKeyLength(t *testing.T) {
secret := testSecret(t)
for _, n := range []int{0, 16, 31, 33, 64} {
if _, err := WrapKey(make([]byte, n), secret); err == nil {
t.Errorf("WrapKey accepted a %d-byte plaintext key", n)
}
}
}
// A v2 blob demands exactly 32 bytes on the read path too, so a caller cannot
// go back to passing a public key.
func TestUnwrapV2RefusesNonPRFSecret(t *testing.T) {
blob, err := WrapKey(testKey(t), testSecret(t))
if err != nil {
t.Fatalf("WrapKey: %v", err)
}
if _, _, err := UnwrapKey(blob, make([]byte, 77)); !errors.Is(err, ErrKeyUnwrap) {
t.Fatalf("err = %v, want ErrKeyUnwrap for a 77-byte secret", err)
}
if _, _, err := UnwrapKey(blob, nil); err == nil {
t.Fatal("an empty secret unwrapped a v2 blob")
}
}
func TestUnwrapRejectsOversizeBlob(t *testing.T) {
if _, _, err := UnwrapKey(make([]byte, maxBlobLen+1), testSecret(t)); !errors.Is(err, ErrBlobTooLong) {
t.Fatalf("err = %v, want ErrBlobTooLong", err)
}
}
+70
View File
@@ -0,0 +1,70 @@
package webauthn
import (
"crypto/sha256"
"encoding/base64"
"errors"
"fmt"
)
// The WebAuthn PRF extension is where cold-start unlock gets its secret
// (Vikunja #14). The authenticator evaluates a keyed PRF over a salt we
// choose and returns 32 bytes that are:
//
// - stable — the same credential and the same salt always give the same
// bytes, which is what lets a blob wrapped today be opened tomorrow;
// - secret — they never leave the authenticator except as this output, so
// unlike the credential public key they are not sitting in passkeys.json;
// - bound to user verification — the assertion that produces them required
// a gesture, so the bytes cannot be harvested silently.
//
// The salt is fixed and public. It is a domain separator, not a secret: it
// makes maven's PRF output different from any other relying party's use of
// the same credential.
// prfSaltInput — the string hashed into the 32-byte evaluation salt. Changing
// it invalidates every wrapped key file in existence, which is why it is a
// constant and not configuration.
const prfSaltInput = "maven-coldstart-unlock-v1"
// PRFSalt returns the fixed 32-byte PRF evaluation salt.
func PRFSalt() []byte {
sum := sha256.Sum256([]byte(prfSaltInput))
return sum[:]
}
// ErrNoPRF is returned when a browser reports no PRF result — either the
// authenticator does not implement the extension, or the platform stripped
// it. Cold-start unlock is unavailable for that credential, and the correct
// response is to say so rather than to fall back to something weaker.
var ErrNoPRF = errors.New("webauthn: authenticator returned no PRF result (cold-start unlock unavailable)")
// DecodePRFResult parses the base64url PRF output the browser read out of
// getClientExtensionResults().prf.results.first and checks it is usable as
// wrapping key material.
//
// The browser is not trusted to send something sensible: a short, empty, or
// all-zero result would silently produce a blob that anyone can open, so all
// three are refused here rather than at the crypto layer.
func DecodePRFResult(b64 string) ([]byte, error) {
if b64 == "" {
return nil, ErrNoPRF
}
secret, err := decodeB64Any(b64)
if err != nil {
return nil, fmt.Errorf("webauthn: prf result: %w", err)
}
if err := checkSecret(secret); err != nil {
return nil, err
}
return secret, nil
}
// decodeB64Any accepts padded or unpadded base64url — browsers differ, and
// the JS helper on the passkey page strips padding.
func decodeB64Any(s string) ([]byte, error) {
if b, err := base64.RawURLEncoding.DecodeString(s); err == nil {
return b, nil
}
return base64.URLEncoding.DecodeString(s)
}
+114
View File
@@ -0,0 +1,114 @@
package webauthn
import (
"bytes"
"encoding/base64"
"encoding/json"
"errors"
"testing"
)
// The salt is the identity of every wrapped key file ever written. If it
// changes, every deployment's blob becomes unopenable, so it is pinned here.
func TestPRFSaltIsStable(t *testing.T) {
salt := PRFSalt()
if len(salt) != 32 {
t.Fatalf("salt is %d bytes, want 32", len(salt))
}
if got := base64.RawURLEncoding.EncodeToString(salt); got != base64.RawURLEncoding.EncodeToString(PRFSalt()) {
t.Fatal("PRFSalt is not deterministic")
}
// Mutating the returned slice must not affect the next caller.
salt[0] ^= 0xff
if bytes.Equal(salt, PRFSalt()) {
t.Fatal("PRFSalt returned shared backing state")
}
}
func TestDecodePRFResult(t *testing.T) {
raw := make([]byte, 32)
for i := range raw {
raw[i] = byte(i + 1)
}
for _, enc := range []string{
base64.RawURLEncoding.EncodeToString(raw),
base64.URLEncoding.EncodeToString(raw),
} {
got, err := DecodePRFResult(enc)
if err != nil {
t.Fatalf("DecodePRFResult(%q): %v", enc, err)
}
if !bytes.Equal(got, raw) {
t.Errorf("decoded %x, want %x", got, raw)
}
}
}
// No PRF must be a distinguishable, named failure — never a silent fallback to
// some other secret.
func TestDecodePRFResultNoPRF(t *testing.T) {
if _, err := DecodePRFResult(""); !errors.Is(err, ErrNoPRF) {
t.Fatalf("err = %v, want ErrNoPRF", err)
}
}
func TestDecodePRFResultRejectsUnusable(t *testing.T) {
zeros := base64.RawURLEncoding.EncodeToString(make([]byte, 32))
short := base64.RawURLEncoding.EncodeToString(make([]byte, 16))
long := base64.RawURLEncoding.EncodeToString(make([]byte, 64))
for name, in := range map[string]string{
"not base64": "!!!!",
"all zero": zeros,
"too short": short,
"too long": long,
} {
t.Run(name, func(t *testing.T) {
if _, err := DecodePRFResult(in); err == nil {
t.Fatalf("accepted a %s PRF result", name)
}
})
}
}
// Both option builders must ask for PRF, or the browser never produces a
// secret and cold-start unlock silently never works.
func TestOptionsRequestPRF(t *testing.T) {
rp := NewRP(Config{Origin: "http://localhost:8080", RPID: "localhost", RPName: "maven"})
create, _, err := rp.CreationOptions([]byte("u"), "u")
if err != nil {
t.Fatalf("CreationOptions: %v", err)
}
if _, ok := extPRF(t, create)["prf"]; !ok {
t.Error("creation options do not request the prf extension")
}
assert, _, err := rp.AssertionOptions()
if err != nil {
t.Fatalf("AssertionOptions: %v", err)
}
prf, ok := extPRF(t, assert)["prf"].(map[string]any)
if !ok {
t.Fatal("assertion options do not request the prf extension")
}
eval, _ := prf["eval"].(map[string]any)
first, _ := eval["first"].(string)
if first != base64.RawURLEncoding.EncodeToString(PRFSalt()) {
t.Errorf("prf.eval.first = %q, want the fixed salt", first)
}
}
func extPRF(t *testing.T, opts any) map[string]any {
t.Helper()
b, err := json.Marshal(opts)
if err != nil {
t.Fatalf("marshal options: %v", err)
}
var m struct {
Extensions map[string]any `json:"extensions"`
}
if err := json.Unmarshal(b, &m); err != nil {
t.Fatalf("unmarshal options: %v", err)
}
return m.Extensions
}
+16
View File
@@ -124,6 +124,13 @@ func (rp *RP) CreationOptions(userID []byte, userName string) (map[string]any, s
"timeout": 60000,
"attestation": "none",
"excludeCredentials": []any{},
// PRF: ask the authenticator at enrollment time whether it can
// produce a per-credential secret. Nothing is wrapped here — the
// browser reports support back and mavweb decides whether cold-start
// unlock is available for this credential. See internal/webauthn/prf.go.
"extensions": map[string]any{
"prf": map[string]any{},
},
}, challengeB64, nil
}
@@ -193,6 +200,15 @@ func (rp *RP) AssertionOptions() (map[string]any, string, error) {
"rpId": rp.cfg.RPID,
"allowCredentials": []any{},
"userVerification": "required",
// PRF evaluation over the fixed cold-start salt. The 32 bytes that
// come back are the ONLY thing that can unwrap the database key.
"extensions": map[string]any{
"prf": map[string]any{
"eval": map[string]any{
"first": base64.RawURLEncoding.EncodeToString(PRFSalt()),
},
},
},
}, challengeB64, nil
}
+67
View File
@@ -0,0 +1,67 @@
#!/usr/bin/env bash
# gen-stt-fixtures.sh — regenerate the golden STT audio fixtures.
#
# The fixtures in cmd/mavsttd/testdata/*.wav are SYNTHESISED, not recorded.
# They come out of the same piper voices maven speaks with, so nothing of the
# owner's voice is committed and every fixture is reproducible from this
# script plus the voice model. They are also small: 16 kHz mono s16le, a
# couple of seconds each.
#
# Usage:
# scripts/gen-stt-fixtures.sh
#
# Voices are picked up from, in order, $PIPER_VOICE_RU / $PIPER_VOICE_EN, then
# the repo's models/tts, then ~/esp-server/voices. The English voice is not
# vendored; if it is missing the English fixture is skipped and the existing
# one is left alone.
set -euo pipefail
root="$(cd "$(dirname "$0")/.." && pwd)"
out="$root/cmd/mavsttd/testdata"
piper="${PIPER_BIN:-$root/deps/piper/piper}"
espeak="${PIPER_ESPEAK:-$root/deps/piper/espeak-ng-data}"
pick_voice() {
for c in "$@"; do
[ -f "$c" ] && { echo "$c"; return 0; }
done
return 1
}
ru="$(pick_voice "${PIPER_VOICE_RU:-}" "$root/models/tts/ru_RU-irina-medium.onnx" "$HOME/esp-server/voices/ru_RU-irina-medium.onnx")" || {
echo "no russian piper voice found" >&2
exit 1
}
en="$(pick_voice "${PIPER_VOICE_EN:-}" "$root/models/tts/en_US-lessac-medium.onnx" "$HOME/esp-server/voices/en_US-lessac-medium.onnx")" || en=""
# synth <voice> <out.wav> <text>
# piper emits raw 22050 Hz s16le on stdout; ffmpeg resamples to the canonical
# 16 kHz mono and writes a plain 44-byte-header WAV (-fflags bitexact keeps
# ffmpeg's encoder LIST chunk out, so the bytes are stable across ffmpeg
# builds and internal/audio.PCMFromWAV reads them without scanning).
synth() {
local voice="$1" dest="$2" text="$3"
printf '%s' "$text" | LD_LIBRARY_PATH="$(dirname "$piper")" "$piper" \
--model "$voice" --config "$voice.json" \
--espeak_data "$espeak" --output_raw --quiet |
ffmpeg -hide_banner -loglevel error -y \
-f s16le -ar 22050 -ac 1 -i - \
-af "adelay=200,apad=pad_dur=0.2" \
-ar 16000 -ac 1 -c:a pcm_s16le -fflags bitexact "$dest"
echo "wrote $dest ($(stat -c%s "$dest") bytes)"
}
synth "$ru" "$out/ru_reminder.wav" "Напомни мне через час позвонить маме."
synth "$ru" "$out/ru_fact.wav" "Отметь, что я выпил воды."
synth "$ru" "$out/ru_query.wav" "Что у меня сегодня по календарю?"
if [ -n "$en" ]; then
# Keep the English line free of words piper spells out letter by letter —
# "nginx" comes out of lessac as "engine X", which is a TTS artefact and
# would make the fixture assert on the wrong thing.
synth "$en" "$out/en_act.wav" "Restart the web server and check the disk space."
else
echo "no english piper voice found — skipping en_act.wav" >&2
fi
echo "fixtures regenerated; expected transcripts live in $out/golden_v1.json"