Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 08f3db318f | |||
| 69e2800ef3 | |||
| a8fcb404be | |||
| dc4c5b7841 | |||
| 33e53ee897 |
@@ -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
|
||||
|
||||
@@ -72,6 +72,7 @@ var praxisCapabilities = []praxisCapability{
|
||||
},
|
||||
},
|
||||
listChangesCapability{},
|
||||
entityAttentionCapability{},
|
||||
}
|
||||
|
||||
// handlePraxisAct — dispatches ecosystem tool acts through the Praxis tools API.
|
||||
@@ -190,6 +191,106 @@ func (listChangesCapability) handle(ctx context.Context, h *reactiveHandler, px
|
||||
return "изменения: " + strings.Join(parts, "; ")
|
||||
}
|
||||
|
||||
// entityAttentionCapability answers "what's going on with X" by resolving X to
|
||||
// a canonical Nexus entity and asking Praxis for that entity's attention items
|
||||
// (Vikunja #272). Unlike listAttentionCapability it is scoped: the entity_id
|
||||
// travels to Praxis as a query parameter instead of Maven filtering an unscoped
|
||||
// list client-side, which is what makes the ref canonical end to end.
|
||||
//
|
||||
// It also folds in what Maven herself knows about the same entity — facts the
|
||||
// enrichment worker has already resolved to that entity_id — so one question
|
||||
// gets one answer across both stores.
|
||||
type entityAttentionCapability struct{}
|
||||
|
||||
func (entityAttentionCapability) aliases() []string {
|
||||
return []string{"entity_attention", "что с", "как дела у", "статус"}
|
||||
}
|
||||
|
||||
func (entityAttentionCapability) handle(ctx context.Context, h *reactiveHandler, px *praxisClient, dec router.Decision) string {
|
||||
subject := dec.Slots.Value
|
||||
if subject == "" {
|
||||
subject = dec.Slots.Text
|
||||
}
|
||||
if subject == "" {
|
||||
return "про что именно спросить?"
|
||||
}
|
||||
if h.ecosystem == nil || h.ecosystem.nexus == nil {
|
||||
// Without Nexus there is no canonical ref to scope by. Say so rather
|
||||
// than quietly answering about something else.
|
||||
return "не могу связать это с сущностью — Nexus не настроен."
|
||||
}
|
||||
|
||||
entityID, displayName, ambiguous, err := h.ecosystem.resolveEntityReference(ctx, subject, nil)
|
||||
if err != nil {
|
||||
log.Printf("ecosystem: entity attention resolve %q: %v", subject, err)
|
||||
return "экосистема недоступна, попробуй ещё раз."
|
||||
}
|
||||
if len(ambiguous) > 0 {
|
||||
return "уточни, что именно: " + strings.Join(ambiguous, ", ") + "?"
|
||||
}
|
||||
if entityID == "" {
|
||||
return "не знаю такой сущности."
|
||||
}
|
||||
if displayName == "" {
|
||||
displayName = subject
|
||||
}
|
||||
|
||||
items, err := px.ListAttentionForEntity(ctx, entityID, 20)
|
||||
if err != nil {
|
||||
log.Printf("ecosystem: praxis attention for %s: %v", entityID, err)
|
||||
return "не могу сейчас узнать, что требует внимания по «" + displayName + "»."
|
||||
}
|
||||
h.recordPraxisTrace(ctx, "entity_attention", map[string]any{
|
||||
"entity_id": entityID, "count": len(items),
|
||||
})
|
||||
|
||||
var parts []string
|
||||
for _, item := range items {
|
||||
title, _ := item["title"].(string)
|
||||
if title == "" {
|
||||
continue
|
||||
}
|
||||
parts = append(parts, title)
|
||||
// Same surfaced != acknowledged rule as the unscoped digest.
|
||||
if id, ok := item["id"].(string); ok && id != "" {
|
||||
if _, err := px.Surface(ctx, id); err != nil {
|
||||
log.Printf("ecosystem: praxis surface %s: %v", id, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
if known := h.localFactsForEntity(ctx, entityID); known != "" {
|
||||
parts = append(parts, known)
|
||||
}
|
||||
if len(parts) == 0 {
|
||||
return "по «" + displayName + "» ничего нет."
|
||||
}
|
||||
return "по «" + displayName + "»: " + strings.Join(parts, "; ")
|
||||
}
|
||||
|
||||
// localFactsForEntity summarises Maven's own facts already resolved to this
|
||||
// canonical entity. Empty when the store is unavailable or nothing matched —
|
||||
// entity-scoped memory is an enrichment of the answer, never a precondition.
|
||||
func (h *reactiveHandler) localFactsForEntity(ctx context.Context, entityID string) string {
|
||||
if h.dataStore == nil || entityID == "" {
|
||||
return ""
|
||||
}
|
||||
facts, err := h.dataStore.FactsByEntity(ctx, entityID, 3)
|
||||
if err != nil {
|
||||
log.Printf("ecosystem: facts by entity %s: %v", entityID, err)
|
||||
return ""
|
||||
}
|
||||
var parts []string
|
||||
for _, f := range facts {
|
||||
if f.Value != "" {
|
||||
parts = append(parts, f.Value)
|
||||
}
|
||||
}
|
||||
if len(parts) == 0 {
|
||||
return ""
|
||||
}
|
||||
return "я помню: " + strings.Join(parts, ", ")
|
||||
}
|
||||
|
||||
// recordPraxisTrace — writes a fact recording a cross-service ecosystem call.
|
||||
// The fact is stored with source "praxis:trace" so the proactive loop can
|
||||
// reference it and the dashboard can display recent ecosystem activity.
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,213 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/router"
|
||||
"github.com/kami/maven/internal/store"
|
||||
)
|
||||
|
||||
// Entity-ref propagation, Maven side (Vikunja #272): the canonical Nexus
|
||||
// entity_id must reach Praxis as a query scope rather than being resolved and
|
||||
// then thrown away, and the enrichment that produces those ids must degrade
|
||||
// visibly instead of silently.
|
||||
|
||||
func entityAttentionDec(subject string) router.Decision {
|
||||
return router.Decision{
|
||||
Intent: router.IntentAct,
|
||||
Slots: router.Slots{Fn: "entity_attention", HasFn: true, Value: subject},
|
||||
}
|
||||
}
|
||||
|
||||
// TestEntityAttention_ScopesPraxisByCanonicalID: the resolved id must travel
|
||||
// to Praxis in the request, not be used for client-side filtering.
|
||||
func TestEntityAttention_ScopesPraxisByCanonicalID(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": "indexer queue is backing up", "importance": 3.0},
|
||||
))
|
||||
h := ecoHandler(t, nexus, praxis, nil)
|
||||
|
||||
reply := h.handlePraxisAct(ctx, entityAttentionDec("muzick indexer"))
|
||||
if !strings.Contains(reply, "indexer queue is backing up") {
|
||||
t.Fatalf("expected the scoped item in the reply, got %q", reply)
|
||||
}
|
||||
|
||||
var scoped bool
|
||||
for _, r := range praxis.Requests() {
|
||||
if r.Method == "GET" && strings.HasPrefix(r.Path, "/api/v1/tools/attention") &&
|
||||
strings.Contains(r.Query, "entity_id=ent_muzick") {
|
||||
scoped = true
|
||||
}
|
||||
}
|
||||
if !scoped {
|
||||
t.Fatalf("expected attention scoped by entity_id, got requests %+v", praxis.Requests())
|
||||
}
|
||||
if praxis.Count("POST", "/api/v1/tools/surface") == 0 {
|
||||
t.Error("a spoken scoped item must be surfaced, like the unscoped digest")
|
||||
}
|
||||
}
|
||||
|
||||
// TestEntityAttention_FoldsInLocalFactsForSameEntity: facts the enrichment
|
||||
// worker already tagged with the same canonical id join the same answer.
|
||||
func TestEntityAttention_FoldsInLocalFactsForSameEntity(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
nexus := newFakeNexus(t, fixtureNexusResolved("ent_espresso", "the espresso machine", "device"))
|
||||
praxis := newFakePraxis(t, fixturePraxisAttentionItems())
|
||||
h := ecoHandler(t, nexus, praxis, nil)
|
||||
|
||||
id, err := h.dataStore.WriteFactAboutSubject(ctx, time.Now(), store.KindEnv,
|
||||
"descaled", "the espresso machine", "descaled in june", "infer:pref", 0.8, sql.NullInt64{})
|
||||
if err != nil {
|
||||
t.Fatalf("WriteFactAboutSubject: %v", err)
|
||||
}
|
||||
if err := h.dataStore.ResolveFactEntity(ctx, id, "ent_espresso", store.ResolutionResolved); err != nil {
|
||||
t.Fatalf("ResolveFactEntity: %v", err)
|
||||
}
|
||||
|
||||
reply := h.handlePraxisAct(ctx, entityAttentionDec("the espresso machine"))
|
||||
if !strings.Contains(reply, "descaled in june") {
|
||||
t.Fatalf("expected entity-scoped local facts in the reply, got %q", reply)
|
||||
}
|
||||
}
|
||||
|
||||
// TestEntityAttention_AmbiguousAsksInsteadOfGuessing.
|
||||
func TestEntityAttention_AmbiguousAsksInsteadOfGuessing(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"},
|
||||
))
|
||||
praxis := newFakePraxis(t, fixturePraxisAttentionItems())
|
||||
h := ecoHandler(t, nexus, praxis, nil)
|
||||
|
||||
reply := h.handlePraxisAct(ctx, entityAttentionDec("muzick"))
|
||||
if !strings.Contains(reply, "Muzick indexer") || !strings.Contains(reply, "Muzick web") {
|
||||
t.Fatalf("ambiguous subject must ask, got %q", reply)
|
||||
}
|
||||
if praxis.Count("GET", "/api/v1/tools/attention") != 0 {
|
||||
t.Fatal("an ambiguous subject must not be queried against praxis")
|
||||
}
|
||||
}
|
||||
|
||||
// TestEntityAttention_MissingAndDegradedAreDistinct: "no such entity" and
|
||||
// "Nexus is down" must not produce the same answer.
|
||||
func TestEntityAttention_MissingAndDegradedAreDistinct(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
nexus := newFakeNexus(t, fixtureNexusNotFound())
|
||||
praxis := newFakePraxis(t, fixturePraxisAttentionItems())
|
||||
h := ecoHandler(t, nexus, praxis, nil)
|
||||
|
||||
missing := h.handlePraxisAct(ctx, entityAttentionDec("нечто"))
|
||||
if missing == "" {
|
||||
t.Fatal("an unknown entity must still get an answer")
|
||||
}
|
||||
|
||||
nexus.SetFault(503)
|
||||
degraded := h.handlePraxisAct(ctx, entityAttentionDec("нечто"))
|
||||
if degraded == missing {
|
||||
t.Fatalf("outage and unknown-entity must not read the same: %q", degraded)
|
||||
}
|
||||
}
|
||||
|
||||
// TestEntityAttention_DelayedNexusDegradesNotHangs: a slow Nexus past the
|
||||
// caller's deadline degrades and never queries Praxis with an empty scope.
|
||||
func TestEntityAttention_DelayedNexusDegradesNotHangs(t *testing.T) {
|
||||
nexus := newFakeNexus(t, fixtureNexusResolved("ent_muzick", "Muzick indexer", "service"))
|
||||
praxis := newFakePraxis(t, fixturePraxisAttentionItems())
|
||||
h := ecoHandler(t, nexus, praxis, nil)
|
||||
nexus.SetDelay(2 * time.Second)
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond)
|
||||
defer cancel()
|
||||
reply := h.handlePraxisAct(ctx, entityAttentionDec("muzick indexer"))
|
||||
if reply == "" {
|
||||
t.Fatal("a delayed resolve must still answer")
|
||||
}
|
||||
if praxis.Count("GET", "/api/v1/tools/attention") != 0 {
|
||||
t.Fatal("praxis must not be queried without a resolved scope")
|
||||
}
|
||||
}
|
||||
|
||||
// TestEntityAttention_WithoutNexusSaysSo: no Nexus means no canonical ref, so
|
||||
// the scoped query is refused rather than answered about something else.
|
||||
func TestEntityAttention_WithoutNexusSaysSo(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)
|
||||
|
||||
reply := h.handlePraxisAct(ctx, entityAttentionDec("muzick indexer"))
|
||||
if strings.Contains(reply, "disk almost full") {
|
||||
t.Fatalf("unscoped items must not be passed off as entity-scoped, got %q", reply)
|
||||
}
|
||||
if praxis.Count("GET", "/api/v1/tools/attention") != 0 {
|
||||
t.Fatal("no canonical ref means no scoped query at all")
|
||||
}
|
||||
}
|
||||
|
||||
// TestEnrichmentBackoff_HoldsAndReleases: repeated Nexus failures back the
|
||||
// fact off instead of hammering, and the fact is retried once the window
|
||||
// elapses. Nothing is ever given up on.
|
||||
func TestEnrichmentBackoff_HoldsAndReleases(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
nexus := newFakeNexus(t, fixtureNexusResolved("ent_espresso", "the espresso machine", "device"))
|
||||
st := newTestStore(t)
|
||||
if _, err := st.WriteFactAboutSubject(ctx, time.Now(), store.KindEnv, "likes",
|
||||
"the espresso machine", `"true"`, "infer:pref", 0.8, sql.NullInt64{}); err != nil {
|
||||
t.Fatalf("WriteFactAboutSubject: %v", err)
|
||||
}
|
||||
|
||||
clock := newFakeClock(time.Date(2026, 8, 1, 3, 0, 0, 0, time.UTC))
|
||||
w := newFactEnrichmentWorker(st, stubEcosystem(nexus.URL, ""), time.Hour)
|
||||
w.now = clock.Now
|
||||
|
||||
nexus.SetFault(503)
|
||||
w.tick(ctx)
|
||||
failedCalls := nexus.Count("POST", "/api/v1/resolve")
|
||||
if failedCalls != 1 {
|
||||
t.Fatalf("expected one resolve attempt, got %d", failedCalls)
|
||||
}
|
||||
|
||||
// Immediately after a failure the fact is in backoff: no second call.
|
||||
w.tick(ctx)
|
||||
if nexus.Count("POST", "/api/v1/resolve") != failedCalls {
|
||||
t.Fatal("a fact in backoff must not be retried on the very next tick")
|
||||
}
|
||||
if s := w.status(ctx); s.Pending != 1 || s.InBackoff != 1 || s.MaxAttempts != 1 {
|
||||
t.Fatalf("degradation must be reported, got %+v", s)
|
||||
}
|
||||
|
||||
// Once the window elapses and Nexus recovers, the fact resolves.
|
||||
clock.Advance(2 * time.Minute)
|
||||
nexus.SetFault(0)
|
||||
w.tick(ctx)
|
||||
facts, err := st.FactsByEntity(ctx, "ent_espresso", 10)
|
||||
if err != nil {
|
||||
t.Fatalf("FactsByEntity: %v", err)
|
||||
}
|
||||
if len(facts) != 1 {
|
||||
t.Fatalf("expected the fact resolved after recovery, got %+v", facts)
|
||||
}
|
||||
if s := w.status(ctx); s.Pending != 0 || s.MaxAttempts != 0 {
|
||||
t.Fatalf("recovery must clear the degradation report, got %+v", s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEnrichmentBackoff_GrowsAndIsCapped(t *testing.T) {
|
||||
if enrichmentBackoff(1) != time.Minute {
|
||||
t.Fatalf("first retry should be a minute, got %v", enrichmentBackoff(1))
|
||||
}
|
||||
if enrichmentBackoff(3) != 4*time.Minute {
|
||||
t.Fatalf("third retry should be four minutes, got %v", enrichmentBackoff(3))
|
||||
}
|
||||
if enrichmentBackoff(50) != time.Hour {
|
||||
t.Fatalf("backoff must cap at an hour, got %v", enrichmentBackoff(50))
|
||||
}
|
||||
}
|
||||
@@ -9,6 +9,7 @@ package main
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/store"
|
||||
@@ -24,10 +25,69 @@ type factEnrichmentWorker struct {
|
||||
eco *ecosystemWiring
|
||||
interval time.Duration
|
||||
batch int // facts resolved per tick; keeps a single slow tick bounded
|
||||
now func() time.Time
|
||||
|
||||
// Retry state for facts whose resolution failed transiently. Kept in
|
||||
// memory rather than in the DB: a restart legitimately retries
|
||||
// everything, and the backoff exists to spare a struggling Nexus, not
|
||||
// to be durable. A fact is never given up on — degraded means slower,
|
||||
// not dropped.
|
||||
mu sync.Mutex
|
||||
attempt map[int64]int // fact id → consecutive failures
|
||||
nextTry map[int64]time.Time // fact id → earliest retry
|
||||
skipped int // facts held back by backoff on the last tick
|
||||
}
|
||||
|
||||
// enrichmentBackoff is the wait before retrying a fact after n consecutive
|
||||
// failures, capped so a long Nexus outage still retries about hourly.
|
||||
func enrichmentBackoff(n int) time.Duration {
|
||||
d := time.Minute
|
||||
for i := 1; i < n && d < time.Hour; i++ {
|
||||
d *= 2
|
||||
}
|
||||
if d > time.Hour {
|
||||
d = time.Hour
|
||||
}
|
||||
return d
|
||||
}
|
||||
|
||||
func newFactEnrichmentWorker(st *store.Store, eco *ecosystemWiring, interval time.Duration) *factEnrichmentWorker {
|
||||
return &factEnrichmentWorker{store: st, eco: eco, interval: interval, batch: 20}
|
||||
return &factEnrichmentWorker{
|
||||
store: st,
|
||||
eco: eco,
|
||||
interval: interval,
|
||||
batch: 20,
|
||||
now: time.Now,
|
||||
attempt: map[int64]int{},
|
||||
nextTry: map[int64]time.Time{},
|
||||
}
|
||||
}
|
||||
|
||||
// enrichmentStatus is what the worker reports about its own health: how many
|
||||
// facts are waiting, how many are currently in backoff, and the worst retry
|
||||
// count seen. Degradation is reported, never hidden — a Nexus that has been
|
||||
// down all day must be visible as a backlog, not as facts that silently
|
||||
// never got tagged.
|
||||
type enrichmentStatus struct {
|
||||
Pending int
|
||||
InBackoff int
|
||||
MaxAttempts int
|
||||
}
|
||||
|
||||
func (w *factEnrichmentWorker) status(ctx context.Context) enrichmentStatus {
|
||||
var st enrichmentStatus
|
||||
if pending, err := w.store.PendingFactResolutions(ctx, 1000); err == nil {
|
||||
st.Pending = len(pending)
|
||||
}
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
st.InBackoff = w.skipped
|
||||
for _, n := range w.attempt {
|
||||
if n > st.MaxAttempts {
|
||||
st.MaxAttempts = n
|
||||
}
|
||||
}
|
||||
return st
|
||||
}
|
||||
|
||||
func (w *factEnrichmentWorker) run(ctx context.Context) {
|
||||
@@ -57,18 +117,50 @@ func (w *factEnrichmentWorker) tick(ctx context.Context) {
|
||||
log.Printf("factenrichment: list pending: %v", err)
|
||||
return
|
||||
}
|
||||
skipped, failed := 0, 0
|
||||
for _, f := range pending {
|
||||
w.resolveOne(ctx, f)
|
||||
if !w.due(f.ID) {
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
if !w.resolveOne(ctx, f) {
|
||||
failed++
|
||||
}
|
||||
}
|
||||
w.mu.Lock()
|
||||
w.skipped = skipped
|
||||
w.mu.Unlock()
|
||||
if failed > 0 {
|
||||
log.Printf("factenrichment: %d/%d resolutions failed this tick, %d held in backoff",
|
||||
failed, len(pending), skipped)
|
||||
}
|
||||
}
|
||||
|
||||
func (w *factEnrichmentWorker) resolveOne(ctx context.Context, f store.Fact) {
|
||||
// due reports whether a fact's backoff window has elapsed.
|
||||
func (w *factEnrichmentWorker) due(id int64) bool {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
next, ok := w.nextTry[id]
|
||||
return !ok || !w.now().Before(next)
|
||||
}
|
||||
|
||||
// resolveOne resolves one pending fact. It returns false when the attempt
|
||||
// failed transiently: the fact stays pending and is retried on a backoff.
|
||||
func (w *factEnrichmentWorker) resolveOne(ctx context.Context, f store.Fact) bool {
|
||||
entityID, _, ambiguous, err := w.eco.resolveEntityReference(ctx, f.Subject, nil)
|
||||
if err != nil {
|
||||
// Transient (Nexus unreachable) — leave pending, retry next tick.
|
||||
// Transient (Nexus unreachable) — leave pending, back off, retry later.
|
||||
log.Printf("factenrichment: resolve fact %d subject %q: %v", f.ID, f.Subject, err)
|
||||
return
|
||||
w.mu.Lock()
|
||||
w.attempt[f.ID]++
|
||||
w.nextTry[f.ID] = w.now().Add(enrichmentBackoff(w.attempt[f.ID]))
|
||||
w.mu.Unlock()
|
||||
return false
|
||||
}
|
||||
w.mu.Lock()
|
||||
delete(w.attempt, f.ID)
|
||||
delete(w.nextTry, f.ID)
|
||||
w.mu.Unlock()
|
||||
state := store.ResolutionNotFound
|
||||
switch {
|
||||
case entityID != "":
|
||||
@@ -78,5 +170,7 @@ func (w *factEnrichmentWorker) resolveOne(ctx context.Context, f store.Fact) {
|
||||
}
|
||||
if err := w.store.ResolveFactEntity(ctx, f.ID, entityID, state); err != nil {
|
||||
log.Printf("factenrichment: record resolution for fact %d: %v", f.ID, err)
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
@@ -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, `{}`),
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -611,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
|
||||
@@ -677,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
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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),
|
||||
|
||||
@@ -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" },
|
||||
|
||||
@@ -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"
|
||||
)
|
||||
@@ -237,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
|
||||
@@ -263,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
|
||||
@@ -884,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.
|
||||
@@ -1113,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 {
|
||||
@@ -1223,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,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