Compare commits

...

7 Commits

Author SHA1 Message Date
kami a8fcb404be Scan the LAN, bounded to configured subnets (#257)
internal/netscan/ discovers hosts on the network Maven is configured to look at:
a TCP-connect scan (net.DialTimeout, no raw sockets, no privileges) plus a read
of the kernel's ARP cache. Wired as a read-only query source, "network", so
"какие устройства в сети?" is answered by a scan instead of by whatever old note
happens to be nearest.

Scanning is a read, but an unbounded scanner on a home LAN is noisy and easy to
point somewhere it should not go, so the package is built around four bounds:

  - Scan takes NO target argument. The range comes from the config block and
    from nowhere else, so there is no exported way to scan an arbitrary prefix
    and nothing an utterance, the router, or a scanned host says can retarget
    it. That is asserted directly: the test watches every address handed to the
    dialer and fails if one falls outside the configured prefix. The ARP cache —
    the one input the network itself populates — is filtered to the configured
    range for the same reason.
  - Every configured CIDR must be private (RFC1918 / CGNAT / link-local) and no
    larger than 1024 addresses. 8.8.8.0/24, 0.0.0.0/0 and 10.0.0.0/8 are refused
    at config load, not after the packets have left.
  - Rate-limited to a configured connections-per-second across the whole scan,
    so it looks like background traffic rather than a portscan.
  - Bounded in total by MaxHosts, a per-connection timeout, a 20s turn budget
    and the context; a canceled scan stops dialing immediately.

Off unless configured: dark without "enabled": true, and applyDefaults
normalises a disabled block to nil. deploy/mavend.json carries it disabled.

BLUETOOTH IS NOT SHIPPED, AND IS BLOCKED, NOT SKIPPED. The plan's other half
(internal/bluetooth/, RSSI presence probes) needs a bluez stack that is not
here: bluetoothctl and hcitool are not installed, bluetoothd is not installed,
the bluetooth unit is inactive, and org.bluez is not on the system bus. hci0
exists as a kernel device and nothing can talk to it. The docker deploy is
further away still — it would need host networking, the D-Bus system socket
passed in, and CAP_NET_ADMIN. Writing an exec wrapper around a binary that does
not exist, against an output format nothing here can produce, would be a guess
dressed as a feature. It needs a decision about privileging the container before
any of it is worth writing.

Vikunja #257
2026-08-01 06:35:10 +04:00
kami dc4c5b7841 Read and control the house through Home Assistant (#256)
A `smarthome` block points Maven at a Home Assistant instance. She reads its
entity states to answer "что включено дома?", and every controllable device
becomes a PROPOSED row in the existing act allowlist — cmd
["smarthome",<entity_id>,<service>], scope smarthome:<domain> — so nothing new
had to be invented for the mutating half. ProposeTool/EnableTool/DisableTool,
tool.Matcher and the confirm turn are untouched; one branch in Executor.Exec
routes such a row to the client instead of exec, and "smarthome" is never run as
a binary. This is the same trick overnight/mcp-tools used for #251, on purpose.

Discovery only ever PROPOSES, and every control row is destructive=true: there
is no read-only way to turn the heating off, so flipping something in his flat
always costs a confirm turn and always had to be enabled by hand on /tools,
behind step-up.

The entity and the service come from the row he enabled, never from the
utterance — Exec drops the spoken tail for a house row. A router that misheard
can pick the wrong lamp; it cannot compose a target of its own. The service is
checked against the domain's table on the way out too, so a hand-edited cmd
column cannot reach an arbitrary Home Assistant service. set_brightness and
set_temperature are deliberately absent: a spoken number the router got wrong is
a wrong act on real hardware, and on/off is the whole of what a voice turn can
defend.

The read side is a query source ("home", before calendar and the recall passes)
so "что нового дома?" is not answered from an old note. Its matcher needs a
house marker plus an ask plus a device word and bails out on weather wording,
because "какая температура на улице?" belongs to the weather source.

Off unless configured: the block is dark without "enabled": true, and
applyDefaults normalises a disabled block to nil so "off" stays in one place.
deploy/mavend.json carries it disabled, with the token as ${HA_TOKEN}.

NOT shipped, and not faked: MQTT / Zigbee2MQTT (plan steps 2 and 5) and the
sensor-to-fact and presence-probe pipelines. There is no broker and no Home
Assistant anywhere on this network — 8123 and 1883 are closed on every host in
192.168.1.0/24 — the module tree is vendored so a paho dependency cannot be
added offline, and Home Assistant already fronts Zigbee2MQTT where it exists.
Writing a sensor pipeline with no sensor to test it against would be a guess.

Vikunja #256
2026-08-01 06:27:39 +04:00
kami 33e53ee897 Add a replayable full-system simulator on a fake clock (#284)
A scenario is a JSON file under cmd/mavend/testdata/scenarios: a start
instant, a script of canned model answers, and a list of steps at "HH:MM".
Each step does one thing — say, audio, signal, arrive, tick, fault — and
then asserts on what she said, what was sent, which ecosystem services were
called, and what landed in the intake journal.

Between those boundaries the real components run: the real router cascade
(stage0, the LLM router over a scripted completer, the classifier
underneath it), the real store, the real reactive handler, the real tick
loop, and the same intake-decorated ipc.CoreAPI the daemon wires. What is
faked is only what a test cannot have: the model, the microphone, the
speaker, the delivery sink, and the ecosystem HTTP services.

Time is a single fakeClock threaded into every reader — the handler, the
intake publish stamp and tick(ctx, now) — so there is no time.Now() on the
replay path and a scenario is reproducible. TestSimulatorIsDeterministic
enforces that by replaying twice and diffing the transcripts byte for byte;
advanceTo refuses a step that goes backwards.

Two scenarios ship. morning_missed replays #284's own description: he
appears at the desk, a feed item, a mail candidate and a relayed
notification arrive through the morning, two ticks pass, and the assertions
are as much about nothing being sent at him unprompted as about what she
said. evening_degraded picks up the tier-2 pipeline case #288 deferred
here — a golden WAV through the STT seam to a written fact — and then puts
the ecosystem into 503 and checks that the proactive loop stays quiet and
that intake keeps working without it.

This is test-only code. Nothing in the production binaries changed, so the
daemon behaves identically when no scenario is running.

`make simulate` runs them verbose so the transcript is readable; `make
test` runs them with everything else.

Vikunja #284
2026-08-01 06:15:21 +04:00
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
61 changed files with 6764 additions and 279 deletions
+20 -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: simulate stt-fixtures test-stt-golden all build build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav clean test fmt-check vet run-stt run-tts run-web download-embedder deps-go eval-router eval-recall eval-phrasing eval-models
all: build
@@ -91,6 +91,14 @@ vet:
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
$(GO) vet ./internal/... ./cmd/...
# simulate — replay every scripted day under cmd/mavend/testdata/scenarios
# through the real router, store, tick loop and intake journal, on a fake clock
# (Vikunja #284). Verbose so the transcript of each scenario lands in the
# terminal. Also runs as part of `make test`; this target is for reading it.
simulate:
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
$(GO) test -v -count=1 -run TestSimulator ./cmd/mavend/
test: fmt-check vet
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
$(GO) test -race -coverprofile=coverage.out ./internal/... ./cmd/...
@@ -139,6 +147,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)
+45
View File
@@ -74,6 +74,18 @@ var querySources = []querySource{
// it by inventing news. Its matcher needs a feed noun plus an ask, so
// "у меня новая лента в инстаграме" is untouched.
{"feeds", (*reactiveHandler).queryFeeds},
// Before "calendar" and before the recall sources: "что включено дома?" is
// a question about the house, and the notes pass would otherwise answer it
// from whatever he once said about the lights. Its matcher needs a house
// marker plus an ask plus a device word, and it bails out on weather
// wording, so "какая температура на улице?" still reaches the weather
// source.
{"home", (*reactiveHandler).queryHome},
// Next to "home" and for the same reason: "какие устройства в сети?" is a
// question about the LAN, and the recall pass would otherwise answer it
// from an old note about the router. Its matcher needs a network word plus
// an ask plus a device noun, so "интернет не работает" is untouched.
{"network", (*reactiveHandler).queryNetwork},
{"calendar", (*reactiveHandler).queryCalendar},
{"weather", (*reactiveHandler).queryWeather},
{"embed", (*reactiveHandler).queryEmbed},
@@ -275,6 +287,39 @@ func (h *reactiveHandler) queryCalendar(ctx context.Context, t *queryTurn) (stri
return f.FormatEntries(entries, date), true
}
// queryHome answers a question about the house. Read-only by construction: it
// calls States and nothing else, so there is no confirm turn here — the only
// way to CHANGE something is an enabled allowlist row through tool.Executor.
func (h *reactiveHandler) queryHome(ctx context.Context, t *queryTurn) (string, bool) {
if !isHomeQuery(t.dec.Utterance) {
return "", false
}
if h.home == nil {
// Claim the turn rather than fall through: "дом не подключён" is true,
// and letting general knowledge answer would be an invented house.
return "дом не подключён — я его не вижу.", true
}
ctxH, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
return h.home.homeSummary(ctxH)
}
// queryNetwork answers a question about the LAN with a bounded scan. There is
// no confirm turn because nothing is changed, and no way to widen the range
// because Scan takes no target — the utterance selects the question, never the
// subnet.
func (h *reactiveHandler) queryNetwork(ctx context.Context, t *queryTurn) (string, bool) {
if !isNetworkQuery(t.dec.Utterance) {
return "", false
}
if h.netscan == nil {
// Claim the turn: "сканирование не настроено" is true, and general
// knowledge would answer with an invented list of devices.
return "сканирование сети не настроено.", true
}
return h.netscan.scanSummary(ctx)
}
func (h *reactiveHandler) queryWeather(ctx context.Context, t *queryTurn) (string, bool) {
if !isWeatherQuery(t.dec.Utterance) {
return "", false
+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")
}
}
+104 -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,12 @@ func run(args []string) error {
go voiceW.mcp.run(ctx)
}
dl.unlock()
// Re-enumerate the house for new devices (nil unless configured).
if voiceW != nil && voiceW.home != nil {
go voiceW.home.run(ctx)
}
dl.unlock(st)
log.Printf("mavend: unlocked via passkey assertion")
return nil
}
@@ -609,6 +682,13 @@ func run(args []string) error {
voiceW.mcp.run(ctx)
}()
}
if voiceW != nil && voiceW.home != nil {
wg.Add(1)
go func() {
defer wg.Done()
voiceW.home.run(ctx)
}()
}
}
<-ctx.Done()
+142
View File
@@ -0,0 +1,142 @@
package main
import (
"context"
"fmt"
"log"
"strings"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/netscan"
)
// scanBudget — the whole spoken scan, end to end. A voice turn that takes
// longer than this has already failed as a turn, so the scan returns whatever
// it found rather than keeping him waiting.
const scanBudget = 20 * time.Second
// scanReadOut — how many hosts she names out loud. The rest are a count: a
// spoken list of twenty IP addresses is not an answer.
const scanReadOut = 6
// netWiring — the LAN scanner, when the `netscan` block is enabled. nil ⇒ Maven
// never puts a discovery packet on the network.
//
// Unlike the house, a scan is a READ, so it is a query source rather than an
// act: there is no allowlist row and no confirm turn, because nothing changes.
// What makes that safe is that the range is not an argument — see
// internal/netscan's package comment.
type netWiring struct {
scanner *netscan.Scanner
subnets []string
}
// wireNetScan builds the scanner. nil unless the block is enabled and valid.
func wireNetScan(cfg *config.Config) *netWiring {
nc, ok := cfg.NetScanner()
if !ok {
return nil
}
if err := netscan.Validate(nc); err != nil {
// config.validate already ran this, so reaching here is a programming
// error rather than a config one. Not fatal: the scanner off is a
// working Maven.
log.Printf("netscan: not wired: %v", err)
return nil
}
return &netWiring{scanner: netscan.New(nc), subnets: nc.Subnets}
}
// scanSummary answers "какие устройства в сети?" in one line.
func (w *netWiring) scanSummary(ctx context.Context) (string, bool) {
if w == nil {
return "", false
}
ctx, cancel := context.WithTimeout(ctx, scanBudget)
defer cancel()
hosts, err := w.scanner.Scan(ctx)
if err != nil {
log.Printf("netscan: scan: %v", err)
return "не получилось просканировать сеть.", true
}
if len(hosts) == 0 {
return "в сети никого не нашла.", true
}
shown := hosts
if len(shown) > scanReadOut {
shown = shown[:scanReadOut]
}
parts := make([]string, 0, len(shown))
for _, h := range shown {
s := h.Addr
if len(h.Ports) > 0 {
ps := make([]string, 0, len(h.Ports))
for _, p := range h.Ports {
ps = append(ps, fmt.Sprintf("%d", p))
}
s += " (" + strings.Join(ps, ", ") + ")"
}
parts = append(parts, s)
}
out := fmt.Sprintf("нашла %d %s: %s", len(hosts), hostWord(len(hosts)), strings.Join(parts, "; "))
if len(hosts) > len(shown) {
out += fmt.Sprintf(" и ещё %d", len(hosts)-len(shown))
}
return out + ".", true
}
// hostWord — Russian counts inflect the noun: 1 устройство, 2-4 устройства,
// 5+ устройств, and the teens are all the last form.
func hostWord(n int) string {
if n%100 >= 11 && n%100 <= 14 {
return "устройств"
}
switch n % 10 {
case 1:
return "устройство"
case 2, 3, 4:
return "устройства"
default:
return "устройств"
}
}
// isNetworkQuery recognises a question about the LAN, narrowly. It needs a
// network word AND an ask: "интернет не работает" is a complaint, not a request
// to scan, and a scan she runs unasked is exactly the noisy behaviour the
// bounds exist to prevent.
func isNetworkQuery(u string) bool {
s := strings.ToLower(strings.TrimSpace(u))
if s == "" {
return false
}
network := false
for _, w := range []string{"в сети", "в сетке", "сеть", "сети", "локальн", "wifi", "wi-fi", "вайфай"} {
if strings.Contains(s, w) {
network = true
break
}
}
if !network {
return false
}
// An explicit ask to scan, or a phrase that can only be about the LAN.
// "кто в сети" carries no device noun but means nothing else.
for _, w := range []string{"просканируй", "сканируй", "скан", "просканир", "кто в сети", "кто в сетке"} {
if strings.Contains(s, w) {
return true
}
}
ask := strings.Contains(s, "?") || homeWord(s, "какие") || homeWord(s, "кто") ||
homeWord(s, "что") || homeWord(s, "сколько") || strings.Contains(s, "покажи")
if !ask {
return false
}
for _, w := range []string{"устройств", "хост", "компьютер", "машин", "адрес"} {
if strings.Contains(s, w) {
return true
}
}
return false
}
+111
View File
@@ -0,0 +1,111 @@
package main
import (
"context"
"strings"
"testing"
"github.com/kami/maven/internal/config"
)
func TestWireNetScanOffUnlessEnabled(t *testing.T) {
for name, cfg := range map[string]*config.Config{
"no block": {},
"written but dark": {NetScan: &config.NetScanConfig{
Subnets: []string{"192.168.1.0/24"},
}},
"enabled but nothing to scan": {NetScan: &config.NetScanConfig{Enabled: true}},
"enabled but public": {NetScan: &config.NetScanConfig{
Subnets: []string{"8.8.8.0/24"}, Enabled: true,
}},
"enabled but far too wide": {NetScan: &config.NetScanConfig{
Subnets: []string{"10.0.0.0/8"}, Enabled: true,
}},
} {
t.Run(name, func(t *testing.T) {
if w := wireNetScan(cfg); w != nil {
t.Fatal("the scanner must not wire for this config")
}
})
}
var w *netWiring
if _, ok := w.scanSummary(context.Background()); ok {
t.Fatal("a nil wiring must not claim a query")
}
ok := wireNetScan(&config.Config{NetScan: &config.NetScanConfig{
Subnets: []string{"192.168.1.0/24"}, Enabled: true,
}})
if ok == nil {
t.Fatal("a valid enabled block should wire")
}
}
// A loopback /32 with nothing listening on the scanned port: the summary must
// come back honest rather than inventing a host. This also exercises the real
// dialer end to end without touching anything outside this box.
func TestScanSummaryOnAnEmptyRange(t *testing.T) {
w := wireNetScan(&config.Config{NetScan: &config.NetScanConfig{
// Port 1 on loopback: nothing listens and the connection is refused
// immediately, so the scan is fast and touches only this machine.
Subnets: []string{"127.0.0.1/32"}, Ports: []int{1}, Rate: 1000, Enabled: true,
}})
if w == nil {
t.Fatal("wireNetScan returned nil")
}
out, claimed := w.scanSummary(context.Background())
if !claimed {
t.Fatal("the summary did not claim the turn")
}
if out == "" {
t.Fatal("empty summary")
}
// Persona: feminine self-reference, informal address, no pet names.
low := strings.ToLower(out)
for _, bad := range []string{"нашёл", "не смог ", "вы ", "ваш", "милый", "дорогой"} {
if strings.Contains(low, bad) {
t.Errorf("persona violation %q in %q", bad, out)
}
}
}
func TestHostWordAgreesWithTheCount(t *testing.T) {
for n, want := range map[int]string{
1: "устройство", 2: "устройства", 4: "устройства", 5: "устройств",
11: "устройств", 12: "устройств", 21: "устройство", 22: "устройства",
25: "устройств", 111: "устройств", 101: "устройство", 0: "устройств",
} {
if got := hostWord(n); got != want {
t.Errorf("hostWord(%d) = %q, want %q", n, got, want)
}
}
}
func TestIsNetworkQuery(t *testing.T) {
yes := []string{
"какие устройства в сети?",
"кто в сети?",
"просканируй сеть",
"покажи устройства в локальной сети",
"сколько машин в сети",
}
no := []string{
"",
"интернет не работает",
"сеть какая-то медленная",
"я в сети инстаграма",
"что включено дома?",
"напомни оплатить интернет",
}
for _, u := range yes {
if !isNetworkQuery(u) {
t.Errorf("isNetworkQuery(%q) = false, want true", u)
}
}
for _, u := range no {
if isNetworkQuery(u) {
t.Errorf("isNetworkQuery(%q) = true, want false", u)
}
}
}
+798
View File
@@ -0,0 +1,798 @@
// mavend/simulator_test.go — the replayable full-system simulator
// (Vikunja #284, 20-07-2026-BACKLOG.md item 7).
//
// # What it is
//
// A scripted day, replayed through the real mavend code paths, with every
// boundary faked and the clock under the scenario's control. A scenario is a
// JSON file in testdata/scenarios; the harness reads it, builds a world, walks
// the steps in order, and asserts on what actually happened:
//
// what Maven SAID — the reply text of every utterance
// what was SENT — every delivery.Sendable the dispatcher emitted
// what ARRIVED — the unified intake journal from #283
// what TOOLS were called — the recorded requests against fake Praxis/Nexis/Hexis
// what did NOT happen — expect_no_send / expect_no_call, first-class
//
// The last one is the point. Maven's hard constraints are mostly negative —
// not a nag, not autonomous, nothing executed without confirmation — and a
// harness that can only assert on things that happened cannot test any of
// them. "Nothing was sent" is an assertion here, not an absence of one.
//
// # Determinism
//
// No time.Now() runs inside a replay. The scenario names a start instant, each
// step names a wall-clock offset from it, and the harness advances a fakeClock
// to that offset before running the step. Every clock reader in the world —
// the handler's `now`, the tick loop's `tick(ctx, now)`, the intake journal's
// publish stamp — is wired to that clock. Two runs of the same file produce
// the same transcript, and a scenario about 08:35 does not behave differently
// at 03:00 in CI.
//
// The tick is driven by the scenario, not by a ticker: tick() already takes
// `now` as an argument, so the only thing the daemon's ticker contributed was
// wall-clock timing, which is exactly what a replay must not have.
//
// # Why this shape and not a binary
//
// Vikunja #288 (golden-audio STT) deferred its tier-2 "audio → STT → router →
// phraser" scenarios to this task, and asked that they reuse a fixture format
// rather than inventing a third. A scenario here can name a WAV from
// cmd/mavsttd/testdata and the harness will feed it through the STT seam. As a
// test it runs under `make test` on every change, which a separate binary
// would not.
//
// # Production is untouched
//
// Every file this task adds is a _test.go file or testdata. There is no
// simulator in the daemon, no flag, no config key, and no code path that
// checks whether a simulation is running. The seams it uses — stt.Transcriber,
// tts.Synthesizer, router.Completer, delivery.Sink, ipc.CoreAPI, the
// event.Bus from #283 — all already existed for the production wiring.
package main
import (
"context"
"encoding/json"
"fmt"
"os"
"path/filepath"
"strings"
"sync"
"testing"
"time"
"github.com/kami/maven/internal/audio"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/delivery"
"github.com/kami/maven/internal/dialogue"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/llm"
"github.com/kami/maven/internal/loop"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/store"
"github.com/kami/maven/internal/tool"
"github.com/kami/maven/internal/voice"
)
// ---------------------------------------------------------------------------
// Scenario format
// ---------------------------------------------------------------------------
// scenario — one scripted day. schema_version matches the convention already
// set by testdata/system_safety_scenarios.json.
type scenario struct {
SchemaVersion int `json:"schema_version"`
Name string `json:"name"`
Description string `json:"description,omitempty"`
// Start — the instant the day begins, RFC3339. Every step offset is
// relative to it, and nothing in the run reads a real clock.
Start string `json:"start"`
// Script — what the resident model answers. The world has no llama-server;
// see scriptedLLM for how an entry is chosen.
Script []scriptEntry `json:"script,omitempty"`
// Praxis / Nexus / Hexis — canned bodies for the ecosystem fakes. Absent ⇒
// that service is not wired at all, which is the default box.
Praxis string `json:"praxis_attention,omitempty"`
Nexus string `json:"nexus_resolve,omitempty"`
Hexis string `json:"hexis_capabilities,omitempty"`
Steps []step `json:"steps"`
}
// scriptEntry — one canned model answer. Match is a substring of the user
// message; the first entry whose Match is contained in it wins, and an entry
// with an empty Match is the catch-all.
//
// Route and Reply are separate because the same model serves both contracts
// (CLAUDE.md, "LLM output contract"): a grammar-constrained call is a routing
// call and gets Route, an unconstrained one is a phrasing call and gets Reply.
type scriptEntry struct {
Match string `json:"match"`
Route string `json:"route,omitempty"`
Reply string `json:"reply,omitempty"`
}
// step — one scripted moment. At is "HH:MM" or "HH:MM:SS", interpreted in the
// start instant's location; the clock is advanced to it before the step runs.
//
// A step does exactly one thing (say / audio / signal / fact / tick / arrive)
// and then asserts. Assertions are evaluated against everything recorded since
// the run began, except expect_no_send and expect_no_call, which are scoped to
// this step — "nothing was sent because of THIS" is the useful question.
type step struct {
At string `json:"at"`
Note string `json:"note,omitempty"`
// --- stimuli (at most one per step) ---
// Say — an utterance, as text, through the same runTurn the IPC chat path
// uses.
Say string `json:"say,omitempty"`
// Audio — a WAV under cmd/mavsttd/testdata, fed through the STT seam. This
// is #288's deferred tier 2. The harness uses the deterministic stt stub
// unless a real transcriber is available, so the assertion a scenario can
// make about an audio step is about the PIPELINE, not about whisper's
// accuracy — that is what cmd/mavsttd/golden_test.go is for.
Audio string `json:"audio,omitempty"`
// Signal — a presence/world fact arriving from a poller or /api/signal.
Signal *signalStep `json:"signal,omitempty"`
// Arrive — an intake write from a module: an ambient notification, a feed
// item, a mail candidate. Goes through the same decorated ipc.CoreAPI the
// daemon gives those callers, so it lands in the journal exactly as it
// would in production.
Arrive *arriveStep `json:"arrive,omitempty"`
// Tick — run one iteration of the proactive loop at this instant.
Tick bool `json:"tick,omitempty"`
// Fault — make every ecosystem fake answer with this HTTP status from now
// on. The degraded-mode lever; ClearFault puts them back.
Fault int `json:"fault,omitempty"`
ClearFault bool `json:"clear_fault,omitempty"`
// --- assertions ---
ExpectReply []string `json:"expect_reply_contains,omitempty"`
ExpectNotReply []string `json:"expect_reply_lacks,omitempty"`
ExpectSent []string `json:"expect_sent_contains,omitempty"`
ExpectNoSend bool `json:"expect_no_send,omitempty"`
ExpectCalled []string `json:"expect_called,omitempty"`
ExpectNotCalled []string `json:"expect_not_called,omitempty"`
ExpectEvents []string `json:"expect_events,omitempty"`
ExpectNoEvents bool `json:"expect_no_events,omitempty"`
}
type signalStep struct {
Key string `json:"key"`
Value string `json:"value"`
Source string `json:"source"`
Kind string `json:"kind,omitempty"`
}
type arriveStep struct {
// Note / Fact / Task — exactly one. Each mirrors the intake seam its real
// caller uses.
Note *arriveNote `json:"note,omitempty"`
Fact *signalStep `json:"fact,omitempty"`
Task *arriveTask `json:"task,omitempty"`
AsOf string `json:"as_of,omitempty"` // "HH:MM" — OccurredAt, when it differs from the step time
Source string `json:"source"`
}
type arriveNote struct {
Text string `json:"text"`
}
type arriveTask struct {
Text string `json:"text"`
Evidence string `json:"evidence,omitempty"`
Status string `json:"status,omitempty"`
}
// ---------------------------------------------------------------------------
// The world
// ---------------------------------------------------------------------------
// simWorld — every faked boundary plus the real components between them.
type simWorld struct {
t *testing.T
clock *fakeClock
loc *time.Location
start time.Time
store *store.Store
api ipc.CoreAPI // the intake-decorated adapter, same as the daemon builds
bus *event.Bus
handler *reactiveHandler
tick *tickLoop
sink *recordingSink
llm *scriptedLLM
praxis *fakeServer
nexus *fakeServer
hexis *fakeServer
// transcript — everything that happened, in order. Printed on failure so a
// broken scenario is diagnosable without a debugger.
transcript []string
replies []string
}
// recordingSink captures every send, mutex-guarded (the tick loop dispatches
// from its own goroutine in production and the race detector is on here).
type recordingSink struct {
mu sync.Mutex
sends []delivery.Sendable
}
func (s *recordingSink) Send(_ context.Context, d delivery.Sendable) error {
s.mu.Lock()
defer s.mu.Unlock()
s.sends = append(s.sends, d)
return nil
}
func (s *recordingSink) all() []delivery.Sendable {
s.mu.Lock()
defer s.mu.Unlock()
out := make([]delivery.Sendable, len(s.sends))
copy(out, s.sends)
return out
}
func (s *recordingSink) count() int {
s.mu.Lock()
defer s.mu.Unlock()
return len(s.sends)
}
// scriptedLLM stands in for llama-server on BOTH contracts the resident model
// serves: grammar-constrained routing and unconstrained phrasing.
//
// It is not a stub that ignores its input — a scenario that scripts an answer
// for "что я пропустил" and gets asked something else must fail, not silently
// return the wrong intent. An unmatched call returns an error, and the router
// then falls through to the classifier cascade exactly as it does in
// production when llama-server is unreachable. That fall-through is itself
// worth exercising: it is the failure floor CLAUDE.md refuses to let rot.
type scriptedLLM struct {
mu sync.Mutex
entries []scriptEntry
calls []llm.Req
}
func (s *scriptedLLM) Complete(_ context.Context, r llm.Req) (string, error) {
s.mu.Lock()
defer s.mu.Unlock()
s.calls = append(s.calls, r)
routing := r.Grammar != ""
for _, e := range s.entries {
if e.Match != "" && !strings.Contains(strings.ToLower(r.User), strings.ToLower(e.Match)) {
continue
}
if routing && e.Route != "" {
return e.Route, nil
}
if !routing && e.Reply != "" {
return e.Reply, nil
}
}
return "", fmt.Errorf("simulator: no scripted %s answer for %q",
map[bool]string{true: "route", false: "reply"}[routing], truncateRunes(r.User, 60))
}
// ---------------------------------------------------------------------------
// Building the world
// ---------------------------------------------------------------------------
func newSimWorld(t *testing.T, sc scenario) *simWorld {
t.Helper()
start, err := time.Parse(time.RFC3339, sc.Start)
if err != nil {
t.Fatalf("scenario %q: bad start %q: %v", sc.Name, sc.Start, err)
}
clock := newFakeClock(start)
st := newTestStore(t)
bus := event.NewBus(512)
// The same decorator the daemon wires, on the same clock: intake in a
// replay is journalled exactly as it is in production.
api := newIntakeAPI(ipc.NewStoreAPI(st), bus, clock.Now)
sink := &recordingSink{}
rules := loop.DefaultRules()
gatherer := loop.NewGatherer(st, rules)
dispatcher := delivery.NewDispatcher(delivery.Config{
Voice: sink, Ntfy: sink, Telegram: sink, Nudges: st, Reminders: st,
})
tl := newTickLoop(st, gatherer, dispatcher, phraser.NewStub(), rules,
time.Minute, 5*time.Minute, 0, nil, nil, nil, nil)
scripted := &scriptedLLM{entries: sc.Script}
w := &simWorld{
t: t, clock: clock, loc: start.Location(), start: start,
store: st, api: api, bus: bus, tick: tl, sink: sink, llm: scripted,
}
// Ecosystem fakes, wired only when the scenario supplies a body — a box
// with no praxis block has no praxis client, and a scenario must be able to
// reproduce that.
eco := &ecosystemWiring{}
if sc.Praxis != "" {
w.praxis = newFakePraxis(t, sc.Praxis)
eco.praxis = newPraxisClient(w.praxis.URL)
}
if sc.Nexus != "" {
w.nexus = newFakeNexus(t, sc.Nexus)
}
if sc.Hexis != "" {
w.hexis = newFakeHexis(t, sc.Hexis, fixtureHexisExecuted("exec_1", "completed"))
}
// The router: the same cascade the daemon builds — stage-0 grammars, the
// LLM router on the scripted model, the classifier underneath. Keeping the
// classifier in is deliberate; it is the failure floor, and a scenario that
// scripts no route for an utterance exercises it.
emb := router.NewHashEmbedder(1024)
matcher := tool.NewMatcher(nil)
rtr := buildRouter(emb, matcher, config.DefaultRouterThreshold, router.NewLLMRouter(scripted))
w.handler = &reactiveHandler{
stt: simTranscriber{},
tts: simSynthesizer{},
router: rtr,
embedder: emb,
api: api,
matcher: matcher,
phraser: phraser.NewStub(),
replier: newLLMReplier(scripted, nil),
now: clock.Now,
memStore: st.VectorMemory(),
dataStore: st,
queryMinScore: config.DefaultQueryMinScore,
queryMinMargin: config.DefaultQueryMinMargin,
timeParser: router.StubDateTimeParser{},
dialogueSessions: dialogue.NewSessionStore(time.Hour),
clarifyStore: dialogue.NewClarifyStore(time.Hour),
clarifyMaxAttempts: dialogue.DefaultMaxAttempts,
ecosystem: eco,
}
return w
}
// simTranscriber — the STT seam. Deterministic by construction: it returns the
// text the harness parked for this step, so the pipeline under test is
// "audio arrives → a turn runs", not "whisper heard correctly". Transcription
// accuracy is cmd/mavsttd/golden_test.go's job (#288 tier 1), and duplicating
// it here would make every scenario depend on a 500 MB model.
type simTranscriber struct{ text string }
func (s simTranscriber) Transcribe(_ context.Context, _ audio.Audio) (string, float64, error) {
return s.text, 1.0, nil
}
// simSynthesizer — the TTS seam. A scenario asserts on what Maven SAID, which
// is the reply text; the waveform is not the artefact under test.
type simSynthesizer struct{}
func (simSynthesizer) Synthesize(_ context.Context, _ string) (audio.Audio, error) {
return audio.Audio{Format: audio.PCM16kMono}, nil
}
// ---------------------------------------------------------------------------
// Running
// ---------------------------------------------------------------------------
func (w *simWorld) logf(format string, args ...any) {
w.transcript = append(w.transcript,
fmt.Sprintf("%s %s", w.clock.Now().In(w.loc).Format("15:04:05"), fmt.Sprintf(format, args...)))
}
// dump prints the whole transcript. Called on any failure — a scenario that
// broke on step 7 is unreadable without the six steps before it.
func (w *simWorld) dump() {
w.t.Logf("--- replay transcript ---\n%s", strings.Join(w.transcript, "\n"))
}
// advanceTo moves the clock to the step's offset. Time only ever moves
// FORWARD: a scenario with steps out of order is a bug in the scenario, and
// silently reordering it would hide the bug.
func (w *simWorld) advanceTo(at string) {
w.t.Helper()
if at == "" {
return
}
target := w.timeOf(at)
now := w.clock.Now()
if target.Before(now) {
w.t.Fatalf("step at %s goes backwards from %s — scenario steps must be in order",
at, now.In(w.loc).Format("15:04:05"))
}
w.clock.Advance(target.Sub(now))
}
// timeOf resolves an "HH:MM" or "HH:MM:SS" step offset against the scenario's
// start day and location.
func (w *simWorld) timeOf(at string) time.Time {
w.t.Helper()
layout := "15:04"
if strings.Count(at, ":") == 2 {
layout = "15:04:05"
}
hm, err := time.Parse(layout, at)
if err != nil {
w.t.Fatalf("bad step time %q: %v", at, err)
}
return time.Date(w.start.Year(), w.start.Month(), w.start.Day(),
hm.Hour(), hm.Minute(), hm.Second(), 0, w.loc)
}
func (w *simWorld) run(sc scenario) {
ctx := context.Background()
for i, s := range sc.Steps {
w.advanceTo(s.At)
if s.Note != "" {
w.logf("# %s", s.Note)
}
sendsBefore := w.sink.count()
callsBefore := w.callCount()
eventsBefore := w.bus.Len()
w.stimulate(ctx, s)
w.assert(i, s, sendsBefore, callsBefore, eventsBefore)
}
}
func (w *simWorld) stimulate(ctx context.Context, s step) {
if s.Fault != 0 || s.ClearFault {
for _, fs := range []*fakeServer{w.praxis, w.nexus, w.hexis} {
if fs != nil {
fs.SetFault(s.Fault)
}
}
w.logf("fault=%d on every ecosystem fake", s.Fault)
}
switch {
case s.Say != "":
reply := w.handler.runTurn(ctx, s.Say)
w.replies = append(w.replies, reply)
w.logf("он: %s", s.Say)
w.logf("она: %s", reply)
case s.Audio != "":
text := w.audioText(s.Audio)
// Swap in a transcriber parked with this step's text, then run the same
// push-to-talk entry point the voice client calls.
w.handler.stt = simTranscriber{text: text}
resp, err := w.handler.HandlePushToTalk(ctx, voicePTT(), 0)
if err != nil {
w.t.Fatalf("push-to-talk on %s: %v", s.Audio, err)
}
w.replies = append(w.replies, resp.ReplyText)
w.logf("[wav %s → %q]", filepath.Base(s.Audio), text)
w.logf("она: %s", resp.ReplyText)
case s.Signal != nil:
w.write(ctx, *s.Signal, w.clock.Now())
w.logf("сигнал: %s=%s (%s)", s.Signal.Key, s.Signal.Value, s.Signal.Source)
case s.Arrive != nil:
w.arrive(ctx, *s.Arrive)
case s.Tick:
w.tick.tick(ctx, w.clock.Now())
w.logf("tick")
}
}
func (w *simWorld) write(ctx context.Context, sig signalStep, ts time.Time) {
w.t.Helper()
kind := sig.Kind
if kind == "" {
kind = "env"
}
if _, err := w.api.WriteFact(ctx, ipc.WriteFactReq{
Ts: ts, Kind: kind, Key: sig.Key, Value: sig.Value, Source: sig.Source, Confidence: 1.0,
}); err != nil {
w.t.Fatalf("write fact %s: %v", sig.Key, err)
}
}
func (w *simWorld) arrive(ctx context.Context, a arriveStep) {
w.t.Helper()
// AsOf is when the thing HAPPENED, which for a feed item or a relayed
// notification is usually earlier than when Maven heard about it. It does
// not move the clock — only the timestamp on the row and the envelope.
ts := w.clock.Now()
if a.AsOf != "" {
ts = w.timeOf(a.AsOf)
}
switch {
case a.Fact != nil:
f := *a.Fact
if f.Source == "" {
f.Source = a.Source
}
w.write(ctx, f, ts)
w.logf("пришло: факт %s=%s (%s)", f.Key, f.Value, f.Source)
case a.Note != nil:
if _, err := w.api.WriteNote(ctx, ts, a.Note.Text, nil, a.Source); err != nil {
w.t.Fatalf("write note from %s: %v", a.Source, err)
}
w.logf("пришло: заметка от %s — %s", a.Source, truncateRunes(a.Note.Text, 60))
case a.Task != nil:
status := a.Task.Status
if status == "" {
status = store.TaskCandidate
}
if _, err := w.api.CaptureTask(ctx, ipc.CaptureTaskReq{
Text: a.Task.Text, Source: a.Source, Evidence: a.Task.Evidence, Status: status, Ts: ts,
}); err != nil {
w.t.Fatalf("capture task from %s: %v", a.Source, err)
}
w.logf("пришло: задача от %s — %s", a.Source, a.Task.Text)
default:
w.t.Fatalf("arrive step from %s carries nothing", a.Source)
}
}
// audioText resolves a scenario's WAV reference to the text the fixture is
// known to contain, by reading cmd/mavsttd's golden manifest (#288's format,
// reused rather than duplicated). An unknown reference fails the scenario
// rather than quietly transcribing to "".
func (w *simWorld) audioText(ref string) string {
w.t.Helper()
manifest := filepath.Join("..", "mavsttd", "testdata", "golden_v1.json")
raw, err := os.ReadFile(manifest)
if err != nil {
w.t.Fatalf("audio step %q: reading %s: %v", ref, manifest, err)
}
var m struct {
Cases []struct {
Name string `json:"name"`
WAV string `json:"wav"`
Text string `json:"text"`
} `json:"cases"`
}
if err := json.Unmarshal(raw, &m); err != nil {
w.t.Fatalf("audio step %q: parsing %s: %v", ref, manifest, err)
}
for _, c := range m.Cases {
if c.Name == ref || c.WAV == ref {
return c.Text
}
}
w.t.Fatalf("audio step %q: no such case in %s", ref, manifest)
return ""
}
func voicePTT() voice.PushToTalkReq {
return voice.PushToTalkReq{Audio: audio.Audio{Format: audio.PCM16kMono}}
}
// callCount — how many requests every wired ecosystem fake has seen.
func (w *simWorld) callCount() int {
n := 0
for _, fs := range []*fakeServer{w.praxis, w.nexus, w.hexis} {
if fs != nil {
n += len(fs.Requests())
}
}
return n
}
func (w *simWorld) callPaths() []string {
var out []string
for _, fs := range []*fakeServer{w.praxis, w.nexus, w.hexis} {
if fs == nil {
continue
}
for _, r := range fs.Requests() {
out = append(out, r.Method+" "+r.Path)
}
}
return out
}
// ---------------------------------------------------------------------------
// Assertions
// ---------------------------------------------------------------------------
func (w *simWorld) assert(i int, s step, sendsBefore, callsBefore, eventsBefore int) {
w.t.Helper()
where := fmt.Sprintf("step %d (%s)", i+1, s.At)
if s.Note != "" {
where += " " + s.Note
}
fail := func(format string, args ...any) {
w.dump()
w.t.Errorf("%s: %s", where, fmt.Sprintf(format, args...))
}
lastReply := ""
if len(w.replies) > 0 {
lastReply = w.replies[len(w.replies)-1]
}
for _, want := range s.ExpectReply {
if !containsFold(lastReply, want) {
fail("reply %q does not contain %q", lastReply, want)
}
}
for _, unwanted := range s.ExpectNotReply {
if containsFold(lastReply, unwanted) {
fail("reply %q contains %q and must not", lastReply, unwanted)
}
}
sent := w.sink.all()
for _, want := range s.ExpectSent {
if !anyContains(sendableTexts(sent), want) {
fail("nothing sent mentions %q; sent so far: %v", want, sendableTexts(sent))
}
}
// Scoped to this step on purpose: "nothing was sent BECAUSE OF THIS" is the
// question a not-a-nag constraint asks.
if s.ExpectNoSend && len(sent) > sendsBefore {
fail("expected nothing to be sent, got %v", sendableTexts(sent[sendsBefore:]))
}
paths := w.callPaths()
for _, want := range s.ExpectCalled {
if !anyContains(paths, want) {
fail("no ecosystem call matches %q; calls so far: %v", want, paths)
}
}
for _, unwanted := range s.ExpectNotCalled {
if anyContains(paths[callsBefore:], unwanted) {
fail("an ecosystem call matched %q and must not have: %v", unwanted, paths[callsBefore:])
}
}
evs := w.bus.Recent(0)
for _, want := range s.ExpectEvents {
if !anyContains(eventLines(evs), want) {
fail("no intake event matches %q; journal: %v", want, eventLines(evs))
}
}
if s.ExpectNoEvents && w.bus.Len() > eventsBefore {
fail("expected nothing to arrive, journal grew to %d", w.bus.Len())
}
}
func sendableTexts(sends []delivery.Sendable) []string {
out := make([]string, 0, len(sends))
for _, s := range sends {
out = append(out, fmt.Sprintf("[%s] %s", s.RuleName, s.Body))
}
return out
}
func eventLines(evs []event.Event) []string {
out := make([]string, 0, len(evs))
for _, e := range evs {
out = append(out, fmt.Sprintf("%s/%s %s %s", e.Source, e.Kind, e.Title, e.Body))
}
return out
}
func containsFold(hay, needle string) bool {
return strings.Contains(strings.ToLower(hay), strings.ToLower(needle))
}
func anyContains(hay []string, needle string) bool {
for _, h := range hay {
if containsFold(h, needle) {
return true
}
}
return false
}
// ---------------------------------------------------------------------------
// The test
// ---------------------------------------------------------------------------
const scenarioDir = "testdata/scenarios"
// TestSimulatorScenarios replays every scenario file. Adding a scenario is
// adding a JSON file — no Go change, which is the property that makes this
// cheap enough to actually use.
func TestSimulatorScenarios(t *testing.T) {
entries, err := os.ReadDir(scenarioDir)
if err != nil {
t.Fatalf("reading %s: %v", scenarioDir, err)
}
var ran int
for _, ent := range entries {
if ent.IsDir() || !strings.HasSuffix(ent.Name(), ".json") {
continue
}
ran++
name := strings.TrimSuffix(ent.Name(), ".json")
t.Run(name, func(t *testing.T) {
sc := loadScenario(t, filepath.Join(scenarioDir, ent.Name()))
w := newSimWorld(t, sc)
w.run(sc)
if testing.Verbose() {
w.dump()
}
})
}
if ran == 0 {
t.Fatalf("no scenarios in %s — the harness would pass vacuously", scenarioDir)
}
}
func loadScenario(t *testing.T, path string) scenario {
t.Helper()
raw, err := os.ReadFile(path)
if err != nil {
t.Fatalf("reading %s: %v", path, err)
}
var sc scenario
dec := json.NewDecoder(strings.NewReader(string(raw)))
dec.DisallowUnknownFields() // a typo'd assertion key must fail, not be ignored
if err := dec.Decode(&sc); err != nil {
t.Fatalf("parsing %s: %v", path, err)
}
if sc.SchemaVersion != 1 {
t.Fatalf("%s: schema_version = %d, want 1", path, sc.SchemaVersion)
}
if sc.Name == "" || sc.Start == "" || len(sc.Steps) == 0 {
t.Fatalf("%s: a scenario needs a name, a start and at least one step", path)
}
return sc
}
// TestSimulatorIsDeterministic replays one scenario twice and requires an
// identical transcript. This is the property the whole task rests on: if a
// time.Now() creeps into a replayed path, two runs diverge and this fails.
func TestSimulatorIsDeterministic(t *testing.T) {
path := filepath.Join(scenarioDir, "morning_missed.json")
sc := loadScenario(t, path)
transcriptOf := func() string {
w := newSimWorld(t, sc)
w.run(sc)
return strings.Join(w.transcript, "\n")
}
first := transcriptOf()
second := transcriptOf()
if first != second {
t.Errorf("two replays of the same scenario diverged:\n--- first ---\n%s\n--- second ---\n%s", first, second)
}
// And the transcript's own timestamps must be the scenario's, not today's.
if strings.Contains(first, time.Now().Format("15:04")) && !strings.Contains(sc.Start, time.Now().Format("15:04")) {
t.Error("transcript carries the wall clock — something in the replay path read time.Now()")
}
}
// TestSimulatorRefusesBackwardsSteps guards the one scenario-authoring mistake
// that would silently produce a meaningless run.
func TestSimulatorRefusesBackwardsSteps(t *testing.T) {
// Not table-driven through run() because advanceTo calls t.Fatalf; this
// checks the ordering arithmetic directly.
sc := scenario{SchemaVersion: 1, Name: "x", Start: "2026-08-01T08:30:00+03:00",
Steps: []step{{At: "09:00"}}}
w := newSimWorld(t, sc)
w.advanceTo("09:00")
if got := w.clock.Now().In(w.loc).Format("15:04"); got != "09:00" {
t.Fatalf("clock at %s after advancing to 09:00", got)
}
w.advanceTo("09:30")
if got := w.clock.Now().In(w.loc).Format("15:04"); got != "09:30" {
t.Fatalf("clock at %s after advancing to 09:30", got)
}
}
+215
View File
@@ -0,0 +1,215 @@
package main
import (
"context"
"log"
"strings"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/smarthome"
"github.com/kami/maven/internal/store"
)
// homeWiring — the Home Assistant client, when the `smarthome` block is present
// AND enabled. nil ⇒ the house is not wired, nothing was proposed, and an
// allowlist row that happens to look like a house row refuses to run.
//
// It lives on the voice wiring for the same reason MCP does: a house control IS
// an act. It goes through tool.Executor, the enabled allowlist and the confirm
// turn, all of which only exist on the voice/chat path.
type homeWiring struct {
client *smarthome.Client
st *store.Store
refresh time.Duration
}
// wireSmartHome builds the client and proposes what it found. It never fails
// the daemon: an instance that is down at boot is logged and retried, because
// Maven starting is not contingent on someone else's process.
func wireSmartHome(cfg *config.Config, st *store.Store) *homeWiring {
hc, ok := cfg.SmartHomeClient()
if !ok || st == nil {
return nil
}
if err := smarthome.Validate(hc); err != nil {
// config.validate already ran this, so reaching here is a programming
// error rather than a config one. Still not fatal: the house off is a
// working Maven.
log.Printf("smarthome: not wired: %v", err)
return nil
}
w := &homeWiring{
client: smarthome.NewClient(hc),
st: st,
refresh: time.Duration(cfg.SmartHome.Refresh),
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
w.propose(ctx)
return w
}
// caller is the tool.HomeCaller seam.
func (w *homeWiring) caller() *smarthome.Client {
if w == nil {
return nil
}
return w.client
}
// propose writes a 'proposed' allowlist row for every controllable device. It
// does NOT enable anything: a reachable house is a place Maven may look, not a
// set of switches she may flip. Kami enables what he wants on /tools, behind
// step-up, which is the same gate a shell tool goes through.
//
// Sensors are read but never proposed — there is nothing to call on them.
func (w *homeWiring) propose(ctx context.Context) {
if w == nil {
return
}
ents, err := w.client.States(ctx)
if err != nil {
log.Printf("smarthome: read states: %v", err)
return
}
now := time.Now()
fresh, devices := 0, 0
for _, e := range ents {
svcs := smarthome.Services(e.Domain)
if len(svcs) == 0 {
continue
}
devices++
for _, s := range svcs {
name := smarthome.LocalName(e.ID, s.Verb)
provenance := "дом: " + s.Name + " → " + e.Name + " (" + e.ID + ")"
ok, err := w.st.ProposeSmartHomeTool(ctx, name, smarthome.Scope(e.Domain),
smarthome.Cmd(e.ID, s.Name), provenance, now)
if err != nil {
log.Printf("smarthome: propose %s: %v", name, err)
continue
}
if ok {
fresh++
}
}
}
log.Printf("smarthome: %d entities, %d controllable", len(ents), devices)
if fresh > 0 {
log.Printf("smarthome: %d new device proposal(s) waiting on /tools", fresh)
}
}
// run re-enumerates the house and picks up devices that appeared, until ctx is
// canceled.
func (w *homeWiring) run(ctx context.Context) {
if w == nil {
return
}
iv := w.refresh
if iv <= 0 {
iv = config.DefaultSmartHomeRefresh
}
t := time.NewTicker(iv)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
w.propose(ctx)
}
}
}
// homeSummary answers "что дома?" — a read of the current entity states, one
// short line. Read-only: it can never call a service, so it needs no confirm
// and no allowlist row.
func (w *homeWiring) homeSummary(ctx context.Context) (string, bool) {
if w == nil {
return "", false
}
ents, err := w.client.States(ctx)
if err != nil {
log.Printf("smarthome: summary: %v", err)
return "не смогла достучаться до дома.", true
}
if len(ents) == 0 {
return "дом ничего не отдаёт.", true
}
var on []string
var sensors []string
for _, e := range ents {
switch {
case e.Domain == "sensor" || e.Domain == "binary_sensor":
if len(sensors) < 3 && e.State != "" && e.State != "unavailable" {
sensors = append(sensors, e.Name+" "+e.State+e.Unit)
}
case e.State == "on" || e.State == "open" || e.State == "unlocked":
on = append(on, e.Name)
}
}
var parts []string
if len(on) > 0 {
if len(on) > 5 {
on = on[:5]
}
parts = append(parts, "включено: "+strings.Join(on, ", "))
} else {
parts = append(parts, "всё выключено")
}
if len(sensors) > 0 {
parts = append(parts, strings.Join(sensors, ", "))
}
return strings.Join(parts, "; ") + ".", true
}
// isHomeQuery recognises a question about the house, narrowly. "дома" on its
// own is not enough — "я дома" is a fact, not a question — so it takes a house
// marker AND an ask AND either a device word or the word "включ…". Weather
// wording bails out first: "какая температура на улице?" belongs to the weather
// source, and both questions contain "температура".
func isHomeQuery(u string) bool {
s := strings.ToLower(strings.TrimSpace(u))
if s == "" {
return false
}
for _, w := range []string{"погод", "на улице", "прогноз"} {
if strings.Contains(s, w) {
return false
}
}
for _, phrase := range []string{"что включено", "что выключено", "умный дом", "что в доме включено"} {
if strings.Contains(s, phrase) {
return true
}
}
house := homeWord(s, "дома") || strings.Contains(s, "в доме") || strings.Contains(s, "в квартире")
if !house {
return false
}
ask := strings.Contains(s, "?") || homeWord(s, "что") || homeWord(s, "какая") ||
homeWord(s, "какой") || homeWord(s, "сколько")
if !ask {
return false
}
for _, w := range []string{"свет", "лампа", "лампы", "розетк", "датчик", "температур", "включ", "выключ"} {
if strings.Contains(s, w) {
return true
}
}
return false
}
// homeWord — whole-token membership, so "дома" does not fire on "домашний".
// Punctuation is trimmed off each token because a spoken question arrives with
// a question mark glued to the last word.
func homeWord(s, w string) bool {
for _, tok := range strings.Fields(s) {
if strings.Trim(tok, ".,!?;:") == w {
return true
}
}
return false
}
+190
View File
@@ -0,0 +1,190 @@
package main
import (
"context"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/config"
)
const haStatesFixture = `[
{"entity_id":"light.living_room","state":"on","attributes":{"friendly_name":"Гостиная"}},
{"entity_id":"switch.kettle","state":"off","attributes":{"friendly_name":"Чайник"}},
{"entity_id":"sensor.bedroom_temp","state":"22.5","attributes":{"friendly_name":"Спальня","unit_of_measurement":"°C"}}
]`
func TestWireSmartHomeOffUnlessEnabled(t *testing.T) {
st := newTestStore(t)
for name, cfg := range map[string]*config.Config{
"no block": {},
"written but dark": {SmartHome: &config.SmartHomeConfig{
URL: "http://ha.lan:8123", Token: "t",
}},
} {
t.Run(name, func(t *testing.T) {
if w := wireSmartHome(cfg, st); w != nil {
t.Fatal("the house must be off unless the block is enabled")
}
})
}
// nil wiring must be safe everywhere it is reachable.
var w *homeWiring
w.propose(context.Background())
w.run(context.Background())
if w.caller() != nil {
t.Fatal("a nil wiring must have no caller")
}
if _, ok := w.homeSummary(context.Background()); ok {
t.Fatal("a nil wiring must not claim a query")
}
}
// An unreachable instance must not stop the daemon and must propose nothing.
func TestWireSmartHomeUnreachableIsNotFatal(t *testing.T) {
st := newTestStore(t)
w := wireSmartHome(&config.Config{SmartHome: &config.SmartHomeConfig{
// Port 1 on loopback: nothing listens, and it fails fast.
URL: "http://127.0.0.1:1", Token: "t", Enabled: true,
}}, st)
if w == nil {
t.Fatal("a configured house should still wire")
}
tools, err := st.ListTools(context.Background(), "")
if err != nil {
t.Fatal(err)
}
if len(tools) != 0 {
t.Fatalf("an instance that never answered must propose nothing, got %+v", tools)
}
}
// Discovery proposes one row per controllable service, always destructive,
// always 'proposed'. A sensor gets no row: there is nothing to call on it.
func TestProposeOnlyProposesControllableDevices(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte(haStatesFixture))
}))
defer srv.Close()
st := newTestStore(t)
w := wireSmartHome(&config.Config{SmartHome: &config.SmartHomeConfig{
URL: srv.URL, Token: "t", Enabled: true,
}}, st)
if w == nil {
t.Fatal("wireSmartHome returned nil for an enabled, reachable house")
}
tools, err := st.ListTools(context.Background(), "")
if err != nil {
t.Fatal(err)
}
got := map[string]bool{}
for _, tl := range tools {
got[tl.Name] = true
if tl.Status != "proposed" {
t.Errorf("%s status = %q: discovery must never enable", tl.Name, tl.Status)
}
if !tl.Destructive {
t.Errorf("%s is not destructive: every house control needs the confirm turn", tl.Name)
}
if len(tl.Cmd) == 0 || tl.Cmd[0] != "smarthome" {
t.Errorf("%s cmd = %v", tl.Name, tl.Cmd)
}
}
for _, want := range []string{
"home_light_living_room_on", "home_light_living_room_off",
"home_switch_kettle_on", "home_switch_kettle_off",
} {
if !got[want] {
t.Errorf("missing proposal %q (have %v)", want, got)
}
}
if len(tools) != 4 {
t.Fatalf("got %d rows, want 4 — the sensor must not be proposed: %+v", len(tools), tools)
}
// A second pass must be idempotent: re-discovery duplicates nothing and
// never rewrites a row Kami already enabled.
if err := st.EnableTool(context.Background(), "home_switch_kettle_on",
[]string{"smarthome", "switch.kettle", "turn_on"}, true, "smarthome:switch", time.Now()); err != nil {
t.Fatal(err)
}
w.propose(context.Background())
again, err := st.ListTools(context.Background(), "")
if err != nil {
t.Fatal(err)
}
if len(again) != 4 {
t.Fatalf("re-discovery duplicated rows: %d", len(again))
}
for _, tl := range again {
if tl.Name == "home_switch_kettle_on" && tl.Status != "enabled" {
t.Errorf("re-discovery un-enabled a device he had enabled: %q", tl.Status)
}
}
}
func TestHomeSummaryReadsState(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte(haStatesFixture))
}))
defer srv.Close()
w := wireSmartHome(&config.Config{SmartHome: &config.SmartHomeConfig{
URL: srv.URL, Token: "t", Enabled: true,
}}, newTestStore(t))
out, ok := w.homeSummary(context.Background())
if !ok {
t.Fatal("summary did not claim the turn")
}
if !strings.Contains(out, "Гостиная") {
t.Errorf("the lamp that is on should be named: %q", out)
}
if strings.Contains(out, "Чайник") {
t.Errorf("a device that is off should not be listed as on: %q", out)
}
if !strings.Contains(out, "22.5") {
t.Errorf("the sensor reading should be there: %q", out)
}
// Persona: no masculine self-reference, no "вы", no pet names.
for _, bad := range []string{"рад ", "готов ", "вы ", "ваш", "милый", "дорогой"} {
if strings.Contains(strings.ToLower(out), bad) {
t.Errorf("persona violation %q in %q", bad, out)
}
}
}
func TestIsHomeQuery(t *testing.T) {
yes := []string{
"что включено дома?",
"что выключено",
"какой свет горит дома",
"свет в доме включен?",
"какая температура в квартире?",
"покажи умный дом",
}
no := []string{
"",
"я дома",
"буду дома в семь",
"какая погода дома", // weather wording wins
"какая температура на улице?",
"домашние дела", // "дома" must not fire on "домашние"
"что мне нужно сделать?",
"напомни выключить чайник в семь", // a reminder, not a house read
}
for _, u := range yes {
if !isHomeQuery(u) {
t.Errorf("isHomeQuery(%q) = false, want true", u)
}
}
for _, u := range no {
if isHomeQuery(u) {
t.Errorf("isHomeQuery(%q) = true, want false", u)
}
}
}
+67
View File
@@ -0,0 +1,67 @@
{
"schema_version": 1,
"name": "evening_degraded",
"description": "The tier-2 pipeline case #288 deferred here, plus degraded mode. A golden WAV goes in at the microphone end and comes out as a written fact, and then the ecosystem starts answering 503 and the proactive loop has to stay quiet instead of falling over. The audio step asserts the PIPELINE — mic to STT seam to router to store to TTS — not whisper's accuracy; cmd/mavsttd/golden_test.go owns accuracy.",
"start": "2026-08-01T21:00:00+03:00",
"praxis_attention": "[{\"id\":\"item_1\",\"title\":\"medicine not taken\",\"importance\":3.0,\"rule\":\"evening_medicine\"}]",
"script": [
{
"match": "выпил воды",
"route": "[{\"intent\":\"fact\",\"key\":\"water\",\"value\":\"выпил\"}]"
},
{
"match": "записала факт: water",
"reply": "{\"response\":\"Записала, что ты выпил воды.\",\"mood\":\"neutral\"}"
},
{
"match": "",
"route": "[{\"intent\":\"chat\",\"text\":\"привет\"}]",
"reply": "{\"response\":\"Я рада тебя слышать.\",\"mood\":\"happy\"}"
}
],
"steps": [
{
"at": "21:00",
"note": "he speaks. The whole voice path runs: push-to-talk, the STT seam parked with the golden transcript, the real router, the real store write, the phrasing contract.",
"audio": "ru_fact",
"expect_reply_contains": ["записала"],
"expect_reply_lacks": ["записал,", "милый", "ваш"],
"expect_events": ["water"]
},
{
"at": "21:05",
"note": "a healthy tick with him just having spoken stays silent",
"tick": true,
"expect_no_send": true
},
{
"at": "21:10",
"note": "the ecosystem goes down",
"fault": 503
},
{
"at": "21:15",
"note": "a tick against a dead ecosystem must degrade, not send half a thought",
"tick": true,
"expect_no_send": true,
"expect_no_events": true
},
{
"at": "21:20",
"note": "intake keeps working while the ecosystem is down — a write does not depend on it",
"arrive": {
"source": "rss:tech",
"note": { "text": "Патч 6.19.1 [tech]\nисправления\nhttps://example.org/b" }
},
"expect_events": ["rss:tech"],
"expect_no_send": true
},
{
"at": "21:25",
"note": "recovery",
"clear_fault": true,
"tick": true,
"expect_no_send": true
}
]
}
+96
View File
@@ -0,0 +1,96 @@
{
"schema_version": 1,
"name": "morning_missed",
"description": "The scenario from Vikunja #284's description, replayed. He appears at 08:30, things arrive through the morning while he is at the desk, and at 08:50 he asks what he missed. The assertions are as much about what did NOT happen — nothing was sent at him unprompted — as about what she said.",
"start": "2026-08-01T08:30:00+03:00",
"praxis_attention": "[{\"id\":\"item_1\",\"title\":\"medicine not taken\",\"importance\":3.0,\"rule\":\"morning_medicine\"}]",
"script": [
{
"match": "выпил воды",
"route": "[{\"intent\":\"fact\",\"key\":\"water\",\"value\":\"выпил\"}]"
},
{
"match": "записала факт: water",
"reply": "{\"response\":\"Записала, что ты выпил воды.\",\"mood\":\"neutral\"}"
},
{
"match": "что я пропустил",
"route": "[{\"intent\":\"query\",\"text\":\"что я пропустил\"}]"
},
{
"match": "",
"route": "[{\"intent\":\"chat\",\"text\":\"привет\"}]",
"reply": "{\"response\":\"Я рада тебя слышать.\",\"mood\":\"happy\"}"
}
],
"steps": [
{
"at": "08:30",
"note": "he appears at the desk",
"signal": { "key": "desk_active", "value": "true", "source": "infer:hyprland" },
"expect_events": ["infer:hyprland"],
"expect_no_send": true
},
{
"at": "08:32",
"note": "a feed item arrives, published half an hour ago",
"arrive": {
"source": "rss:tech",
"as_of": "08:02",
"note": { "text": "Вышло ядро 6.19 [tech]\nкраткое содержание\nhttps://example.org/a" }
},
"expect_events": ["rss:tech"],
"expect_no_send": true
},
{
"at": "08:35",
"note": "the mail reader extracts a candidate — a candidate is never spoken",
"arrive": {
"source": "email:inbox",
"task": { "text": "продлить домен", "evidence": "Домен истекает через 7 дней" }
},
"expect_events": ["email:inbox", "продлить домен"],
"expect_no_send": true
},
{
"at": "08:40",
"note": "the work calendar signal — a relayed notification, below full confidence",
"arrive": {
"source": "ambient:notif",
"fact": {
"key": "calendar_event_20260801_планёрка",
"value": "10:00-11:00 планёрка"
}
},
"expect_events": ["ambient:notif", "планёрка"],
"expect_no_send": true
},
{
"at": "08:45",
"note": "a tick with him present and nothing wrong must stay silent",
"tick": true,
"expect_no_send": true
},
{
"at": "08:50",
"note": "he asks. The query path answers from local recall only: nothing stored clears the score gate, so she refuses rather than inventing a morning summary, and the replier is never reached. That refusal is the no-hallucination floor and this step pins it.",
"say": "что я пропустил?",
"expect_reply_contains": ["не знаю"],
"expect_reply_lacks": ["рад ", "милый", "ваш"]
},
{
"at": "08:55",
"note": "stating a fact writes it and says so, in the feminine",
"say": "я выпил воды",
"expect_reply_contains": ["записала"],
"expect_reply_lacks": ["записал,", "милый"],
"expect_events": ["water"]
},
{
"at": "09:00",
"note": "a second tick, still nothing unprompted",
"tick": true,
"expect_no_send": true
}
]
}
+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) {
+11
View File
@@ -92,6 +92,17 @@ type reactiveHandler struct {
// instead of "ничего нового", which are different truths.
feedsOn bool
// home — the Home Assistant client (Vikunja #256). nil ⇒ the house is not
// configured, which is the default: no `smarthome` block, no reads, no
// switches. Control does not go through this field — it goes through the
// act allowlist and tool.Executor, like every other mutating act.
home *homeWiring
// netscan — the LAN scanner (Vikunja #257). nil ⇒ off, which is the
// default. A scan is a read, so it has no allowlist row; what keeps it
// safe is that its range comes from config and from nowhere else.
netscan *netWiring
weatherProvider weather.Provider
weatherLocation string // default location for weather queries
+19
View File
@@ -49,6 +49,13 @@ type voiceWiring struct {
// server (Vikunja #251). Its tools land in the same allowlist as every
// other act, so nothing else here has to know about it.
mcp *mcpWiring
// home — the Home Assistant client, nil unless the `smarthome` block is
// enabled (Vikunja #256). Its devices land in the same allowlist as every
// other act, so nothing else here has to know about it.
home *homeWiring
// netscan — the LAN scanner, nil unless the `netscan` block is enabled
// (Vikunja #257).
netscan *netWiring
}
// close releases the listener + worker conns. Safe to call on nil (when
@@ -150,6 +157,16 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
if w.mcp != nil {
exec = exec.WithMCP(w.mcp.caller())
}
// The house (Vikunja #256): same story as MCP. Discovery PROPOSES a row per
// controllable device, always destructive, and Kami enables the ones he
// wants on /tools. Off unless the `smarthome` block is enabled.
w.home = wireSmartHome(cfg, dataStore)
if w.home != nil {
exec = exec.WithHome(w.home.caller())
}
// The LAN scanner (Vikunja #257): a read, bounded to the configured
// subnets and rate-limited. Off unless the `netscan` block is enabled.
w.netscan = wireNetScan(cfg)
matcher := tool.NewMatcher(coreAPI)
// ----- weather provider (Open-Meteo when configured, Stub otherwise) -----
@@ -233,6 +250,8 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
phraser: phr,
now: time.Now,
feedsOn: cfg.Feeds != nil,
home: w.home,
netscan: w.netscan,
// nil unless `crawl.on_demand` is on: reading a page he names is a
// capability, and capabilities are off unless configured.
crawler: onDemandCrawler(cfg),
+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)
}
}
}
+20
View File
@@ -45,6 +45,26 @@
]
},
"smarthome": {
"provider": "homeassistant",
"url": "http://192.168.1.50:8123",
"token": "${HA_TOKEN}",
"domains": ["light", "switch", "sensor"],
"max_entities": 40,
"timeout": "10s",
"refresh": "15m",
"enabled": false
},
"netscan": {
"subnets": ["192.168.1.0/24"],
"ports": [22, 80, 443, 8080],
"timeout": "400ms",
"rate": 50,
"max_hosts": 256,
"enabled": false
},
"nexus": { "url": "http://nexus:9740" },
"praxis": { "url": "http://praxis:8989" },
"hexis": { "url": "http://hexis:9741" },
+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
+170 -1
View File
@@ -25,6 +25,8 @@ import (
"github.com/kami/maven/internal/delivery/telegramsink"
"github.com/kami/maven/internal/mcp"
"github.com/kami/maven/internal/morning"
"github.com/kami/maven/internal/netscan"
"github.com/kami/maven/internal/smarthome"
"github.com/kami/maven/internal/update"
"github.com/robfig/cron/v3"
)
@@ -169,6 +171,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.
@@ -223,6 +239,16 @@ type Config struct {
// box. She is a client here, never a server: nothing exposes her own
// capabilities to an outside caller. See MCPConfig.
MCP *MCPConfig `json:"mcp,omitempty"`
// SmartHome — the Home Assistant instance (Vikunja #256). nil / absent /
// disabled ⇒ Maven neither reads the house nor touches it, and no house row
// exists in the act allowlist. See SmartHomeConfig.
SmartHome *SmartHomeConfig `json:"smarthome,omitempty"`
// NetScan — the LAN scanner (Vikunja #257). nil / absent / disabled ⇒
// Maven never puts a packet on the network looking for hosts. See
// NetScanConfig.
NetScan *NetScanConfig `json:"netscan,omitempty"`
}
// MCPConfig — the MCP client block. Servers are dark until one has
@@ -249,6 +275,106 @@ type MCPConfig struct {
MaxBytes int64 `json:"max_bytes,omitempty"`
}
// SmartHomeConfig — the Home Assistant block (Vikunja #256). Dark until
// `"enabled": true`, and even then a discovered device is only ever PROPOSED
// into the act allowlist: Kami enables it on /tools, behind step-up, exactly as
// he would a shell tool. Finding a switch on the network is not the same as
// being allowed to flip it.
type SmartHomeConfig struct {
// Provider — only "homeassistant" is implemented. MQTT / Zigbee2MQTT are
// not: Home Assistant already fronts them, and a broker client is a
// dependency this vendored module tree cannot take on tonight.
Provider string `json:"provider,omitempty"`
// URL — the instance base, "http://192.168.1.50:8123".
URL string `json:"url,omitempty"`
// Token — a long-lived access token. Use ${HA_TOKEN} and keep the value in
// the gitignored env file, like the telegram credentials.
Token string `json:"token,omitempty"`
// Domains — entity domains to take. Empty ⇒ the controllable domains
// (light, switch, fan, cover, lock) plus sensor and binary_sensor for
// reads. Narrow it when the instance is large: a tool name the 1.7B
// half-remembers is a wrong act.
Domains []string `json:"domains,omitempty"`
// MaxEntities — cap on the proposal catalogue. 0 ⇒ 40.
MaxEntities int `json:"max_entities,omitempty"`
// Timeout — per-call budget. 0 ⇒ 10s.
Timeout Duration `json:"timeout,omitempty"`
// Refresh — how often the entity list is re-read and new devices proposed.
// 0 ⇒ 15m. Discovery is idempotent, so this only ever adds rows.
Refresh Duration `json:"refresh,omitempty"`
// Enabled — false (the default) keeps a written block dark, so it can be
// reviewed before the house is wired to a voice.
Enabled bool `json:"enabled,omitempty"`
}
// SmartHomeClient maps the config block onto the smarthome package's own type.
// Returns ok=false when nothing is configured or it is disabled, so validation
// and daemon wiring cannot drift on the mapping.
func (c *Config) SmartHomeClient() (smarthome.Config, bool) {
if c.SmartHome == nil || !c.SmartHome.Enabled {
return smarthome.Config{}, false
}
return smarthome.Config{
URL: c.SmartHome.URL,
Token: c.SmartHome.Token,
Domains: c.SmartHome.Domains,
MaxEntities: c.SmartHome.MaxEntities,
Timeout: time.Duration(c.SmartHome.Timeout),
}, true
}
// NetScanConfig — the LAN scanner block (Vikunja #257). Dark until
// `"enabled": true`.
//
// The important field is Subnets, and it is the ONLY source of a scan target.
// Nothing an utterance, a router or a scanned host says can widen or move the
// range: internal/netscan.Scanner.Scan takes no target argument at all. Each
// subnet must be private and no larger than netscan.MaxPrefixHosts addresses
// (a /22), enforced at config load rather than at the first spoken scan.
type NetScanConfig struct {
// Subnets — CIDRs to scan, "192.168.1.0/24".
Subnets []string `json:"subnets,omitempty"`
// Ports — TCP ports to try per host. Empty ⇒ 22, 80, 443, 8080.
Ports []int `json:"ports,omitempty"`
// Timeout — per-connection budget. 0 ⇒ 400ms.
Timeout Duration `json:"timeout,omitempty"`
// Rate — connections per second across the whole scan. 0 ⇒ 50. Low on
// purpose: a scan should look like background traffic, not a portscan.
Rate int `json:"rate,omitempty"`
// MaxHosts — cap on addresses probed per scan. 0 ⇒ 256.
MaxHosts int `json:"max_hosts,omitempty"`
// Enabled — false (the default) keeps a written block dark.
Enabled bool `json:"enabled,omitempty"`
}
// NetScanner maps the config block onto the netscan package's own type.
// ok=false when absent or disabled, so validation and daemon wiring cannot
// drift on the mapping.
func (c *Config) NetScanner() (netscan.Config, bool) {
if c.NetScan == nil || !c.NetScan.Enabled {
return netscan.Config{}, false
}
return netscan.Config{
Subnets: c.NetScan.Subnets,
Ports: c.NetScan.Ports,
Timeout: time.Duration(c.NetScan.Timeout),
Rate: c.NetScan.Rate,
MaxHosts: c.NetScan.MaxHosts,
}, true
}
// MCPServerConfig — one MCP server.
type MCPServerConfig struct {
// Name — the local handle. It prefixes every tool this server contributes
@@ -870,6 +996,11 @@ type EmailConfig struct {
// DefaultEmailTimeout — extraction budget per message.
const DefaultEmailTimeout = 2 * time.Minute
// DefaultSmartHomeRefresh — how often the house is re-enumerated for new
// devices. Slow on purpose: discovery only adds proposals, and a flat does not
// grow a new lamp every minute.
const DefaultSmartHomeRefresh = 15 * time.Minute
// PhraserConfig — the LLM-backed phraser seam. The daemon spawns llama-server
// as a managed subprocess and sends chat-completion requests to phrase nudge
// and reminder messages. nil ⇒ the template-based Stub is used instead.
@@ -968,7 +1099,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 +1153,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)
}
@@ -1092,6 +1230,20 @@ func (c *Config) applyDefaults() {
c.MCP = nil
}
// Same rule for the house: a block that is not enabled is the same as no
// block at all, so "off" stays in one place.
if c.SmartHome != nil && !c.SmartHome.Enabled {
c.SmartHome = nil
}
if c.SmartHome != nil && c.SmartHome.Refresh <= 0 {
c.SmartHome.Refresh = Duration(DefaultSmartHomeRefresh)
}
// Same rule for the scanner.
if c.NetScan != nil && !c.NetScan.Enabled {
c.NetScan = nil
}
// Same rule for the crawler: a block that neither answers on demand nor
// watches anything has nothing to do, so it is normalised to "off".
if c.Crawl != nil && !c.Crawl.OnDemand && len(c.Crawl.Watches) == 0 {
@@ -1202,6 +1354,23 @@ func (c *Config) validate() error {
if err := mcp.Validate(c.MCPServers()); err != nil {
return err
}
// Same for the house: a missing token or a bare hostname fails at startup,
// not at the first "выключи свет".
if hc, ok := c.SmartHomeClient(); ok {
if p := c.SmartHome.Provider; p != "" && p != "homeassistant" {
return fmt.Errorf("smarthome: provider %q: only \"homeassistant\" is implemented", p)
}
if err := smarthome.Validate(hc); err != nil {
return err
}
}
// A scanner pointed at the public internet, or at a /8, fails here rather
// than after the packets have already left.
if nc, ok := c.NetScanner(); ok {
if err := netscan.Validate(nc); err != nil {
return err
}
}
if len(c.MorningRoutines) > 0 {
if err := morning.Validate(morningRoutinesFromConfig(c.MorningRoutines)); err != nil {
return err
+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
+356
View File
@@ -0,0 +1,356 @@
// Package netscan discovers hosts on the LAN Maven is configured to look at
// (Vikunja #257, docs/plans/12-bluetooth-network-scan.md).
//
// A scan is a read, but an unbounded scanner on a home network is noisy and is
// trivially pointed somewhere it should not go, so the whole package is built
// around four rules:
//
// - The target range NEVER comes from an utterance, a router, an LLM or a
// device. Scan takes no target argument at all: it reads only the CIDRs in
// the config block. There is deliberately no exported way to scan an
// arbitrary range, so no amount of prompt injection or a rogue reply from a
// scanned host can retarget it.
// - Every configured CIDR must be private (RFC1918 / CGNAT / link-local) and
// no larger than MaxPrefixHosts addresses. Scanning the public internet
// from his flat is not a thing Maven does, and /8 is not a home LAN.
// - Rate-limited. Connections leave at a fixed rate, so a scan looks like
// background traffic rather than a portscan to anything watching.
// - Bounded in total. MaxHosts, a per-connection timeout and the caller's
// context all cap the work; a scan that runs long returns what it has.
//
// It is a TCP-connect scan (net.DialTimeout) and an ARP-table read. No raw
// sockets, no SYN scan, no privileges: mavend does not run as root and this
// does not ask it to.
package netscan
import (
"bufio"
"context"
"errors"
"fmt"
"io"
"net"
"net/netip"
"os"
"sort"
"strings"
"sync"
"time"
)
// DefaultPorts — what a scan looks at when the config names nothing. Chosen to
// answer "what is this box" on a home network, not to find a way in.
var DefaultPorts = []int{22, 80, 443, 8080}
const (
// DefaultTimeout — per-connection budget. Short: on a LAN a live host
// answers in single-digit milliseconds, and a filtered port never answers.
DefaultTimeout = 400 * time.Millisecond
// DefaultRate — connections per second across the whole scan.
DefaultRate = 50
// DefaultMaxHosts — cap on addresses probed in one scan.
DefaultMaxHosts = 256
// MaxPrefixHosts — the largest CIDR that may be configured, in addresses.
// 1024 is a /22: generous for a flat, and far short of anything that would
// take minutes or wake up a neighbour's IDS.
MaxPrefixHosts = 1024
// maxParallel — in-flight dials. The rate limiter is the real throttle;
// this only stops a slow subnet from piling up file descriptors.
maxParallel = 16
// arpFile — the kernel's ARP cache. Reading it is free and needs no packet.
arpFile = "/proc/net/arp"
)
var (
// ErrNotConfigured — no netscan block, or it is disabled.
ErrNotConfigured = errors.New("netscan: not configured")
// ErrNoSubnets — enabled with nothing to scan.
ErrNoSubnets = errors.New("netscan: no subnets configured")
)
// Config — the bounds of every scan. There is nothing here that can be
// overridden at call time.
type Config struct {
// Subnets — the ONLY ranges that are ever probed, as CIDRs. Each must be
// private and no bigger than MaxPrefixHosts.
Subnets []string
// Ports — TCP ports to try on each host. Empty ⇒ DefaultPorts.
Ports []int
// Timeout — per-connection budget. 0 ⇒ DefaultTimeout.
Timeout time.Duration
// Rate — connections per second. 0 ⇒ DefaultRate.
Rate int
// MaxHosts — cap on addresses probed per scan. 0 ⇒ DefaultMaxHosts.
MaxHosts int
}
// Host is one machine the scan saw.
type Host struct {
// Addr — the IP.
Addr string
// MAC — from the ARP cache, empty when the kernel has no entry.
MAC string
// Ports — open TCP ports, ascending.
Ports []int
}
// Up reports whether anything at all answered for this host.
func (h Host) Up() bool { return len(h.Ports) > 0 || h.MAC != "" }
// Validate rejects a block that cannot safely run, at config-load time rather
// than at the first spoken scan. This is the guard the whole package rides on:
// if it passes, every later scan is inside these bounds by construction.
func Validate(c Config) error {
if len(c.Subnets) == 0 {
return ErrNoSubnets
}
for _, s := range c.Subnets {
p, err := netip.ParsePrefix(strings.TrimSpace(s))
if err != nil {
return fmt.Errorf("netscan: subnet %q: %w", s, err)
}
if !p.Addr().Is4() {
return fmt.Errorf("netscan: subnet %q: only IPv4 is scanned", s)
}
if !isPrivate(p.Addr()) {
return fmt.Errorf("netscan: subnet %q is not a private range: Maven does not scan the public internet", s)
}
if n := prefixHosts(p); n > MaxPrefixHosts {
return fmt.Errorf("netscan: subnet %q covers %d addresses, limit is %d: narrow the prefix", s, n, MaxPrefixHosts)
}
}
for _, port := range c.Ports {
if port < 1 || port > 65535 {
return fmt.Errorf("netscan: port %d out of range", port)
}
}
if c.Rate < 0 || c.MaxHosts < 0 || c.Timeout < 0 {
return errors.New("netscan: rate, max_hosts and timeout must not be negative")
}
return nil
}
// isPrivate — RFC1918, CGNAT and link-local. Loopback counts: scanning this box
// is harmless and is how the tests run.
func isPrivate(a netip.Addr) bool {
if a.IsLoopback() || a.IsPrivate() || a.IsLinkLocalUnicast() {
return true
}
// 100.64.0.0/10, the carrier-grade NAT range Tailscale hands out.
cgnat := netip.MustParsePrefix("100.64.0.0/10")
return cgnat.Contains(a)
}
// prefixHosts — addresses covered by a v4 prefix.
func prefixHosts(p netip.Prefix) int {
bits := 32 - p.Bits()
if bits >= 31 {
return MaxPrefixHosts + 1
}
return 1 << bits
}
// Scanner probes the configured subnets. Build it with New; the config it holds
// is the config it was validated with, and nothing mutates it afterwards.
type Scanner struct {
cfg Config
// dial is the connect seam; tests swap it.
dial func(ctx context.Context, addr string, timeout time.Duration) bool
// arp is the ARP-cache seam; tests swap it.
arp func() (map[string]string, error)
}
// New builds a scanner. Validate first — this does not.
func New(cfg Config) *Scanner {
if len(cfg.Ports) == 0 {
cfg.Ports = append([]int(nil), DefaultPorts...)
}
if cfg.Timeout <= 0 {
cfg.Timeout = DefaultTimeout
}
if cfg.Rate <= 0 {
cfg.Rate = DefaultRate
}
if cfg.MaxHosts <= 0 {
cfg.MaxHosts = DefaultMaxHosts
}
return &Scanner{cfg: cfg, dial: dialTCP, arp: readARP}
}
// targets expands the configured subnets into addresses, skipping the network
// and broadcast address of each, capped at MaxHosts. Deterministic order, so
// two scans of an unchanged network read the same.
func (s *Scanner) targets() []netip.Addr {
var out []netip.Addr
for _, cidr := range s.cfg.Subnets {
p, err := netip.ParsePrefix(strings.TrimSpace(cidr))
if err != nil {
continue
}
p = p.Masked()
first := p.Addr()
for a := first; p.Contains(a); a = a.Next() {
if len(out) >= s.cfg.MaxHosts {
return out
}
// Skip the network address; the broadcast address is skipped by
// looking one ahead.
if a == first && p.Bits() < 31 {
continue
}
if p.Bits() < 31 && !p.Contains(a.Next()) {
continue
}
out = append(out, a)
}
}
return out
}
// Scan probes every configured address and returns the hosts that answered.
//
// It takes no target: the range is the configured one, always. Callers pass a
// context and nothing else, which is the point — see the package comment.
func (s *Scanner) Scan(ctx context.Context) ([]Host, error) {
if len(s.cfg.Subnets) == 0 {
return nil, ErrNoSubnets
}
arp, err := s.arp()
if err != nil {
// A missing /proc/net/arp costs MAC addresses, not the scan.
arp = map[string]string{}
}
// One token per connection, at Rate per second, shared by every worker.
interval := time.Second / time.Duration(s.cfg.Rate)
if interval <= 0 {
interval = time.Millisecond
}
tick := time.NewTicker(interval)
defer tick.Stop()
type result struct {
addr string
ports []int
}
targets := s.targets()
results := make(chan result, len(targets))
sem := make(chan struct{}, maxParallel)
var wg sync.WaitGroup
scan:
for _, a := range targets {
addr := a.String()
for _, port := range s.cfg.Ports {
// Checked before the select as well as inside it: select picks
// randomly among ready cases, so at a high rate the ticker would
// sometimes win over an already-canceled context and let one more
// probe out.
if ctx.Err() != nil {
break scan
}
select {
case <-ctx.Done():
break scan
case <-tick.C:
}
sem <- struct{}{}
wg.Add(1)
go func(addr string, port int) {
defer wg.Done()
defer func() { <-sem }()
if s.dial(ctx, net.JoinHostPort(addr, itoa(port)), s.cfg.Timeout) {
results <- result{addr: addr, ports: []int{port}}
}
}(addr, port)
}
}
wg.Wait()
close(results)
byAddr := map[string]*Host{}
for r := range results {
h := byAddr[r.addr]
if h == nil {
h = &Host{Addr: r.addr}
byAddr[r.addr] = h
}
h.Ports = append(h.Ports, r.ports...)
}
// A host in the ARP cache is up even with every port closed — it answered
// an ARP request, which is the cheapest liveness signal there is.
for _, a := range targets {
addr := a.String()
mac, ok := arp[addr]
if !ok {
continue
}
if byAddr[addr] == nil {
byAddr[addr] = &Host{Addr: addr}
}
byAddr[addr].MAC = mac
}
out := make([]Host, 0, len(byAddr))
for _, h := range byAddr {
sort.Ints(h.Ports)
out = append(out, *h)
}
sort.Slice(out, func(i, j int) bool {
ai, _ := netip.ParseAddr(out[i].Addr)
aj, _ := netip.ParseAddr(out[j].Addr)
return ai.Less(aj)
})
return out, nil
}
func itoa(n int) string { return fmt.Sprintf("%d", n) }
func dialTCP(ctx context.Context, addr string, timeout time.Duration) bool {
d := net.Dialer{Timeout: timeout}
ctx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
c, err := d.DialContext(ctx, "tcp", addr)
if err != nil {
return false
}
_ = c.Close()
return true
}
func readARP() (map[string]string, error) {
f, err := os.Open(arpFile)
if err != nil {
return nil, err
}
defer f.Close()
return parseARP(f)
}
// parseARP reads the kernel's ARP table. Incomplete entries (all-zero MAC,
// flags 0x0) are dropped: they mean "we asked and nobody answered", which is
// the opposite of a discovered host.
func parseARP(r io.Reader) (map[string]string, error) {
out := map[string]string{}
sc := bufio.NewScanner(r)
first := true
for sc.Scan() {
if first { // header row
first = false
continue
}
f := strings.Fields(sc.Text())
if len(f) < 4 {
continue
}
ip, flags, mac := f[0], f[2], f[3]
if flags == "0x0" || mac == "00:00:00:00:00:00" {
continue
}
if _, err := netip.ParseAddr(ip); err != nil {
continue
}
out[ip] = mac
}
return out, sc.Err()
}
+198
View File
@@ -0,0 +1,198 @@
package netscan
import (
"context"
"errors"
"net/netip"
"strings"
"sync"
"testing"
"time"
)
func TestValidateBounds(t *testing.T) {
ok := []Config{
{Subnets: []string{"192.168.1.0/24"}},
{Subnets: []string{"10.0.0.0/24", "172.16.5.0/28"}, Ports: []int{22, 80}},
{Subnets: []string{"127.0.0.1/32"}},
{Subnets: []string{"100.64.1.0/24"}}, // CGNAT / tailnet
}
for _, c := range ok {
if err := Validate(c); err != nil {
t.Errorf("Validate(%v) = %v, want nil", c.Subnets, err)
}
}
bad := map[string]Config{
"nothing to scan": {},
"public range": {Subnets: []string{"8.8.8.0/24"}},
"whole internet": {Subnets: []string{"0.0.0.0/0"}},
"a slash-8 is not a flat": {Subnets: []string{"10.0.0.0/8"}},
"a /16 is too big": {Subnets: []string{"192.168.0.0/16"}},
"not a cidr": {Subnets: []string{"192.168.1.1"}},
"ipv6": {Subnets: []string{"fd00::/120"}},
"garbage": {Subnets: []string{"выключи свет"}},
"bad port": {Subnets: []string{"192.168.1.0/24"}, Ports: []int{0}},
"huge port": {Subnets: []string{"192.168.1.0/24"}, Ports: []int{70000}},
"negative rate": {Subnets: []string{"192.168.1.0/24"}, Rate: -1},
}
for name, c := range bad {
if err := Validate(c); err == nil {
t.Errorf("Validate(%s) = nil, want an error", strings.ReplaceAll(name, "\n", " "))
}
}
if !errors.Is(Validate(Config{}), ErrNoSubnets) {
t.Error("an empty block should report ErrNoSubnets")
}
}
// The whole safety story: a scanner probes its configured range and nothing
// else. There is no API that takes a target, so this test asserts the negative
// by watching every address the dialer was handed.
func TestScanOnlyTouchesConfiguredSubnet(t *testing.T) {
s := New(Config{Subnets: []string{"192.168.9.0/29"}, Ports: []int{80}, Rate: 10000})
inside := netip.MustParsePrefix("192.168.9.0/29")
var mu sync.Mutex
var seen []string
s.dial = func(_ context.Context, addr string, _ time.Duration) bool {
mu.Lock()
seen = append(seen, addr)
mu.Unlock()
return addr == "192.168.9.3:80"
}
s.arp = func() (map[string]string, error) { return map[string]string{}, nil }
hosts, err := s.Scan(context.Background())
if err != nil {
t.Fatalf("Scan: %v", err)
}
if len(hosts) != 1 || hosts[0].Addr != "192.168.9.3" || len(hosts[0].Ports) != 1 {
t.Fatalf("hosts = %+v", hosts)
}
// A /29 is 8 addresses; network (.0) and broadcast (.7) are skipped.
if len(seen) != 6 {
t.Errorf("probed %d addresses, want 6 (a /29 minus network and broadcast): %v", len(seen), seen)
}
for _, a := range seen {
host, _, _ := strings.Cut(a, ":")
ip, err := netip.ParseAddr(host)
if err != nil || !inside.Contains(ip) {
t.Errorf("probed %q, which is outside the configured subnet", a)
}
}
}
func TestScanHonoursMaxHosts(t *testing.T) {
s := New(Config{Subnets: []string{"192.168.9.0/24"}, Ports: []int{80}, Rate: 10000, MaxHosts: 3})
var mu sync.Mutex
n := 0
s.dial = func(_ context.Context, _ string, _ time.Duration) bool {
mu.Lock()
n++
mu.Unlock()
return false
}
s.arp = func() (map[string]string, error) { return nil, nil }
if _, err := s.Scan(context.Background()); err != nil {
t.Fatal(err)
}
if n != 3 {
t.Errorf("dialed %d times, want 3 (MaxHosts)", n)
}
}
// The rate limiter must actually gate: 6 probes at 200/s cannot finish in less
// than ~25ms. Asserted loosely, since a CI box is not a stopwatch.
func TestScanIsRateLimited(t *testing.T) {
s := New(Config{Subnets: []string{"192.168.9.0/29"}, Ports: []int{80}, Rate: 200})
s.dial = func(context.Context, string, time.Duration) bool { return false }
s.arp = func() (map[string]string, error) { return nil, nil }
start := time.Now()
if _, err := s.Scan(context.Background()); err != nil {
t.Fatal(err)
}
if el := time.Since(start); el < 20*time.Millisecond {
t.Errorf("6 probes at 200/s took %v: the rate limiter is not gating", el)
}
}
func TestScanStopsOnCanceledContext(t *testing.T) {
s := New(Config{Subnets: []string{"192.168.9.0/24"}, Ports: []int{80}, Rate: 10000})
ctx, cancel := context.WithCancel(context.Background())
cancel()
s.dial = func(context.Context, string, time.Duration) bool {
t.Error("a canceled scan still dialed")
return false
}
s.arp = func() (map[string]string, error) { return nil, nil }
if _, err := s.Scan(ctx); err != nil {
t.Fatal(err)
}
}
// A host with every port closed but an ARP entry is still up. A host outside
// the configured range must not be reported even if the kernel knows it —
// otherwise the ARP cache, which is populated by the network rather than by
// Maven, would widen the answer past what he configured.
func TestARPFillsMACWithinTheConfiguredRangeOnly(t *testing.T) {
s := New(Config{Subnets: []string{"192.168.9.0/29"}, Ports: []int{80}, Rate: 10000})
s.dial = func(context.Context, string, time.Duration) bool { return false }
s.arp = func() (map[string]string, error) {
return map[string]string{
"192.168.9.2": "aa:bb:cc:dd:ee:ff",
"10.9.9.9": "11:22:33:44:55:66",
}, nil
}
hosts, err := s.Scan(context.Background())
if err != nil {
t.Fatal(err)
}
if len(hosts) != 1 {
t.Fatalf("hosts = %+v", hosts)
}
if hosts[0].Addr != "192.168.9.2" || hosts[0].MAC != "aa:bb:cc:dd:ee:ff" {
t.Errorf("host = %+v", hosts[0])
}
if !hosts[0].Up() {
t.Error("an ARP entry with no open port is still a live host")
}
}
const arpFixture = `IP address HW type Flags HW address Mask Device
192.168.1.1 0x1 0x2 3c:84:6a:11:22:33 * wlp1s0
192.168.1.50 0x1 0x2 b8:27:eb:44:55:66 * wlp1s0
192.168.1.77 0x1 0x0 00:00:00:00:00:00 * wlp1s0
not-an-ip 0x1 0x2 de:ad:be:ef:00:01 * wlp1s0
short line
`
func TestParseARP(t *testing.T) {
got, err := parseARP(strings.NewReader(arpFixture))
if err != nil {
t.Fatal(err)
}
if len(got) != 2 {
t.Fatalf("got %d entries, want 2: %v", len(got), got)
}
if got["192.168.1.1"] != "3c:84:6a:11:22:33" || got["192.168.1.50"] != "b8:27:eb:44:55:66" {
t.Errorf("entries = %v", got)
}
if _, ok := got["192.168.1.77"]; ok {
t.Error("an incomplete ARP entry (flags 0x0) is not a discovered host")
}
}
func TestNewAppliesDefaults(t *testing.T) {
s := New(Config{Subnets: []string{"192.168.1.0/24"}})
if len(s.cfg.Ports) != len(DefaultPorts) || s.cfg.Rate != DefaultRate ||
s.cfg.MaxHosts != DefaultMaxHosts || s.cfg.Timeout != DefaultTimeout {
t.Errorf("defaults not applied: %+v", s.cfg)
}
// The defaults must not alias the package slice, or a second scanner could
// rewrite DefaultPorts through it.
s.cfg.Ports[0] = 9999
if DefaultPorts[0] == 9999 {
t.Error("New aliased DefaultPorts")
}
}
+231
View File
@@ -0,0 +1,231 @@
package smarthome
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"sort"
"strings"
"time"
)
// DefaultTimeout — per-call budget. A house that takes longer than this to
// answer is not usable in a spoken turn.
const DefaultTimeout = 10 * time.Second
// DefaultMaxEntities — cap on how many entities become allowlist proposals.
// The resident model is a 1.7B with a 4096-token context: a tool name it
// half-remembers is a wrong act, so a bounded, deliberate catalogue beats a
// complete one.
const DefaultMaxEntities = 40
// maxBody — cap on one /api/states response. A Home Assistant with hundreds of
// entities would otherwise stream megabytes into a daemon that wants forty
// names.
const maxBody = 4 << 20
// Config — what a Home Assistant instance needs to be reachable.
type Config struct {
// URL — the base, "http://homeassistant.local:8123". No trailing path.
URL string
// Token — a long-lived access token. Sent as a bearer header and never
// logged.
Token string
// Domains — the entity domains to take. Empty ⇒ every domain in the
// controllable table plus sensor/binary_sensor for reads.
Domains []string
// MaxEntities — 0 ⇒ DefaultMaxEntities.
MaxEntities int
// Timeout — 0 ⇒ DefaultTimeout.
Timeout time.Duration
}
// Validate rejects a block that cannot work, at config-load time rather than at
// the first spoken act.
func Validate(c Config) error {
if c.URL == "" {
return errors.New("smarthome: url is required")
}
u, err := url.Parse(c.URL)
if err != nil {
return fmt.Errorf("smarthome: url: %w", err)
}
if u.Scheme != "http" && u.Scheme != "https" {
return fmt.Errorf("smarthome: url scheme %q: want http or https", u.Scheme)
}
if u.Host == "" {
return errors.New("smarthome: url has no host")
}
if c.Token == "" {
return errors.New("smarthome: token is required")
}
return nil
}
// Client is a Home Assistant REST client. Read (States) and one write
// (CallService); no WebSocket, because a spoken turn is request/response and an
// event stream is a second failure mode for no gain yet.
type Client struct {
cfg Config
http *http.Client
}
// NewClient builds a client. Validate first — this does not.
func NewClient(cfg Config) *Client {
if cfg.Timeout <= 0 {
cfg.Timeout = DefaultTimeout
}
if cfg.MaxEntities <= 0 {
cfg.MaxEntities = DefaultMaxEntities
}
return &Client{cfg: cfg, http: &http.Client{Timeout: cfg.Timeout}}
}
// SetHTTPClient swaps the transport. Tests use it; nothing else should.
func (c *Client) SetHTTPClient(h *http.Client) { c.http = h }
// wanted reports whether an entity's domain is one Maven takes. The config list
// wins when set; otherwise every controllable domain plus the two read-only
// sensor domains.
func (c *Client) wanted(domain string) bool {
if len(c.cfg.Domains) > 0 {
for _, d := range c.cfg.Domains {
if d == domain {
return true
}
}
return false
}
if _, ok := controllable[domain]; ok {
return true
}
return domain == "sensor" || domain == "binary_sensor"
}
type haState struct {
EntityID string `json:"entity_id"`
State string `json:"state"`
Attributes json.RawMessage `json:"attributes"`
}
type haAttrs struct {
FriendlyName string `json:"friendly_name"`
Unit string `json:"unit_of_measurement"`
}
// States reads every entity Maven cares about, sorted by id and capped at
// MaxEntities so the catalogue is deterministic across restarts — a proposal
// list that reshuffles itself would make /tools unreadable.
func (c *Client) States(ctx context.Context) ([]Entity, error) {
body, err := c.do(ctx, http.MethodGet, "/api/states", nil)
if err != nil {
return nil, err
}
var raw []haState
if err := json.Unmarshal(body, &raw); err != nil {
return nil, fmt.Errorf("smarthome: decode states: %w", err)
}
out := make([]Entity, 0, len(raw))
for _, s := range raw {
domain := DomainOf(s.EntityID)
if domain == "" || !c.wanted(domain) {
continue
}
e := Entity{ID: s.EntityID, Domain: domain, Name: s.EntityID, State: s.State}
if len(s.Attributes) > 0 {
var a haAttrs
// Attributes are free-form per integration; a shape we cannot read
// costs the friendly name, not the entity.
if err := json.Unmarshal(s.Attributes, &a); err == nil {
if a.FriendlyName != "" {
e.Name = a.FriendlyName
}
e.Unit = a.Unit
}
}
out = append(out, e)
}
sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
if len(out) > c.cfg.MaxEntities {
out = out[:c.cfg.MaxEntities]
}
return out, nil
}
// CallService performs one service call against one entity and returns a short
// Russian confirmation.
//
// The entity id and service are NOT taken from the utterance: they come from
// the allowlist row that Kami enabled, so the router can only pick a row, never
// compose a target. That is the whole reason control is encoded in the cmd
// column instead of parsed out of speech.
func (c *Client) CallService(ctx context.Context, entityID, service string) (string, error) {
domain := DomainOf(entityID)
if domain == "" {
return "", ErrUnknownEntity
}
svcs := Services(domain)
if len(svcs) == 0 {
return "", ErrNotControllable
}
known := false
for _, s := range svcs {
if s.Name == service {
known = true
break
}
}
if !known {
return "", fmt.Errorf("%w: %s has no service %q", ErrNotControllable, domain, service)
}
payload, err := json.Marshal(map[string]string{"entity_id": entityID})
if err != nil {
return "", fmt.Errorf("smarthome: encode call: %w", err)
}
path := "/api/services/" + url.PathEscape(domain) + "/" + url.PathEscape(service)
if _, err := c.do(ctx, http.MethodPost, path, payload); err != nil {
return "", err
}
return "готово", nil
}
// do issues one authenticated request and returns the (capped) body.
func (c *Client) do(ctx context.Context, method, path string, body []byte) ([]byte, error) {
if c.cfg.URL == "" || c.cfg.Token == "" {
return nil, ErrNotConfigured
}
target := strings.TrimRight(c.cfg.URL, "/") + path
var rdr io.Reader
if body != nil {
rdr = bytes.NewReader(body)
}
req, err := http.NewRequestWithContext(ctx, method, target, rdr)
if err != nil {
return nil, fmt.Errorf("smarthome: request: %w", err)
}
req.Header.Set("Authorization", "Bearer "+c.cfg.Token)
req.Header.Set("Accept", "application/json")
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
resp, err := c.http.Do(req)
if err != nil {
return nil, fmt.Errorf("smarthome: %s %s: %w", method, path, err)
}
defer resp.Body.Close()
out, err := io.ReadAll(io.LimitReader(resp.Body, maxBody))
if err != nil {
return nil, fmt.Errorf("smarthome: read %s: %w", path, err)
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
// The body of an error can contain the instance's own detail; the token
// never appears in it, but keep it to one line anyway.
return nil, fmt.Errorf("smarthome: %s %s: http %d", method, path, resp.StatusCode)
}
return out, nil
}
+121
View File
@@ -0,0 +1,121 @@
// Package smarthome talks to a Home Assistant instance so Maven can read what
// the house is doing and change it (Vikunja #256,
// docs/plans/11-smarthome-integration.md).
//
// The shape of this package is copied deliberately from internal/mcp: a
// controllable entity becomes a PROPOSED row in the existing act allowlist,
// encoded in the columns that already exist — cmd
// ["smarthome", "<entity_id>", "<service>"], scope "smarthome:<domain>". So
// ProposeTool/EnableTool/DisableTool, tool.Matcher and the confirm turn need no
// change, and turning a light off in his flat goes through exactly the same
// gate as `restart nginx`.
//
// Two rules that are not negotiable here:
//
// - Discovery only ever PROPOSES. Finding a switch on the network is not the
// same as being allowed to flip it; Kami enables it on /tools, behind
// step-up.
// - Every control row is destructive=true. There is no read-only way to turn
// the heating off. That means a spoken act always gets the confirm turn,
// which is the point.
//
// MQTT / Zigbee2MQTT (steps 2 and 5 of the plan) are NOT here: they need a
// broker client dependency and the module cache in this repo is vendored, and
// there is no broker on this network to test one against. Home Assistant's REST
// API is stdlib-only and already fronts Zigbee2MQTT when it is present.
package smarthome
import (
"errors"
"strings"
)
var (
// ErrNotConfigured — no smarthome block, or it is disabled.
ErrNotConfigured = errors.New("smarthome: not configured")
// ErrUnknownEntity — the entity vanished between discovery and the call.
ErrUnknownEntity = errors.New("smarthome: unknown entity")
// ErrNotControllable — the entity's domain has no service Maven will call.
ErrNotControllable = errors.New("smarthome: entity is not controllable")
)
// cmdPrefix marks an allowlist row as a Home Assistant service call rather than
// a process. It is never run as a binary — tool.Executor branches on it before
// it ever reaches exec.
const cmdPrefix = "smarthome"
// Entity is one thing in the house, as Home Assistant sees it.
type Entity struct {
// ID — the Home Assistant entity_id, "light.living_room".
ID string
// Domain — the part before the dot. Decides which services apply.
Domain string
// Name — friendly_name when the instance has one, else ID.
Name string
// State — "on", "off", "22.5", …
State string
// Unit — unit_of_measurement, for sensors.
Unit string
}
// Service is one thing Maven can do to an entity.
type Service struct {
// Name — the Home Assistant service, "turn_on".
Name string
// Verb — the local suffix used to build the allowlist row name.
Verb string
}
// controllable maps a domain to the services Maven will expose for it. A domain
// that is not in this table gets no control row at all — the list is an
// allowlist, not a default, so a new HA integration cannot quietly hand her a
// verb nobody reviewed. set_temperature and set_brightness take a value and are
// deliberately absent: a spoken number that the router got wrong is a wrong act
// on real hardware, and on/off is the whole of what a voice turn can defend.
var controllable = map[string][]Service{
"light": {{Name: "turn_on", Verb: "on"}, {Name: "turn_off", Verb: "off"}},
"switch": {{Name: "turn_on", Verb: "on"}, {Name: "turn_off", Verb: "off"}},
"fan": {{Name: "turn_on", Verb: "on"}, {Name: "turn_off", Verb: "off"}},
"cover": {{Name: "open_cover", Verb: "open"}, {Name: "close_cover", Verb: "close"}},
"lock": {{Name: "lock", Verb: "lock"}, {Name: "unlock", Verb: "unlock"}},
}
// Services returns the services exposed for an entity, nil when its domain is
// not controllable (a sensor, a person, a weather entity: readable, not
// flippable).
func Services(domain string) []Service { return controllable[domain] }
// DomainOf splits "light.living_room" into "light". Empty when the id has no
// dot, which Home Assistant guarantees it does.
func DomainOf(entityID string) string {
i := strings.IndexByte(entityID, '.')
if i <= 0 {
return ""
}
return entityID[:i]
}
// LocalName is the allowlist row name for one entity+service. Prefixed so a
// house row is recognisable on /tools without opening the config, and so it
// cannot collide with a shell tool Kami named himself.
func LocalName(entityID, verb string) string {
return "home_" + strings.ReplaceAll(entityID, ".", "_") + "_" + verb
}
// Scope is the store scope for an entity's domain.
func Scope(domain string) string { return cmdPrefix + ":" + domain }
// Cmd is the allowlist cmd column for an entity+service.
func Cmd(entityID, service string) []string { return []string{cmdPrefix, entityID, service} }
// ParseCmd recognises a Home Assistant row. ok=false ⇒ an ordinary process row,
// and the caller execs it as it always did.
func ParseCmd(cmd []string) (entityID, service string, ok bool) {
if len(cmd) != 3 || cmd[0] != cmdPrefix {
return "", "", false
}
if cmd[1] == "" || cmd[2] == "" {
return "", "", false
}
return cmd[1], cmd[2], true
}
+211
View File
@@ -0,0 +1,211 @@
package smarthome
import (
"context"
"errors"
"net/http"
"net/http/httptest"
"strings"
"testing"
)
// statesFixture is a trimmed /api/states response from a Home Assistant with
// one light, one switch, one sensor and two entities Maven must ignore.
const statesFixture = `[
{"entity_id":"light.living_room","state":"on","attributes":{"friendly_name":"Гостиная"}},
{"entity_id":"switch.kettle","state":"off","attributes":{"friendly_name":"Чайник"}},
{"entity_id":"sensor.bedroom_temp","state":"22.5","attributes":{"unit_of_measurement":"°C"}},
{"entity_id":"person.kami","state":"home","attributes":{}},
{"entity_id":"automation.wake","state":"on","attributes":[]}
]`
func newTestClient(t *testing.T, h http.HandlerFunc) (*Client, *httptest.Server) {
t.Helper()
srv := httptest.NewServer(h)
t.Cleanup(srv.Close)
c := NewClient(Config{URL: srv.URL, Token: "tok"})
return c, srv
}
func TestStatesFiltersAndNames(t *testing.T) {
var auth string
c, _ := newTestClient(t, func(w http.ResponseWriter, r *http.Request) {
auth = r.Header.Get("Authorization")
if r.URL.Path != "/api/states" {
t.Errorf("path = %q", r.URL.Path)
}
_, _ = w.Write([]byte(statesFixture))
})
got, err := c.States(context.Background())
if err != nil {
t.Fatalf("States: %v", err)
}
if auth != "Bearer tok" {
t.Errorf("Authorization = %q", auth)
}
// person and automation are neither controllable nor sensors.
want := []string{"light.living_room", "sensor.bedroom_temp", "switch.kettle"}
if len(got) != len(want) {
t.Fatalf("got %d entities, want %d: %+v", len(got), len(want), got)
}
for i, id := range want {
if got[i].ID != id {
t.Errorf("entity %d = %q, want %q (sorted by id)", i, got[i].ID, id)
}
}
if got[0].Name != "Гостиная" || got[0].Domain != "light" || got[0].State != "on" {
t.Errorf("light = %+v", got[0])
}
if got[1].Unit != "°C" {
t.Errorf("sensor unit = %q", got[1].Unit)
}
// An attributes value of the wrong shape must not lose the entity.
if got[2].Name != "Чайник" {
t.Errorf("switch name = %q", got[2].Name)
}
}
func TestStatesRespectsConfiguredDomainsAndCap(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte(statesFixture))
}))
defer srv.Close()
c := NewClient(Config{URL: srv.URL, Token: "t", Domains: []string{"switch"}})
got, err := c.States(context.Background())
if err != nil {
t.Fatalf("States: %v", err)
}
if len(got) != 1 || got[0].ID != "switch.kettle" {
t.Fatalf("domain filter: %+v", got)
}
c = NewClient(Config{URL: srv.URL, Token: "t", MaxEntities: 2})
got, err = c.States(context.Background())
if err != nil {
t.Fatalf("States: %v", err)
}
if len(got) != 2 {
t.Fatalf("cap: got %d entities, want 2", len(got))
}
}
func TestCallServicePostsEntityID(t *testing.T) {
var path, body string
c, _ := newTestClient(t, func(w http.ResponseWriter, r *http.Request) {
path = r.URL.Path
b := make([]byte, 256)
n, _ := r.Body.Read(b)
body = string(b[:n])
_, _ = w.Write([]byte(`[]`))
})
out, err := c.CallService(context.Background(), "light.living_room", "turn_off")
if err != nil {
t.Fatalf("CallService: %v", err)
}
if out != "готово" {
t.Errorf("out = %q", out)
}
if path != "/api/services/light/turn_off" {
t.Errorf("path = %q", path)
}
if !strings.Contains(body, `"entity_id":"light.living_room"`) {
t.Errorf("body = %q", body)
}
}
// A service that is not in the domain's table never leaves the box. The
// allowlist is the gate, and it is enforced on the way out too, so a corrupted
// or hand-edited cmd column cannot reach an arbitrary Home Assistant service.
func TestCallServiceRefusesUnknownServiceAndDomain(t *testing.T) {
called := false
c, _ := newTestClient(t, func(w http.ResponseWriter, r *http.Request) {
called = true
_, _ = w.Write([]byte(`[]`))
})
for _, tc := range []struct {
entity, service string
want error
}{
{"light.living_room", "delete_everything", ErrNotControllable},
{"sensor.bedroom_temp", "turn_on", ErrNotControllable},
{"nodot", "turn_on", ErrUnknownEntity},
} {
if _, err := c.CallService(context.Background(), tc.entity, tc.service); !errors.Is(err, tc.want) {
t.Errorf("CallService(%q,%q) err = %v, want %v", tc.entity, tc.service, err, tc.want)
}
}
if called {
t.Error("a refused call still reached the network")
}
}
func TestHTTPErrorIsAnError(t *testing.T) {
c, _ := newTestClient(t, func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusUnauthorized)
})
if _, err := c.States(context.Background()); err == nil {
t.Fatal("want error on 401")
}
}
func TestUnconfiguredClientRefuses(t *testing.T) {
c := NewClient(Config{})
if _, err := c.States(context.Background()); !errors.Is(err, ErrNotConfigured) {
t.Fatalf("err = %v, want ErrNotConfigured", err)
}
}
func TestValidate(t *testing.T) {
ok := Config{URL: "http://ha.lan:8123", Token: "t"}
if err := Validate(ok); err != nil {
t.Fatalf("Validate(ok): %v", err)
}
for name, c := range map[string]Config{
"no url": {Token: "t"},
"no token": {URL: "http://ha.lan:8123"},
"bad scheme": {URL: "ftp://ha.lan", Token: "t"},
"no host": {URL: "http://", Token: "t"},
"not a url": {URL: "://x", Token: "t"},
"bare string": {URL: "ha.lan:8123", Token: "t"},
} {
if err := Validate(c); err == nil {
t.Errorf("Validate(%s) = nil, want error", name)
}
}
}
func TestAllowlistEncoding(t *testing.T) {
cmd := Cmd("light.living_room", "turn_off")
id, svc, ok := ParseCmd(cmd)
if !ok || id != "light.living_room" || svc != "turn_off" {
t.Fatalf("ParseCmd(%v) = %q,%q,%v", cmd, id, svc, ok)
}
// Anything that is not exactly a three-element smarthome row stays a
// process row, or the executor would swallow a real shell tool.
for _, bad := range [][]string{
nil,
{"smarthome"},
{"smarthome", "light.x"},
{"smarthome", "light.x", "turn_on", "extra"},
{"smarthome", "", "turn_on"},
{"smarthome", "light.x", ""},
{"systemctl", "restart", "nginx"},
} {
if _, _, ok := ParseCmd(bad); ok {
t.Errorf("ParseCmd(%v) = ok, want not a smarthome row", bad)
}
}
if got := LocalName("light.living_room", "off"); got != "home_light_living_room_off" {
t.Errorf("LocalName = %q", got)
}
if got := Scope("light"); got != "smarthome:light" {
t.Errorf("Scope = %q", got)
}
if Services("light") == nil || Services("sensor") != nil {
t.Error("Services: light must be controllable and sensor must not")
}
if DomainOf("light.x") != "light" || DomainOf("nodot") != "" || DomainOf(".x") != "" {
t.Error("DomainOf")
}
}
+13
View File
@@ -95,6 +95,19 @@ func (s *Store) ProposeMCPTool(ctx context.Context, name, scope string, cmd []st
return n > 0, nil
}
// ProposeSmartHomeTool is ProposeTool for a controllable device discovered on
// the Home Assistant instance (Vikunja #256). Like ProposeMCPTool the proposal
// already knows what it would run, so cmd is written with it and Kami only has
// to press enable.
//
// It is still a PROPOSAL, and destructive is not a parameter: there is no
// read-only way to turn a lamp off, so every house row carries the confirm
// turn. Re-discovery on every refresh is idempotent — an existing row is never
// touched, so a device he disabled stays disabled.
func (s *Store) ProposeSmartHomeTool(ctx context.Context, name, scope string, cmd []string, utterance string, ts time.Time) (bool, error) {
return s.ProposeMCPTool(ctx, name, scope, cmd, true, utterance, ts)
}
// EnableTool fills cmd + destructive and flips status to 'enabled'. This is the
// human "enable" act (the authed surface calls it); it upserts so enabling a
// name that was never proposed still works. An empty cmd is refused — an
+37
View File
@@ -13,6 +13,11 @@
// - Args are passed as argv, NEVER through a shell. STT text lands as
// positional arguments to Cmd; there is no `sh -c`, so "restart nginx;
// rm -rf" can't inject — the tail is one argv element to the named binary.
// - An enabled row whose cmd is ["smarthome", "<entity_id>", "<service>"] is
// a Home Assistant service call instead of a process (Vikunja #256), by
// exactly the same trick and under exactly the same rules. Control rows are
// always destructive, so flipping something in his flat always costs a
// confirm turn.
// - An enabled row whose cmd is ["mcp", "<server>", "<tool>"] is a call to a
// configured MCP server instead of a process (Vikunja #251). It goes
// through every rule above unchanged — enabled, and confirmed if it
@@ -38,6 +43,7 @@ import (
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/mcp"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/smarthome"
)
// API — the narrow slice of ipc.CoreAPI the executor and matcher need. Backed
@@ -64,6 +70,14 @@ type MCPCaller interface {
CallPositional(ctx context.Context, server, tool string, args []string) (string, error)
}
// HomeCaller is the seam for an act that is a Home Assistant service call
// rather than a process (Vikunja #256). internal/smarthome.Client satisfies it.
// nil ⇒ the house is not configured, and a house row refuses to run rather than
// silently doing nothing.
type HomeCaller interface {
CallService(ctx context.Context, entityID, service string) (string, error)
}
// Executor runs enabled tools. run is the exec seam (default: real process);
// tests swap it. timeout bounds each invocation.
type Executor struct {
@@ -71,6 +85,7 @@ type Executor struct {
timeout time.Duration
run func(ctx context.Context, argv []string) (string, error)
mcp MCPCaller
home HomeCaller
}
// NewExecutor builds the executor. timeout<=0 defaults to 30s.
@@ -88,6 +103,14 @@ func (e *Executor) WithMCP(m MCPCaller) *Executor {
return e
}
// WithHome attaches the Home Assistant caller. Called once at wiring time when
// the smarthome block is enabled; without it, a row whose cmd is
// ["smarthome", …] refuses.
func (e *Executor) WithHome(h HomeCaller) *Executor {
e.home = h
return e
}
// Exec looks up name in the store and runs Cmd+args as argv (no shell).
// confirmed=true is the second turn of a destructive act (the user said "да");
// it bypasses the ErrNeedsConfirm gate. Non-enabled ⇒ ErrNotEnabled; a
@@ -117,6 +140,20 @@ func (e *Executor) Exec(ctx context.Context, name string, args []string, confirm
defer cancel()
return e.mcp.CallPositional(ctx, server, remote, args)
}
// A house row is a Home Assistant service call, not a process (Vikunja
// #256). Same story: enabled, and confirmed — every control row is
// destructive, because there is no read-only way to turn the heating off.
// The spoken args are dropped on purpose: the entity and the service come
// from the row Kami enabled, so a router that misheard can pick the wrong
// row but can never compose a target of its own.
if entityID, service, ok := smarthome.ParseCmd(t.Cmd); ok {
if e.home == nil {
return "", ErrNotEnabled
}
ctx, cancel := context.WithTimeout(ctx, e.timeout)
defer cancel()
return e.home.CallService(ctx, entityID, service)
}
argv := append(append([]string(nil), t.Cmd...), args...)
if len(argv) == 0 {
return "", ErrNotEnabled
+77
View File
@@ -185,3 +185,80 @@ func TestExecMCPRowWithoutCallerRefuses(t *testing.T) {
t.Fatal(`"mcp" must never be run as a binary`)
}
}
// fakeHome records what the executor asked the house to do.
type fakeHome struct {
entity, service string
calls int
}
func (f *fakeHome) CallService(_ context.Context, entityID, service string) (string, error) {
f.calls++
f.entity, f.service = entityID, service
return "готово", nil
}
// A house row goes through the same allowlist and the same confirm turn as any
// other act, and it is never exec'd as a binary (Vikunja #256).
func TestExecSmartHomeRow(t *testing.T) {
api := fakeAPI{tools: map[string]ipc.Tool{
"home_light_x_off": {
Name: "home_light_x_off", Scope: "smarthome:light",
Cmd: []string{"smarthome", "light.x", "turn_off"}, Destructive: true, Status: "enabled",
},
"home_draft": {
Name: "home_draft", Scope: "smarthome:light",
Cmd: []string{"smarthome", "light.y", "turn_on"}, Destructive: true, Status: "proposed",
},
}}
ran := false
newExec := func(h HomeCaller) *Executor {
e := NewExecutor(api, time.Second)
e.run = func(context.Context, []string) (string, error) { ran = true; return "", nil }
if h != nil {
e = e.WithHome(h)
}
return e
}
// No house configured ⇒ the row refuses rather than being exec'd.
if _, err := newExec(nil).Exec(context.Background(), "home_light_x_off", nil, true); !errors.Is(err, ErrNotEnabled) {
t.Fatalf("unconfigured house: err = %v, want ErrNotEnabled", err)
}
if ran {
t.Fatal(`"smarthome" was run as a binary`)
}
// Configured, but not confirmed ⇒ the confirm turn, before any call.
fh := &fakeHome{}
if _, err := newExec(fh).Exec(context.Background(), "home_light_x_off", nil, false); !errors.Is(err, ErrNeedsConfirm) {
t.Fatalf("err = %v, want ErrNeedsConfirm", err)
}
if fh.calls != 0 {
t.Fatal("an unconfirmed house act reached the house")
}
// A merely proposed row never runs, confirmed or not.
if _, err := newExec(fh).Exec(context.Background(), "home_draft", nil, true); !errors.Is(err, ErrNotEnabled) {
t.Fatalf("proposed row: err = %v, want ErrNotEnabled", err)
}
if fh.calls != 0 {
t.Fatal("a proposed house row reached the house")
}
// Confirmed ⇒ the service call, with the entity from the ROW and the
// spoken tail dropped.
out, err := newExec(fh).Exec(context.Background(), "home_light_x_off", []string{"light.somewhere_else"}, true)
if err != nil {
t.Fatalf("Exec: %v", err)
}
if out != "готово" {
t.Errorf("out = %q", out)
}
if fh.entity != "light.x" || fh.service != "turn_off" {
t.Errorf("called %s/%s: the target must come from the enabled row, never from the utterance", fh.entity, fh.service)
}
if ran {
t.Fatal(`"smarthome" was run as a binary`)
}
}
+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"