Compare commits

...

5 Commits

Author SHA1 Message Date
kami 69e2800ef3 Cover ecosystem degraded modes with a shared fault-injection harness (#276)
Extend the fake Nexus/Praxis/Hexis harness with request header and query
capture, a malformed-body lever, a response delay lever, and a request
counter, then add a degraded-mode suite on top of it: independent outages,
malformed and drifted contracts, cancellation, execution failure vs
transport failure, ambiguous targets, no autonomous Praxis to Hexis
chaining, confirmation for mutating capabilities, and recovery without a
restart.
2026-08-01 06:47:55 +04:00
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
41 changed files with 4766 additions and 23 deletions
+9 -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: 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
.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/...
+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
+312
View File
@@ -0,0 +1,312 @@
package main
import (
"context"
"net/http"
"strings"
"testing"
"time"
hexisclient "github.com/kami/hexis/pkg/client"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/store"
)
// Phase-5 hardening suite (Vikunja #276). Everything here drives the shared
// fake ecosystem (fakeecosystem_test.go) rather than one-off inline handlers,
// so the same fault levers — SetFault, SetBody, SetDelay — cover every
// service. What is asserted is the degraded-mode contract:
//
// - services degrade independently: one outage never mutes the others,
// - a degraded reply is never silent, never fabricated, never "success",
// - contract drift (old shape, unknown fields, garbage) is survivable,
// - Maven never acts on an ambiguous target and never chains
// Praxis observation into Hexis execution on its own.
// ecoHandler wires a handler against whichever of the three fakes is given
// (pass nil to leave a service unconfigured, which is a different state from
// "configured but down").
func ecoHandler(t *testing.T, nexus, praxis, hexis *fakeServer) *reactiveHandler {
t.Helper()
st := newTestStore(t)
clock := newFakeClock(time.Date(2026, 8, 1, 9, 0, 0, 0, time.UTC))
w := &ecosystemWiring{}
if nexus != nil {
w.nexus = newNexusClient(nexus.URL)
}
if praxis != nil {
w.praxis = newPraxisClient(praxis.URL)
}
if hexis != nil {
w.hexis = hexisclient.New(hexis.URL)
}
return &reactiveHandler{
api: ipc.NewStoreAPI(st),
dataStore: st,
now: clock.Now,
ecosystem: w,
}
}
func traceFacts(t *testing.T, h *reactiveHandler) []store.Fact {
t.Helper()
facts, err := h.dataStore.RecentFacts(context.Background(), 50)
if err != nil {
t.Fatalf("read facts: %v", err)
}
var out []store.Fact
for _, f := range facts {
if f.Source == "praxis:trace" {
out = append(out, f)
}
}
return out
}
func restartCaps() string {
return fixtureHexisCapabilities(map[string]any{
"id": "cap_restart", "name": "restart", "read_only": true,
})
}
// TestEcosystem_OutagesAreIndependent: Praxis being down must not disable the
// Nexus+Hexis action path, and vice versa. A shared "ecosystem is broken"
// mode would take away working capability for no reason.
func TestEcosystem_OutagesAreIndependent(t *testing.T) {
ctx := context.Background()
nexus := newFakeNexus(t, fixtureNexusResolved("ent_muzick", "Muzick indexer", "service"))
praxis := newFakePraxis(t, fixturePraxisAttentionItems(map[string]any{
"id": "item_1", "title": "disk almost full", "importance": 3.0,
}))
hexis := newFakeHexis(t, restartCaps(), fixtureHexisExecuted("exec_1", "succeeded"))
h := ecoHandler(t, nexus, praxis, hexis)
praxis.SetFault(503)
if reply := h.handleHexisAct(ctx, actDec("muzick indexer")); !strings.Contains(reply, "выполнена") {
t.Fatalf("praxis outage must not block the hexis path, got %q", reply)
}
praxis.SetFault(0)
hexis.SetFault(503)
nexus.SetFault(503)
reply := h.handlePraxisAct(ctx, praxisActDec("list_attention"))
if !strings.Contains(reply, "disk almost full") {
t.Fatalf("nexus/hexis outage must not block the praxis digest, got %q", reply)
}
}
// TestEcosystem_MalformedNexusResponseFailsClosed: a 200 carrying garbage is a
// dependency failure, not "no such entity". It must stop before Hexis.
func TestEcosystem_MalformedNexusResponseFailsClosed(t *testing.T) {
ctx := context.Background()
nexus := newFakeNexus(t, fixtureNexusResolved("ent_muzick", "Muzick indexer", "service"))
hexis := newFakeHexis(t, restartCaps(), fixtureHexisExecuted("exec_1", "succeeded"))
h := ecoHandler(t, nexus, nil, hexis)
nexus.SetBody(`{"status":"resolved","entity":`)
reply := h.handleHexisAct(ctx, actDec("muzick indexer"))
if reply == "" || strings.Contains(reply, "выполнена") {
t.Fatalf("malformed nexus body must degrade, got %q", reply)
}
if hexis.Count("", "/api/v1") != 0 {
t.Fatal("hexis must not be contacted after a malformed nexus response")
}
}
// TestEcosystem_UnknownContractFieldsTolerated: a newer Nexus adding fields
// must not break an older Maven. Same for the older flat resolve shape.
func TestEcosystem_UnknownContractFieldsTolerated(t *testing.T) {
ctx := context.Background()
for name, body := range map[string]string{
"future": fixtureNexusResolvedFuture("ent_muzick", "Muzick indexer", "service"),
"flat": fixtureNexusResolvedFlat("ent_muzick", "Muzick indexer", "service"),
} {
t.Run(name, func(t *testing.T) {
nexus := newFakeNexus(t, body)
hexis := newFakeHexis(t, restartCaps(), fixtureHexisExecuted("exec_1", "succeeded"))
h := ecoHandler(t, nexus, nil, hexis)
if reply := h.handleHexisAct(ctx, actDec("muzick indexer")); !strings.Contains(reply, "выполнена") {
t.Fatalf("%s contract shape must still resolve and execute, got %q", name, reply)
}
})
}
}
// TestEcosystem_CancelledContextDegrades: a caller hanging up (turn abandoned,
// deadline hit) must surface as degradation, never as a fabricated result.
func TestEcosystem_CancelledContextDegrades(t *testing.T) {
nexus := newFakeNexus(t, fixtureNexusResolved("ent_muzick", "Muzick indexer", "service"))
hexis := newFakeHexis(t, restartCaps(), fixtureHexisExecuted("exec_1", "succeeded"))
h := ecoHandler(t, nexus, nil, hexis)
nexus.SetDelay(2 * time.Second)
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond)
defer cancel()
reply := h.handleHexisAct(ctx, actDec("muzick indexer"))
if reply == "" || strings.Contains(reply, "выполнена") {
t.Fatalf("cancelled resolve must degrade, got %q", reply)
}
if hexis.Count("", "/api/v1") != 0 {
t.Fatal("hexis must not be contacted after a cancelled resolve")
}
}
// TestEcosystem_ExecutionFailureIsNotSuccess: Hexis answering 200 with
// status=failed is a partial failure — the call worked, the command did not.
// Maven must report it as a failure and must not write a success trace.
func TestEcosystem_ExecutionFailureIsNotSuccess(t *testing.T) {
ctx := context.Background()
nexus := newFakeNexus(t, fixtureNexusResolved("ent_muzick", "Muzick indexer", "service"))
hexis := newFakeHexis(t, restartCaps(), fixtureHexisExecutionFailed("exec_1", "unit not found"))
h := ecoHandler(t, nexus, nil, hexis)
reply := h.handleHexisAct(ctx, actDec("muzick indexer"))
if strings.Contains(reply, "выполнена") {
t.Fatalf("failed execution must not read as success, got %q", reply)
}
if reply == "" {
t.Fatal("failed execution must say something")
}
for _, f := range traceFacts(t, h) {
if strings.HasPrefix(f.Key, "praxis:hexis:") {
t.Fatalf("failed execution must not write a success trace: %+v", f)
}
}
}
// TestEcosystem_AmbiguousTargetBlocksExecution: ambiguity blocks mutation, and
// the clarification must name the candidates rather than pick one.
func TestEcosystem_AmbiguousTargetBlocksExecution(t *testing.T) {
ctx := context.Background()
nexus := newFakeNexus(t, fixtureNexusAmbiguous(
map[string]string{"entity_id": "ent_a", "display_name": "Muzick indexer"},
map[string]string{"entity_id": "ent_b", "display_name": "Muzick web"},
))
hexis := newFakeHexis(t, restartCaps(), fixtureHexisExecuted("exec_1", "succeeded"))
h := ecoHandler(t, nexus, nil, hexis)
reply := h.handleHexisAct(ctx, actDec("muzick"))
if !strings.Contains(reply, "Muzick indexer") || !strings.Contains(reply, "Muzick web") {
t.Fatalf("ambiguous resolve must list candidates, got %q", reply)
}
if hexis.Count("POST", "/api/v1/execute") != 0 {
t.Fatal("ambiguous target must never execute")
}
}
// TestEcosystem_NoAutonomousPraxisToHexis: reading the attention digest is an
// observation. Maven must never turn an observed problem into a Hexis command
// by herself — she is not autonomous.
func TestEcosystem_NoAutonomousPraxisToHexis(t *testing.T) {
ctx := context.Background()
praxis := newFakePraxis(t, fixturePraxisAttentionItems(
map[string]any{"id": "item_1", "title": "muzick indexer is down", "importance": 4.0, "rule": "service_down"},
))
nexus := newFakeNexus(t, fixtureNexusResolved("ent_muzick", "Muzick indexer", "service"))
hexis := newFakeHexis(t, restartCaps(), fixtureHexisExecuted("exec_1", "succeeded"))
h := ecoHandler(t, nexus, praxis, hexis)
_ = h.handlePraxisAct(ctx, praxisActDec("list_attention"))
if hexis.Count("", "/api/v1") != 0 {
t.Fatal("attention digest must not contact hexis on its own")
}
if nexus.Count("", "/api/v1/resolve") != 0 {
t.Fatal("attention digest must not resolve targets for autonomous action")
}
}
// TestEcosystem_MutatingCapabilityWaitsForConfirmation: a non-read-only
// capability parks for an explicit spoken confirm bound to capability+target.
func TestEcosystem_MutatingCapabilityWaitsForConfirmation(t *testing.T) {
ctx := context.Background()
nexus := newFakeNexus(t, fixtureNexusResolved("ent_muzick", "Muzick indexer", "service"))
caps := fixtureHexisCapabilities(map[string]any{"id": "cap_restart", "name": "restart", "read_only": false})
hexis := newFakeHexis(t, caps, fixtureHexisExecuted("exec_1", "succeeded"))
h := ecoHandler(t, nexus, nil, hexis)
reply := h.handleHexisAct(ctx, actDec("restart"))
if !strings.Contains(reply, "restart") || !strings.Contains(reply, "да") {
t.Fatalf("mutating capability must ask for confirmation, got %q", reply)
}
if hexis.Count("POST", "/api/v1/execute") != 0 {
t.Fatal("mutating capability must not execute before confirmation")
}
h.mu.Lock()
pending := h.pendingHexis
h.mu.Unlock()
if pending == nil || pending.capabilityID != "cap_restart" || pending.entityID != "ent_muzick" {
t.Fatalf("confirmation must be bound to capability+target, got %+v", pending)
}
}
// TestEcosystem_SurfaceFailureStillDelivers: surfacing is bookkeeping. If the
// surface call fails the digest must still be spoken — a partial failure
// downgrades bookkeeping, not the answer.
func TestEcosystem_SurfaceFailureStillDelivers(t *testing.T) {
ctx := context.Background()
praxis := newFakeServer(t, map[string]http.HandlerFunc{
"GET /api/v1/tools/attention": jsonHandler(200, fixturePraxisAttentionItems(
map[string]any{"id": "item_1", "title": "disk almost full", "importance": 3.0},
)),
"POST /api/v1/tools/surface": jsonHandler(500, `{"error":"boom"}`),
})
h := ecoHandler(t, nil, praxis, nil)
reply := h.handlePraxisAct(ctx, praxisActDec("list_attention"))
if !strings.Contains(reply, "disk almost full") {
t.Fatalf("failed surface must not swallow the digest, got %q", reply)
}
if praxis.Count("POST", "/api/v1/tools/surface") == 0 {
t.Fatal("expected the surface attempt")
}
}
// TestEcosystem_TotalOutageSaysSoForEveryPath: with all three down, every
// entry point degrades explicitly instead of returning empty or inventing.
func TestEcosystem_TotalOutageSaysSoForEveryPath(t *testing.T) {
ctx := context.Background()
nexus := newFakeNexus(t, fixtureNexusResolved("ent_muzick", "Muzick indexer", "service"))
praxis := newFakePraxis(t, fixturePraxisAttentionItems())
hexis := newFakeHexis(t, restartCaps(), fixtureHexisExecuted("exec_1", "succeeded"))
for _, fs := range []*fakeServer{nexus, praxis, hexis} {
fs.SetFault(503)
}
h := ecoHandler(t, nexus, praxis, hexis)
for name, reply := range map[string]string{
"hexis act": h.handleHexisAct(ctx, actDec("muzick indexer")),
"attention": h.handlePraxisAct(ctx, praxisActDec("list_attention")),
"changes": h.handlePraxisAct(ctx, praxisActDec("list_changes")),
"acknowledge": h.handlePraxisAct(ctx, praxisActDec("acknowledge_item")),
} {
if reply == "" {
t.Errorf("%s: total outage must not answer with silence", name)
}
if strings.Contains(reply, "выполнена") {
t.Errorf("%s: total outage must not claim success: %q", name, reply)
}
}
if len(traceFacts(t, h)) != 0 {
t.Fatal("a total outage must not leave success traces behind")
}
}
// TestEcosystem_RecoveryAfterOutageNeedsNoRestart: once the dependency comes
// back the very next turn works — no cached failure state, no restart.
func TestEcosystem_RecoveryAfterOutageNeedsNoRestart(t *testing.T) {
ctx := context.Background()
praxis := newFakePraxis(t, fixturePraxisAttentionItems(
map[string]any{"id": "item_1", "title": "disk almost full", "importance": 3.0},
))
h := ecoHandler(t, nil, praxis, nil)
praxis.SetFault(503)
if reply := h.handlePraxisAct(ctx, praxisActDec("list_attention")); strings.Contains(reply, "disk") {
t.Fatalf("outage must not serve content, got %q", reply)
}
praxis.SetFault(0)
if reply := h.handlePraxisAct(ctx, praxisActDec("list_attention")); !strings.Contains(reply, "disk almost full") {
t.Fatalf("recovery must work on the next turn, got %q", reply)
}
}
+96 -4
View File
@@ -14,7 +14,9 @@ import (
type capturedRequest struct {
Method string
Path string
Query string
Body []byte
Header http.Header
}
// fakeServer is the common shell behind fakeNexus/fakePraxis/fakeHexis: an
@@ -27,7 +29,9 @@ type fakeServer struct {
mu sync.Mutex
requests []capturedRequest
fault int // non-zero: every request gets this HTTP status instead of routing
fault int // non-zero: every request gets this HTTP status instead of routing
garbage string // non-empty: returned 200 verbatim instead of routing (malformed-contract lever)
delay time.Duration
}
// newFakeServer starts a server dispatching to routes keyed by "METHOD
@@ -47,14 +51,34 @@ func newFakeServer(t *testing.T, routes map[string]http.HandlerFunc) *fakeServer
}
}
fs.mu.Lock()
fs.requests = append(fs.requests, capturedRequest{Method: r.Method, Path: r.URL.Path, Body: body})
fs.requests = append(fs.requests, capturedRequest{
Method: r.Method,
Path: r.URL.Path,
Query: r.URL.RawQuery,
Body: body,
Header: r.Header.Clone(),
})
fault := fs.fault
garbage := fs.garbage
delay := fs.delay
fs.mu.Unlock()
if delay > 0 {
select {
case <-time.After(delay):
case <-r.Context().Done():
return
}
}
if fault != 0 {
http.Error(w, "injected fault", fault)
return
}
if garbage != "" {
w.Header().Set("Content-Type", "application/json")
w.Write([]byte(garbage))
return
}
for key, handler := range routes {
method, prefix := splitRouteKey(key)
@@ -90,6 +114,36 @@ func (fs *fakeServer) SetFault(status int) {
fs.fault = status
}
// SetBody makes every subsequent request answer 200 with the given body,
// bypassing the route table. Used to serve a malformed or contract-violating
// payload where the transport itself is healthy. Pass "" to clear it.
func (fs *fakeServer) SetBody(body string) {
fs.mu.Lock()
defer fs.mu.Unlock()
fs.garbage = body
}
// SetDelay stalls every subsequent request for d before answering, so callers
// can drive client timeouts and context cancellation deterministically. The
// delay is abandoned as soon as the client hangs up.
func (fs *fakeServer) SetDelay(d time.Duration) {
fs.mu.Lock()
defer fs.mu.Unlock()
fs.delay = d
}
// Count returns how many captured requests used the given method and path
// prefix. "" matches any method.
func (fs *fakeServer) Count(method, prefix string) int {
n := 0
for _, r := range fs.Requests() {
if (method == "" || r.Method == method) && hasPrefix(r.Path, prefix) {
n++
}
}
return n
}
// Requests returns a snapshot of captured requests, in arrival order.
func (fs *fakeServer) Requests() []capturedRequest {
fs.mu.Lock()
@@ -118,6 +172,32 @@ func fixtureNexusResolved(entityID, displayName, entityType string) string {
})
}
// fixtureNexusResolvedFlat is the flat resolve shape documented in
// ECOSYSTEM-SPEC.md §1.5 (entity_id/entity_type/display_name at the top
// level) rather than the nested "entity" object — the older of the two
// wire shapes Maven must keep accepting.
func fixtureNexusResolvedFlat(entityID, displayName, entityType string) string {
return mustJSON(map[string]any{
"status": "resolved",
"entity_id": entityID,
"entity_type": entityType,
"display_name": displayName,
})
}
// fixtureNexusResolvedFuture is a resolved response from a hypothetical newer
// Nexus: same required fields plus unknown ones. Decoding must ignore the
// extras, not fail — forward compatibility is what lets the ecosystem be
// upgraded one service at a time.
func fixtureNexusResolvedFuture(entityID, displayName, entityType string) string {
return mustJSON(map[string]any{
"status": "resolved",
"entity": map[string]any{"id": entityID, "display_name": displayName, "type": entityType, "tenant": "home"},
"provenance": map[string]any{"resolver": "v3", "graph_epoch": 42},
"score_breakdown": []any{map[string]any{"signal": "alias", "weight": 0.9}},
})
}
func fixtureNexusNotFound() string {
return `{"status":"not_found"}`
}
@@ -138,6 +218,13 @@ func fixtureHexisExecuted(id, status string) string {
return mustJSON(map[string]any{"id": id, "status": status})
}
// fixtureHexisExecutionFailed is a well-formed Hexis response reporting that
// the command itself failed: the call succeeded, the execution did not. Maven
// must distinguish this from a transport failure and from success.
func fixtureHexisExecutionFailed(id, message string) string {
return mustJSON(map[string]any{"id": id, "status": "failed", "error": message})
}
func fixturePraxisAttentionItems(items ...map[string]any) string {
return mustJSON(items)
}
@@ -191,8 +278,13 @@ func newFakeNexus(t *testing.T, resolveBody string) *fakeServer {
// fault is injected via SetFault.
func newFakePraxis(t *testing.T, attentionBody string) *fakeServer {
return newFakeServer(t, map[string]http.HandlerFunc{
"GET /api/v1/tools/attention": jsonHandler(http.StatusOK, attentionBody),
"POST /api/v1/tools/surface": jsonHandler(http.StatusOK, `{}`),
"GET /api/v1/tools/attention": jsonHandler(http.StatusOK, attentionBody),
"GET /api/v1/tools/changes": jsonHandler(http.StatusOK, `[]`),
"POST /api/v1/tools/surface": jsonHandler(http.StatusOK, `{}`),
"POST /api/v1/tools/acknowledge": jsonHandler(http.StatusOK, `{}`),
"POST /api/v1/tools/resolve": jsonHandler(http.StatusOK, `{}`),
"POST /api/v1/tools/ignore": jsonHandler(http.StatusOK, `{}`),
"POST /api/v1/tools/pin": jsonHandler(http.StatusOK, `{}`),
})
}
+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")
}
}
+34 -10
View File
@@ -204,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)
@@ -247,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)
}
@@ -297,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)
@@ -361,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.
@@ -482,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)
}
@@ -526,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)
@@ -599,6 +611,11 @@ func run(args []string) error {
go voiceW.mcp.run(ctx)
}
// 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
@@ -665,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),
+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 {
+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())
}
}
+25
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 ---
+9
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
@@ -593,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 {
+16
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
}
@@ -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 {
+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
}
+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`)
}
}