Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 69e2800ef3 | |||
| a8fcb404be | |||
| dc4c5b7841 | |||
| 33e53ee897 | |||
| 45b5e16eff |
@@ -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/...
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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, `{}`),
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -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, ""
|
||||
}
|
||||
@@ -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
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
@@ -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()
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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>
|
||||
@@ -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, "<script>") {
|
||||
t.Error("intake title is missing from the page entirely")
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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" },
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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())
|
||||
}
|
||||
}
|
||||
@@ -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 ---
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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`)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user