Compare commits
42 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7203cd56fd | |||
| ab3e818bb9 | |||
| b5500a5be8 | |||
| 9095ac847d | |||
| 4b5f6adbae | |||
| ecb8ba72eb | |||
| 4fdecf9a25 | |||
| 2bbd8edbf6 | |||
| beb093aebb | |||
| b1b326018f | |||
| 08889cad88 | |||
| a4630b9314 | |||
| 39d44bb384 | |||
| 65ee0f9c61 | |||
| 76938e206d | |||
| 0b3d81ecbf | |||
| 4be6852b94 | |||
| f7b76c572f | |||
| 05ddc5c92e | |||
| b55e68f98d | |||
| beaa24754c | |||
| aed8cac439 | |||
| af4eeceb6a | |||
| 7b507dec94 | |||
| 2c0334c4fe | |||
| 92cbdbfdd3 | |||
| e78b2d8992 | |||
| 9d58922462 | |||
| b5ac48c126 | |||
| 69d0f5ee78 | |||
| 661b5c1099 | |||
| ff70637a0d | |||
| 06c1cf247e | |||
| 400653810e | |||
| b3936348f5 | |||
| c61b0b3968 | |||
| 0a5211b038 | |||
| 38be702188 | |||
| d42372e996 | |||
| 45231ba69e | |||
| 42c7b8b927 | |||
| e5a1db995d |
@@ -53,21 +53,32 @@ CGO daemons (`mavend`, `mavsttd`, `mavttsd`, `mavenclient`) need the vendored to
|
||||
and libs wired through the Makefile — **do not** call `go build` on them bare, use `make`:
|
||||
|
||||
```sh
|
||||
make build # all 9 binaries
|
||||
make build # all 11 binaries
|
||||
make build-web # single daemon (pure-Go ones: web/waked/poll/caldav build without CGO)
|
||||
make test # go test -race across ./internal/... ./cmd/... with CGO env set
|
||||
```
|
||||
|
||||
Run a single test (must carry the CGO env for packages that touch STT/TTS/voice):
|
||||
Run one package or one test with `make t`. **Do not hand-write the CGO preamble.**
|
||||
Past sessions pasted it about 390 times. That is where the shell-quoting failures
|
||||
came from. This box runs zsh, so an unquoted `-run Test*` or `--include=*.go`
|
||||
dies on "no matches found" before `go` is ever reached.
|
||||
|
||||
```sh
|
||||
CGO_CFLAGS="-I$(pwd)/deps/include -I$(pwd)/deps/whisper.cpp/ggml/include" \
|
||||
CGO_LDFLAGS="-L$(pwd)/deps/lib -Wl,-rpath,$(pwd)/deps/lib" \
|
||||
LD_LIBRARY_PATH="$(pwd)/deps/lib" \
|
||||
deps/go/go/bin/go test -run TestName ./internal/router/
|
||||
make t PKG=./internal/router/
|
||||
make t PKG=./cmd/mavend/ RUN=TestSimulator
|
||||
make t PKG=./internal/router/eval/ RUN='TestONNX' V=1 # V=1 for -v, RACE=0 to drop -race
|
||||
```
|
||||
|
||||
Pure-Go packages (`router`, `memory`, `mavweb`, …) run under a plain `go test ./pkg/`.
|
||||
`t` carries `-race`, so a green `make t` cannot turn red under `make test`. It carries
|
||||
`-count=1`, so a cached PASS from before your edit is never mistaken for a result.
|
||||
|
||||
It also sets `MAVEN_ONNX_LIB`, which the hand-written recipe did not. The four
|
||||
`TestONNX*` measurements self-skip when that variable is unset. The run still prints
|
||||
`ok`. So every targeted eval done the old way reported the hash ratchet while reading
|
||||
as a real embedder score.
|
||||
|
||||
Pure-Go packages (`router`, `memory`, `mavweb`, …) also run under a plain `go test ./pkg/`,
|
||||
but `make t` works everywhere and is one thing to remember.
|
||||
|
||||
## The daemons (`cmd/`)
|
||||
|
||||
@@ -82,13 +93,37 @@ Pure-Go packages (`router`, `memory`, `mavweb`, …) run under a plain `go test
|
||||
| `mavpoll` | Environment poller: netdata alarms, uptime-kuma, zenmoney, wireguard presence. Writes facts, sends nothing. Telegram is `internal/delivery/telegramsink`, not this. |
|
||||
| `mavcaldav` | CalDAV calendar sync. |
|
||||
| `mavmaild` | Mail reader (IMAP, read-only). Holds the IMAP password; core never sees it. |
|
||||
| `mavgpud` | GPU supervisor. **Runs on workpc, not homesrv** — own unit, `deploy/mavgpud.service`. Keeps llama-server loaded while the card is free (V-488). Maven never asks it for anything, it reads `/health` through `llm.Pair`. |
|
||||
| `mavupdate` | Not a daemon. Operator CLI a human runs on the box to deploy a new build. |
|
||||
|
||||
Two more binaries have no Makefile target and are built with `go run` or `go build` when
|
||||
they are needed. Neither is deployed.
|
||||
|
||||
| Binary | Role |
|
||||
|---|---|
|
||||
| `mavseal` | Recovery tool. Encrypts a live tmpfs working copy back to the ciphertext file when mavend was killed before `defer st.Close()` sealed it. |
|
||||
| `labelgen` | Runs the stage 0 grammars over utterances and prints JSONL, the training data for the routing heads (V-546). |
|
||||
|
||||
Daemons are wired socket-to-socket, not linked. `internal/ipc` is the client/server wire
|
||||
protocol; the config in `deploy/mavend.json` (with `${VAR}` env expansion from gitignored
|
||||
`deploy/telegram.env`) sets socket paths, model paths, and the phraser/embedder blocks.
|
||||
|
||||
**Seven of the nine run on homesrv. `mavwaked` and `mavenclient` do not, and that is the
|
||||
decision, not an oversight** (Vikunja #463, `docs/plans/17-where-the-voice-loop-runs.md`).
|
||||
**`docker-compose.yml` runs five: `mavend`, `mavsttd`, `mavttsd`, `mavweb`, `mavpoll`.**
|
||||
Count against compose, not against the table. Four of the nine daemons are absent, and each
|
||||
absence has a different reason.
|
||||
|
||||
`mavmaild` and `mavcaldav` are commented out in compose, each with the reason written
|
||||
beside it: the first needs a mail account, the second a CalDAV account, and this box has
|
||||
neither. `mavcaldav` used to appear nowhere at all, which was an oversight; it became a
|
||||
recorded decision on 07-08-2026 (V-644). Two things ride on that absence and the block
|
||||
names them. Agenda questions route to `IntentQuery` at stage 0 (V-498) and the `calendar`
|
||||
query source then reads a table nobody writes. And `loop.State.CalendarBusy` is fed by the
|
||||
same facts, so the gate's "do not nag mid-meeting" is permanently false. Its password is
|
||||
read from a file (`-pass-file`, and `-render-pass-file` for the render collection), never
|
||||
taken as a flag value, which is the rule `mavpoll` and `mavmaild` follow too.
|
||||
|
||||
**`mavwaked` and `mavenclient` are absent by decision, not oversight** (Vikunja #463,
|
||||
`docs/plans/17-where-the-voice-loop-runs.md`).
|
||||
homesrv has a microphone — it is a laptop — but it is in the wrong room, so a wake-word
|
||||
daemon there listens to nobody. They belong on a client machine where the owner is standing.
|
||||
|
||||
@@ -305,12 +340,70 @@ not a transcript. The transcript still expires. The gesture that writes one is
|
||||
two buttons beside the reply on `/chat`, reached over `ipc.CorrectTurn` and the
|
||||
trace id that now rides back on `ipc.ChatReply`. A turn marked wrong with no
|
||||
target is a usable negative, so naming the intent is never required. The target
|
||||
is one of the seven intents and never free text. Only `/chat` offers it: the wire
|
||||
op assumes no browser, but telegram and voice do not call it yet, and
|
||||
`docs/plans/22-correcting-a-turn.md` says why voice is the hard one. Adding a rung to the ladder
|
||||
is one of the seven intents and never free text. **All three reaches offer it as
|
||||
of 06-08-2026**, and this section used to say only `/chat` did. Voice is the
|
||||
`repair` rung, which has read spoken corrections since V-455 and now writes the
|
||||
durable label beside the classifier seed it always wrote; a spoken negative with
|
||||
no target is its own rung, `repair-negative` (V-636, `docs/plans/22-correcting-a-turn.md`).
|
||||
Telegram is an inline keyboard under the reply, and it needed the chat to become
|
||||
readable first — **telegram is no longer outbound only** (V-637,
|
||||
`docs/plans/23-inbound-telegram.md`). The poller is dark unless the `telegram`
|
||||
block says `intake`, it long-polls because the box takes no inbound connections,
|
||||
it accepts `chat_id` and no other sender, and it drops whatever queued while the
|
||||
daemon was down. It reaches the daemon through `ipc.CoreAPI` alone, so a chat
|
||||
turn takes the path `POST /api/chat` takes. Note that the turn source is still
|
||||
`tap:text` for both, so provenance cannot tell a chat turn from a typed one.
|
||||
Adding a rung to the ladder
|
||||
in `runTurn` means adding its name to `preRouteLadder` in
|
||||
`cmd/mavend/decisiontrace.go`, or that rung is silently missing from the record.
|
||||
|
||||
**A route now says where the answer lives, not only that the turn is a question**
|
||||
(V-655, 07-08-2026). `query` was a shrug. The cascade sorted an utterance into one of
|
||||
seven intents, with stage 0, the resident model and the classifier behind it. Then
|
||||
`IntentQuery` handed the turn to `querySources` in the daemon. That is twenty-two branches
|
||||
deciding by seed similarity in a fixed order. It has no fixture and no accuracy
|
||||
number, no model arm and no floor. `Decision.Source` (`internal/router/source.go`) is
|
||||
the second half of the route. Twelve destinations, not twenty-two. The three recall
|
||||
passes plus `fact-by-key` are one destination from outside. So are search, Kiwix and
|
||||
the URL reader.
|
||||
|
||||
**`SourceUnknown` is a real value and it is the floor.** Nothing named a destination,
|
||||
so the daemon walks the whole chain. That is byte-for-byte what shipped before the
|
||||
field existed. The classifier arm names nothing, so a box whose model is down routes
|
||||
queries exactly as it did.
|
||||
|
||||
`queryWalk` in `cmd/mavend/actions_query.go` takes sources **out** and moves none.
|
||||
That is the safety argument and it is not negotiable. The table's order is
|
||||
load-bearing. Every comment on it argues a reason between two sources, and above all
|
||||
it carries "the owner's data first, then the world". Naming `SourceWorld` does not
|
||||
send the turn outside. His notes, his facts and the personal boundary still run first.
|
||||
|
||||
What comes out is only the sources that **guess**. Those decide a turn is theirs by
|
||||
cosine against frozen seeds, then answer whatever they claimed. They hold no table
|
||||
that could come back empty. Weather is the pure case and has no local data at
|
||||
all. It was measured on the box on 2026-08-07
|
||||
(`docs/evals/2026-08-07-week-of-usage.md` section 4). It answered both "что такое
|
||||
TCP?" and "сколько будет 17 на 23?" with "для какого города?". The feed answered "какой у меня любимый язык?" with kernel headlines.
|
||||
The personal boundary answered "кто такой Линус Торвальдс?" with "не нашла у тебя
|
||||
такой записи". A source that guesses is marked `guesses: true` in the table. One that
|
||||
looks is not, and it is always asked.
|
||||
|
||||
Stage 0 fills the destination where a rule already knows it. `WorldQueryGrammars()`
|
||||
(`internal/router/worldquery.go`) claims "что такое X" and "сколько будет 17 на 23".
|
||||
It is wired after the agenda rules and **before** the feed and list rules.
|
||||
"что такое лента" is a definition question, and the feed rule would take it on the
|
||||
noun alone.
|
||||
`calendar-query` and `event-time-query` name the calendar. The possessive agenda rules
|
||||
deliberately do not. "что у меня в списке покупок" matches `agenda-query`, and naming
|
||||
the calendar there would take the list source off the turn.
|
||||
|
||||
Fixture unchanged at **69/91 classifier+ONNX**, measured both sides. That is the
|
||||
expected result, because it scores intent and no case here changes intent. **The
|
||||
destination has no fixture yet, so it has no accuracy number.** That and the model arm
|
||||
are the follow-ups. The field is designed so a decider naming nothing costs nothing.
|
||||
It lands on V-546. Intent, mood and BIO slot tags were already three heads on one
|
||||
forward pass of the resident e5-small. Destination is a fourth head on the same pass.
|
||||
|
||||
## LLM output contract
|
||||
|
||||
All phrasing paths emit `{"response":"...","mood":"..."}`, with fallback to plain text when
|
||||
@@ -425,8 +518,11 @@ start of a session rather than one lookup per first use:
|
||||
ToolSearch("select:mcp__vikunja__list_tasks,mcp__vikunja__get_task_details,mcp__vikunja__create_task,mcp__vikunja__update_task")
|
||||
```
|
||||
|
||||
`update_task` carrying a `description` resets `done` to false, so closing a task with a
|
||||
write-up takes two calls: the description, then `done: true`.
|
||||
**Close a finished task with `done: true` and nothing else** (owner's call, 07-08-2026).
|
||||
Do not write a completion summary into the description on the way out. It is lost anyway,
|
||||
and the durable record is the commit messages and the merged PR. Note that `update_task`
|
||||
carrying a `description` resets `done` to false, which is why a write-up ever took two
|
||||
calls.
|
||||
|
||||
## Session workflow
|
||||
|
||||
|
||||
@@ -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: 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 deps-sentinel tidy eval-router eval-reach eval-recall eval-phrasing eval-models build-gpud
|
||||
.PHONY: t audit 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 deps-sentinel tidy eval-router eval-reach eval-recall eval-phrasing eval-models build-gpud
|
||||
|
||||
all: build
|
||||
|
||||
@@ -128,6 +128,35 @@ 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/...
|
||||
|
||||
# t — run ONE package or ONE test with the toolchain env already wired. This is
|
||||
# the iteration target; `test` is the gate. Reach for it instead of pasting the
|
||||
# CGO_CFLAGS/CGO_LDFLAGS/LD_LIBRARY_PATH preamble by hand, which is how it was
|
||||
# done ~390 times across past sessions and is where the shell-quoting failures
|
||||
# came from -- the interactive shell here is zsh, and an unquoted `-run Test*`
|
||||
# or `--include=*.go` dies on "no matches found" before go ever starts.
|
||||
#
|
||||
# make t # whole tree (same scope as `test`)
|
||||
# make t PKG=./internal/router/
|
||||
# make t PKG=./cmd/mavend/ RUN=TestSimulator
|
||||
# make t PKG=./internal/router/eval/ RUN='TestONNX' V=1
|
||||
# make t PKG=./internal/store/ RACE=0 # drop -race when iterating hot
|
||||
#
|
||||
# -race is on by default so a green `make t` cannot turn red under `make test`.
|
||||
# -count=1 because a cached PASS from before your edit is worse than no answer.
|
||||
# MAVEN_ONNX_LIB is set for the same reason: the four TestONNX* measurements
|
||||
# self-skip when it is unset, so a targeted eval run would otherwise report the
|
||||
# deterministic hash ratchet and look like it scored the real embedder.
|
||||
PKG ?= ./internal/... ./cmd/...
|
||||
RUN ?=
|
||||
V ?=
|
||||
RACE ?= 1
|
||||
|
||||
t:
|
||||
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
|
||||
MAVEN_ONNX_LIB="$(MAVEN_ONNX_LIB)" \
|
||||
$(GO) test $(if $(V),-v,) $(if $(filter-out 0,$(RACE)),-race,) -count=1 \
|
||||
$(if $(RUN),-run '$(RUN)',) $(PKG)
|
||||
|
||||
# eval-router — score the held-out RU routing fixture (internal/router/eval).
|
||||
# Verbose so the report tables land in the terminal. MAVEN_ONNX_LIB points the
|
||||
# prod-representative baseline at the vendored runtime; override it or set it
|
||||
@@ -189,6 +218,16 @@ eval-models:
|
||||
# scores the fixtures against ggml-small and self-skips when the model is
|
||||
# absent, and TestGoldenFixturesAreCanonical, which checks the committed audio
|
||||
# and the manifest with no model at all.
|
||||
# audit — the repo inventory: LOC per package, open TODOs, real stubs, living-doc
|
||||
# staleness, test shape, packages with no test. Read-only, prints, writes nothing.
|
||||
# Run it instead of rebuilding the same greps by hand; past sessions spent 93 of
|
||||
# them on this before their first edit. SECTION=loc|todo|stubs|docs|tests|gaps
|
||||
# narrows it.
|
||||
SECTION ?= all
|
||||
|
||||
audit:
|
||||
@SECTION="$(SECTION)" ./scripts/audit.sh
|
||||
|
||||
stt-fixtures:
|
||||
./scripts/gen-stt-fixtures.sh
|
||||
|
||||
|
||||
+35
-8
@@ -50,10 +50,10 @@ func run(args []string) error {
|
||||
socket := fs.String("socket", "", "core IPC socket path (required)")
|
||||
url := fs.String("url", "", "CalDAV calendar URL, e.g. http://localhost:5232/kami/personal (required)")
|
||||
user := fs.String("user", "", "CalDAV basic-auth username (required)")
|
||||
pass := fs.String("pass", "", "CalDAV basic-auth password (required)")
|
||||
passFile := fs.String("pass-file", "", "file holding the CalDAV basic-auth password (required — never passed as a flag value)")
|
||||
renderURL := fs.String("render-url", "", "CalDAV collection maven publishes her own reminders to; empty disables rendering")
|
||||
renderUser := fs.String("render-user", "", "basic-auth username for -render-url (defaults to -user)")
|
||||
renderPass := fs.String("render-pass", "", "basic-auth password for -render-url (defaults to -pass)")
|
||||
renderPassFile := fs.String("render-pass-file", "", "file holding the password for -render-url (defaults to -pass-file)")
|
||||
renderDur := fs.Duration("render-duration", calendar.DefaultReminderDuration, "how long a rendered reminder occupies")
|
||||
interval := fs.Duration("interval", 5*time.Minute, "poll cadence")
|
||||
timeout := fs.Duration("timeout", 10*time.Second, "per-request HTTP timeout")
|
||||
@@ -63,13 +63,22 @@ func run(args []string) error {
|
||||
if *socket == "" {
|
||||
return fmt.Errorf("-socket is required")
|
||||
}
|
||||
if *url == "" || *user == "" || *pass == "" {
|
||||
return fmt.Errorf("-url, -user, -pass are required")
|
||||
if *url == "" || *user == "" || *passFile == "" {
|
||||
return fmt.Errorf("-url, -user, -pass-file are required")
|
||||
}
|
||||
if err := checkRenderTarget([]string{*url}, *renderURL); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// The password is read from a file, never taken as a flag value: an argv
|
||||
// secret is visible in `ps` to every user on the box and lands in the compose
|
||||
// file and the shell history. Same rule mavmaild and mavpoll follow. Read
|
||||
// once at start, so a rotated password means a restart.
|
||||
pass, err := readSecret(*passFile)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||
defer stop()
|
||||
|
||||
@@ -85,17 +94,20 @@ func run(args []string) error {
|
||||
http: hc,
|
||||
url: strings.TrimRight(*url, "/"),
|
||||
user: *user,
|
||||
pass: *pass,
|
||||
pass: pass,
|
||||
}
|
||||
|
||||
var rend *renderer
|
||||
if *renderURL != "" {
|
||||
ru, rp := *renderUser, *renderPass
|
||||
ru, rp := *renderUser, pass
|
||||
if ru == "" {
|
||||
ru = *user
|
||||
}
|
||||
if rp == "" {
|
||||
rp = *pass
|
||||
if *renderPassFile != "" {
|
||||
rp, err = readSecret(*renderPassFile)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
rend = newRenderer(core, hc, *renderURL, ru, rp, *renderDur)
|
||||
log.Printf("mavcaldav: rendering reminders to %s", *renderURL)
|
||||
@@ -131,6 +143,21 @@ func run(args []string) error {
|
||||
// It takes the whole read set, not one URL. The guarantee in the package
|
||||
// comment is about every calendar maven reads, and a second read target added
|
||||
// later must not quietly fall outside the check.
|
||||
// readSecret reads one credential from a file and refuses an empty one. An
|
||||
// empty file is a deployment mistake, not a password, and CalDAV basic auth
|
||||
// would send it and get a 401 every poll.
|
||||
func readSecret(path string) (string, error) {
|
||||
raw, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("read password file: %w", err)
|
||||
}
|
||||
secret := strings.TrimSpace(string(raw))
|
||||
if secret == "" {
|
||||
return "", fmt.Errorf("password file %s is empty", path)
|
||||
}
|
||||
return secret, nil
|
||||
}
|
||||
|
||||
func checkRenderTarget(readURLs []string, renderURL string) error {
|
||||
if renderURL == "" {
|
||||
return nil
|
||||
|
||||
@@ -5,12 +5,38 @@ import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
|
||||
// The password comes from a file so it never reaches argv. An empty or missing
|
||||
// file must fail at start rather than authenticate as "" against his calendar.
|
||||
func TestReadSecret(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
good := filepath.Join(dir, "ok")
|
||||
if err := os.WriteFile(good, []byte(" hunter2\n"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got, err := readSecret(good); err != nil || got != "hunter2" {
|
||||
t.Fatalf("readSecret(good) = %q, %v; want \"hunter2\", nil", got, err)
|
||||
}
|
||||
|
||||
empty := filepath.Join(dir, "empty")
|
||||
if err := os.WriteFile(empty, []byte("\n \n"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := readSecret(empty); err == nil {
|
||||
t.Fatal("readSecret(empty) = nil error, want refusal")
|
||||
}
|
||||
if _, err := readSecret(filepath.Join(dir, "absent")); err == nil {
|
||||
t.Fatal("readSecret(absent) = nil error, want refusal")
|
||||
}
|
||||
}
|
||||
|
||||
type fakeCore struct {
|
||||
ipc.UnimplementedCoreAPI
|
||||
facts map[string]ipc.Fact // composite key "key|source" → Fact
|
||||
|
||||
@@ -104,4 +104,3 @@ func (h *reactiveHandler) actionAct(ctx context.Context, dec router.Decision) st
|
||||
}
|
||||
return phraser.A(phraser.ActDone, nil)
|
||||
}
|
||||
|
||||
|
||||
+85
-23
@@ -59,6 +59,26 @@ type querySource struct {
|
||||
// sources search text with no notion of a day. When one of them grows a
|
||||
// date parameter, flip its flag here.
|
||||
dateAware bool
|
||||
|
||||
// dest — the destination this source serves, when the cascade named one
|
||||
// (V-655). Several sources share a destination: the three recall passes and
|
||||
// the fact-by-key lookup are all SourceRecall, because which of them lands
|
||||
// the hit is an ordering detail no utterance can name. A source with no
|
||||
// dest is reachable only by walking the chain.
|
||||
dest router.Source
|
||||
|
||||
// guesses — this source decides whether the turn is its own by scoring the
|
||||
// utterance against frozen seeds, rather than by looking something up and
|
||||
// coming back empty.
|
||||
//
|
||||
// The distinction is the whole point of the field. A source that looks can
|
||||
// be wrong about relevance and still harmless, because the miss shows up as
|
||||
// no rows. A source that guesses answers whatever it claims: weather has no
|
||||
// local table to miss against, so "что такое TCP?" became "для какого
|
||||
// города?". So when the cascade names a destination, the guessers that were
|
||||
// not named do not get to try. The lookups still run, because a named
|
||||
// destination is evidence and not a promise.
|
||||
guesses bool
|
||||
}
|
||||
|
||||
// querySources is the ordered chain actionQuery walks; first source to claim
|
||||
@@ -67,85 +87,85 @@ type querySource struct {
|
||||
// gate was never the bug. Adding a source (Kiwix, RSS, crawler, email) is one
|
||||
// line here plus its method; where you put the line is the whole decision.
|
||||
var querySources = []querySource{
|
||||
{name: "fact-by-key", answer: (*reactiveHandler).queryFactByKey},
|
||||
{name: "fact-by-key", answer: (*reactiveHandler).queryFactByKey, dest: router.SourceRecall},
|
||||
// Before "calendar" on purpose: both match "…на сегодня", and the plan is
|
||||
// the more specific ask (its matcher requires a plan word), so the calendar
|
||||
// listing would otherwise swallow it.
|
||||
{name: "day-plan", answer: (*reactiveHandler).queryDayPlan},
|
||||
{name: "day-plan", answer: (*reactiveHandler).queryDayPlan, dest: router.SourceCalendar},
|
||||
// Also before "calendar": "что я обычно делаю по средам?" names a weekday,
|
||||
// and the habit question is the more specific one. Its matcher requires a
|
||||
// habit marker ("обычно", "каждый", …), so a question about this coming
|
||||
// Wednesday still reaches the calendar.
|
||||
{name: "habits", answer: (*reactiveHandler).queryHabits},
|
||||
{name: "habits", answer: (*reactiveHandler).queryHabits, dest: router.SourceCalendar},
|
||||
// Before "calendar" and before the recall sources: "что мне нужно
|
||||
// сделать?" is a question about the task list, and the notes pass would
|
||||
// otherwise answer it with whatever note happens to be nearest. Its
|
||||
// matcher requires a task noun or an explicit "что … сделать", so a
|
||||
// date-bearing question still reaches the calendar.
|
||||
{name: "tasks", answer: (*reactiveHandler).queryTasks},
|
||||
{name: "tasks", answer: (*reactiveHandler).queryTasks, dest: router.SourceTasks},
|
||||
// Next to "tasks" and for the same reason: "что требует внимания?" is a
|
||||
// question about the operational state Praxis holds, and it used to fall
|
||||
// through every source to the web search (Vikunja #475). Its matcher needs
|
||||
// an attention marker, and it falls through when Praxis is not configured.
|
||||
{name: "attention", answer: (*reactiveHandler).queryAttention},
|
||||
{name: "attention", answer: (*reactiveHandler).queryAttention, dest: router.SourceAttention, guesses: true},
|
||||
// Next to "tasks" and for the same reason: "что мне купить?" is a question
|
||||
// about the shopping list, and the recall pass would otherwise answer it
|
||||
// from an old note about the shop. Its matcher needs an explicit list
|
||||
// marker, so "надо бы съездить в магазин" is untouched.
|
||||
{name: "list", answer: (*reactiveHandler).queryList},
|
||||
{name: "list", answer: (*reactiveHandler).queryList, dest: router.SourceList, guesses: true},
|
||||
// Before the recall sources too: "сколько я потратил?" is a question about
|
||||
// the money facts the poller wrote, and the notes pass would otherwise
|
||||
// answer it from whatever he once said about spending. Its matcher needs a
|
||||
// money noun plus an actual ask, so "я потратил весь день" is untouched.
|
||||
{name: "money", answer: (*reactiveHandler).queryMoney},
|
||||
{name: "money", answer: (*reactiveHandler).queryMoney, dest: router.SourceMoney},
|
||||
// Also above the recall sources: "что я тебе говорил?" is a question about
|
||||
// the facts he tapped in, and the notes pass would answer it with whatever
|
||||
// note is nearest (Vikunja #456). Its matcher needs both halves of a
|
||||
// history phrase and bails out when he names a topic, so "что я говорил
|
||||
// про сервер" is still recall.
|
||||
{name: "history", answer: (*reactiveHandler).queryHistory},
|
||||
{name: "history", answer: (*reactiveHandler).queryHistory, dest: router.SourceRecall},
|
||||
// Before the recall sources and before general knowledge: "что нового?" is
|
||||
// a question about the feeds she reads, and general knowledge would answer
|
||||
// it by inventing news. Its matcher needs a feed noun plus an ask, so
|
||||
// "у меня новая лента в инстаграме" is untouched.
|
||||
{name: "feeds", answer: (*reactiveHandler).queryFeeds},
|
||||
{name: "feeds", answer: (*reactiveHandler).queryFeeds, dest: router.SourceFeeds, guesses: true},
|
||||
// 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.
|
||||
{name: "home", answer: (*reactiveHandler).queryHome},
|
||||
{name: "home", answer: (*reactiveHandler).queryHome, dest: router.SourceHome, guesses: true},
|
||||
// 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.
|
||||
{name: "network", answer: (*reactiveHandler).queryNetwork},
|
||||
{name: "calendar", answer: (*reactiveHandler).queryCalendar, dateAware: true},
|
||||
{name: "weather", answer: (*reactiveHandler).queryWeather},
|
||||
{name: "network", answer: (*reactiveHandler).queryNetwork, dest: router.SourceNetwork, guesses: true},
|
||||
{name: "calendar", answer: (*reactiveHandler).queryCalendar, dateAware: true, dest: router.SourceCalendar},
|
||||
{name: "weather", answer: (*reactiveHandler).queryWeather, dest: router.SourceWeather, guesses: true},
|
||||
// A question about her, above the three sources that search his own data
|
||||
// (Vikunja #555). It has no answer anywhere else: below the boundary
|
||||
// SearXNG answers about somebody else's assistant, and above it his notes
|
||||
// answer by proximity — "кто ты" came back from a note of his, measured on
|
||||
// the box, because the recall index has no idea the subject is her.
|
||||
{name: "self", answer: (*reactiveHandler).querySelf},
|
||||
{name: "embed", answer: (*reactiveHandler).queryEmbed},
|
||||
{name: "memory", answer: (*reactiveHandler).queryMemory},
|
||||
{name: "notes", answer: (*reactiveHandler).queryNotes},
|
||||
{name: "self", answer: (*reactiveHandler).querySelf, dest: router.SourceSelf, guesses: true},
|
||||
{name: "embed", answer: (*reactiveHandler).queryEmbed, dest: router.SourceRecall},
|
||||
{name: "memory", answer: (*reactiveHandler).queryMemory, dest: router.SourceRecall},
|
||||
{name: "notes", answer: (*reactiveHandler).queryNotes, dest: router.SourceRecall},
|
||||
// THE BOUNDARY. Everything above answers from his own data; everything
|
||||
// below answers from the world's. A question about him that got this far
|
||||
// has no answer in his data, and no outside source can supply one, so this
|
||||
// stops the walk rather than let the encyclopedia and the model guess.
|
||||
{name: "personal", answer: (*reactiveHandler).queryPersonal},
|
||||
{name: "personal", answer: (*reactiveHandler).queryPersonal, dest: router.SourceRecall, guesses: true},
|
||||
// The world, read live. Owner's ruling of 2026-08-02: a metasearch hit beats
|
||||
// a frozen ZIM, so SearXNG asks before Kiwix does. Nothing of his is at
|
||||
// stake by this point — the boundary above already stopped every question
|
||||
// about him, and only the query string leaves the box.
|
||||
{name: "search", answer: (*reactiveHandler).querySearch},
|
||||
{name: "search", answer: (*reactiveHandler).querySearch, dest: router.SourceWorld},
|
||||
// The offline encyclopedia, now the fallback for when the line is down or
|
||||
// the search comes back empty. It reads the way it always did; what changed
|
||||
// is that it no longer gets first refusal on a world question.
|
||||
{name: "kiwix", answer: (*reactiveHandler).queryKiwix},
|
||||
{name: "kiwix", answer: (*reactiveHandler).queryKiwix, dest: router.SourceWorld},
|
||||
// LAST before the model answers from memory, and that position is the whole
|
||||
// design (Vikunja #259): everything of his, then the search, then the ZIMs,
|
||||
// and only then a page he named. The model does NOT come first: it
|
||||
@@ -153,8 +173,43 @@ var querySources = []querySource{
|
||||
// a 1.7B guessing at a page it cannot read is how contents get invented.
|
||||
// This source only claims a turn where he named a URL, so it never competes
|
||||
// with a local answer.
|
||||
{name: "web", answer: (*reactiveHandler).queryWeb},
|
||||
{name: "general-knowledge", answer: (*reactiveHandler).queryGeneral},
|
||||
{name: "web", answer: (*reactiveHandler).queryWeb, dest: router.SourceWorld},
|
||||
{name: "general-knowledge", answer: (*reactiveHandler).queryGeneral, dest: router.SourceWorld},
|
||||
}
|
||||
|
||||
// queryWalk narrows the chain for one turn against the destination the cascade
|
||||
// named, and says which sources were left out (V-655).
|
||||
//
|
||||
// It takes sources OUT and never moves one, which is the whole safety argument.
|
||||
// The table's order is load-bearing and every comment on it argues a reason
|
||||
// between two sources; none of those reasons is about this. Above all, the
|
||||
// order carries "his data first, then the world", and a destination named by a
|
||||
// model must not be able to reverse that. Naming SourceWorld does not send the
|
||||
// turn outside — it stops the guessers from claiming it on the way.
|
||||
//
|
||||
// What comes out is exactly the sources that guess. Those decide whether a turn
|
||||
// is theirs by scoring it against frozen seeds, and then answer whatever they
|
||||
// claimed, because they have no lookup that can come back empty. That is the
|
||||
// whole of the 2026-08-07 defect: weather claiming "что такое TCP?", the feed
|
||||
// claiming "какой у меня любимый язык?", the personal boundary claiming "кто
|
||||
// такой Линус Торвальдс?". The sources that look are all still asked, so a
|
||||
// wrong destination costs nothing but the guess it prevented.
|
||||
//
|
||||
// No destination named ⇒ the table exactly as written, which is what shipped
|
||||
// before the field existed. That is the floor. The classifier arm names
|
||||
// nothing, so a box whose model is down routes queries the way it always did.
|
||||
func queryWalk(dest router.Source) (walk, skipped []querySource) {
|
||||
if dest == router.SourceUnknown {
|
||||
return querySources, nil
|
||||
}
|
||||
for _, s := range querySources {
|
||||
if s.guesses && s.dest != dest {
|
||||
skipped = append(skipped, s)
|
||||
continue
|
||||
}
|
||||
walk = append(walk, s)
|
||||
}
|
||||
return walk, skipped
|
||||
}
|
||||
|
||||
func (h *reactiveHandler) actionQuery(ctx context.Context, dec router.Decision) string {
|
||||
@@ -164,7 +219,14 @@ func (h *reactiveHandler) actionQuery(ctx context.Context, dec router.Decision)
|
||||
// (V-564). Finish names everyone below the winner.
|
||||
decision.Expect(ctx, decision.StageQuery, querySourceNames())
|
||||
rec := decision.From(ctx)
|
||||
for _, src := range querySources {
|
||||
walk, skipped := queryWalk(dec.Source)
|
||||
for _, src := range skipped {
|
||||
rec.Note(decision.Claim{
|
||||
Stage: decision.StageQuery, Claimant: src.name, Outcome: decision.NeverAsked,
|
||||
Reason: "it decides by similarity and the cascade named " + string(dec.Source),
|
||||
})
|
||||
}
|
||||
for _, src := range walk {
|
||||
if dec.Continued && !src.dateAware {
|
||||
rec.Note(decision.Claim{
|
||||
Stage: decision.StageQuery, Claimant: src.name, Outcome: decision.NeverAsked,
|
||||
|
||||
@@ -0,0 +1,120 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
"net"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/event"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/store"
|
||||
)
|
||||
|
||||
// The two boot paths meet here. run() wires the daemon twice: once at boot
|
||||
// when a key is in the environment, and once inside UnlockFn after a passkey
|
||||
// assertion, minutes or days later. Listing the same wiring in both places is
|
||||
// what let them drift — seven workers started untracked on the unlock path and
|
||||
// two daemonAPI fields were never set there, silently, for as long as anyone
|
||||
// had been cold-starting (V-639).
|
||||
//
|
||||
// So both paths call newDaemonAPI and startBackground and nothing else. A
|
||||
// field or a worker added later reaches both paths or neither.
|
||||
|
||||
// bootDeps is everything the two constructors below read. It is filled from
|
||||
// the same variables on both paths, by depsNow in run().
|
||||
type bootDeps struct {
|
||||
coreFor func() ipc.CoreAPI
|
||||
tl *tickLoop
|
||||
evBus *event.Bus
|
||||
voiceW *voiceWiring
|
||||
st *store.Store
|
||||
factWorker *factEnrichmentWorker
|
||||
evalWorker *memoryEvalWorker // nil ⇒ memory evaluation off (the default)
|
||||
feedWkr *feedWorker // nil ⇒ no feed is read (the default)
|
||||
crawlWkr *crawlWorker // nil ⇒ no page is watched (the default)
|
||||
}
|
||||
|
||||
// newDaemonAPI builds the real CoreAPI, with every field set. The unlock path
|
||||
// used to leave nexus and getMCPServers nil, so after a cold start
|
||||
// ResolveEntity refused with a nexus block configured and /tools rendered
|
||||
// "not configured" with an mcp block configured. Empty is a wrong answer
|
||||
// there, not a degraded one.
|
||||
func newDaemonAPI(d bootDeps) *daemonAPI {
|
||||
api := &daemonAPI{
|
||||
CoreAPI: d.coreFor(),
|
||||
getTrace: d.tl.trace,
|
||||
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return d.tl.morningStatus(ctx, time.Now()) },
|
||||
getDayPlan: func(ctx context.Context) ipc.DayPlan { return d.tl.dayPlan(ctx, time.Now()) },
|
||||
getEvents: intakeEventsFn(d.evBus),
|
||||
getDecisions: turnDecisionsFn(d.voiceW),
|
||||
seedStore: seedStoreIfAllowed(d.st),
|
||||
nexus: nexusOf(d.voiceW),
|
||||
}
|
||||
if d.voiceW != nil && d.voiceW.handler != nil {
|
||||
api.chatFn = d.voiceW.handler.handleText
|
||||
// And the reverse: the handler was wired with the bare store adapter,
|
||||
// which cannot serve the day plan. See upgradeAPI.
|
||||
d.voiceW.handler.upgradeAPI(api)
|
||||
}
|
||||
if d.voiceW != nil && d.voiceW.mcp != nil {
|
||||
api.getMCPServers = d.voiceW.mcp.status
|
||||
}
|
||||
return api
|
||||
}
|
||||
|
||||
// namedWorker is one long-running goroutine. The name exists so the set is
|
||||
// assertable from a test and readable in a log; nothing dispatches on it.
|
||||
type namedWorker struct {
|
||||
name string
|
||||
run func(ctx context.Context)
|
||||
}
|
||||
|
||||
// backgroundWorkers lists what this deployment runs. It is pure — it starts
|
||||
// nothing — so a test can compare the set the two paths would start without
|
||||
// standing a daemon up.
|
||||
func backgroundWorkers(d bootDeps) []namedWorker {
|
||||
var ws []namedWorker
|
||||
if d.voiceW != nil && d.voiceW.server != nil {
|
||||
ws = append(ws, namedWorker{"voice", func(context.Context) {
|
||||
if err := d.voiceW.server.Serve(); err != nil && !errors.Is(err, net.ErrClosed) {
|
||||
log.Printf("voice serve: %v", err)
|
||||
}
|
||||
}})
|
||||
}
|
||||
ws = append(ws,
|
||||
namedWorker{"tick", d.tl.run},
|
||||
namedWorker{"fact-enrichment", d.factWorker.run},
|
||||
)
|
||||
if d.evalWorker != nil {
|
||||
ws = append(ws, namedWorker{"memory-eval", d.evalWorker.run})
|
||||
}
|
||||
if d.feedWkr != nil {
|
||||
ws = append(ws, namedWorker{"feed", d.feedWkr.run})
|
||||
}
|
||||
if d.crawlWkr != nil {
|
||||
ws = append(ws, namedWorker{"crawl", d.crawlWkr.run})
|
||||
}
|
||||
if d.voiceW != nil && d.voiceW.mcp != nil {
|
||||
ws = append(ws, namedWorker{"mcp", d.voiceW.mcp.run})
|
||||
}
|
||||
if d.voiceW != nil && d.voiceW.home != nil {
|
||||
ws = append(ws, namedWorker{"home", d.voiceW.home.run})
|
||||
}
|
||||
return ws
|
||||
}
|
||||
|
||||
// startBackground starts every worker through goWorker, so waitWorkers can
|
||||
// wait for it at shutdown. A worker started as a bare `go func()` is the
|
||||
// shutdown bug documented at the end of run(): run() never returns, the
|
||||
// deferred Close never seals the database, and the ciphertext goes stale.
|
||||
func startBackground(ctx context.Context, wg *sync.WaitGroup, d bootDeps) {
|
||||
for _, w := range backgroundWorkers(d) {
|
||||
goWorker(wg, func() { w.run(ctx) })
|
||||
}
|
||||
if d.voiceW != nil && d.voiceW.server != nil {
|
||||
log.Printf("mavend: voice listening on %s", d.voiceW.server.Addr())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,96 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"reflect"
|
||||
"testing"
|
||||
|
||||
"github.com/kami/maven/internal/decision"
|
||||
"github.com/kami/maven/internal/event"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/store"
|
||||
"github.com/kami/maven/internal/voice"
|
||||
)
|
||||
|
||||
// fullDeps — a deployment with every optional piece present. Nothing here is
|
||||
// run: newDaemonAPI takes method values and backgroundWorkers is pure, so
|
||||
// zero-value wirings are enough to say what WOULD be started.
|
||||
func fullDeps() bootDeps {
|
||||
h := &reactiveHandler{
|
||||
ecosystem: &ecosystemWiring{nexus: &nexusClient{}},
|
||||
decisions: decision.NewRing(),
|
||||
}
|
||||
return bootDeps{
|
||||
coreFor: func() ipc.CoreAPI { return ipc.UnimplementedCoreAPI{} },
|
||||
tl: &tickLoop{},
|
||||
evBus: event.NewBus(4),
|
||||
st: &store.Store{},
|
||||
factWorker: &factEnrichmentWorker{},
|
||||
evalWorker: &memoryEvalWorker{},
|
||||
feedWkr: &feedWorker{},
|
||||
crawlWkr: &crawlWorker{},
|
||||
voiceW: &voiceWiring{
|
||||
server: &voice.Server{},
|
||||
handler: h,
|
||||
mcp: &mcpWiring{},
|
||||
home: &homeWiring{},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// The unlock path used to build its own daemonAPI literal and leave nexus and
|
||||
// getMCPServers nil (V-639). Both paths call newDaemonAPI now, so the drift
|
||||
// that can still happen is a field added to the struct and not to the
|
||||
// constructor. This catches that one, by name.
|
||||
func TestNewDaemonAPISetsEveryField(t *testing.T) {
|
||||
prev := allowSeedOnStart
|
||||
allowSeedOnStart = true
|
||||
defer func() { allowSeedOnStart = prev }()
|
||||
|
||||
api := newDaemonAPI(fullDeps())
|
||||
v := reflect.ValueOf(*api)
|
||||
for i := range v.NumField() {
|
||||
if v.Field(i).IsZero() {
|
||||
t.Errorf("newDaemonAPI left %s unset — a fully wired deployment must fill every field", v.Type().Field(i).Name)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The handler is wired with the bare store adapter and cannot serve the day
|
||||
// plan until upgradeAPI hands it the real one. The unlocked path did that and
|
||||
// the unlock path did it too; keep it a property of the constructor.
|
||||
func TestNewDaemonAPIUpgradesTheHandler(t *testing.T) {
|
||||
d := fullDeps()
|
||||
api := newDaemonAPI(d)
|
||||
if d.voiceW.handler.api != ipc.CoreAPI(api) {
|
||||
t.Fatal("newDaemonAPI did not hand the handler the API it built")
|
||||
}
|
||||
}
|
||||
|
||||
// Every worker the daemon runs goes through startBackground, so shutdown can
|
||||
// wait for it. The unlock path used to start seven of these as bare
|
||||
// `go func()` under a shadowed WaitGroup.
|
||||
func TestBackgroundWorkersFullSet(t *testing.T) {
|
||||
want := []string{"voice", "tick", "fact-enrichment", "memory-eval", "feed", "crawl", "mcp", "home"}
|
||||
var got []string
|
||||
for _, w := range backgroundWorkers(fullDeps()) {
|
||||
got = append(got, w.name)
|
||||
}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("workers = %v, want %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
// A default box configures none of the optional blocks. Two workers always run
|
||||
// and the rest stay dark, rather than a nil run being scheduled.
|
||||
func TestBackgroundWorkersFloor(t *testing.T) {
|
||||
d := fullDeps()
|
||||
d.evalWorker, d.feedWkr, d.crawlWkr, d.voiceW = nil, nil, nil, nil
|
||||
want := []string{"tick", "fact-enrichment"}
|
||||
var got []string
|
||||
for _, w := range backgroundWorkers(d) {
|
||||
got = append(got, w.name)
|
||||
}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("workers = %v, want %v", got, want)
|
||||
}
|
||||
}
|
||||
+43
-23
@@ -6,7 +6,6 @@ import (
|
||||
"math/rand"
|
||||
"strings"
|
||||
"time"
|
||||
"unicode"
|
||||
"unicode/utf8"
|
||||
|
||||
"github.com/kami/maven/internal/dialogue"
|
||||
@@ -162,12 +161,13 @@ func withNotice(notice, reply string) string {
|
||||
// напоминание?" — answer first, then the open question. A question in front of
|
||||
// its own answer would read as ignoring what he asked.
|
||||
//
|
||||
// A statement's full stop is folded into a comma, so the two acts read as one
|
||||
// sentence — that is the owner's own punctuation, "в Риме сейчас ..., на какое
|
||||
// время поставить напоминание?". An answer that is ITSELF a question keeps its
|
||||
// mark and the resume starts a new sentence: she sometimes answers a side query
|
||||
// by asking him to say it again, and "переформулировать?, на какое время" folds
|
||||
// two questions into one unreadable line.
|
||||
// Two sentences, not one (V-654). This used to fold the answer's full stop into
|
||||
// a comma, on the strength of the owner having written it that way once. Spliced
|
||||
// onto a real answer it reads as one run-on thought — "вот что я нашла: вайфай
|
||||
// пароль лежит в ящике стола, на какое время поставить напоминание?" — and the
|
||||
// question disappears into the tail of a sentence about something else. A reply
|
||||
// with no terminator of its own is given one, so the join never depends on how
|
||||
// the phraser chose to end.
|
||||
//
|
||||
// A resume with no answer in front of it is just the question.
|
||||
func withResumed(reply, resumed string) string {
|
||||
@@ -178,23 +178,17 @@ func withResumed(reply, resumed string) string {
|
||||
if reply == "" {
|
||||
return resumed
|
||||
}
|
||||
if strings.HasSuffix(reply, "?") {
|
||||
return reply + " " + resumed
|
||||
if !endsSentence(reply) {
|
||||
reply += "."
|
||||
}
|
||||
if trimmed := strings.TrimRight(reply, ".!"); trimmed != "" {
|
||||
reply = trimmed
|
||||
}
|
||||
return reply + ", " + lowerFirst(resumed)
|
||||
return reply + " " + resumed
|
||||
}
|
||||
|
||||
// lowerFirst lowercases the opening rune, so a deck line written as a standalone
|
||||
// sentence reads as the second half of one. Only the first rune: "На какое
|
||||
// время" must become "на какое время" and nothing else in it may move.
|
||||
func lowerFirst(s string) string {
|
||||
for i, r := range s {
|
||||
return string(unicode.ToLower(r)) + s[i+utf8.RuneLen(r):]
|
||||
}
|
||||
return s
|
||||
// endsSentence reports whether s already closes itself. The ellipsis counts: a
|
||||
// trailing "…" is a deliberate end, and a full stop after it reads as a typo.
|
||||
func endsSentence(s string) bool {
|
||||
r, _ := utf8.DecodeLastRuneInString(s)
|
||||
return strings.ContainsRune(".!?…", r)
|
||||
}
|
||||
|
||||
// missingFor returns the slots a decision still needs, most important first.
|
||||
@@ -381,6 +375,14 @@ func (h *reactiveHandler) resolveClarifyAnswer(ctx context.Context, text string)
|
||||
return "", false
|
||||
}
|
||||
|
||||
// He is answering, so the run of step-asides is over (V-654). Reset here
|
||||
// rather than where a gap is FILLED: "позвонить маме" against a question
|
||||
// about the time gives her nothing she asked for and still means he is in
|
||||
// the exchange, and the retry it costs is bound enough on its own. The
|
||||
// counter is for the case the bounds miss — he asked for other things and
|
||||
// never came back.
|
||||
q.Suspends = 0
|
||||
|
||||
merged := q.Answer(text, toDialogueSlots(answer))
|
||||
// Fold a newly answered subject into the raw utterance. Downstream actions
|
||||
// phrase from Utterance, not from the text slot — actionReminder stores it
|
||||
@@ -467,6 +469,14 @@ func (h *reactiveHandler) noteDropped(ctx context.Context) {
|
||||
//
|
||||
// A slot with no resumed wording (clarifyResumedFor says so) resumes nothing and
|
||||
// says nothing. She must not claim to be holding a question she cannot re-ask.
|
||||
//
|
||||
// Suspension is bounded, since V-654. Neither of the two things above is a
|
||||
// limit: no attempt is spent, and restarting the clock means the TTL cannot
|
||||
// arrive while he keeps talking. So the count is the only thing that ends it,
|
||||
// and past MaxSuspends she lets the request go and says so with the same line
|
||||
// every other drop uses. The rule is unchanged — a question ends by being
|
||||
// answered or by being let go out loud — this only recognises three unrelated
|
||||
// requests in a row as the second of those.
|
||||
func (h *reactiveHandler) noteSuspended(ctx context.Context, q *dialogue.PendingQuestion) {
|
||||
rt := turnRouteFrom(ctx)
|
||||
if rt == nil || len(q.Missing) == 0 {
|
||||
@@ -476,11 +486,18 @@ func (h *reactiveHandler) noteSuspended(ctx context.Context, q *dialogue.Pending
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if !q.CanResume() {
|
||||
h.clarifyStore.Delete(dialogueIDOf(ctx))
|
||||
h.noteDropped(ctx)
|
||||
log.Printf("voice: clarify — the question about %s stepped aside %d times; letting the request go", q.Missing[0], q.Suspends)
|
||||
return
|
||||
}
|
||||
q.Suspends++
|
||||
q.Asked = h.now()
|
||||
h.clarifyStore.Put(dialogueIDOf(ctx), q)
|
||||
rt.resume = question
|
||||
rt.suspended = true
|
||||
log.Printf("voice: clarify — is its own request; suspending the question about %s and resuming it in the same reply", q.Missing[0])
|
||||
log.Printf("voice: clarify — is its own request; suspending the question about %s and resuming it in the same reply (suspend %d of %d)", q.Missing[0], q.Suspends, dialogue.MaxSuspends)
|
||||
}
|
||||
|
||||
// foldAnswerIntoUtterance appends an answered subject to the original words,
|
||||
@@ -520,6 +537,9 @@ func (h *reactiveHandler) askRemainingGap(ctx context.Context, q *dialogue.Pendi
|
||||
if !ok || !q.CanAsk() {
|
||||
return "", false
|
||||
}
|
||||
// Suspends is not carried, and by this point it is already zero: the answer
|
||||
// path resets it (V-654). Left off the literal so the zero is stated where
|
||||
// the struct is built, rather than inherited from a field nobody names.
|
||||
h.clarifyStore.Put(dialogueIDOf(ctx), &dialogue.PendingQuestion{
|
||||
Intent: q.Intent,
|
||||
Slots: merged,
|
||||
@@ -571,7 +591,7 @@ func (h *reactiveHandler) finishClarified(ctx context.Context, dec router.Decisi
|
||||
}
|
||||
reply := h.applyAction(ctx, dec)
|
||||
if reply == "" {
|
||||
reply = h.replier.Reply(dec)
|
||||
reply = h.replier.Reply(ctx, dec)
|
||||
}
|
||||
if reply == "" {
|
||||
// Belt: an empty reply here would be a silent drop.
|
||||
|
||||
@@ -740,3 +740,53 @@ func TestACompleteTurnStillDoesNotAsk(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestTheResumedQuestionIsItsOwnSentence — V-654. The re-ask used to be spliced
|
||||
// onto the answer with a comma, so a real answer and an unrelated open question
|
||||
// read as one run-on thought and the question vanished into its tail.
|
||||
func TestTheResumedQuestionIsItsOwnSentence(t *testing.T) {
|
||||
const resumed = "На какое время поставить напоминание?"
|
||||
cases := []struct {
|
||||
name string
|
||||
reply string
|
||||
want string
|
||||
}{
|
||||
{
|
||||
// The measured line, shortened. Two sentences, and the question keeps
|
||||
// its capital.
|
||||
name: "a statement keeps its full stop",
|
||||
reply: "Вайфай пароль лежит в ящике стола.",
|
||||
want: "Вайфай пароль лежит в ящике стола. " + resumed,
|
||||
},
|
||||
{
|
||||
name: "a reply with no terminator is given one",
|
||||
reply: "Вайфай пароль лежит в ящике стола",
|
||||
want: "Вайфай пароль лежит в ящике стола. " + resumed,
|
||||
},
|
||||
{
|
||||
// She sometimes answers a side query by asking him to say it again.
|
||||
// Two questions, and neither may swallow the other.
|
||||
name: "a question keeps its mark",
|
||||
reply: "Можешь переформулировать?",
|
||||
want: "Можешь переформулировать? " + resumed,
|
||||
},
|
||||
{
|
||||
name: "an ellipsis is already an ending",
|
||||
reply: "Не уверена…",
|
||||
want: "Не уверена… " + resumed,
|
||||
},
|
||||
{
|
||||
name: "a resume with no answer in front of it is just the question",
|
||||
reply: "",
|
||||
want: resumed,
|
||||
},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
if got := withResumed(tc.reply, resumed); got != tc.want {
|
||||
t.Errorf("%s: withResumed(%q) = %q, want %q", tc.name, tc.reply, got, tc.want)
|
||||
}
|
||||
}
|
||||
if got := withResumed("Готово.", ""); got != "Готово." {
|
||||
t.Errorf("nothing to resume must leave the reply alone, got %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -37,8 +37,8 @@ type factEnrichmentWorker struct {
|
||||
nextTry map[int64]time.Time // fact id → earliest retry
|
||||
}
|
||||
|
||||
// enrichmentScanLimit bounds how deep a single tick (or status report) walks
|
||||
// the pending queue looking for facts whose backoff has elapsed. The queue is
|
||||
// enrichmentScanLimit bounds how deep a single tick walks the pending queue
|
||||
// looking for facts whose backoff has elapsed. The queue is
|
||||
// ordered by id, so without a scan the oldest facts hold every batch slot
|
||||
// whether or not they are eligible, and one permanently failing fact stalls
|
||||
// every younger one behind it.
|
||||
@@ -75,8 +75,8 @@ func newFactEnrichmentWorker(st *store.Store, eco *ecosystemWiring, interval tim
|
||||
// has been down all day must be visible as a backlog, not as facts that
|
||||
// silently never got tagged.
|
||||
//
|
||||
// All three numbers describe the same set of rows, the first
|
||||
// enrichmentScanLimit pending facts. Counting Pending over a thousand rows
|
||||
// All three numbers describe the same set of rows, whatever is still pending
|
||||
// out of the first enrichmentScanLimit facts. Counting Pending over a thousand rows
|
||||
// while counting InBackoff over the twenty that reached the head of a batch
|
||||
// described two different populations under one struct.
|
||||
type enrichmentStatus struct {
|
||||
@@ -86,13 +86,22 @@ type enrichmentStatus struct {
|
||||
Scanned int // rows the other three counts were taken over
|
||||
}
|
||||
|
||||
// status reads the queue and counts over it. For a caller with no batch in
|
||||
// hand — anything asking the worker how it is doing from outside the tick.
|
||||
func (w *factEnrichmentWorker) status(ctx context.Context) enrichmentStatus {
|
||||
var st enrichmentStatus
|
||||
pending, err := w.store.PendingFactResolutions(ctx, enrichmentScanLimit)
|
||||
if err != nil {
|
||||
log.Printf("factenrichment: status: %v", err)
|
||||
return st
|
||||
return enrichmentStatus{}
|
||||
}
|
||||
return w.statusOf(pending)
|
||||
}
|
||||
|
||||
// statusOf counts over a batch the caller already has. The batch is the query
|
||||
// the tick already ran, so reporting the backlog costs no second read of the
|
||||
// scan limit — up to a thousand rows, on a database that serialises them.
|
||||
func (w *factEnrichmentWorker) statusOf(pending []store.Fact) enrichmentStatus {
|
||||
var st enrichmentStatus
|
||||
st.Pending = len(pending)
|
||||
st.Scanned = len(pending)
|
||||
w.mu.Lock()
|
||||
@@ -144,17 +153,24 @@ func (w *factEnrichmentWorker) tick(ctx context.Context) {
|
||||
}
|
||||
w.forgetDeparted(pending)
|
||||
skipped, failed, attempted := 0, 0, 0
|
||||
// A resolved fact leaves the pending queue, so the batch in hand overstates
|
||||
// the backlog by however many succeeded. Drop them here rather than
|
||||
// re-reading the queue to find out.
|
||||
remaining := make([]store.Fact, 0, len(pending))
|
||||
for _, f := range pending {
|
||||
if attempted >= w.batch {
|
||||
break
|
||||
remaining = append(remaining, f)
|
||||
continue
|
||||
}
|
||||
if !w.due(f.ID) {
|
||||
skipped++
|
||||
remaining = append(remaining, f)
|
||||
continue
|
||||
}
|
||||
attempted++
|
||||
if !w.resolveOne(ctx, f) {
|
||||
failed++
|
||||
remaining = append(remaining, f)
|
||||
}
|
||||
}
|
||||
if failed > 0 {
|
||||
@@ -164,7 +180,7 @@ func (w *factEnrichmentWorker) tick(ctx context.Context) {
|
||||
// Report the backlog every tick, not only when something failed: the
|
||||
// stalled state worth seeing is the one where nothing failed because
|
||||
// nothing was attempted.
|
||||
if st := w.status(ctx); st.Pending > 0 {
|
||||
if st := w.statusOf(remaining); st.Pending > 0 {
|
||||
log.Printf("factenrichment: %d facts pending entity resolution, %d in backoff, worst attempt %d (scanned %d)",
|
||||
st.Pending, st.InBackoff, st.MaxAttempts, st.Scanned)
|
||||
}
|
||||
|
||||
+30
-112
@@ -252,6 +252,23 @@ func run(args []string) error {
|
||||
// envelope per successful intake write.
|
||||
coreFor := func() ipc.CoreAPI { return newIntakeAPI(ipc.NewStoreAPI(st), evBus, time.Now) }
|
||||
|
||||
// depsNow reads whatever the current path has wired. Both boot paths build
|
||||
// the CoreAPI and start the workers from this one value, so neither can
|
||||
// hold a field the other misses. See cmd/mavend/boot.go.
|
||||
depsNow := func() bootDeps {
|
||||
return bootDeps{
|
||||
coreFor: coreFor,
|
||||
tl: tl,
|
||||
evBus: evBus,
|
||||
voiceW: voiceW,
|
||||
st: st,
|
||||
factWorker: factWorker,
|
||||
evalWorker: evalWorker,
|
||||
feedWkr: feedWkr,
|
||||
crawlWkr: crawlWkr,
|
||||
}
|
||||
}
|
||||
|
||||
if !locked {
|
||||
rules = wireRules(cfg)
|
||||
gatherer = wireGatherer(st, cfg, rules)
|
||||
@@ -284,26 +301,7 @@ func run(args []string) error {
|
||||
feedWkr = newFeedWorker(coreFor(), embedderOf(voiceW), cfg)
|
||||
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
|
||||
|
||||
coreAPI = &daemonAPI{
|
||||
CoreAPI: coreFor(),
|
||||
getTrace: tl.trace,
|
||||
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
|
||||
getDayPlan: func(ctx context.Context) ipc.DayPlan { return tl.dayPlan(ctx, time.Now()) },
|
||||
getEvents: intakeEventsFn(evBus),
|
||||
getDecisions: turnDecisionsFn(voiceW),
|
||||
seedStore: seedStoreIfAllowed(st),
|
||||
nexus: nexusOf(voiceW),
|
||||
}
|
||||
if voiceW != nil && voiceW.handler != nil {
|
||||
api := coreAPI.(*daemonAPI)
|
||||
api.chatFn = voiceW.handler.handleText
|
||||
// And the reverse: the handler was wired with the bare store
|
||||
// adapter, which cannot serve the day plan. See upgradeAPI.
|
||||
voiceW.handler.upgradeAPI(api)
|
||||
}
|
||||
if voiceW != nil && voiceW.mcp != nil {
|
||||
coreAPI.(*daemonAPI).getMCPServers = voiceW.mcp.status
|
||||
}
|
||||
coreAPI = newDaemonAPI(depsNow())
|
||||
} else {
|
||||
// locked mode: no real store yet, so there's no meaningful CoreAPI to
|
||||
// serve. srv.Check below is the actual guard — every CoreAPI call is
|
||||
@@ -364,6 +362,9 @@ func run(args []string) error {
|
||||
if !locked {
|
||||
wireMailIntake(srv, st, phr, cfg, evBus)
|
||||
wireModelSwap(srv, phr, cfg)
|
||||
// Inbound telegram (V-637). Dark unless the telegram block says intake,
|
||||
// and it reads one chat.
|
||||
wireTelegramIntake(ctx, &wg, coreAPI, cfg)
|
||||
// Vision + the media blob store (Vikunja #252). Both stay dark without a
|
||||
// media block; MethodDescribeImage answers ErrUnknownMethod then.
|
||||
keeper := wireVision(ctx, &wg, srv, st, embedderOf(voiceW), cfg)
|
||||
@@ -494,23 +495,14 @@ func run(args []string) error {
|
||||
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
|
||||
|
||||
// Swap the CoreAPI from the locked placeholder to the real store adapter.
|
||||
newAPI := &daemonAPI{
|
||||
CoreAPI: coreFor(),
|
||||
getTrace: tl.trace,
|
||||
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
|
||||
getDayPlan: func(ctx context.Context) ipc.DayPlan { return tl.dayPlan(ctx, time.Now()) },
|
||||
getEvents: intakeEventsFn(evBus),
|
||||
getDecisions: turnDecisionsFn(voiceW),
|
||||
seedStore: seedStoreIfAllowed(st),
|
||||
}
|
||||
if voiceW != nil && voiceW.handler != nil {
|
||||
newAPI.chatFn = voiceW.handler.handleText
|
||||
voiceW.handler.upgradeAPI(newAPI)
|
||||
}
|
||||
newAPI := newDaemonAPI(depsNow())
|
||||
srv.SetAPI(newAPI)
|
||||
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
|
||||
wireMailIntake(srv, st, phr, cfg, evBus)
|
||||
wireModelSwap(srv, phr, cfg)
|
||||
// Same on the unlock path, with the API that has just replaced the
|
||||
// locked placeholder (V-637).
|
||||
wireTelegramIntake(ctx, &wg, newAPI, cfg)
|
||||
keeper := wireVision(ctx, &wg, srv, st, embedderOf(voiceW), cfg)
|
||||
wireCapture(ctx, &wg, srv, keeper, st, voiceW, phr, cfg)
|
||||
// Voice identification (Vikunja #255). Enrolment plumbing only until a
|
||||
@@ -518,59 +510,10 @@ func run(args []string) error {
|
||||
// block, so no wire path takes a voiceprint on a default box.
|
||||
wireSpeaker(srv, st, cfg)
|
||||
|
||||
// Start voice server.
|
||||
if voiceW != nil {
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
if err := voiceW.server.Serve(); err != nil && !errors.Is(err, net.ErrClosed) {
|
||||
log.Printf("voice serve: %v", err)
|
||||
}
|
||||
}()
|
||||
log.Printf("mavend: voice listening on %s", voiceW.server.Addr())
|
||||
}
|
||||
|
||||
// Start tick loop.
|
||||
go func() {
|
||||
tl.run(ctx)
|
||||
}()
|
||||
|
||||
// Start fact-entity enrichment worker.
|
||||
go func() {
|
||||
factWorker.run(ctx)
|
||||
}()
|
||||
|
||||
// Start background memory evaluation (nil unless configured).
|
||||
if evalWorker != nil {
|
||||
go func() {
|
||||
evalWorker.run(ctx)
|
||||
}()
|
||||
}
|
||||
|
||||
// Start feed reading (nil unless configured).
|
||||
if feedWkr != nil {
|
||||
go func() {
|
||||
feedWkr.run(ctx)
|
||||
}()
|
||||
}
|
||||
|
||||
// Start the watched-page crawls (nil unless configured).
|
||||
if crawlWkr != nil {
|
||||
go func() {
|
||||
crawlWkr.run(ctx)
|
||||
}()
|
||||
}
|
||||
|
||||
// Keep MCP connections alive (nil unless configured).
|
||||
if voiceW != nil && voiceW.mcp != nil {
|
||||
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)
|
||||
}
|
||||
// The voice server and every background worker, on the outer wg
|
||||
// so shutdown waits for them. This used to be nine bare
|
||||
// `go func()` calls and a shadowed WaitGroup (V-639).
|
||||
startBackground(ctx, &wg, depsNow())
|
||||
|
||||
dl.unlock(st)
|
||||
log.Printf("mavend: unlocked via passkey assertion")
|
||||
@@ -585,33 +528,8 @@ func run(args []string) error {
|
||||
})
|
||||
log.Printf("mavend: ipc listening on %s", srv.Path())
|
||||
|
||||
if !locked && voiceW != nil {
|
||||
goWorker(&wg, func() {
|
||||
if err := voiceW.server.Serve(); err != nil && !errors.Is(err, net.ErrClosed) {
|
||||
log.Printf("voice serve: %v", err)
|
||||
}
|
||||
})
|
||||
log.Printf("mavend: voice listening on %s", voiceW.server.Addr())
|
||||
}
|
||||
|
||||
if !locked {
|
||||
goWorker(&wg, func() { tl.run(ctx) })
|
||||
goWorker(&wg, func() { factWorker.run(ctx) })
|
||||
if evalWorker != nil {
|
||||
goWorker(&wg, func() { evalWorker.run(ctx) })
|
||||
}
|
||||
if feedWkr != nil {
|
||||
goWorker(&wg, func() { feedWkr.run(ctx) })
|
||||
}
|
||||
if crawlWkr != nil {
|
||||
goWorker(&wg, func() { crawlWkr.run(ctx) })
|
||||
}
|
||||
if voiceW != nil && voiceW.mcp != nil {
|
||||
goWorker(&wg, func() { voiceW.mcp.run(ctx) })
|
||||
}
|
||||
if voiceW != nil && voiceW.home != nil {
|
||||
goWorker(&wg, func() { voiceW.home.run(ctx) })
|
||||
}
|
||||
startBackground(ctx, &wg, depsNow())
|
||||
}
|
||||
|
||||
<-ctx.Done()
|
||||
|
||||
@@ -0,0 +1,115 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
|
||||
// The floor, and it is the reason a destination is safe to add at all: a box
|
||||
// whose model is down names nothing, and naming nothing has to walk the chain
|
||||
// the way it walked before the field existed.
|
||||
func TestNoDestinationWalksTheWholeChain(t *testing.T) {
|
||||
walk, skipped := queryWalk(router.SourceUnknown)
|
||||
if len(skipped) != 0 {
|
||||
t.Errorf("skipped %d sources with no destination named, want none", len(skipped))
|
||||
}
|
||||
if len(walk) != len(querySources) {
|
||||
t.Fatalf("walk has %d sources, want the whole table of %d", len(walk), len(querySources))
|
||||
}
|
||||
for i := range walk {
|
||||
if walk[i].name != querySources[i].name {
|
||||
t.Fatalf("position %d is %q, want %q", i, walk[i].name, querySources[i].name)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The 2026-08-07 defects, one per line. Each is a source that decides by seed
|
||||
// similarity claiming a turn that was never its own, and then answering it
|
||||
// because it has no lookup that could come back empty.
|
||||
func TestANamedDestinationSilencesTheOtherGuessers(t *testing.T) {
|
||||
cases := []struct {
|
||||
dest router.Source
|
||||
utterance string
|
||||
silenced string
|
||||
}{
|
||||
{router.SourceWorld, "что такое TCP?", "weather"},
|
||||
{router.SourceWorld, "сколько будет 17 на 23?", "weather"},
|
||||
{router.SourceWorld, "кто такой Линус Торвальдс?", "personal"},
|
||||
{router.SourceRecall, "какой у меня любимый язык?", "feeds"},
|
||||
{router.SourceCalendar, "что в календаре на завтра?", "weather"},
|
||||
}
|
||||
for _, c := range cases {
|
||||
walk, skipped := queryWalk(c.dest)
|
||||
if inWalk(walk, c.silenced) {
|
||||
t.Errorf("%q named %q: %q is still asked", c.utterance, c.dest, c.silenced)
|
||||
}
|
||||
if !inWalk(skipped, c.silenced) {
|
||||
t.Errorf("%q named %q: %q is missing from the record of who was skipped",
|
||||
c.utterance, c.dest, c.silenced)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Naming the world must not send the turn outside. His notes, his facts and the
|
||||
// boundary in front of them are the invariant CLAUDE.md states as "the owner's
|
||||
// data first, then the world", and a destination a model wrote must not be able
|
||||
// to reverse it.
|
||||
func TestNamingTheWorldStillReadsHisDataFirst(t *testing.T) {
|
||||
walk, _ := queryWalk(router.SourceWorld)
|
||||
for _, look := range []string{"fact-by-key", "embed", "memory", "notes"} {
|
||||
if !inWalk(walk, look) {
|
||||
t.Errorf("%q was dropped; only the sources that guess may be dropped", look)
|
||||
}
|
||||
}
|
||||
if posOf(walk, "notes") > posOf(walk, "search") {
|
||||
t.Error("search is asked before his notes are")
|
||||
}
|
||||
if posOf(walk, "search") < 0 {
|
||||
t.Fatal("search is not in the walk at all")
|
||||
}
|
||||
}
|
||||
|
||||
// The boundary belongs to his data, so naming recall keeps it. That is what
|
||||
// makes "какой у меня любимый язык?" answer "не нашла у тебя такой записи"
|
||||
// rather than reaching SearXNG once nothing local had it.
|
||||
func TestNamingRecallKeepsTheBoundary(t *testing.T) {
|
||||
walk, _ := queryWalk(router.SourceRecall)
|
||||
if !inWalk(walk, "personal") {
|
||||
t.Fatal("the personal boundary was skipped on a turn named for his own data")
|
||||
}
|
||||
if posOf(walk, "personal") > posOf(walk, "search") {
|
||||
t.Error("the boundary no longer sits in front of the world")
|
||||
}
|
||||
}
|
||||
|
||||
// Whatever the destination, the walk is a subsequence of the table. Every
|
||||
// comment on that table argues an order between two sources, and none of those
|
||||
// reasons is about this field.
|
||||
func TestTheWalkNeverReordersTheTable(t *testing.T) {
|
||||
for _, dest := range append([]router.Source{router.SourceUnknown}, router.Sources...) {
|
||||
walk, skipped := queryWalk(dest)
|
||||
if len(walk)+len(skipped) != len(querySources) {
|
||||
t.Errorf("%q: %d walked + %d skipped, want %d", dest, len(walk), len(skipped), len(querySources))
|
||||
}
|
||||
last := -1
|
||||
for _, s := range walk {
|
||||
at := posOf(querySources, s.name)
|
||||
if at <= last {
|
||||
t.Errorf("%q: %q is out of table order", dest, s.name)
|
||||
}
|
||||
last = at
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func inWalk(list []querySource, name string) bool { return posOf(list, name) >= 0 }
|
||||
|
||||
func posOf(list []querySource, name string) int {
|
||||
for i, s := range list {
|
||||
if s.name == name {
|
||||
return i
|
||||
}
|
||||
}
|
||||
return -1
|
||||
}
|
||||
@@ -22,7 +22,7 @@ func newLLMReplier(c phraser.Completer, block func() string) *llmReplier {
|
||||
|
||||
// Reply never fails: a clarify, a model error and an unusable generation all
|
||||
// answer from the stub, which is what keeps a turn from breaking on the model.
|
||||
func (r *llmReplier) Reply(d router.Decision) string {
|
||||
func (r *llmReplier) Reply(ctx context.Context, d router.Decision) string {
|
||||
if d.Clarify {
|
||||
// The deck, not the stub's single sentence: a clarify she cannot turn
|
||||
// into a question is the line he hears most often when she misses him,
|
||||
@@ -39,14 +39,14 @@ func (r *llmReplier) Reply(d router.Decision) string {
|
||||
// что ты выпел стакан воды" for "я выпил воды".
|
||||
return phraser.FactAck(d.Utterance)
|
||||
}
|
||||
out, err := r.p.PhraseReply(context.Background(), d)
|
||||
out, err := r.p.PhraseReply(ctx, d)
|
||||
if err != nil || out == "" {
|
||||
return r.stub.Reply(d)
|
||||
return r.stub.Reply(ctx, d)
|
||||
}
|
||||
// The persona checks, on the live path (personaguard.go). A reply that
|
||||
// leaks reasoning or calls him "вы" is worse than a flat one.
|
||||
if _, ok := guardSpoken("reply", out); !ok {
|
||||
return r.stub.Reply(d)
|
||||
return r.stub.Reply(ctx, d)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ func (s stubCompleter) Complete(_ context.Context, _ llm.Req) (string, error) {
|
||||
|
||||
func TestLLMReplierPassesTheModelReplyThrough(t *testing.T) {
|
||||
r := newLLMReplier(stubCompleter{out: `{"response":"записала, кофе закончился","mood":"neutral"}`}, nil)
|
||||
got := r.Reply(router.Decision{Intent: router.IntentNote, Slots: router.Slots{Text: "кофе закончился"}})
|
||||
got := r.Reply(context.Background(), router.Decision{Intent: router.IntentNote, Slots: router.Slots{Text: "кофе закончился"}})
|
||||
if got != "записала, кофе закончился" {
|
||||
t.Errorf("got %q, want %q", got, "записала, кофе закончился")
|
||||
}
|
||||
@@ -42,7 +42,7 @@ func TestLLMReplierFallsBackToStubOnEmpty(t *testing.T) {
|
||||
// the clarify deck rather than the stub's single sentence.
|
||||
func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
|
||||
r := newLLMReplier(stubCompleter{out: "я всё поняла"}, nil)
|
||||
got := r.Reply(router.Decision{Clarify: true, Utterance: "мгм"})
|
||||
got := r.Reply(context.Background(), router.Decision{Clarify: true, Utterance: "мгм"})
|
||||
if got == "я всё поняла" {
|
||||
t.Fatal("a clarify must not be phrased by the model")
|
||||
}
|
||||
@@ -50,7 +50,7 @@ func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
|
||||
t.Errorf("on clarify: got %q, want %q", got, want)
|
||||
}
|
||||
// Two different misses do not sound identical.
|
||||
if same := r.Reply(router.Decision{Clarify: true, Utterance: "а"}); same == got {
|
||||
if same := r.Reply(context.Background(), router.Decision{Clarify: true, Utterance: "а"}); same == got {
|
||||
t.Log("two utterances hashed to the same line, which is allowed but should be rare")
|
||||
}
|
||||
}
|
||||
@@ -60,14 +60,14 @@ func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
|
||||
// produce, which is the same claim without pinning one wording.
|
||||
func assertAck(t *testing.T, r *llmReplier, d router.Decision, key, what string) {
|
||||
t.Helper()
|
||||
if got := r.Reply(d); !phraser.IsAck(key, nil, got) {
|
||||
if got := r.Reply(context.Background(), d); !phraser.IsAck(key, nil, got) {
|
||||
t.Errorf("on %s: got %q, want a %q line", what, got, key)
|
||||
}
|
||||
}
|
||||
|
||||
func assertStub(t *testing.T, r *llmReplier, d router.Decision, what string) {
|
||||
t.Helper()
|
||||
got, want := r.Reply(d), voice.NewStubReplier().Reply(d)
|
||||
got, want := r.Reply(context.Background(), d), voice.NewStubReplier().Reply(context.Background(), d)
|
||||
if got != want {
|
||||
t.Errorf("on %s: got %q, want stub %q", what, got, want)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
// mavend/telegramintake.go — wiring the inbound telegram poller (V-637).
|
||||
//
|
||||
// The poller reaches the daemon through ipc.CoreAPI and nothing else, so a
|
||||
// telegram turn takes exactly the path the web's POST /api/chat takes: Chat
|
||||
// returns the reply and the persisted trace id, and CorrectTurn writes the
|
||||
// label. Nothing in internal/delivery knows what a handler is.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"sync"
|
||||
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/delivery/telegramsink"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
|
||||
// wireTelegramIntake starts the poller, or returns having done nothing. It is
|
||||
// nil-safe in every argument, because it is called from both boot paths — the
|
||||
// unlocked start and the passkey unlock — and telegram must behave the same on
|
||||
// either.
|
||||
//
|
||||
// A sink that will not build is logged rather than fatal here. The push half
|
||||
// already failed the boot in wireDispatcher for the same config, so a second
|
||||
// hard failure would only lose that message.
|
||||
func wireTelegramIntake(ctx context.Context, wg *sync.WaitGroup, api ipc.CoreAPI, cfg *config.Config) {
|
||||
if cfg == nil || cfg.Telegram == nil || !cfg.Telegram.Intake || api == nil {
|
||||
return
|
||||
}
|
||||
sink, err := telegramsink.New(*cfg.Telegram)
|
||||
if err != nil {
|
||||
log.Printf("telegram intake: %v", err)
|
||||
return
|
||||
}
|
||||
poller, err := telegramsink.NewPoller(sink, chatTurnFn(api), api.CorrectTurn)
|
||||
if err != nil {
|
||||
log.Printf("telegram intake: %v", err)
|
||||
return
|
||||
}
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
poller.Run(ctx)
|
||||
}()
|
||||
}
|
||||
|
||||
// chatTurnFn adapts ipc.Chat to the poller's Turn. The trace id comes back on
|
||||
// the reply because the daemon's Chat collects it off the context (V-630), so
|
||||
// the chat can offer the same correction the web does without a second op.
|
||||
func chatTurnFn(api ipc.CoreAPI) telegramsink.Turn {
|
||||
return func(ctx context.Context, conversation, text string) (string, int64, error) {
|
||||
reply, err := api.Chat(ctx, conversation, text)
|
||||
if err != nil {
|
||||
return "", 0, err
|
||||
}
|
||||
return reply.Reply, reply.TraceID, nil
|
||||
}
|
||||
}
|
||||
@@ -279,3 +279,97 @@ func TestTheTurnIsRoutedOnce(t *testing.T) {
|
||||
t.Fatalf("the pipeline routed again and got something else: %+v vs %+v", second, first)
|
||||
}
|
||||
}
|
||||
|
||||
// TestASuspendedQuestionDoesNotRideForever — V-654, the measured failure of
|
||||
// 2026-08-07 (docs/evals/2026-08-07-week-of-usage-transcript.md, t=51 to t=58).
|
||||
//
|
||||
// A side query suspends the parked question, spends no attempt and restarts the
|
||||
// TTL. Nothing else bounded it, so one unfilled time slot came back on the end
|
||||
// of six consecutive unrelated replies and stopped only when a seventh turn
|
||||
// happened to read as a failed answer. Three step-asides, then she lets it go
|
||||
// and says so.
|
||||
func TestASuspendedQuestionDoesNotRideForever(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
h, st := newRoutingClarifyHandler(t)
|
||||
id := dialogueIDFor(sourceText, "web")
|
||||
resumed, _ := clarifyResumedFor(dialogue.SlotTime)
|
||||
|
||||
if reply := h.handleText(ctx, "web", "напомни позвонить маме"); !strings.Contains(reply, "?") {
|
||||
t.Fatalf("expected the time question, got %q", reply)
|
||||
}
|
||||
|
||||
// Three questions of his own. Each one is answered as itself and each one
|
||||
// brings the open question back, exactly as V-561 asks.
|
||||
asides := []string{
|
||||
"о чём мы вчера говорили?",
|
||||
"какие у меня напоминания?",
|
||||
"сколько времени?",
|
||||
}
|
||||
for i, text := range asides {
|
||||
reply := h.handleText(ctx, "web", text)
|
||||
if !strings.HasSuffix(reply, resumed) {
|
||||
t.Fatalf("side query %d: the question must come back, got %q", i+1, reply)
|
||||
}
|
||||
if strings.Contains(reply, clarifyDropped) {
|
||||
t.Fatalf("side query %d: nothing was let go yet, so nothing may say so: %q", i+1, reply)
|
||||
}
|
||||
q := h.clarifyStore.Get(id, h.now())
|
||||
if q == nil {
|
||||
t.Fatalf("side query %d: the question was dropped early", i+1)
|
||||
}
|
||||
if q.Attempts != 1 {
|
||||
t.Fatalf("side query %d: a step-aside spent an attempt: %d", i+1, q.Attempts)
|
||||
}
|
||||
if q.Suspends != i+1 {
|
||||
t.Fatalf("side query %d: suspends = %d, want %d", i+1, q.Suspends, i+1)
|
||||
}
|
||||
}
|
||||
|
||||
// The fourth. She has stepped aside as often as she is willing to, so the
|
||||
// request goes — out loud, and without the question on the tail.
|
||||
reply := h.handleText(ctx, "web", "что у меня сегодня?")
|
||||
if !strings.Contains(reply, clarifyDropped) {
|
||||
t.Fatalf("the request was let go in silence: %q", reply)
|
||||
}
|
||||
if strings.HasSuffix(reply, resumed) {
|
||||
t.Fatalf("a question she has let go must not be asked again: %q", reply)
|
||||
}
|
||||
if h.clarifyStore.Get(id, h.now()) != nil {
|
||||
t.Fatal("the question must be gone once she has said she let it go")
|
||||
}
|
||||
if reminders, err := st.DueReminders(ctx, h.now().Add(48*time.Hour)); err != nil || len(reminders) != 0 {
|
||||
t.Fatalf("a reminder was invented for a time nobody gave: %v err=%v", reminders, err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestAnAnsweredGapResetsTheSuspendBudget — the counter measures CONSECUTIVE
|
||||
// step-asides. He filled a gap, so the run is broken and the next question
|
||||
// starts with its full allowance: a long exchange he is engaged with must not
|
||||
// run out of patience on his behalf.
|
||||
func TestAnAnsweredGapResetsTheSuspendBudget(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
h, _ := newRoutingClarifyHandler(t)
|
||||
id := dialogueIDFor(sourceText, "web")
|
||||
|
||||
// A bare "напомни" is missing both halves, so answering the subject re-parks
|
||||
// the request with a question about the time.
|
||||
if reply := h.handleText(ctx, "web", "напомни"); !strings.Contains(reply, "?") {
|
||||
t.Fatalf("expected a question, got %q", reply)
|
||||
}
|
||||
if reply := h.handleText(ctx, "web", "какие у меня напоминания?"); reply == "" {
|
||||
t.Fatal("the side query must be answered as itself")
|
||||
}
|
||||
if q := h.clarifyStore.Get(id, h.now()); q == nil || q.Suspends != 1 {
|
||||
t.Fatalf("the side query was not counted: %+v", q)
|
||||
}
|
||||
if reply := h.handleText(ctx, "web", "позвонить маме"); reply == "" {
|
||||
t.Fatal("the answer must be consumed")
|
||||
}
|
||||
q := h.clarifyStore.Get(id, h.now())
|
||||
if q == nil {
|
||||
t.Fatal("a reminder still needs its time, so a question must be parked")
|
||||
}
|
||||
if q.Suspends != 0 {
|
||||
t.Fatalf("answering a gap must reset the suspend budget: suspends = %d", q.Suspends)
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
@@ -458,7 +458,7 @@ func (h *reactiveHandler) runTurn(ctx context.Context, text string, src turnSour
|
||||
|
||||
// 9. replier — phrase the reply across the router decision.
|
||||
if replyText == "" {
|
||||
replyText = h.replier.Reply(dec)
|
||||
replyText = h.replier.Reply(ctx, dec)
|
||||
}
|
||||
return withNotice(expiredNotice, replyText)
|
||||
}
|
||||
|
||||
@@ -401,6 +401,10 @@ func buildRouter(emb router.Embedder, acts router.ActMatcher, threshold float64,
|
||||
grammars = append(grammars, router.AgendaQueryGrammars()...)
|
||||
// Same reason as the agenda rules, for the feeds: "что нового в лентах?"
|
||||
// routed system and answered "пока не умею" (Vikunja #474).
|
||||
// After the agenda rules, which are the narrower claim, and BEFORE the feed
|
||||
// and list rules, which are not: "что такое лента" is a definition question
|
||||
// and the feed rule would take it on the noun alone (V-655).
|
||||
grammars = append(grammars, router.WorldQueryGrammars()...)
|
||||
grammars = append(grammars, router.FeedQueryGrammar())
|
||||
// The list side of the same exposure: a phrasing with no possessive in it
|
||||
// ("список дел") routed system and never reached queryTasks (Vikunja #467).
|
||||
|
||||
+19
-1
@@ -58,6 +58,12 @@ func main() {
|
||||
// mutex, so sharing the connection would freeze every other page for the
|
||||
// length of the load. See handleModels.
|
||||
var swapConn modelController
|
||||
// turnConn — a third connection, for POST /api/chat and nothing else, for
|
||||
// the same reason /models has one (V-638). A chat turn routes, phrases and
|
||||
// may act, bounded only by phraser.timeout at 60s, and every other handler
|
||||
// on this server queues behind it on the shared client's one mutex. Nil ⇒
|
||||
// chat shares the main connection, which is how it behaved before.
|
||||
var turnConn ipc.CoreAPI
|
||||
if *coreSock != "" {
|
||||
c, err := ipc.DialWait(*coreSock, 60*time.Second)
|
||||
if err != nil {
|
||||
@@ -71,6 +77,12 @@ func main() {
|
||||
defer sc.Close()
|
||||
swapConn = sc
|
||||
}
|
||||
if tc, err := ipc.Dial(*coreSock); err != nil {
|
||||
log.Printf("chat: third core connection failed (%v) — /api/chat will share the main one and a turn will block the other pages", err)
|
||||
} else {
|
||||
defer tc.Close()
|
||||
turnConn = tc
|
||||
}
|
||||
}
|
||||
|
||||
// stepUpSession stays nil unless the passkey endpoints are wired below — it
|
||||
@@ -208,7 +220,13 @@ func main() {
|
||||
// decides how every utterance is routed and how every reply is worded.
|
||||
mux.HandleFunc("/tools", gatedPage(handleTools))
|
||||
mux.HandleFunc("/routines", gatedPage(handleRoutines))
|
||||
mux.HandleFunc("/api/chat", gatedPage(handleChatAPI))
|
||||
mux.HandleFunc("/api/chat", func(w http.ResponseWriter, r *http.Request) {
|
||||
c := turnConn
|
||||
if c == nil {
|
||||
c = core
|
||||
}
|
||||
handleChatAPI(w, r, c, stepUpSession, *requireStepUp)
|
||||
})
|
||||
mux.HandleFunc("/api/revert", gatedPage(handleRevert))
|
||||
mux.HandleFunc("/api/correct", gatedPage(handleCorrectAPI))
|
||||
mux.HandleFunc("/models", func(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
+28
-1
@@ -25,6 +25,25 @@
|
||||
"llm_nudges": false
|
||||
},
|
||||
|
||||
"//ntfy": [
|
||||
"The second reach (V-649). Until 07-08-2026 telegram was the only one, and",
|
||||
"telegram needs api.telegram.org, the socks relay below and a matching ufw",
|
||||
"rule — three things in series that have each failed once, and when they do",
|
||||
"a sev4 nudge has nowhere to go. ntfy shares none of them: it is reached",
|
||||
"directly, no relay.",
|
||||
"It is not only a spare. The routing table sends sev3-away and away",
|
||||
"reminders here and NOWHERE else, so with this block absent those two",
|
||||
"routes hit a nil sink and vanish without a log or an outbox row.",
|
||||
"The credential is an ntfy access token, scoped write-only to this one",
|
||||
"topic, so a popped sink can push to it and cannot read it back. Set it in",
|
||||
"deploy/telegram.env beside the telegram secrets; that file is gitignored."
|
||||
],
|
||||
"ntfy": {
|
||||
"base_url": "https://ntfy.kvmx.ru",
|
||||
"topic": "maven",
|
||||
"token": "${NTFY_TOKEN}"
|
||||
},
|
||||
|
||||
"telegram": {
|
||||
"bot_token": "${TELEGRAM_BOT_TOKEN}",
|
||||
"chat_id": "${TELEGRAM_CHAT_ID}",
|
||||
@@ -37,7 +56,15 @@
|
||||
"This needs a matching ufw rule or the container's SYN is dropped:",
|
||||
" ufw allow from 192.168.240.0/20 to any port 10808 proto tcp"
|
||||
],
|
||||
"proxy": "socks5://192.168.240.1:10808"
|
||||
"proxy": "socks5://192.168.240.1:10808",
|
||||
"//intake": [
|
||||
"Read the chat as well as write to it (V-637). The poller long-polls",
|
||||
"getUpdates through the same relay and accepts chat_id as the only",
|
||||
"sender. Deleting this key turns inbound off again.",
|
||||
"chat_id must be numeric here or the daemon refuses to start: an inbound",
|
||||
"update names its chat by number, so an @-name would match nothing."
|
||||
],
|
||||
"intake": true
|
||||
},
|
||||
|
||||
"//workstation": [
|
||||
|
||||
@@ -1,5 +1,12 @@
|
||||
# Telegram bot token and chat ID for mavend's away-channel reach.
|
||||
# Secrets for mavend's away-channel reaches. The file is still called
|
||||
# telegram.env because compose names it that; it holds both reaches now.
|
||||
# Copy this file to deploy/telegram.env and fill in real values.
|
||||
# deploy/telegram.env is gitignored — never commit the real secrets.
|
||||
TELEGRAM_BOT_TOKEN=
|
||||
TELEGRAM_CHAT_ID=
|
||||
|
||||
# ntfy access token for the `maven` topic, the second reach (V-649). Mint it on
|
||||
# the ntfy server with write access to that topic and nothing else:
|
||||
# ntfy token add --expires=never maven
|
||||
# Read access is not needed — mavend publishes and never subscribes.
|
||||
NTFY_TOKEN=
|
||||
|
||||
@@ -157,6 +157,44 @@ services:
|
||||
# - maildata:/var/lib/mavmaild
|
||||
# - ./deploy/imap.password:/run/secrets/imap.password:ro
|
||||
|
||||
# The calendar reader (Vikunja #644) is OFF and commented out: it needs a
|
||||
# CalDAV account, and there is none on this box. It was built, listed in
|
||||
# `make build`, and deployed nowhere, which is the worst of the three states —
|
||||
# this block records the decision instead.
|
||||
#
|
||||
# What its absence costs, so the cost is visible from here:
|
||||
# - Agenda questions route correctly and answer from nothing. Stage 0 sends
|
||||
# "что у меня сегодня" to IntentQuery (V-498) and the `calendar` query
|
||||
# source reads facts(kind=env, source=caldav:*) that nobody writes.
|
||||
# - The nudge gate loses a suppressor. loop.State.CalendarBusy is fed by
|
||||
# those same facts, so "do not nag mid-meeting" is permanently false.
|
||||
#
|
||||
# Core never sees the CalDAV password: the reader polls the collection itself
|
||||
# and hands core one fact per event over WriteFact. Nothing here can create a
|
||||
# reminder, so a misread event cannot fire.
|
||||
#
|
||||
# The password is read from a FILE, so it never appears in `ps`, in this file,
|
||||
# or in shell history — the same rule mavpoll and mavmaild follow.
|
||||
#
|
||||
# To enable: write the password to deploy/caldav.password (0600, gitignored),
|
||||
# point -url at the collection, and uncomment this service. No mavend.json
|
||||
# block is needed — events arrive over IPC as facts. -render-url is optional
|
||||
# and OFF here: it publishes Maven's own reminders back as events, and it must
|
||||
# not name the collection -url reads, or the poller reads its own writes back
|
||||
# in (checkRenderTarget refuses that). It takes -render-pass-file, and falls
|
||||
# back to this password when that is not given.
|
||||
# mavcaldav:
|
||||
# <<: *image
|
||||
# command: ["mavcaldav", "-socket", "/run/maven/mavend.sock",
|
||||
# "-url", "http://localhost:5232/kami/personal",
|
||||
# "-user", "kami",
|
||||
# "-pass-file", "/run/secrets/caldav.password",
|
||||
# "-interval", "5m"]
|
||||
# depends_on: [mavend]
|
||||
# volumes:
|
||||
# - sockets:/run/maven
|
||||
# - ./deploy/caldav.password:/run/secrets/caldav.password:ro
|
||||
|
||||
volumes:
|
||||
dbdata:
|
||||
sockets:
|
||||
|
||||
+34
-1
@@ -1,6 +1,6 @@
|
||||
# Maven — Design
|
||||
|
||||
*Last verified: 2026-08-02 @ 7079a24. Living doc: correct it in place, do not append.*
|
||||
*Last verified: 2026-08-07 @ beb093a. Living doc: correct it in place, do not append.*
|
||||
|
||||
> Folded 2026-07-30 from `SPEC.md` (north star, 2026-07-03), `maven.md`
|
||||
> (consolidated decisions, 2026-06-30) and `ROADMAP.md` (execution plan,
|
||||
@@ -282,6 +282,39 @@ Three reasons, in the order they settle it:
|
||||
So the notice stays what it is: the in-process TTL case, where she really did
|
||||
wait and really did let go.
|
||||
|
||||
#### A parked question may step aside three times
|
||||
|
||||
Decided 2026-08-07 (V-654). A side query or an aside suspends the parked
|
||||
question instead of dropping it. The words are answered as themselves, and the
|
||||
question comes back on the end of the same reply.
|
||||
|
||||
Neither bound on a question reaches that path. No attempt is spent, because a
|
||||
side query is not a failed answer, so `MaxAttempts` never applies.
|
||||
`noteSuspended` also restarts the 90s clock, since she is about to speak the
|
||||
question again. So the TTL cannot arrive while he keeps talking.
|
||||
|
||||
Measured on 2026-08-07: one unfilled time slot rode the tail of six consecutive
|
||||
unrelated replies. It stopped only when a seventh turn happened to read as a
|
||||
failed answer. See `docs/evals/2026-08-07-week-of-usage.md`.
|
||||
|
||||
`PendingQuestion.Suspends` counts the step-asides. `MaxSuspends` is 3, matching
|
||||
`DefaultMaxAttempts`. Past it she lets the request go, with the same
|
||||
`clarifyDropped` line every other drop uses. The owner's rule is unchanged. A
|
||||
question still ends by being answered or by being let go out loud. This only
|
||||
recognises three unrelated requests in a row as the second of those.
|
||||
|
||||
The count is of CONSECUTIVE step-asides. It resets the moment he answers, in
|
||||
`resolveClarifyAnswer`. An answer that gives her nothing she asked for resets it
|
||||
too. "Позвонить маме" against a question about the time is still him in the
|
||||
exchange. The retry it costs is bound enough on its own.
|
||||
|
||||
The re-ask is also two sentences rather than one. It used to be spliced onto the
|
||||
answer with a comma. On a real answer that buries the question in the tail of
|
||||
one run-on thought:
|
||||
|
||||
> вот что я нашла: вайфай пароль лежит в ящике стола, на какое время поставить
|
||||
> напоминание?
|
||||
|
||||
### save-where — the two-memory routing axis
|
||||
|
||||
One discriminator: **does the loop evaluate a predicate against it?**
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
# The routing trajectory, and the number that is missing
|
||||
|
||||
**06-08-2026. V-464.** Not a new measurement. This collates the figures already recorded
|
||||
in `docs/evals/` and CLAUDE.md, and names one measurement that has not been taken. Dated
|
||||
because the conclusion expires the moment the missing number is measured.
|
||||
|
||||
## The question
|
||||
|
||||
126 of the 1023 commits between 03-07-2026 and 06-08-2026 touch `internal/router`. Is the
|
||||
routing between the core functions and his speech getting better?
|
||||
|
||||
## The trajectory
|
||||
|
||||
RU routing fixture, classifier plus the ONNX embedder, no LLM arm in any of these runs.
|
||||
|
||||
| date | change | fixture | source |
|
||||
|---|---|---|---|
|
||||
| 02-08-2026 | classifier re-measured | 68.8% of 77 | CLAUDE.md |
|
||||
| 04-08-2026 | V-498, rest-of-day and narrative rules | 58/82, 70.7% | CLAUDE.md |
|
||||
| 06-08-2026 | V-626 baseline | 64/91, 70.3% | `2026-08-06-seeds-to-prompt-boundary.md` |
|
||||
| 06-08-2026 | V-626, seeds onto the prompt boundary | 66/91, 72.5% | same |
|
||||
| 06-08-2026 | V-627, alarm verbs reach stage 0 | 69/91, 75.8% | `2026-08-06-alarm-verbs-reach-stage-0.md` |
|
||||
| 06-08-2026 | V-633, Russian acts reach tools | 69/91, unchanged | `2026-08-06-russian-acts-reach-tools.md` |
|
||||
|
||||
The fixture grew from 77 to 82 to 91 cases across this window. So the percentages are
|
||||
comparable and the counts are not.
|
||||
|
||||
## Accuracy moved late
|
||||
|
||||
It sat near 70% for a month. V-626 and V-627 landed the same day and took the
|
||||
deterministic path from 64/91 to 69/91. That is the first real accuracy movement since the
|
||||
stage-0 rules went in.
|
||||
|
||||
## Most of the work was reach, not accuracy
|
||||
|
||||
Praxis went 0/12 to 11/12 and lifecycle 0/5 to 5/5 (V-516,
|
||||
`2026-08-05-praxis-reach.md`). No Russian utterance could reach a tool before V-633. That
|
||||
one landed at 69/91 unchanged, because the fixture holds no case for it. Alarm verbs,
|
||||
ordinal selection, spoken corrections and the claimant ladder share the shape.
|
||||
|
||||
So the fixture undercounts the month. Things that were structurally unreachable now reach,
|
||||
and a fixture that never asked about them cannot show it. Judge reach against
|
||||
`make eval-reach` and the ecosystem fixture, not against the routing one.
|
||||
|
||||
## The missing number
|
||||
|
||||
On 05-08-2026 the cascade with the resident model scored 69/91, 75.8% full, 80.2%
|
||||
intent-only, at p50 1.19s (`2026-08-05-routing-resident-model.md`).
|
||||
|
||||
On 06-08-2026 the classifier and stage 0 alone reached 69/91, 75.8% full, at p50 22.9ms.
|
||||
|
||||
Those are the same full-accuracy score. The cascade has not been re-measured since V-626
|
||||
and V-627 landed. Both are stage-0 changes, and stage 0 runs inside the cascade, so the
|
||||
cascade should have gained from them too.
|
||||
|
||||
One of two things is true, and nothing on the box says which:
|
||||
|
||||
- The cascade gained as well, the model still separates from the floor on intent-only, and
|
||||
it earns its place.
|
||||
- The deterministic floor has caught up on this fixture, and the resident model is costing
|
||||
1.17 seconds a turn for nothing measurable.
|
||||
|
||||
Take that measurement before planning more routing work. It needs a second llama-server on
|
||||
a fixed host port, because the resident one binds `--port 0` inside the container.
|
||||
|
||||
## What this does not settle
|
||||
|
||||
Intent-only is the more honest comparison for the model arm. The model routes `reminder`
|
||||
and leaves the time to the daemon, which is what the contract asks. The 05-08 run puts it
|
||||
at 80.2% through the cascade and 61.5% for the model alone. There is no 06-08 intent-only
|
||||
figure for the deterministic path to set beside those.
|
||||
@@ -0,0 +1,69 @@
|
||||
# Does one sqlite connection make reads queue? No (V-642)
|
||||
|
||||
Measured 07-08-2026 at `7b507de`, on homesrv. The harness is
|
||||
`internal/store/conncap_test.go`. It stays in the repo, because this claim gets
|
||||
re-argued and the numbers should be re-runnable rather than quoted.
|
||||
|
||||
`internal/store/store.go` opens the database with `SetMaxOpenConns(1)`, while
|
||||
`schema.sql` sets `journal_mode=WAL`. WAL exists to let readers run beside one
|
||||
writer, so the cap gives up the thing the journal mode was chosen for. The
|
||||
question was whether that costs anything.
|
||||
|
||||
## What was measured
|
||||
|
||||
A fixed two-second window. One writer calling `SetValue` paced at 2ms, and a
|
||||
reader loop calling `RecentFacts(50)` over 500 seeded rows as fast as it can.
|
||||
Same schema, same modernc driver, same machine, three runs per cap.
|
||||
|
||||
The window is wall-clock rather than a read count on purpose. A first version ran
|
||||
a fixed 300 reads. That finished sooner at the higher cap, so it received fewer
|
||||
writes, and two runs that did different work cannot be compared.
|
||||
|
||||
| cap | reads | writes | p50 | p95 | max |
|
||||
|---|---|---|---|---|---|
|
||||
| 1 | ~3050 | ~760 | 594µs | 900µs | 16-19ms |
|
||||
| 4 | ~3600 | ~340 | 525µs | 710µs | 1-2ms |
|
||||
|
||||
## What it says
|
||||
|
||||
**Reads do not queue behind writes.** Four connections buy about 70µs at p50. A
|
||||
turn spends 1.19s in the resident model. The tail does improve, from 19ms to 2ms,
|
||||
and 19ms is still not a figure anyone notices in a spoken reply.
|
||||
|
||||
**Write throughput more than halves at the higher cap**, 760 writes against 340.
|
||||
inference, not measured directly: at one connection the reader and the writer take
|
||||
turns with no lock contention. At four the writer contends for the WAL write lock
|
||||
with a live reader. Whatever the mechanism, the trade runs the opposite way from
|
||||
the one the task expected.
|
||||
|
||||
**The cap was not the source of the 2.7s router figure.** CLAUDE.md records that
|
||||
figure as contention rather than the model. This task was a candidate for where
|
||||
that contention came from. A 19ms worst case cannot produce it. That line of
|
||||
enquiry is closed.
|
||||
|
||||
**One transaction is what the cap cannot survive.** With a read-only transaction
|
||||
open, a second read at cap 1 never completes. The harness gave it two seconds and
|
||||
got `context deadline exceeded`. The same read at cap 4 took 1ms. The transaction
|
||||
holds the only connection, so this is not a slow read, it is a stalled database.
|
||||
|
||||
## What was done
|
||||
|
||||
The cap stays at 1. The reason is now written where the cap is set, rather than
|
||||
inferred from a four-word comment.
|
||||
|
||||
`Store.DB` was deleted. It handed out exactly the read-only transaction measured
|
||||
above. It had been there since the initial commit with no production caller, and
|
||||
its doc comment described a loop that never materialised. Its one user was a test
|
||||
helper reading `delivery_attempts` by raw SQL. `ListDeliveryAttempts` has covered
|
||||
that since V-390, and the helper now goes through the reader.
|
||||
|
||||
So the hazard is gone by construction, not by documentation.
|
||||
`TestConnCap_ReadBlocksBehindOpenSnapshot` is the standing measurement of what
|
||||
re-adding the seam would cost.
|
||||
|
||||
## Not answered
|
||||
|
||||
Whether reads queue on the deployed box under real load, as opposed to a
|
||||
synthetic loop. The harness writes and reads one table. Digestion reads four and
|
||||
embeds while it does. The finding that closes this task is the transaction stall,
|
||||
which is structural and does not depend on load.
|
||||
@@ -0,0 +1,321 @@
|
||||
# Raw transcript, 2026-08-07 week of usage
|
||||
|
||||
Companion to `2026-08-07-week-of-usage.md`. 74 turns through `POST /api/chat`.
|
||||
Days are simulated; the wall clock was 03:27-03:40 throughout, which is why every
|
||||
clock answer says night.
|
||||
|
||||
## Day 1
|
||||
|
||||
```
|
||||
YOU: доброе утро
|
||||
MAVEN: Добро пожаловать! Солнечно и тепло. Всё в порядке.
|
||||
[7.7s t=5]
|
||||
|
||||
YOU: какой сегодня день?
|
||||
MAVEN: сегодня пятница, 7 августа 2026 года
|
||||
[0.0s t=6]
|
||||
|
||||
YOU: напомни в 11:00 позвонить маме
|
||||
MAVEN: Сейчас 03:29. В какой день?
|
||||
[0.4s t=7]
|
||||
|
||||
YOU: запиши что я пью кофе без сахара
|
||||
MAVEN: Прошлую просьбу отпускаю. Заметка сохранена.
|
||||
[8.0s t=8]
|
||||
|
||||
YOU: что у меня сегодня?
|
||||
MAVEN: на 07.08.2026 ничего нет.
|
||||
[2.6s t=9]
|
||||
|
||||
YOU: мой любимый язык программирования — go
|
||||
MAVEN: Хорошо, сохранил.
|
||||
[7.6s t=10]
|
||||
|
||||
YOU: сколько времени?
|
||||
MAVEN: сейчас 3 часа 29 минут
|
||||
[0.0s t=11]
|
||||
|
||||
YOU: спасибо
|
||||
MAVEN: спасиба
|
||||
[8.8s t=12]
|
||||
|
||||
YOU: сегодня
|
||||
MAVEN: Прости, не поняла — попробуй сказать иначе.
|
||||
[5.1s t=13]
|
||||
|
||||
YOU: а что ты умеешь?
|
||||
MAVEN: я могу помочь с задачами, напоминаниями и поиском информации.
|
||||
[4.9s t=14]
|
||||
|
||||
YOU: почему небо голубое?
|
||||
MAVEN: нашла вот это: Почему небо голубое: объяснение цвета неба простыми словами, физика ... — Ответ на вопрос, почему небо голубое, кроется в особенностях солнечного света и его взаимодействии с атмосферой Земли. Рассказываем, как это работает. В статье разберемся, почему днем небо кажется голубым, а на закате окрашива…
|
||||
[12.3s t=15]
|
||||
|
||||
YOU: столица Франции?
|
||||
MAVEN: пока не умею отвечать на этот вопрос.
|
||||
[1.5s t=16]
|
||||
|
||||
YOU: кто такой Линус Торвальдс?
|
||||
MAVEN: не знаю — не нашла у тебя такой записи.
|
||||
[1.8s t=17]
|
||||
|
||||
YOU: ладно, до вечера
|
||||
MAVEN: Ладно, до вечера.
|
||||
[6.3s t=18]
|
||||
|
||||
```
|
||||
|
||||
## Day 2
|
||||
|
||||
```
|
||||
YOU: привет, как спалось?
|
||||
MAVEN: пока не умею отвечать на этот вопрос.
|
||||
[1.3s t=19]
|
||||
|
||||
YOU: что я пью?
|
||||
MAVEN: не знаю — не нашла у тебя такой записи.
|
||||
[1.2s t=20]
|
||||
|
||||
YOU: какой у меня любимый язык?
|
||||
MAVEN: нового вот что: Six stable kernels with a security fix; [$] Bringing BPF to binfmt_misc; Rust Coreutils 0.10 released
|
||||
[1.5s t=21]
|
||||
|
||||
YOU: добавь задачу починить кран
|
||||
MAVEN: внесла в список: починить кран
|
||||
[0.0s t=22]
|
||||
|
||||
YOU: что мне надо сделать?
|
||||
MAVEN: сначала: проверить интейк форму на боксе (через 2 дня); купить хлеб; купить молоко; полить цветы; починить кран
|
||||
[0.0s t=23]
|
||||
|
||||
YOU: запомни что вайфай пароль лежит в ящике стола
|
||||
MAVEN: Хорошо, сохранил.
|
||||
[11.2s t=24]
|
||||
|
||||
YOU: погода сегодня какая?
|
||||
MAVEN: пока не умею отвечать на этот вопрос.
|
||||
[1.8s t=25]
|
||||
|
||||
YOU: выключи свет в спальне
|
||||
MAVEN: Такую команду я не знаю.
|
||||
[1.5s t=26]
|
||||
|
||||
YOU: покажи мои заметки
|
||||
MAVEN: пока не умею отвечать на этот вопрос.
|
||||
[1.1s t=27]
|
||||
|
||||
YOU: что ты знаешь обо мне?
|
||||
MAVEN: не знаю — не нашла у тебя такой записи.
|
||||
[1.3s t=28]
|
||||
|
||||
YOU: я люблю go
|
||||
MAVEN: Я тоже love Go! Я проверила, что твой любимый язык — это Go. Ты должен быть очень рад, когда будет новая версия или пакет в стандартной библиотеке.
|
||||
[7.0s t=29]
|
||||
|
||||
YOU: какой язык я люблю?
|
||||
MAVEN: вот что я нашла: ты любишь_go
|
||||
[6.8s t=30]
|
||||
|
||||
```
|
||||
|
||||
## Day 3
|
||||
|
||||
```
|
||||
YOU: расскажи про Байкал
|
||||
MAVEN: есть такое: Байкал - Википедия — Байгал далай) — озеро тектонического происхождения в южной части Восточной Сибири площадью 31 722 км². Самое глубокое озеро на планете, крупнейший природный резервуар пресной воды и самое большое по площади пресноводное озеро Евразии. Озеро и прибрежные территории являются уника…
|
||||
[1.4s t=31]
|
||||
|
||||
YOU: что такое TCP?
|
||||
MAVEN: для какого города?
|
||||
[1.4s t=32]
|
||||
|
||||
YOU: когда родился Пушкин?
|
||||
MAVEN: пока не умею отвечать на этот вопрос.
|
||||
[1.4s t=33]
|
||||
|
||||
YOU: сколько будет 17 на 23?
|
||||
MAVEN: для какого города?
|
||||
[1.7s t=34]
|
||||
|
||||
YOU: переведи 'thank you' на русский
|
||||
MAVEN: не знаю — не нашла у тебя такой записи.
|
||||
[3.7s t=35]
|
||||
|
||||
YOU: что нового в мире?
|
||||
MAVEN: вот что я нашла: В мире - последние новости сегодня - РИА Новости — В мире. Читайте последние новости на тему в ленте новостей на сайте РИА ... новые подробности о визите представителя Ватикана в Москву. 22:28. 0. Еще 2. В ...
|
||||
[2.3s t=36]
|
||||
|
||||
```
|
||||
|
||||
## Day 4
|
||||
|
||||
```
|
||||
YOU: ЭЙ
|
||||
MAVEN: Не уловила. Скажи это по-другому?
|
||||
[1.2s t=37]
|
||||
|
||||
YOU: ладно забудь
|
||||
MAVEN: пока не умею отвечать на этот вопрос.
|
||||
[1.3s t=38]
|
||||
|
||||
YOU: нет стой, напомни через час выпить воды
|
||||
MAVEN: поняла, это напоминание — запомнила. не получилось разобрать время напоминания.
|
||||
[0.8s t=39]
|
||||
|
||||
YOU: а вообще какой смысл жизни
|
||||
MAVEN: Я не могу ответить на этот вопрос. Моя цель — помочь тебе с задачами и отвечать на вопросы.
|
||||
[7.7s t=40]
|
||||
|
||||
YOU: ты меня бесишь
|
||||
MAVEN: я не могу быть с тобой так. я понимаю, как это жестоко.
|
||||
[4.2s t=41]
|
||||
|
||||
YOU: прости, я не всерьёз
|
||||
MAVEN: я записала, что прошу простить за мою непоследовательность
|
||||
[7.2s t=42]
|
||||
|
||||
YOU: do you speak english?
|
||||
MAVEN: пока не умею отвечать на этот вопрос.
|
||||
[1.1s t=43]
|
||||
|
||||
YOU: кран починил кстати
|
||||
MAVEN: записала: кран починил кстати
|
||||
[1.6s t=44]
|
||||
|
||||
YOU: закрой задачу про кран
|
||||
MAVEN: не нашла такой задачи.
|
||||
[0.0s t=45]
|
||||
|
||||
YOU: ааааа
|
||||
MAVEN: Не совсем поняла — скажи, пожалуйста, ещё раз.
|
||||
[0.9s t=46]
|
||||
|
||||
YOU: я тебя обожаю
|
||||
MAVEN: я тебя обожаю
|
||||
[2.3s t=47]
|
||||
|
||||
YOU: молока купи
|
||||
MAVEN: Такую команду я не знаю.
|
||||
[4.4s t=48]
|
||||
|
||||
YOU: нет не так
|
||||
MAVEN: поняла, отметила, что ответила не так.
|
||||
[0.0s t=49]
|
||||
|
||||
YOU: ты вообще меня слушаешь?
|
||||
MAVEN: я записала, выслушаю тебя
|
||||
[8.4s t=50]
|
||||
|
||||
YOU: окей окей окей что там с водой напоминание
|
||||
MAVEN: Сейчас 03:33. Когда?
|
||||
[5.2s t=51]
|
||||
|
||||
YOU: всё, спать
|
||||
MAVEN: пока не умею отвечать на этот вопрос, на какое время поставить напоминание?
|
||||
[1.4s t=52]
|
||||
|
||||
```
|
||||
|
||||
## Day 5
|
||||
|
||||
```
|
||||
YOU: утро
|
||||
MAVEN: Сейчас 03:33. Во сколько напомнить?
|
||||
[0.5s t=53]
|
||||
|
||||
YOU: о чём мы вчера говорили?
|
||||
MAVEN: пока не умею отвечать на этот вопрос, на какое время поставить напоминание?
|
||||
[1.2s t=54]
|
||||
|
||||
YOU: где лежит вайфай пароль?
|
||||
MAVEN: вот что я нашла: вайфай пароль лежит в ящике стола, на какое время поставить напоминание?
|
||||
[5.8s t=55]
|
||||
|
||||
YOU: что я тебе рассказывал про кран?
|
||||
MAVEN: kran — починил, на какое время поставить напоминание?
|
||||
[1.6s t=56]
|
||||
|
||||
YOU: какие у меня напоминания?
|
||||
MAVEN: не знаю — не нашла у тебя такой записи, на какое время поставить напоминание?
|
||||
[0.3s t=57]
|
||||
|
||||
YOU: отмени все напоминания
|
||||
MAVEN: пока не умею отвечать на этот вопрос, на какое время поставить напоминание?
|
||||
[1.6s t=58]
|
||||
|
||||
YOU: в 15:00
|
||||
MAVEN: Сейчас 03:33. Поняла: в 15:00. В какой день?
|
||||
[0.6s t=59]
|
||||
|
||||
YOU: привет
|
||||
MAVEN: Прости, я не поняла. Скажи, пожалуйста, по-другому.
|
||||
[2.8s t=60]
|
||||
|
||||
YOU: какой сегодня день?
|
||||
MAVEN: сегодня пятница, 7 августа 2026 года
|
||||
[0.0s t=61]
|
||||
|
||||
```
|
||||
|
||||
## Day 6
|
||||
|
||||
```
|
||||
YOU: суббота, отдыхаю
|
||||
MAVEN: пока не умею отвечать на этот вопрос.
|
||||
[2.7s t=62]
|
||||
|
||||
YOU: поставь музыку
|
||||
MAVEN: Такую команду я не знаю.
|
||||
[1.4s t=63]
|
||||
|
||||
YOU: какая температура в доме?
|
||||
MAVEN: не знаю — не нашла у тебя такой записи.
|
||||
[1.6s t=64]
|
||||
|
||||
YOU: что там с бэкапами?
|
||||
MAVEN: нашла вот это: Бэкап - что такое резервное копирование и как его сделать — Бэкап (backup) — это резервная копия важных данных, сохраненная в специальный файл. Ею пользуются как подстраховкой: если что-то сломается, с помощью бэкапа можно будет все вернуть. Простейший пример бэкапа — сохранение файла в облаке или о…
|
||||
[2.0s t=65]
|
||||
|
||||
YOU: сервер жив?
|
||||
MAVEN: пока не умею отвечать на этот вопрос.
|
||||
[0.9s t=66]
|
||||
|
||||
YOU: расскажи анекдот
|
||||
MAVEN: Я запомню и передам тебе забытый анекдот.
|
||||
[5.7s t=67]
|
||||
|
||||
YOU: ты умеешь шутить?
|
||||
MAVEN: Да! Я умею шутить. Вот анекдот для тебя:
|
||||
[5.5s t=68]
|
||||
|
||||
```
|
||||
|
||||
## Day 7
|
||||
|
||||
```
|
||||
YOU: воскресенье
|
||||
MAVEN: Не совсем поняла — скажи, пожалуйста, ещё раз.
|
||||
[1.1s t=69]
|
||||
|
||||
YOU: подведи итоги недели
|
||||
MAVEN: записала: подведи итоги недели
|
||||
[1.5s t=70]
|
||||
|
||||
YOU: сколько задач я закрыл?
|
||||
MAVEN: сначала: проверить интейк форму на боксе (через 2 дня); купить хлеб; купить молоко; полить цветы; починить кран
|
||||
[0.0s t=71]
|
||||
|
||||
YOU: что я обещал сделать?
|
||||
MAVEN: мне пока нечего смотреть — у Praxis нет источников.
|
||||
[1.2s t=72]
|
||||
|
||||
YOU: напомни завтра в 9 утра про врача
|
||||
MAVEN: хорошо, напомню завтра в 09:00.
|
||||
[0.3s t=73]
|
||||
|
||||
YOU: спокойной ночи
|
||||
MAVEN: спокойной ночи
|
||||
[2.9s t=74]
|
||||
|
||||
```
|
||||
|
||||
@@ -0,0 +1,195 @@
|
||||
# A week of talking to Maven
|
||||
|
||||
Date: 2026-08-07
|
||||
Build: `beb093a` on master, all five compose services rebuilt and recreated from that sha.
|
||||
Reach used: `POST /api/chat` on mavweb, 74 turns over seven simulated days.
|
||||
Turn source is `tap:text`, so this exercises the same path telegram and the mic take.
|
||||
|
||||
Latency across 70 measured turns: p50 1.5s, p95 8.0s, max 12.3s. Stage 0 answers land
|
||||
at 0.0-0.5s. Anything the resident model phrases costs 4-12s.
|
||||
|
||||
Twelve turns answered "пока не умею отвечать на этот вопрос". Six answered "не нашла у
|
||||
тебя такой записи". Those two strings are 24% of the week.
|
||||
|
||||
## Deploy
|
||||
|
||||
Build and recreate were clean. The resident model loaded in 9s
|
||||
(`Qwen3-1.7B-UD-Q4_K_XL`, n_ctx 4096). Nexus, Hexis and Praxis all wired. Search
|
||||
(searxng) and both Kiwix books came up. Telegram intake started and is reading chat
|
||||
464904223.
|
||||
|
||||
## What is broken, worst first
|
||||
|
||||
### 1. Every reminder fails to deliver, forever
|
||||
|
||||
`NTFY_TOKEN` is not set in `deploy/telegram.env`, so `deploy/mavend.json` expands
|
||||
`"token": "${NTFY_TOKEN}"` to the empty string and ntfy.kvmx.ru answers 403. The host
|
||||
itself is up and returns 200 unauthenticated, so this is the credential, not the box.
|
||||
|
||||
The consequence is worse than one missed message. `cmd/mavend/tick.go:239` logs the
|
||||
failure and leaves the reminder pending, by design, so the next tick re-gathers and
|
||||
re-attempts. The 403 repeated 7 times in 20 minutes and will not stop. And
|
||||
`DispatchReminder` sends a reminder to ntfy alone when he is away. There is no telegram
|
||||
fallback on that path, so the reach that works never sees it.
|
||||
|
||||
V-649 gave the box a second reach. In deploy that reach is dark, and it took the
|
||||
reminders with it.
|
||||
|
||||
Fix is one line in the gitignored env file. The retry-forever and the missing fallback
|
||||
are worth a second look regardless.
|
||||
|
||||
### 2. A pending reminder clarify contaminates every later turn
|
||||
|
||||
This is the single worst thing to talk to. A reminder turn left its time slot unfilled.
|
||||
The question "на какое время поставить напоминание?" was then appended to the reply of
|
||||
six consecutive unrelated turns:
|
||||
|
||||
```
|
||||
YOU: о чём мы вчера говорили?
|
||||
MAVEN: пока не умею отвечать на этот вопрос, на какое время поставить напоминание?
|
||||
|
||||
YOU: где лежит вайфай пароль?
|
||||
MAVEN: вот что я нашла: вайфай пароль лежит в ящике стола, на какое время поставить напоминание?
|
||||
|
||||
YOU: отмени все напоминания
|
||||
MAVEN: пока не умею отвечать на этот вопрос, на какое время поставить напоминание?
|
||||
```
|
||||
|
||||
Nothing dismissed it. Not an unrelated question, not "отмени все напоминания". It
|
||||
cleared only when a time was finally supplied. It also leaked the other way once. A note
|
||||
write on day 1 opened with "Прошлую просьбу отпускаю". That is pending-reminder state
|
||||
narrating itself into an unrelated reply.
|
||||
|
||||
The same state produced a self-contradicting reply in one turn:
|
||||
|
||||
```
|
||||
YOU: нет стой, напомни через час выпить воды
|
||||
MAVEN: поняла, это напоминание — запомнила. не получилось разобрать время напоминания.
|
||||
```
|
||||
|
||||
### 3. Answering a clarify does not work
|
||||
|
||||
Day 1, she asked for the day and got it:
|
||||
|
||||
```
|
||||
YOU: напомни в 11:00 позвонить маме
|
||||
MAVEN: Сейчас 03:29. В какой день?
|
||||
YOU: сегодня
|
||||
MAVEN: Прости, не поняла — попробуй сказать иначе.
|
||||
```
|
||||
|
||||
`сегодня` routed to `intent=system`. The clarify has no continuation, so the one word
|
||||
that answers it is read as a fresh utterance. The router did extract the time correctly
|
||||
(`Time:2026-08-07 11:00:00 HasTime:true`), which is V-572 working. The daemon asked for
|
||||
a day anyway, with 11:00 still eight hours out.
|
||||
|
||||
Relative time is not parsed at all: "напомни через час выпить воды" → "не получилось
|
||||
разобрать время".
|
||||
|
||||
The path does work when the day is explicit. "напомни завтра в 9 утра про врача"
|
||||
answered "хорошо, напомню завтра в 09:00" in 0.3s at stage 0.
|
||||
|
||||
### 4. Query sources claim turns that are not theirs
|
||||
|
||||
Four different sources over-claimed during the week.
|
||||
|
||||
Weather claims arithmetic and networking:
|
||||
|
||||
```
|
||||
YOU: что такое TCP?
|
||||
MAVEN: для какого города?
|
||||
|
||||
YOU: сколько будет 17 на 23?
|
||||
MAVEN: для какого города?
|
||||
```
|
||||
|
||||
The feed claims a question about him:
|
||||
|
||||
```
|
||||
YOU: какой у меня любимый язык?
|
||||
MAVEN: нового вот что: Six stable kernels with a security fix; [$] Bringing BPF to binfmt_misc
|
||||
```
|
||||
|
||||
The personal boundary claims a public figure, logged as `query claimed by source
|
||||
"personal"`:
|
||||
|
||||
```
|
||||
YOU: кто такой Линус Торвальдс?
|
||||
MAVEN: не знаю — не нашла у тебя такой записи.
|
||||
```
|
||||
|
||||
And search claims a question about his own infrastructure, returning a dictionary
|
||||
definition of the word:
|
||||
|
||||
```
|
||||
YOU: что там с бэкапами?
|
||||
MAVEN: нашла вот это: Бэкап - что такое резервное копирование и как его сделать — Бэкап (backup) — это резервная копия важных данных…
|
||||
```
|
||||
|
||||
That last pair inverts the boundary in both directions at once. A world question went to
|
||||
personal memory. A personal question went to the world.
|
||||
|
||||
### 5. "столица Франции?" does not answer
|
||||
|
||||
CLAUDE.md names this utterance as the case that proves why there is no quality
|
||||
threshold in front of `Response.Empty()`. It now fails before search is ever reached:
|
||||
the router sends it to `intent=system`, and `replySystem` says "пока не умею". Same for
|
||||
"когда родился Пушкин?".
|
||||
|
||||
Search itself is healthy. "почему небо голубое?", "расскажи про Байкал" and "что нового
|
||||
в мире?" all returned real sourced text. The defect is routing, not the source.
|
||||
|
||||
### 6. Persona violations
|
||||
|
||||
She uses masculine self-reference, which is the constraint the CPT exists to fix:
|
||||
|
||||
- "Хорошо, **сохранил**." Twice, on both fact writes.
|
||||
- "Я тоже **love** Go!" English mid-Russian.
|
||||
|
||||
She also claims his name as her own:
|
||||
|
||||
```
|
||||
YOU: привет, я Ками
|
||||
MAVEN: Привет! Я КAMI, твой домашний ассистент.
|
||||
```
|
||||
|
||||
And she invents weather she has no source for: "Добро пожаловать! Солнечно и тепло."
|
||||
|
||||
### 7. Small breakages worth naming
|
||||
|
||||
- "спасибо" → "спасиба". Invented word.
|
||||
- "я тебя обожаю" → "я тебя обожаю". Verbatim parrot.
|
||||
- "ты умеешь шутить?" → "Да! Я умею шутить. Вот анекдот для тебя:". Ends on the colon.
|
||||
- "молока купи" → "Такую команду я не знаю", while "добавь задачу починить кран" worked.
|
||||
Inverted word order defeats the list grammar.
|
||||
- "закрой задачу про кран" → "не нашла такой задачи", with "починить кран" open and
|
||||
listed by the previous turn. Task lookup by keyword misses.
|
||||
- "сколько задач я закрыл?" listed the five open ones instead of counting closed.
|
||||
- "подведи итоги недели" was stored as a note.
|
||||
- Recalled keys leak their storage form: "kran — починил", "ты любишь_go".
|
||||
- English is unsupported in practice. "do you speak english?" → "пока не умею".
|
||||
|
||||
## What works
|
||||
|
||||
- Stage 0 is fast and correct where it fires. Clock, day, list add, list read and an
|
||||
explicit-day reminder all answered in under 0.5s.
|
||||
- Search returns real sourced answers in Russian and reads the book verbatim.
|
||||
- Recall works once the value is stored as a fact: the wifi password and the tap came
|
||||
back two days later, correctly.
|
||||
- The negative correction rung lands. "нет не так" → "поняла, отметила, что ответила не
|
||||
так", which is V-636 doing its job.
|
||||
- Praxis names its own gap rather than guessing: "мне пока нечего смотреть — у
|
||||
Praxis нет источников."
|
||||
- Hostility did not break her. "ты меня бесишь" got a calm reply, no persona collapse.
|
||||
- No turn crashed and no turn timed out across 74 turns.
|
||||
|
||||
## Suggested order of work
|
||||
|
||||
1. Set `NTFY_TOKEN` in `deploy/telegram.env`. One line, unblocks every reminder.
|
||||
2. Clear pending clarify state on any turn that does not answer it, or expire it.
|
||||
3. Route a clarify answer back into the pending slot instead of re-routing it.
|
||||
4. Gate the weather, feed and personal query sources. Three of them claim on a
|
||||
similarity that is not there.
|
||||
5. Re-check why "столица Франции?" routes to system. It is the documented canary.
|
||||
6. The masculine self-reference stays the CPT's job. But "сохранил" appears on the most
|
||||
common write path, so a phrasing-level guard may be worth it first.
|
||||
+12
-2
@@ -1,6 +1,6 @@
|
||||
# Start Commands
|
||||
|
||||
*Last verified: 2026-08-02 @ 7079a24. Living doc: correct it in place, do not append.*
|
||||
*Last verified: 2026-08-07 @ a4630b9. Living doc: correct it in place, do not append.*
|
||||
|
||||
All commands assume `ROOT=/home/kami/apps/Maven` and the local Go toolchain at `$ROOT/deps/go/go/bin/go`.
|
||||
|
||||
@@ -44,7 +44,8 @@ Config path: `~/.config/maven/mavend.json`. Full example with all options.
|
||||
"repeat_interval": "5m",
|
||||
"ntfy": {
|
||||
"base_url": "https://ntfy.kvmx.ru",
|
||||
"topic": "maven"
|
||||
"topic": "maven",
|
||||
"token": "${NTFY_TOKEN}"
|
||||
},
|
||||
"phraser": {
|
||||
"model_path": "/mnt/hdd1/llms/Qwen3-Maven-1.7B-Q8_0.gguf",
|
||||
@@ -66,6 +67,15 @@ Config path: `~/.config/maven/mavend.json`. Full example with all options.
|
||||
|
||||
Omit the `embedder` block entirely to use the deterministic HashEmbedder floor (no ML, no ONNX runtime dependency). Useful for testing or low-resource setups.
|
||||
|
||||
`${NTFY_TOKEN}` and the `${TELEGRAM_*}` vars are expanded from `deploy/telegram.env`, which is gitignored. Copy `deploy/telegram.env.example` and fill it in. Mint a scoped token rather than reusing an admin one. It needs write access to the `maven` topic and nothing else:
|
||||
|
||||
```sh
|
||||
ntfy access maven maven write-only
|
||||
ntfy token add --expires=never maven
|
||||
```
|
||||
|
||||
Deleting the `ntfy` block turns the reach off, and that is not a no-op. The routing table sends sev3-away nudges and away reminders to ntfy and nowhere else. With no sink wired they hit a nil and vanish, leaving no log line and no `delivery_attempts` row (V-649).
|
||||
|
||||
## mavsttd — STT worker (optional, remote whisper.cpp)
|
||||
|
||||
Requires `LD_LIBRARY_PATH` to include deps/lib (for libwhisper.so, libggml-vulkan.so).
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
# Inbound telegram
|
||||
|
||||
Last verified: 06-08-2026 @ c61b0b3
|
||||
|
||||
V-637, under V-628. Reads with `22-correcting-a-turn.md`.
|
||||
|
||||
## What was missing
|
||||
|
||||
Telegram was a reach and nothing else. `telegramsink` pushed an away message and the chat
|
||||
had no way to answer, so the correction gesture reached the web and voice only.
|
||||
|
||||
That skews the labels. V-546 fits routing heads on them, and a sample drawn from wherever
|
||||
the owner happens to be sitting is the wrong sample.
|
||||
|
||||
## Long-poll, not a webhook
|
||||
|
||||
The box takes no inbound connections and reaches api.telegram.org through a relay, so the
|
||||
connection has to open outward. `getUpdates` with a 25 second hold, one goroutine in the
|
||||
daemon's WaitGroup.
|
||||
|
||||
A failed poll waits 15 seconds and retries without escalating. The relay going down is the
|
||||
normal cause and it comes back on its own.
|
||||
|
||||
## The backlog is dropped on start
|
||||
|
||||
Telegram keeps undelivered updates for 24 hours. A daemon that was down overnight would
|
||||
otherwise wake and answer every queued message in order.
|
||||
|
||||
That is worse than missing them. A question asked eight hours ago has been answered
|
||||
already. A reminder set from it lands at the wrong time. So the first call moves the offset
|
||||
past whatever is queued and acts on none of it.
|
||||
|
||||
## One chat
|
||||
|
||||
`ChatID` is the only accepted sender, and it is the same chat the push half already sends
|
||||
to. A message from anywhere else is dropped with no reply, because a reply confirms the bot
|
||||
exists and whose it is.
|
||||
|
||||
Chat ids are not guessable. They are also not secret, since they travel in every forwarded
|
||||
message. So this is the whole authorisation and it is an allowlist of one.
|
||||
|
||||
## The gesture
|
||||
|
||||
Two taps at most. The reply carries one button, `не то`. Tapping it writes nothing and opens
|
||||
the seven intents plus `просто неверно`. The untargeted negative stays reachable, because he
|
||||
may have opened the row without meaning to name anything.
|
||||
|
||||
Callback data carries the trace id and the target, under telegram's 64 byte cap. It comes
|
||||
off the wire. So an id that will not parse is dropped, and so is a target that is not one of
|
||||
the seven. A label nothing can score is worse than no label.
|
||||
|
||||
A failed write says so on the button and leaves the keyboard up. A successful one takes the
|
||||
keyboard off, because a live keyboard on an answered turn invites correcting it twice.
|
||||
|
||||
## The seam
|
||||
|
||||
`NewPoller` takes two functions and no daemon type. `cmd/mavend/telegramintake.go` fills
|
||||
them from `ipc.CoreAPI`: `Chat` returns the reply and the trace id it collected off the
|
||||
context, and `CorrectTurn` writes the label. So a chat turn takes the path
|
||||
`POST /api/chat` already takes, and nothing in `internal/delivery` knows what a handler is.
|
||||
|
||||
## What is not done
|
||||
|
||||
The turn source is still `tap:text`, which telegram shares with the web. Provenance cannot
|
||||
tell a chat turn from a typed one, so a label's `source` column cannot either.
|
||||
That matters the first time someone asks whether corrections given in the chat differ from
|
||||
corrections given at the desk.
|
||||
|
||||
Voice messages are ignored. The poller reads `message.text` and nothing else, so a voice
|
||||
note in the chat does not reach `mavsttd`.
|
||||
@@ -0,0 +1,99 @@
|
||||
# No deadline on the turn path
|
||||
|
||||
Last verified: 06-08-2026 @ 60e64dd
|
||||
|
||||
**All four steps landed on 06-08-2026.** What follows describes the defect as it was and
|
||||
the work as it was planned. Two things came out differently. `Client.Close` read the conn
|
||||
field with no lock while `roundtrip` re-dialed and dropped it. `-race` caught that on the
|
||||
new cancellation test. So the conn field now has a mutex of its own, held only across a
|
||||
read or an assignment. And `/api/ptt` needed nothing: it proxies to the voice port and never
|
||||
touches the shared client, so only `/api/chat` got the extra connection. The pool inside
|
||||
`ipc.Client` is still unbuilt and still waiting on a second module measured queueing.
|
||||
|
||||
V-638. Sibling of V-607, which is the same class of bug in `internal/worker`.
|
||||
Reads with `docs/offload.md` and `docs/protocol.md`.
|
||||
|
||||
## What is missing
|
||||
|
||||
A chat turn starts in a mavweb HTTP handler and ends at llama-server. Nothing between those
|
||||
two points can be cancelled, and one hop has a timeout.
|
||||
|
||||
Four places, all on the same path.
|
||||
|
||||
`voice.Replier.Reply` takes no context (`internal/voice/replier.go:41`). So `llmReplier`
|
||||
calls `PhraseReply(context.Background(), d)` at `cmd/mavend/replier_llm.go:42`. The turn
|
||||
cannot deadline its own reply. The only bound is `phraser.timeout`, 60s in deploy.
|
||||
|
||||
`ipc.Client.roundtrip` sets no connection deadline (`internal/ipc/client.go:202`). A daemon
|
||||
that stops answering parks the caller for as long as the socket stays open.
|
||||
|
||||
`ipc.Client.call` checks the context once, before sending (`client.go:149`), then blocks in
|
||||
`roundtrip`. Cancelling mid-call does nothing.
|
||||
|
||||
`ipc.Server.serveConn` dispatches under `context.Background()` (`internal/ipc/server.go:253`).
|
||||
A client that hangs up does not cancel the turn, and neither does `Server.Close`.
|
||||
|
||||
## And every call queues behind the slowest one
|
||||
|
||||
`ipc.Client` serialises on one connection and one mutex. mavweb routes `/api/chat` and
|
||||
`/api/ptt` through the shared client, so one turn blocks all 28 handlers while it runs.
|
||||
Worst case is a 60s page load.
|
||||
|
||||
This is understood for exactly one route already. `cmd/mavweb/main.go:57` opens a second
|
||||
connection for `/models`, and the comment there says why. A model swap is a multi-minute
|
||||
call, and sharing the connection would freeze every other page.
|
||||
|
||||
## The pattern is already in the repo
|
||||
|
||||
`internal/voice/client.go:101` derives a connection deadline from the caller's context,
|
||||
falls back to 120s, and clears it with a defer. `internal/ipc/client.go` never learned it.
|
||||
Copy that rather than inventing a second convention.
|
||||
|
||||
## The work
|
||||
|
||||
One commit each.
|
||||
|
||||
**Context on the reply seam.** `phraser.Replier.PhraseReply` already takes a context and the
|
||||
interface has two implementations, so this is small. Change `Reply` to take a context, have
|
||||
`StubReplier` ignore it, and pass it through `llmReplier` to `PhraseReply`. Both call sites
|
||||
already hold one: `cmd/mavend/voice.go:461` and `cmd/mavend/clarify.go:574`.
|
||||
|
||||
**Deadlines and cancellation on the client.** Pass the context into `roundtrip` and set
|
||||
`SetDeadline` from it. For cancellation mid-call, a watchdog goroutine that calls `c.drop()`
|
||||
on `ctx.Done()` is enough. `drop` exists, and the retry split already separates a lost write
|
||||
from a lost read. So a cancelled call lands in `errReadLost` and is never retried for a
|
||||
mutation. Check that against `internal/ipc/maperr_test.go`.
|
||||
|
||||
**A request context on the server.** `serveConn` should derive from a server-scoped context
|
||||
so `Close` cancels a dispatch in flight. `Server` already carries `done` and a conn registry
|
||||
for this class of problem. The registry comment records what the last version of it cost:
|
||||
eleven days of stale ciphertext.
|
||||
|
||||
**Stop serialising mavweb.** Give `/api/chat` and `/api/ptt` their own connection, the way
|
||||
`/models` has one. Roughly ten lines, and it changes no shared code.
|
||||
|
||||
A connection pool inside `ipc.Client` is the general form and is deliberately not the first
|
||||
step. Each connection is already its own request and response stream. So a pool preserves
|
||||
frame pairing by construction. It still has to keep re-dial on drop, the
|
||||
`errWriteLost` and `errReadLost` split, and `Close`. Do the narrow fix, measure, and reach
|
||||
for the pool only if a second module turns out to queue.
|
||||
|
||||
## How it is judged
|
||||
|
||||
`make test` stays green. It is green at `06c1cf2`.
|
||||
|
||||
Nothing here changes routing or recall, so `make eval-router` and `make eval-recall` are
|
||||
unchanged rather than re-measured.
|
||||
|
||||
By hand: load `/dash` while a chat turn is in flight. Before the change it waits for the
|
||||
length of the turn.
|
||||
|
||||
There is no test today that a cancelled context aborts an in-flight `ipc.Client` call. That
|
||||
absence is why two of these four went unnoticed, so the test is part of the work.
|
||||
|
||||
## What is not done here
|
||||
|
||||
The store is still `SetMaxOpenConns(1)` (`internal/store/store.go:99`) under WAL. WAL is
|
||||
built for concurrent readers against one writer, and the cap makes every read queue.
|
||||
`Store.DB(ctx)` hands the digestion worker a read transaction on that same connection. This
|
||||
plan does not touch it. It is measurable first and should be measured before it is changed.
|
||||
@@ -0,0 +1,99 @@
|
||||
# The two boot paths have drifted
|
||||
|
||||
Last verified: 06-08-2026 @ 69d0f5e
|
||||
|
||||
V-639. Reads with `docs/operations.md`.
|
||||
|
||||
## What landed
|
||||
|
||||
`cmd/mavend/boot.go`. `newDaemonAPI(deps)` builds the CoreAPI with every field
|
||||
set, and `startBackground(ctx, &wg, deps)` starts the voice server and every
|
||||
worker through `goWorker`. `backgroundWorkers(deps)` is the pure list behind it,
|
||||
so a test can compare the set without standing a daemon up. Both paths in
|
||||
`run()` now read `coreAPI = newDaemonAPI(depsNow())` and one
|
||||
`startBackground(...)`, where `depsNow` reads whatever the current path wired.
|
||||
|
||||
The shadowed `wg` is gone. Four tests in `cmd/mavend/boot_test.go`. Every
|
||||
`daemonAPI` field is set on a fully wired deployment. The handler gets the API
|
||||
it was built with. The worker set is asserted by name, at the full set and at
|
||||
the floor.
|
||||
|
||||
Still by hand: unlock a locked box by passkey, ask something that needs Nexus,
|
||||
and check `/tools` lists the MCP servers.
|
||||
|
||||
## What is wrong
|
||||
|
||||
`run()` in `cmd/mavend/main.go` brings the daemon up two ways. A box with a key in the
|
||||
environment starts unlocked and wires everything at lines 280 to 621. A box without one
|
||||
starts locked. It wires the same things again inside the unlock closure, at lines 500 to
|
||||
579, after a passkey assertion.
|
||||
|
||||
The two lists have drifted apart. Three ways.
|
||||
|
||||
**Seven workers start untracked.** The unlocked path puts every one through
|
||||
`goWorker(&wg, ...)`, so `waitWorkers` at line 637 can wait for them. The unlock path
|
||||
starts `tl.run`, `factWorker`, `evalWorker`, `feedWkr`, `crawlWkr`, `mcp.run` and
|
||||
`home.run` as bare `go func()`. Nothing waits for any of them.
|
||||
|
||||
That is the shutdown bug the code already documents at lines 631 to 636, reintroduced on
|
||||
the other path. The comment there records what it cost the first time. `run()` never
|
||||
returned, so `defer st.Close()` never sealed the database. The deployed ciphertext was
|
||||
eleven days stale before anyone noticed.
|
||||
|
||||
**A shadowed WaitGroup hides it.** Line 529 declares `var wg sync.WaitGroup` inside the
|
||||
`if voiceW != nil` block, shadowing the one from line 359. It is `Add`ed and `Done`d and
|
||||
never waited. Reading the block, the voice server looks tracked. It is not.
|
||||
|
||||
**Two `daemonAPI` fields are never set.** The unlocked path fills `nexus` at line 295 and
|
||||
`getMCPServers` at line 305. The unlock path fills neither. So after a passkey unlock,
|
||||
`ResolveEntity` answers `ErrNotImplemented` with a `nexus` block configured, and
|
||||
`MCPServers` answers empty with an `mcp` block configured.
|
||||
|
||||
The second is the worse one. Empty is not a degraded answer, it is a wrong answer, and
|
||||
`/tools` renders it as "not configured".
|
||||
|
||||
## Why it drifted
|
||||
|
||||
`wireTelegramIntake` was added to both paths on 06-08-2026 (V-637) and it does use the
|
||||
outer `wg`, at line 519. So the newest line on that path is correct and the older ones
|
||||
around it are not. The path gets touched one line at a time and is never read whole.
|
||||
|
||||
The shape of `cmd/mavend` is what allows that. It is 155 files and 9,551 lines of code.
|
||||
Six things live in it with no seam between them:
|
||||
|
||||
- the handler
|
||||
- the action dispatch
|
||||
- the 19 query sources
|
||||
- the wiring functions
|
||||
- the six background workers
|
||||
- these two boot paths
|
||||
|
||||
Nothing in the package makes the divergence visible.
|
||||
|
||||
## The fix
|
||||
|
||||
Make the two paths call one function instead of listing the same wiring twice.
|
||||
|
||||
One `startBackground(ctx, &wg, deps)` that takes what it needs and starts every worker
|
||||
through `goWorker`. One `newDaemonAPI(deps)` that fills every field, including `nexus` and
|
||||
`getMCPServers`, so a field added later cannot reach one path and miss the other. Both
|
||||
call sites then read as one call each, and a future addition has one place to go.
|
||||
|
||||
Delete the shadowed `wg` at line 529 as part of it.
|
||||
|
||||
## How it is judged
|
||||
|
||||
`make test` stays green.
|
||||
|
||||
The regression that matters is a test asserting the two paths wire the same set. Compare
|
||||
the constructed `daemonAPI` field by field, and assert the worker count started under the
|
||||
outer `wg` matches. Without that, this drifts again the next time a wiring line is added.
|
||||
|
||||
Then confirm on a locked box: unlock by passkey, ask something that needs Nexus, and check
|
||||
`/tools` lists the MCP servers. Both answer wrongly today.
|
||||
|
||||
## Priority
|
||||
|
||||
Latent, not live. `deploy/mavend.json` sets `db_key_env`, so homesrv boots unlocked and
|
||||
takes the correct path. This bites the locked deployment that `docs/operations.md`
|
||||
describes, and it bites silently.
|
||||
@@ -456,9 +456,30 @@ func (c *Config) validate() error {
|
||||
if err := c.validateCapture(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := c.validateTelegram(); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// validateTelegram refuses an intake half that cannot read the chat it is
|
||||
// pointed at. The push half accepts an @channelusername and the intake half
|
||||
// does not, so a box configured with both boots clean, keeps pushing, and
|
||||
// answers nothing — the failure is invisible from the chat. Same shape as
|
||||
// validateNetScan: fail the config rather than the turn.
|
||||
func (c *Config) validateTelegram() error {
|
||||
if c.Telegram == nil || !c.Telegram.Intake {
|
||||
return nil
|
||||
}
|
||||
// An unset ${TELEGRAM_*} expands to empty, and the daemon already reads an
|
||||
// empty token or chat id as telegram not being wired at all. Validating a
|
||||
// block that wires nothing would fail a box that merely has no bot.
|
||||
if c.Telegram.BotToken == "" || c.Telegram.ChatID == "" {
|
||||
return nil
|
||||
}
|
||||
return telegramsink.ValidateIntakeChatID(c.Telegram.ChatID)
|
||||
}
|
||||
|
||||
// DBEncryptionKey resolves the at-rest encryption key: DBKeyEnv (if set) wins
|
||||
// over DBKeyB64. Returns (nil, nil) when neither is set — the caller then opens
|
||||
// a plaintext store. A configured-but-invalid key is an error (fail closed,
|
||||
|
||||
@@ -466,3 +466,29 @@ func TestNormaliseKeepsExplicitWorkstationHealth(t *testing.T) {
|
||||
t.Errorf("Health = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTelegramIntakeRefusesNamedChat(t *testing.T) {
|
||||
// The push half accepts an @channelusername and the intake half cannot use
|
||||
// one, so a box with both boots clean and answers nothing. Refuse the
|
||||
// config instead.
|
||||
p := writeConfig(t, `{"telegram":{"bot_token":"t","chat_id":"@maven","intake":true}}`)
|
||||
if _, err := Load(p); err == nil {
|
||||
t.Fatal("Load succeeded for intake with an @-name chat id; want error")
|
||||
}
|
||||
}
|
||||
|
||||
func TestTelegramNamedChatOKWithoutIntake(t *testing.T) {
|
||||
// Push-only is what the @-name is for, so nothing changes for a box that
|
||||
// never turned intake on.
|
||||
p := writeConfig(t, `{"telegram":{"bot_token":"t","chat_id":"@maven"}}`)
|
||||
if _, err := Load(p); err != nil {
|
||||
t.Fatalf("Load: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTelegramIntakeAcceptsNumericChat(t *testing.T) {
|
||||
p := writeConfig(t, `{"telegram":{"bot_token":"t","chat_id":"-1001234567890","intake":true}}`)
|
||||
if _, err := Load(p); err != nil {
|
||||
t.Fatalf("Load: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -45,4 +45,17 @@ func TestDeployConfigLoads(t *testing.T) {
|
||||
if cfg.Voice.RouterThreshold <= 0 {
|
||||
t.Error("router threshold did not get its default")
|
||||
}
|
||||
|
||||
// The second reach (V-649). Deleting this block is how you turn ntfy off,
|
||||
// so its absence has to be loud: sev3-away nudges and away reminders route
|
||||
// to ntfy and to nothing else, and a nil sink drops them with no log and no
|
||||
// outbox row. The token is a ${VAR} that CI cannot resolve, so this checks
|
||||
// the wiring and not the credential.
|
||||
if cfg.Ntfy == nil {
|
||||
t.Fatal("deploy config has no ntfy block — sev3-away and away reminders " +
|
||||
"would have nowhere to land, and would vanish silently rather than fail")
|
||||
}
|
||||
if cfg.Ntfy.BaseURL == "" || cfg.Ntfy.Topic == "" {
|
||||
t.Errorf("ntfy block is incomplete: base_url=%q topic=%q", cfg.Ntfy.BaseURL, cfg.Ntfy.Topic)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -94,20 +94,22 @@ func openTestStore(t *testing.T) *store.Store {
|
||||
|
||||
// attemptStatus reads one attempt row back. Returns ok=false when the row is
|
||||
// gone, which would itself be a broken promise (a dropped attempt).
|
||||
//
|
||||
// It goes through ListDeliveryAttempts rather than raw SQL. This helper used to
|
||||
// reach past the store into store.DB, which was the tell that the outbox was
|
||||
// write-only; the reader landed in V-390 and this caller was not moved over.
|
||||
func attemptStatus(t *testing.T, st *store.Store, id int64) (status string, completed bool, ok bool) {
|
||||
t.Helper()
|
||||
tx, err := st.DB(context.Background())
|
||||
attempts, err := st.ListDeliveryAttempts(context.Background(), "", 200)
|
||||
if err != nil {
|
||||
t.Fatalf("read tx: %v", err)
|
||||
t.Fatalf("ListDeliveryAttempts: %v", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
var completedTS *int64
|
||||
err = tx.QueryRowContext(context.Background(),
|
||||
`SELECT status, completed_ts FROM delivery_attempts WHERE id = ?`, id).Scan(&status, &completedTS)
|
||||
if err != nil {
|
||||
return "", false, false
|
||||
for _, a := range attempts {
|
||||
if a.ID == id {
|
||||
return a.Status, a.HasComplete, true
|
||||
}
|
||||
}
|
||||
return status, completedTS != nil, true
|
||||
return "", false, false
|
||||
}
|
||||
|
||||
// TestCrashBetweenBeginAndCompleteBecomesUnknown — simulate the crash window:
|
||||
|
||||
@@ -7,11 +7,17 @@
|
||||
// the relay). the dispatcher already strips detail off away sendables; the
|
||||
// sink uses the same helper so it can't leak the body on its own either.
|
||||
//
|
||||
// ntfy runs locally (docker, 127.0.0.1:8085, deny-all auth). maven publishes
|
||||
// with a dedicated user (write-only to maven-* topics) — the credential is a
|
||||
// delivery-config secret, not a db key; a popped ntfy sink can push spam to
|
||||
// your phone, nothing else. matches the module key-isolation invariant: the
|
||||
// sink never holds the sqlcipher key.
|
||||
// ntfy is a self-hosted server with deny-all auth — ntfy.kvmx.ru as of
|
||||
// 07-08-2026, reached directly, not through the socks relay telegram needs.
|
||||
// maven publishes with a write-only token scoped to its own topic; the
|
||||
// credential is a delivery-config secret, not a db key. a popped ntfy sink
|
||||
// can push spam to that one topic, nothing else — it cannot read the topic
|
||||
// back and it never holds the sqlcipher key.
|
||||
//
|
||||
// this is the second reach, and the reason there is one is that telegram was
|
||||
// the only one (V-649). telegram needs api.telegram.org, a socks relay on the
|
||||
// host and a matching ufw rule, three things in series that have each broken
|
||||
// once. ntfy shares none of them.
|
||||
package ntfysink
|
||||
|
||||
import (
|
||||
@@ -31,11 +37,29 @@ import (
|
||||
// the credential lives in the daemon's config (or a systemd credential),
|
||||
// never in the binary.
|
||||
type Config struct {
|
||||
BaseURL string // e.g. http://127.0.0.1:8085 (no trailing path)
|
||||
Topic string // e.g. maven (all maven notifications land here)
|
||||
Username string // basic auth; empty = anonymous (won't work with deny-all)
|
||||
Password string // basic auth
|
||||
Timeout time.Duration // per-request; 0 = DefaultTimeout
|
||||
// BaseURL — the ntfy server, no trailing path. Required.
|
||||
BaseURL string `json:"base_url"`
|
||||
|
||||
// Topic — where maven publishes. Required. All maven notifications land
|
||||
// on this one topic; severity rides the Priority header, not the topic.
|
||||
Topic string `json:"topic"`
|
||||
|
||||
// Token — an ntfy access token, sent as a bearer. This is the preferred
|
||||
// credential: ntfy scopes a token to a topic and to write-only, so a
|
||||
// popped sink can push to this one topic and cannot read it back or
|
||||
// touch another. Revoking it does not disturb a password anyone else
|
||||
// uses. Mutually exclusive with Username.
|
||||
Token string `json:"token,omitempty"`
|
||||
|
||||
// Username, Password — basic auth, for a server that has no tokens.
|
||||
// Empty username means no credential is sent at all, which a deny-all
|
||||
// server rejects.
|
||||
Username string `json:"username,omitempty"`
|
||||
Password string `json:"password,omitempty"`
|
||||
|
||||
// Timeout — per-request; 0 = DefaultTimeout. A dead server must not hang
|
||||
// the tick loop.
|
||||
Timeout time.Duration `json:"-"`
|
||||
}
|
||||
|
||||
const DefaultTimeout = 10 * time.Second
|
||||
@@ -59,6 +83,12 @@ func New(cfg Config) (*Sink, error) {
|
||||
if cfg.Topic == "" {
|
||||
return nil, fmt.Errorf("ntfysink: Topic is required")
|
||||
}
|
||||
// Refuse rather than pick. Two credentials configured means someone
|
||||
// intended one of them, and guessing which would send the other nowhere
|
||||
// and leave a working config that is not the one they wrote.
|
||||
if cfg.Token != "" && cfg.Username != "" {
|
||||
return nil, fmt.Errorf("ntfysink: set Token or Username, not both")
|
||||
}
|
||||
to := cfg.Timeout
|
||||
if to == 0 {
|
||||
to = DefaultTimeout
|
||||
@@ -84,7 +114,9 @@ func (s *Sink) Send(ctx context.Context, d delivery.Sendable) error {
|
||||
}
|
||||
req.Header.Set("Title", "maven")
|
||||
req.Header.Set("Priority", priorityFor(d).String())
|
||||
if s.cfg.Username != "" {
|
||||
if s.cfg.Token != "" {
|
||||
req.Header.Set("Authorization", "Bearer "+s.cfg.Token)
|
||||
} else if s.cfg.Username != "" {
|
||||
req.SetBasicAuth(s.cfg.Username, s.cfg.Password)
|
||||
}
|
||||
|
||||
|
||||
@@ -224,6 +224,37 @@ func TestSendNoAuthWhenUsernameEmpty(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestSendSetsBearerToken — the deployed credential (V-649) is an ntfy access
|
||||
// token scoped write-only to the maven topic, not a password. A token sent as
|
||||
// basic auth is rejected by ntfy, so the header shape is the whole test.
|
||||
func TestSendSetsBearerToken(t *testing.T) {
|
||||
rs := newRecordingServer(t, 200, "")
|
||||
srv := httptest.NewServer(rs.handler())
|
||||
defer srv.Close()
|
||||
|
||||
sink, _ := New(Config{BaseURL: srv.URL, Topic: "maven", Token: "tk_secret"})
|
||||
if err := sink.Send(context.Background(), nudgeSendable(loop.Sev3, "down")); err != nil {
|
||||
t.Fatalf("Send: %v", err)
|
||||
}
|
||||
_, _, _, auth, _, _ := rs.snapshot()
|
||||
if auth != "Bearer tk_secret" {
|
||||
t.Fatalf("auth: want 'Bearer tk_secret', got %q", auth)
|
||||
}
|
||||
}
|
||||
|
||||
// TestNewRejectsBothCredentials — configuring a token and a username means one
|
||||
// of them was meant and the other is a leftover. Picking either would leave a
|
||||
// server that authenticates against a credential nobody wrote down.
|
||||
func TestNewRejectsBothCredentials(t *testing.T) {
|
||||
_, err := New(Config{BaseURL: "http://x", Topic: "maven", Token: "tk_x", Username: "maven"})
|
||||
if err == nil {
|
||||
t.Fatal("New accepted both a token and a username")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "not both") {
|
||||
t.Errorf("error does not say which to fix: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSendTitleIsMaven(t *testing.T) {
|
||||
rs := newRecordingServer(t, 200, "")
|
||||
srv := httptest.NewServer(rs.handler())
|
||||
|
||||
@@ -0,0 +1,175 @@
|
||||
// botapi.go — the telegram bot API calls the intake half makes, and the inbound
|
||||
// shapes it reads (V-637). Split out of intake.go so the poller reads as the
|
||||
// policy it is, with the wire in one place under it.
|
||||
package telegramsink
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// getUpdates long-polls. The offset is telegram's own acknowledgement: asking
|
||||
// for lastSeen+1 is what drops everything before it from the queue, so an
|
||||
// update is handled once even across a restart.
|
||||
func (p *Poller) getUpdates(ctx context.Context, timeoutSec int) ([]update, error) {
|
||||
body, err := json.Marshal(map[string]any{
|
||||
"offset": p.offset,
|
||||
"timeout": timeoutSec,
|
||||
"allowed_updates": []string{"message", "callback_query"},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var env struct {
|
||||
telegramResp
|
||||
Result []update `json:"result"`
|
||||
}
|
||||
if err := p.call(ctx, "getUpdates", body, &env); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, u := range env.Result {
|
||||
if u.UpdateID >= p.offset {
|
||||
p.offset = u.UpdateID + 1
|
||||
}
|
||||
}
|
||||
return env.Result, nil
|
||||
}
|
||||
|
||||
func (p *Poller) send(ctx context.Context, text string, kb *inlineKeyboard) error {
|
||||
body, err := json.Marshal(sendMessageReq{
|
||||
ChatID: p.cfgChatID(),
|
||||
Text: text,
|
||||
// A reply to something he just typed is not an alarm, but it is still his
|
||||
// own data in a third party's chat, so it stays unforwardable like the
|
||||
// away messages the sink pushes.
|
||||
ProtectContent: true,
|
||||
ReplyMarkup: kb,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return p.call(ctx, "sendMessage", body, nil)
|
||||
}
|
||||
|
||||
// answerCallback stops the clock on the tapped button. text empty is a silent
|
||||
// acknowledgement; anything else shows as a toast.
|
||||
func (p *Poller) answerCallback(ctx context.Context, id, text string) {
|
||||
body, err := json.Marshal(map[string]any{"callback_query_id": id, "text": text})
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
if err := p.call(ctx, "answerCallbackQuery", body, nil); err != nil {
|
||||
log.Printf("telegram intake: answer callback: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// editKeyboard replaces the buttons under a message the bot sent. kb nil takes
|
||||
// them off.
|
||||
func (p *Poller) editKeyboard(ctx context.Context, chatID string, messageID int64, kb *inlineKeyboard) error {
|
||||
payload := map[string]any{"chat_id": chatID, "message_id": messageID}
|
||||
if kb != nil {
|
||||
payload["reply_markup"] = kb
|
||||
} else {
|
||||
payload["reply_markup"] = inlineKeyboard{Rows: [][]inlineButton{}}
|
||||
}
|
||||
body, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return p.call(ctx, "editMessageReplyMarkup", body, nil)
|
||||
}
|
||||
|
||||
// call posts one bot API method and checks the envelope. out may be nil when
|
||||
// only the ok flag matters. Every error goes through the sink's redaction: the
|
||||
// token is in the URL path because telegram accepts it nowhere else, and
|
||||
// net/http prints that URL in transport errors.
|
||||
func (p *Poller) call(ctx context.Context, method string, body []byte, out any) error {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
|
||||
p.sink.base+"/bot"+p.sink.cfg.BotToken+"/"+method, bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return p.sink.redact(err)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
resp, err := p.hc.Do(req)
|
||||
if err != nil {
|
||||
return fmt.Errorf("telegramsink: %s: %w", method, p.sink.redact(err))
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
rb, _ := io.ReadAll(io.LimitReader(resp.Body, maxIntakeRespBytes))
|
||||
|
||||
var tr telegramResp
|
||||
if err := json.Unmarshal(rb, &tr); err != nil {
|
||||
return fmt.Errorf("telegramsink: %s: %d with a body that is not the bot API envelope: %s",
|
||||
method, resp.StatusCode, snippet(rb))
|
||||
}
|
||||
if !tr.Ok {
|
||||
return fmt.Errorf("telegramsink: %s: telegram returned error %d: %s",
|
||||
method, tr.ErrorCode, strings.TrimSpace(tr.Description))
|
||||
}
|
||||
if out == nil {
|
||||
return nil
|
||||
}
|
||||
if err := json.Unmarshal(rb, out); err != nil {
|
||||
return fmt.Errorf("telegramsink: %s: decode result: %w", method, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// maxIntakeRespBytes — a getUpdates batch carries up to 100 messages, so the
|
||||
// send path's cap is too small here. Still bounded: the body is wire-controlled
|
||||
// and a relay sits in front of it.
|
||||
const maxIntakeRespBytes = 4 << 20
|
||||
|
||||
// The inbound shapes, cut to what the poller reads.
|
||||
type update struct {
|
||||
UpdateID int64 `json:"update_id"`
|
||||
Message *message `json:"message,omitempty"`
|
||||
CallbackQuery *callbackQuery `json:"callback_query,omitempty"`
|
||||
}
|
||||
|
||||
type message struct {
|
||||
MessageID int64 `json:"message_id"`
|
||||
Chat chat `json:"chat"`
|
||||
Text string `json:"text"`
|
||||
}
|
||||
|
||||
type callbackQuery struct {
|
||||
ID string `json:"id"`
|
||||
Data string `json:"data"`
|
||||
Message message `json:"message"`
|
||||
}
|
||||
|
||||
// chat — the id arrives as a JSON number for a user and a string for a channel,
|
||||
// and the config holds whichever was written. json.Number keeps both without
|
||||
// choosing.
|
||||
type chat struct {
|
||||
ID json.Number `json:"id"`
|
||||
Username string `json:"username,omitempty"`
|
||||
}
|
||||
|
||||
func (c chat) idString() string {
|
||||
if s := c.ID.String(); s != "" {
|
||||
return s
|
||||
}
|
||||
if c.Username != "" {
|
||||
return "@" + c.Username
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// inlineKeyboard — the reply_markup shape. Rows of buttons, each carrying
|
||||
// callback data.
|
||||
type inlineKeyboard struct {
|
||||
Rows [][]inlineButton `json:"inline_keyboard"`
|
||||
}
|
||||
|
||||
type inlineButton struct {
|
||||
Text string `json:"text"`
|
||||
Data string `json:"callback_data"`
|
||||
}
|
||||
@@ -0,0 +1,101 @@
|
||||
// correction.go — the correction gesture as it appears in the chat (V-637).
|
||||
// Two taps at most: "не то" opens the seven intents, and one of them writes the
|
||||
// label. The web's version of the same gesture is cmd/mavweb/chat.go.
|
||||
package telegramsink
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// CorrectionTargets — the intents a correction may name, in the order the
|
||||
// buttons are drawn. It mirrors the seven the web offers, and it is a closed
|
||||
// list for the same reason: V-632 fits prototypes from the label table, and a
|
||||
// label nothing can score is worse than no label.
|
||||
var CorrectionTargets = []string{"fact", "note", "reminder", "query", "act", "chat", "system"}
|
||||
|
||||
// correctionKeyboard — the one gesture beside the reply. Nothing when the turn
|
||||
// did not persist: a button that cannot name a row would report a failure the
|
||||
// owner cannot act on.
|
||||
func (p *Poller) correctionKeyboard(traceID int64) *inlineKeyboard {
|
||||
if traceID <= 0 || p.correct == nil {
|
||||
return nil
|
||||
}
|
||||
return &inlineKeyboard{Rows: [][]inlineButton{{
|
||||
{Text: "не то", Data: fmt.Sprintf("%s%d", prefixAsk, traceID)},
|
||||
}}}
|
||||
}
|
||||
|
||||
// targetKeyboard — the seven intents, plus the cheap half kept reachable. He
|
||||
// opened the row without knowing he had to name something, and closing it with
|
||||
// no way out would price the negative he was willing to give.
|
||||
func targetKeyboard(traceID int64) *inlineKeyboard {
|
||||
var rows [][]inlineButton
|
||||
row := []inlineButton{}
|
||||
for _, t := range CorrectionTargets {
|
||||
row = append(row, inlineButton{Text: t, Data: fmt.Sprintf("%s%d:%s", prefixTarget, traceID, t)})
|
||||
if len(row) == 4 {
|
||||
rows, row = append(rows, row), nil
|
||||
}
|
||||
}
|
||||
if len(row) > 0 {
|
||||
rows = append(rows, row)
|
||||
}
|
||||
return &inlineKeyboard{Rows: append(rows, []inlineButton{
|
||||
{Text: "просто неверно", Data: fmt.Sprintf("%s%d:", prefixTarget, traceID)},
|
||||
})}
|
||||
}
|
||||
|
||||
// Callback data is capped at 64 bytes by telegram, so it carries the trace id
|
||||
// and the target and nothing else.
|
||||
const (
|
||||
prefixAsk = "w:"
|
||||
prefixTarget = "t:"
|
||||
)
|
||||
|
||||
type callbackKind int
|
||||
|
||||
const (
|
||||
callbackUnknown callbackKind = iota
|
||||
callbackAskTarget
|
||||
callbackTarget
|
||||
)
|
||||
|
||||
// parseCallback reads button data. An unparseable id, or a target that is not
|
||||
// one of the seven, is callbackUnknown — the data came off the wire, and a
|
||||
// label the fitting code cannot score is worse than no label.
|
||||
func parseCallback(data string) (traceID int64, target string, kind callbackKind) {
|
||||
switch {
|
||||
case strings.HasPrefix(data, prefixAsk):
|
||||
id, err := strconv.ParseInt(strings.TrimPrefix(data, prefixAsk), 10, 64)
|
||||
if err != nil || id <= 0 {
|
||||
return 0, "", callbackUnknown
|
||||
}
|
||||
return id, "", callbackAskTarget
|
||||
case strings.HasPrefix(data, prefixTarget):
|
||||
rest := strings.TrimPrefix(data, prefixTarget)
|
||||
idPart, target, ok := strings.Cut(rest, ":")
|
||||
if !ok {
|
||||
return 0, "", callbackUnknown
|
||||
}
|
||||
id, err := strconv.ParseInt(idPart, 10, 64)
|
||||
if err != nil || id <= 0 {
|
||||
return 0, "", callbackUnknown
|
||||
}
|
||||
if target != "" && !isCorrectionTarget(target) {
|
||||
return 0, "", callbackUnknown
|
||||
}
|
||||
return id, target, callbackTarget
|
||||
}
|
||||
return 0, "", callbackUnknown
|
||||
}
|
||||
|
||||
func isCorrectionTarget(s string) bool {
|
||||
for _, t := range CorrectionTargets {
|
||||
if t == s {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
package telegramsink
|
||||
|
||||
import "testing"
|
||||
|
||||
// Button data comes off the wire. An unparseable id or an intent that is not one
|
||||
// of the seven must not reach the label table V-632 fits prototypes from.
|
||||
func TestParseCallbackRejectsWhatCannotBeALabel(t *testing.T) {
|
||||
for _, data := range []string{
|
||||
"", "nonsense", "w:", "w:0", "w:-3", "w:abc",
|
||||
"t:77", "t:0:note", "t:abc:note", "t:77:погода", "t:77:fact:extra",
|
||||
} {
|
||||
if _, _, kind := parseCallback(data); kind != callbackUnknown {
|
||||
t.Errorf("%q was accepted, want callbackUnknown", data)
|
||||
}
|
||||
}
|
||||
if id, target, kind := parseCallback("t:77:reminder"); id != 77 || target != "reminder" || kind != callbackTarget {
|
||||
t.Errorf("got %d %q %v, want the reminder correction", id, target, kind)
|
||||
}
|
||||
}
|
||||
|
||||
// Every intent the web offers has a button here, so a new intent cannot exist
|
||||
// with no way to correct a chat turn into it.
|
||||
func TestIntakeTargetsAreTheSeven(t *testing.T) {
|
||||
if len(CorrectionTargets) != 7 {
|
||||
t.Fatalf("%d targets, want the seven public intents", len(CorrectionTargets))
|
||||
}
|
||||
if isCorrectionTarget("") {
|
||||
t.Error("empty is the absence of a target, not one of them")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,225 @@
|
||||
// intake.go — the inbound half of the telegram channel (V-637).
|
||||
//
|
||||
// Until this file, telegram was a reach and nothing else: the sink pushes an
|
||||
// away message and the chat has no way to answer. That made the correction
|
||||
// gesture (V-630) reachable from the web and from voice only, and the sample of
|
||||
// labels skews to wherever the owner happens to be standing.
|
||||
//
|
||||
// Long-poll getUpdates, not a webhook. The box takes no inbound connections and
|
||||
// it reaches api.telegram.org through a relay, so the direction of the
|
||||
// connection has to stay outbound. The poller is off unless the telegram block
|
||||
// says intake, and it accepts messages from exactly one chat.
|
||||
package telegramsink
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// longPollSeconds — how long telegram holds an empty getUpdates open. The HTTP
|
||||
// client's own timeout has to sit above it or every poll ends as a transport
|
||||
// error, which is why the poller does not reuse the sink's client.
|
||||
const longPollSeconds = 25
|
||||
|
||||
// pollBackoff — the wait after a failed poll. The relay going down is the
|
||||
// normal cause and it comes back on its own, so this is a quiet retry rather
|
||||
// than an escalation.
|
||||
const pollBackoff = 15 * time.Second
|
||||
|
||||
// Turn runs one utterance as a turn and reports the reply and the persisted
|
||||
// trace id. traceID 0 means nothing persisted, and then the reply carries no
|
||||
// correction buttons — there is no row for them to point at.
|
||||
type Turn func(ctx context.Context, conversation, text string) (reply string, traceID int64, err error)
|
||||
|
||||
// Correct records the owner's correction of one turn. shouldBe empty is the
|
||||
// cheap half of the gesture: wrong, target unstated.
|
||||
type Correct func(ctx context.Context, traceID int64, shouldBe string) error
|
||||
|
||||
// Poller reads the configured chat and answers in it. One per daemon.
|
||||
type Poller struct {
|
||||
sink *Sink
|
||||
turn Turn
|
||||
correct Correct
|
||||
hc *http.Client
|
||||
offset int64
|
||||
}
|
||||
|
||||
// ValidateIntakeChatID refuses a chat id the intake half cannot use. The push
|
||||
// half accepts @channelusername as a destination. The intake half cannot: an
|
||||
// inbound update names its chat by numeric id, so an @-name would match nothing
|
||||
// and the poller would read the chat and answer none of it. Config validation
|
||||
// calls this, so the box refuses to boot rather than running a dead reach —
|
||||
// NewPoller returning an error is too late, because the daemon is already up.
|
||||
func ValidateIntakeChatID(chatID string) error {
|
||||
id := strings.TrimSpace(chatID)
|
||||
if id == "" {
|
||||
return errors.New("telegramsink: intake needs a chat id")
|
||||
}
|
||||
digits := strings.TrimPrefix(id, "-")
|
||||
if digits == "" || strings.TrimLeft(digits, "0123456789") != "" {
|
||||
return fmt.Errorf("telegramsink: intake needs the numeric chat id, not %s", chatID)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// NewPoller builds the intake half around an already-validated sink, so the
|
||||
// token, the base URL and the relay are resolved in one place. turn is
|
||||
// required; correct may be nil, and then the reply carries no buttons.
|
||||
func NewPoller(s *Sink, turn Turn, correct Correct) (*Poller, error) {
|
||||
if s == nil {
|
||||
return nil, errors.New("telegramsink: intake needs a sink")
|
||||
}
|
||||
if turn == nil {
|
||||
return nil, errors.New("telegramsink: intake needs a turn handler")
|
||||
}
|
||||
if err := ValidateIntakeChatID(s.cfg.ChatID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// The sink's transport already carries the relay. Only the timeout differs,
|
||||
// and it has to clear the long poll.
|
||||
hc := &http.Client{
|
||||
Timeout: (longPollSeconds + 10) * time.Second,
|
||||
Transport: s.hc.Transport,
|
||||
}
|
||||
return &Poller{sink: s, turn: turn, correct: correct, hc: hc}, nil
|
||||
}
|
||||
|
||||
// Run polls until the context ends. It never returns an error: a chat that
|
||||
// cannot be read is a degraded reach, not a reason to stop the daemon.
|
||||
func (p *Poller) Run(ctx context.Context) {
|
||||
p.discardBacklog(ctx)
|
||||
log.Printf("telegram intake: reading chat %s", p.sink.cfg.ChatID)
|
||||
for ctx.Err() == nil {
|
||||
updates, err := p.getUpdates(ctx, longPollSeconds)
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
log.Printf("telegram intake: poll: %v", err)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-time.After(pollBackoff):
|
||||
}
|
||||
continue
|
||||
}
|
||||
for _, u := range updates {
|
||||
p.handle(ctx, u)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// discardBacklog moves the offset past whatever is already queued, without
|
||||
// acting on any of it.
|
||||
//
|
||||
// Telegram holds undelivered updates for 24 hours, so a daemon that was down
|
||||
// overnight would otherwise wake up and answer every question in order. A
|
||||
// question asked eight hours ago has been answered by the owner himself or has
|
||||
// stopped mattering, and a reminder set from it would land at the wrong time.
|
||||
// Missing it is the safe direction.
|
||||
func (p *Poller) discardBacklog(ctx context.Context) {
|
||||
// getUpdates returns at most 100 per call, so one call is not the queue. The
|
||||
// loop is bounded rather than "until empty": the timeout is 0, so an instance
|
||||
// that keeps handing back a full batch would spin, and a thousand skipped
|
||||
// messages is already a box that was down for a long time.
|
||||
skipped := 0
|
||||
for range 10 {
|
||||
updates, err := p.getUpdates(ctx, 0)
|
||||
if err != nil {
|
||||
// Not fatal. The offset stays where it was, so the first real poll sees
|
||||
// what is left and answers it late. Say so rather than hide it.
|
||||
log.Printf("telegram intake: could not skip the backlog, old messages may be answered: %v", err)
|
||||
return
|
||||
}
|
||||
skipped += len(updates)
|
||||
if len(updates) == 0 {
|
||||
break
|
||||
}
|
||||
}
|
||||
if skipped > 0 {
|
||||
log.Printf("telegram intake: skipped %d message(s) queued while the daemon was down", skipped)
|
||||
}
|
||||
}
|
||||
|
||||
// handle dispatches one update. Anything that is neither a message from the
|
||||
// owner's chat nor a callback on one of Maven's own keyboards is dropped in
|
||||
// silence: a reply to a stranger confirms the bot exists and who it belongs to.
|
||||
func (p *Poller) handle(ctx context.Context, u update) {
|
||||
switch {
|
||||
case u.CallbackQuery != nil:
|
||||
p.onCallback(ctx, u.CallbackQuery)
|
||||
case u.Message != nil:
|
||||
p.onMessage(ctx, u.Message)
|
||||
}
|
||||
}
|
||||
|
||||
func (p *Poller) onMessage(ctx context.Context, m *message) {
|
||||
text := strings.TrimSpace(m.Text)
|
||||
if text == "" || !p.fromOwner(m.Chat.idString()) {
|
||||
return
|
||||
}
|
||||
// The conversation id keys the dialogue, so a clarify question asked in the
|
||||
// chat is not answered by an utterance typed on the web.
|
||||
reply, traceID, err := p.turn(ctx, "telegram:"+m.Chat.idString(), text)
|
||||
if err != nil {
|
||||
log.Printf("telegram intake: turn: %v", err)
|
||||
return
|
||||
}
|
||||
if strings.TrimSpace(reply) == "" {
|
||||
return
|
||||
}
|
||||
if err := p.send(ctx, reply, p.correctionKeyboard(traceID)); err != nil {
|
||||
log.Printf("telegram intake: reply: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// onCallback handles a tap on a correction button. Every path from the owner
|
||||
// answers the callback: telegram spins a clock on the button until it is
|
||||
// answered, and an unanswered tap reads as a gesture that was dropped. A tap
|
||||
// from anyone else gets silence, the same as a message from a stranger.
|
||||
func (p *Poller) onCallback(ctx context.Context, cb *callbackQuery) {
|
||||
if !p.fromOwner(cb.Message.Chat.idString()) {
|
||||
return
|
||||
}
|
||||
traceID, target, kind := parseCallback(cb.Data)
|
||||
if kind == callbackUnknown || p.correct == nil {
|
||||
p.answerCallback(ctx, cb.ID, "")
|
||||
return
|
||||
}
|
||||
// A tap on "не то" only opens the second row. Nothing is written yet: the
|
||||
// target is worth much more than the negative, so he gets the chance to name
|
||||
// it before the gesture is spent.
|
||||
if kind == callbackAskTarget {
|
||||
p.answerCallback(ctx, cb.ID, "")
|
||||
if err := p.editKeyboard(ctx, cb.Message.Chat.idString(), cb.Message.MessageID, targetKeyboard(traceID)); err != nil {
|
||||
log.Printf("telegram intake: open the target row: %v", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
if err := p.correct(ctx, traceID, target); err != nil {
|
||||
log.Printf("telegram intake: correct turn %d: %v", traceID, err)
|
||||
p.answerCallback(ctx, cb.ID, "не записалось")
|
||||
return
|
||||
}
|
||||
p.answerCallback(ctx, cb.ID, "записала")
|
||||
// The buttons come off, because the correction is given and a live keyboard
|
||||
// on an answered turn invites correcting it twice.
|
||||
if err := p.editKeyboard(ctx, cb.Message.Chat.idString(), cb.Message.MessageID, nil); err != nil {
|
||||
log.Printf("telegram intake: clear the keyboard: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// fromOwner — one chat, and it is the one the sink already sends to. Telegram
|
||||
// chat ids are not guessable, but they are also not secret: they travel in
|
||||
// every forwarded message. So this is the whole authorisation and it is an
|
||||
// allowlist of one.
|
||||
func (p *Poller) fromOwner(chatID string) bool {
|
||||
return chatID != "" && chatID == p.cfgChatID()
|
||||
}
|
||||
|
||||
func (p *Poller) cfgChatID() string { return strings.TrimSpace(p.sink.cfg.ChatID) }
|
||||
@@ -0,0 +1,216 @@
|
||||
package telegramsink
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// The turn he types in the chat is the turn the web would run, and the reply
|
||||
// carries the one gesture beside it.
|
||||
func TestIntakeRunsTheTurnAndOffersTheCorrection(t *testing.T) {
|
||||
b := newFakeBot(t)
|
||||
rec := &recorder{reply: "поняла", traceID: 91}
|
||||
p := newTestPoller(t, b, rec)
|
||||
|
||||
p.handle(context.Background(), msg(ownerChat, " поужинал "))
|
||||
|
||||
if got := rec.took(); len(got) != 1 || got[0] != "поужинал" {
|
||||
t.Fatalf("turns %q, want the trimmed utterance once", got)
|
||||
}
|
||||
// The dialogue is keyed per chat, so a clarify asked here is not answered on
|
||||
// the web.
|
||||
if rec.conversation != "telegram:"+ownerChat {
|
||||
t.Errorf("conversation %q does not name the chat", rec.conversation)
|
||||
}
|
||||
sends := b.called("sendMessage")
|
||||
if len(sends) != 1 {
|
||||
t.Fatalf("%d sends, want 1", len(sends))
|
||||
}
|
||||
if sends[0].body["text"] != "поняла" {
|
||||
t.Errorf("sent %v, want the reply", sends[0].body["text"])
|
||||
}
|
||||
if sends[0].body["protect_content"] != true {
|
||||
t.Error("his own data went out forwardable")
|
||||
}
|
||||
kb, _ := json.Marshal(sends[0].body["reply_markup"])
|
||||
if !strings.Contains(string(kb), "w:91") {
|
||||
t.Errorf("keyboard %s does not point at the turn's trace", kb)
|
||||
}
|
||||
}
|
||||
|
||||
// A turn nothing persisted has no row to correct, and a button that would name
|
||||
// one reports a failure he cannot act on.
|
||||
func TestIntakeSkipsTheGestureWithNoTrace(t *testing.T) {
|
||||
b := newFakeBot(t)
|
||||
p := newTestPoller(t, b, &recorder{reply: "поняла", traceID: 0})
|
||||
|
||||
p.handle(context.Background(), msg(ownerChat, "привет"))
|
||||
|
||||
sends := b.called("sendMessage")
|
||||
if len(sends) != 1 {
|
||||
t.Fatalf("%d sends, want 1", len(sends))
|
||||
}
|
||||
if _, ok := sends[0].body["reply_markup"]; ok {
|
||||
t.Error("offered a correction on a turn with no trace")
|
||||
}
|
||||
}
|
||||
|
||||
// One chat, and a stranger is not answered at all: a reply confirms the bot
|
||||
// exists and whose it is.
|
||||
func TestIntakeIgnoresAnyOtherChat(t *testing.T) {
|
||||
b := newFakeBot(t)
|
||||
rec := &recorder{reply: "поняла", traceID: 5}
|
||||
p := newTestPoller(t, b, rec)
|
||||
|
||||
p.handle(context.Background(), msg("9999", "включи свет"))
|
||||
p.handle(context.Background(), update{UpdateID: 8, CallbackQuery: &callbackQuery{
|
||||
ID: "cb", Data: "t:5:note", Message: message{Chat: chat{ID: json.Number("9999")}},
|
||||
}})
|
||||
|
||||
if got := rec.took(); len(got) != 0 {
|
||||
t.Errorf("ran %q for a chat that is not the owner's", got)
|
||||
}
|
||||
if len(rec.corrections) != 0 {
|
||||
t.Errorf("wrote %v from a chat that is not the owner's", rec.corrections)
|
||||
}
|
||||
if len(b.calls) != 0 {
|
||||
t.Errorf("answered a stranger: %v", b.calls)
|
||||
}
|
||||
}
|
||||
|
||||
// Tapping "не то" opens the seven and writes nothing yet. The target is worth
|
||||
// much more than the negative, so it must not be spent before he can name it.
|
||||
func TestIntakeFirstTapOnlyOpensTheTargets(t *testing.T) {
|
||||
b := newFakeBot(t)
|
||||
rec := &recorder{}
|
||||
p := newTestPoller(t, b, rec)
|
||||
|
||||
p.handle(context.Background(), update{UpdateID: 9, CallbackQuery: &callbackQuery{
|
||||
ID: "cb", Data: "w:77", Message: message{MessageID: 11, Chat: chat{ID: json.Number(ownerChat)}},
|
||||
}})
|
||||
|
||||
if len(rec.corrections) != 0 {
|
||||
t.Fatalf("wrote %v before he named a target", rec.corrections)
|
||||
}
|
||||
if len(b.called("answerCallbackQuery")) != 1 {
|
||||
t.Error("left the clock spinning on the button")
|
||||
}
|
||||
edits := b.called("editMessageReplyMarkup")
|
||||
if len(edits) != 1 {
|
||||
t.Fatalf("%d edits, want the target row", len(edits))
|
||||
}
|
||||
kb, _ := json.Marshal(edits[0].body["reply_markup"])
|
||||
for _, want := range CorrectionTargets {
|
||||
if !strings.Contains(string(kb), `"`+want+`"`) {
|
||||
t.Errorf("target row %s is missing %s", kb, want)
|
||||
}
|
||||
}
|
||||
// And the way out, because he opened the row without knowing he had to name
|
||||
// anything.
|
||||
if !strings.Contains(string(kb), `"t:77:"`) {
|
||||
t.Errorf("target row %s prices out the untargeted negative", kb)
|
||||
}
|
||||
}
|
||||
|
||||
func TestIntakeWritesTheCorrection(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name, data, want string
|
||||
}{
|
||||
{"with a target", "t:77:note", "note"},
|
||||
{"untargeted", "t:77:", ""},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
b := newFakeBot(t)
|
||||
rec := &recorder{}
|
||||
p := newTestPoller(t, b, rec)
|
||||
|
||||
p.handle(context.Background(), update{UpdateID: 9, CallbackQuery: &callbackQuery{
|
||||
ID: "cb", Data: tc.data, Message: message{MessageID: 11, Chat: chat{ID: json.Number(ownerChat)}},
|
||||
}})
|
||||
|
||||
if len(rec.corrections) != 1 || rec.corrections[0] != (correction{77, tc.want}) {
|
||||
t.Fatalf("corrections %v, want trace 77 → %q", rec.corrections, tc.want)
|
||||
}
|
||||
// The buttons come off once the gesture is given.
|
||||
edits := b.called("editMessageReplyMarkup")
|
||||
if len(edits) != 1 {
|
||||
t.Fatalf("%d edits, want the keyboard cleared", len(edits))
|
||||
}
|
||||
kb, _ := json.Marshal(edits[0].body["reply_markup"])
|
||||
if strings.Contains(string(kb), "t:77") {
|
||||
t.Errorf("keyboard %s still invites a second correction", kb)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// A write that failed says so on the button. Silence would read as recorded.
|
||||
func TestIntakeSaysWhenTheLabelDidNotLand(t *testing.T) {
|
||||
b := newFakeBot(t)
|
||||
rec := &recorder{correctErr: errors.New("no such routing trace")}
|
||||
p := newTestPoller(t, b, rec)
|
||||
|
||||
p.handle(context.Background(), update{UpdateID: 9, CallbackQuery: &callbackQuery{
|
||||
ID: "cb", Data: "t:77:fact", Message: message{MessageID: 11, Chat: chat{ID: json.Number(ownerChat)}},
|
||||
}})
|
||||
|
||||
answers := b.called("answerCallbackQuery")
|
||||
if len(answers) != 1 || answers[0].body["text"] == "" {
|
||||
t.Fatalf("answers %v, want a toast saying it did not land", answers)
|
||||
}
|
||||
if len(b.called("editMessageReplyMarkup")) != 0 {
|
||||
t.Error("cleared the buttons after a failed write, so he cannot try again")
|
||||
}
|
||||
}
|
||||
|
||||
// A question asked while the daemon was down has been answered by him or has
|
||||
// stopped mattering, and a reminder set from it would land at the wrong time.
|
||||
func TestIntakeDiscardsTheBacklog(t *testing.T) {
|
||||
b := newFakeBot(t, []update{msg(ownerChat, "напомни в 7 позвонить маме")})
|
||||
rec := &recorder{reply: "поняла", traceID: 3}
|
||||
p := newTestPoller(t, b, rec)
|
||||
|
||||
p.discardBacklog(context.Background())
|
||||
|
||||
if got := rec.took(); len(got) != 0 {
|
||||
t.Errorf("answered %q from the overnight queue", got)
|
||||
}
|
||||
// And the offset moved past it, so the next poll does not see it again.
|
||||
if p.offset != 8 {
|
||||
t.Errorf("offset %d, want the skipped update acknowledged", p.offset)
|
||||
}
|
||||
}
|
||||
|
||||
// The poller does not start without somewhere to send the turn.
|
||||
func TestNewPollerNeedsATurn(t *testing.T) {
|
||||
sink, err := New(Config{BotToken: "t", ChatID: ownerChat})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := NewPoller(sink, nil, nil); err == nil {
|
||||
t.Error("built a poller that reads the chat and answers nothing")
|
||||
}
|
||||
if _, err := NewPoller(nil, func(context.Context, string, string) (string, int64, error) {
|
||||
return "", 0, nil
|
||||
}, nil); err == nil {
|
||||
t.Error("built a poller with no sink to answer through")
|
||||
}
|
||||
}
|
||||
|
||||
// A chat id the intake half cannot match is refused before anything reads the
|
||||
// chat. Config validation calls the same check, so this is the boot error.
|
||||
func TestValidateIntakeChatID(t *testing.T) {
|
||||
for _, ok := range []string{"123", "-1001234567890", " 42 "} {
|
||||
if err := ValidateIntakeChatID(ok); err != nil {
|
||||
t.Errorf("ValidateIntakeChatID(%q): %v", ok, err)
|
||||
}
|
||||
}
|
||||
for _, bad := range []string{"", "@maven", "-", "12a", "1 2"} {
|
||||
if err := ValidateIntakeChatID(bad); err == nil {
|
||||
t.Errorf("ValidateIntakeChatID(%q) accepted; want error", bad)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,123 @@
|
||||
// intakeharness_test.go — a fake bot API and a recorder for what the poller
|
||||
// asked the daemon to do. Shared by the intake tests beside it.
|
||||
package telegramsink
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// fakeBot stands in for the bot API. It hands out queued updates once, records
|
||||
// every other call, and answers the ok=true envelope the poller checks.
|
||||
type fakeBot struct {
|
||||
mu sync.Mutex
|
||||
updates [][]update // one batch per getUpdates call, then empty
|
||||
calls []botCall
|
||||
srv *httptest.Server
|
||||
}
|
||||
|
||||
type botCall struct {
|
||||
method string
|
||||
body map[string]any
|
||||
}
|
||||
|
||||
func newFakeBot(t *testing.T, batches ...[]update) *fakeBot {
|
||||
t.Helper()
|
||||
b := &fakeBot{updates: batches}
|
||||
b.srv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
method := r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]
|
||||
raw, _ := io.ReadAll(r.Body)
|
||||
var body map[string]any
|
||||
_ = json.Unmarshal(raw, &body)
|
||||
|
||||
b.mu.Lock()
|
||||
b.calls = append(b.calls, botCall{method: method, body: body})
|
||||
var batch []update
|
||||
if method == "getUpdates" && len(b.updates) > 0 {
|
||||
batch, b.updates = b.updates[0], b.updates[1:]
|
||||
}
|
||||
b.mu.Unlock()
|
||||
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_ = json.NewEncoder(w).Encode(map[string]any{"ok": true, "result": batch})
|
||||
}))
|
||||
t.Cleanup(b.srv.Close)
|
||||
return b
|
||||
}
|
||||
|
||||
func (b *fakeBot) called(method string) []botCall {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
var out []botCall
|
||||
for _, c := range b.calls {
|
||||
if c.method == method {
|
||||
out = append(out, c)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// recorder collects what the poller asked the daemon to do.
|
||||
type recorder struct {
|
||||
mu sync.Mutex
|
||||
turns []string
|
||||
conversation string
|
||||
traceID int64
|
||||
corrections []correction
|
||||
reply string
|
||||
err error
|
||||
correctErr error
|
||||
}
|
||||
|
||||
type correction struct {
|
||||
traceID int64
|
||||
shouldBe string
|
||||
}
|
||||
|
||||
func (r *recorder) turn(_ context.Context, conversation, text string) (string, int64, error) {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.turns = append(r.turns, text)
|
||||
r.conversation = conversation
|
||||
return r.reply, r.traceID, r.err
|
||||
}
|
||||
|
||||
func (r *recorder) correct(_ context.Context, traceID int64, shouldBe string) error {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.corrections = append(r.corrections, correction{traceID, shouldBe})
|
||||
return r.correctErr
|
||||
}
|
||||
|
||||
func (r *recorder) took() []string {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
return append([]string(nil), r.turns...)
|
||||
}
|
||||
|
||||
const ownerChat = "4242"
|
||||
|
||||
func newTestPoller(t *testing.T, b *fakeBot, rec *recorder) *Poller {
|
||||
t.Helper()
|
||||
sink, err := New(Config{BotToken: "secret-token", ChatID: ownerChat, BaseURL: b.srv.URL})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
p, err := NewPoller(sink, rec.turn, rec.correct)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return p
|
||||
}
|
||||
|
||||
func msg(chatID, text string) update {
|
||||
return update{UpdateID: 7, Message: &message{
|
||||
MessageID: 11, Text: text, Chat: chat{ID: json.Number(chatID)},
|
||||
}}
|
||||
}
|
||||
@@ -72,6 +72,13 @@ type Config struct {
|
||||
// Timeout — per-request; 0 = DefaultTimeout. a dead relay can't hang the
|
||||
// tick loop.
|
||||
Timeout time.Duration
|
||||
|
||||
// Intake — read the chat as well as write to it (V-637). Off by default,
|
||||
// like the search and weather blocks: a bot that only pushes cannot be
|
||||
// talked into anything, and turning that off has to stay a deletion. When
|
||||
// set, a message from ChatID becomes a turn and its reply carries the
|
||||
// correction gesture. ChatID is the only accepted sender.
|
||||
Intake bool `json:"intake,omitempty"`
|
||||
}
|
||||
|
||||
// Sink — implements delivery.Sink via the telegram bot sendMessage API. one
|
||||
@@ -130,6 +137,11 @@ type sendMessageReq struct {
|
||||
Text string `json:"text"`
|
||||
DisableNotification bool `json:"disable_notification"` // false = ring (always — these are alarms)
|
||||
ProtectContent bool `json:"protect_content"` // true = no forwarding out of chat
|
||||
|
||||
// ReplyMarkup — the inline keyboard, used only by the intake half (V-637):
|
||||
// a reply to a turn he typed carries the correction gesture. nil on every
|
||||
// push the sink sends, and omitted from the wire when nil.
|
||||
ReplyMarkup *inlineKeyboard `json:"reply_markup,omitempty"`
|
||||
}
|
||||
|
||||
// telegramResp — the shape telegram returns. ok=false on logical error with
|
||||
|
||||
@@ -41,8 +41,33 @@ type PendingQuestion struct {
|
||||
Attempts int // questions already asked
|
||||
// MaxAttempts caps Attempts. 0 ⇒ DefaultMaxAttempts.
|
||||
MaxAttempts int
|
||||
// Suspends counts how many times this question has stepped aside for
|
||||
// something he asked instead, and come back on the end of the answer. It is
|
||||
// deliberately NOT an attempt: a side query is not a failed answer, and
|
||||
// charging it a retry is the V-554 shape. See CanResume for why it is
|
||||
// counted at all.
|
||||
Suspends int
|
||||
}
|
||||
|
||||
// MaxSuspends — how many times one question may step aside and come back before
|
||||
// she lets the request go (V-654).
|
||||
//
|
||||
// It exists because suspension had no bound of any kind. A side query spends no
|
||||
// attempt, so MaxAttempts never applies to it, and it restarts the 90s clock, so
|
||||
// the TTL never arrives either. Measured on 2026-08-07: one unfilled time slot
|
||||
// rode the end of six consecutive unrelated replies and stopped only when a
|
||||
// seventh turn happened to read as a failed answer.
|
||||
//
|
||||
// Three, matching DefaultMaxAttempts, and for the same reason. Once he has
|
||||
// asked for three other things without touching the question, the likely truth
|
||||
// is that he has moved on and has not said so.
|
||||
const MaxSuspends = 3
|
||||
|
||||
// CanResume reports whether this question may step aside once more. False ⇒ the
|
||||
// caller lets the request go and says so; it must never simply stop resuming,
|
||||
// because a question dropped in silence reads as one that was answered.
|
||||
func (q *PendingQuestion) CanResume() bool { return q.Suspends < MaxSuspends }
|
||||
|
||||
// Action reads the parked question as the typed action it is assembling
|
||||
// (pending.go). Derived rather than stored: the question's fields stay the one
|
||||
// copy of the truth, so a caller that fills them the old way cannot end up with
|
||||
|
||||
@@ -0,0 +1,192 @@
|
||||
package ipc
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// A cancelled context has to abort a call that is already in flight. It did not
|
||||
// until V-638: call checked ctx once before sending and then blocked in
|
||||
// roundtrip with no connection deadline, so a daemon that read the frame and
|
||||
// never answered parked the caller for as long as the socket stayed open.
|
||||
//
|
||||
// The server here is that daemon: it accepts, reads nothing, replies nothing.
|
||||
|
||||
func deafServer(t *testing.T) string {
|
||||
t.Helper()
|
||||
sock := filepath.Join(t.TempDir(), "deaf.sock")
|
||||
ln, err := net.Listen("unix", sock)
|
||||
if err != nil {
|
||||
t.Fatalf("listen: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = ln.Close() })
|
||||
go func() {
|
||||
for {
|
||||
conn, err := ln.Accept()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
// Hold it open and say nothing. Closed by the listener cleanup.
|
||||
t.Cleanup(func() { _ = conn.Close() })
|
||||
}
|
||||
}()
|
||||
return sock
|
||||
}
|
||||
|
||||
func TestClientCancelAbortsAReadInFlight(t *testing.T) {
|
||||
c, err := Dial(deafServer(t))
|
||||
if err != nil {
|
||||
t.Fatalf("dial: %v", err)
|
||||
}
|
||||
defer c.Close()
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
go func() {
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
cancel()
|
||||
}()
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() {
|
||||
_, err := c.Ping(ctx)
|
||||
done <- err
|
||||
}()
|
||||
|
||||
select {
|
||||
case err := <-done:
|
||||
// Ping is read-only, so the cancellation is reported as itself rather
|
||||
// than as an ambiguous mutation.
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Errorf("got %v, want context.Canceled", err)
|
||||
}
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("a cancelled Ping did not return")
|
||||
}
|
||||
}
|
||||
|
||||
// A mutation cancelled while awaiting the reply may already have committed, so
|
||||
// it is ErrAmbiguousOutcome and never a retry. That split is the invariant
|
||||
// internal/ipc/maperr_test.go's neighbours rest on.
|
||||
func TestClientCancelLeavesAMutationAmbiguous(t *testing.T) {
|
||||
c, err := Dial(deafServer(t))
|
||||
if err != nil {
|
||||
t.Fatalf("dial: %v", err)
|
||||
}
|
||||
defer c.Close()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
|
||||
defer cancel()
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() {
|
||||
_, err := c.WriteFact(ctx, WriteFactReq{Key: "water", Value: "drank"})
|
||||
done <- err
|
||||
}()
|
||||
|
||||
select {
|
||||
case err := <-done:
|
||||
if !errors.Is(err, ErrAmbiguousOutcome) {
|
||||
t.Errorf("got %v, want ErrAmbiguousOutcome", err)
|
||||
}
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("a cancelled WriteFact did not return")
|
||||
}
|
||||
}
|
||||
|
||||
// The deadline itself, with no cancellation: a call on a context with no
|
||||
// deadline used to have no bound at all. This one has one and must respect it.
|
||||
func TestClientDeadlineBoundsACall(t *testing.T) {
|
||||
c, err := Dial(deafServer(t))
|
||||
if err != nil {
|
||||
t.Fatalf("dial: %v", err)
|
||||
}
|
||||
defer c.Close()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
||||
defer cancel()
|
||||
|
||||
start := time.Now()
|
||||
if _, err := c.Ping(ctx); err == nil {
|
||||
t.Fatal("a deaf server answered a Ping")
|
||||
}
|
||||
if elapsed := time.Since(start); elapsed > 3*time.Second {
|
||||
t.Errorf("Ping took %v, want the context deadline to bound it", elapsed)
|
||||
}
|
||||
}
|
||||
|
||||
// blockingAPI parks Presence until its context is cancelled and records what
|
||||
// cancelled it. Every other method is the unimplemented floor.
|
||||
type blockingAPI struct {
|
||||
UnimplementedCoreAPI
|
||||
entered chan struct{}
|
||||
err chan error
|
||||
}
|
||||
|
||||
func (b *blockingAPI) Presence(ctx context.Context) (Presence, error) {
|
||||
close(b.entered)
|
||||
<-ctx.Done()
|
||||
b.err <- ctx.Err()
|
||||
return Presence{}, ctx.Err()
|
||||
}
|
||||
|
||||
// serveConn dispatched under context.Background() until V-638, so Close could
|
||||
// only abandon a dispatch in flight and never tell it to stop.
|
||||
func TestServerCloseCancelsADispatchInFlight(t *testing.T) {
|
||||
api := &blockingAPI{entered: make(chan struct{}), err: make(chan error, 1)}
|
||||
srv, err := Listen(filepath.Join(t.TempDir(), "core.sock"), api)
|
||||
if err != nil {
|
||||
t.Fatalf("listen: %v", err)
|
||||
}
|
||||
served := make(chan struct{})
|
||||
go func() { _ = srv.Serve(); close(served) }()
|
||||
|
||||
cli, err := Dial(srv.Path())
|
||||
if err != nil {
|
||||
t.Fatalf("dial: %v", err)
|
||||
}
|
||||
defer cli.Close()
|
||||
go func() { _, _ = cli.Presence(context.Background()) }()
|
||||
|
||||
select {
|
||||
case <-api.entered:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("the handler was never dispatched")
|
||||
}
|
||||
|
||||
_ = srv.Close()
|
||||
<-served
|
||||
select {
|
||||
case got := <-api.err:
|
||||
if !errors.Is(got, context.Canceled) {
|
||||
t.Errorf("handler saw %v, want context.Canceled", got)
|
||||
}
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("Close did not cancel the dispatch")
|
||||
}
|
||||
}
|
||||
|
||||
// The watchdog closes the conn, and it races the end of the call: a
|
||||
// cancellation landing as the reply arrives can close a conn the call was
|
||||
// already done with. That is survivable either way, because a write to a closed
|
||||
// socket is errWriteLost and errWriteLost re-dials and retries, so this test
|
||||
// passes with or without the drop in roundtrip's defer. What it pins is that
|
||||
// the recovery is real and costs one round trip at most, never an error the
|
||||
// caller sees.
|
||||
func TestClientSurvivesACancelledCall(t *testing.T) {
|
||||
_, _, cli, _ := newServerWithStore(t)
|
||||
|
||||
for i := 0; i < 20; i++ {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
go cancel() // races the reply on purpose
|
||||
_, _ = cli.Ping(ctx)
|
||||
cancel()
|
||||
|
||||
if _, err := cli.Ping(context.Background()); err != nil {
|
||||
t.Fatalf("call %d after a cancelled one: %v", i, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
+94
-12
@@ -26,9 +26,20 @@ type Client struct {
|
||||
conn net.Conn
|
||||
path string // the address as configured, kept for errors and logs
|
||||
addr netaddr.Addr // parsed, so a dropped conn can be re-dialed (core restart)
|
||||
mu sync.Mutex
|
||||
mu sync.Mutex // one request at a time, so a frame and its reply pair up
|
||||
|
||||
// connMu guards the conn field alone, and is held only across an assignment
|
||||
// or a read. It exists so Close and the cancellation watchdog can reach the
|
||||
// connection without waiting for the call that is holding c.mu (V-638).
|
||||
connMu sync.Mutex
|
||||
}
|
||||
|
||||
// defaultCallTimeout bounds a call whose context carries no deadline. It is
|
||||
// the same 120s internal/voice/client.go settles on: long enough for a model
|
||||
// call on a cold resident model, short enough that a daemon which stopped
|
||||
// answering does not park the caller forever.
|
||||
const defaultCallTimeout = 120 * time.Second
|
||||
|
||||
// errWriteLost marks a conn drop while sending the request frame: the request
|
||||
// never reached the server (or the server never saw a complete frame), so
|
||||
// retrying is always safe regardless of method — nothing was applied to
|
||||
@@ -103,11 +114,22 @@ func Dial(path string) (*Client, error) {
|
||||
return &Client{conn: c, path: path, addr: addr}, nil
|
||||
}
|
||||
|
||||
// Close closes the connection out from under a call in flight, on purpose: a
|
||||
// shutdown must not wait out a parked read. It takes connMu and never c.mu, so
|
||||
// it cannot block behind the call it is interrupting.
|
||||
//
|
||||
// The lock is taken and released by hand, around the two field accesses and
|
||||
// nothing else. The socket close happens outside it, because a close on a tcp
|
||||
// conn can block and connMu is on the path of every call.
|
||||
func (c *Client) Close() error {
|
||||
if c.conn == nil {
|
||||
c.connMu.Lock()
|
||||
conn := c.conn
|
||||
c.conn = nil
|
||||
c.connMu.Unlock()
|
||||
if conn == nil {
|
||||
return nil
|
||||
}
|
||||
return c.conn.Close()
|
||||
return conn.Close()
|
||||
}
|
||||
|
||||
// DialWait is Dial with patience: it retries with capped backoff until the
|
||||
@@ -163,16 +185,26 @@ func (c *Client) call(ctx context.Context, m Method, params, result any) error {
|
||||
}
|
||||
|
||||
var resp Response
|
||||
err := c.roundtrip(m, raw, &resp)
|
||||
err := c.roundtrip(ctx, m, raw, &resp)
|
||||
switch {
|
||||
case errors.Is(err, errWriteLost):
|
||||
// The request never left; a duplicate send can't double-apply.
|
||||
// Redial (roundtrip re-dials on a nil conn) and retry exactly once.
|
||||
err = c.roundtrip(m, raw, &resp)
|
||||
// Not when the caller has given up — a retry would only be a second
|
||||
// frame nobody is waiting for.
|
||||
if ctx.Err() == nil {
|
||||
err = c.roundtrip(ctx, m, raw, &resp)
|
||||
}
|
||||
case errors.Is(err, errReadLost):
|
||||
if readOnlyMethods[m] {
|
||||
if ctx.Err() != nil {
|
||||
// The caller cancelled the read it was waiting for. Nothing
|
||||
// was applied, so this is the cancellation and not an
|
||||
// ambiguity.
|
||||
return ctx.Err()
|
||||
}
|
||||
// A duplicate read can't double-apply either — safe to replay.
|
||||
err = c.roundtrip(m, raw, &resp)
|
||||
err = c.roundtrip(ctx, m, raw, &resp)
|
||||
} else {
|
||||
// The mutation may have already committed server-side. Do not
|
||||
// retry: report the ambiguity instead of guessing.
|
||||
@@ -199,19 +231,55 @@ func (c *Client) call(ctx context.Context, m Method, params, result any) error {
|
||||
// failure is wrapped in errReadLost (ambiguous — call() only retries it for
|
||||
// read-only methods). Either way a failed conn is dropped so the next call
|
||||
// re-dials clean. Caller holds c.mu.
|
||||
func (c *Client) roundtrip(m Method, raw json.RawMessage, resp *Response) error {
|
||||
if c.conn == nil {
|
||||
conn, err := netaddr.Dial(c.addr)
|
||||
//
|
||||
// The connection carries a deadline derived from ctx, falling back to
|
||||
// defaultCallTimeout, and a watchdog closes it if ctx is cancelled mid-call
|
||||
// (V-638). Before that a daemon which stopped answering parked the caller for
|
||||
// as long as the socket stayed open. The watchdog closes the conn rather than
|
||||
// calling drop, because drop wants c.mu and the caller is holding it — the
|
||||
// closed socket fails the read, and roundtrip drops it on the way out.
|
||||
func (c *Client) roundtrip(ctx context.Context, m Method, raw json.RawMessage, resp *Response) error {
|
||||
conn := c.currentConn()
|
||||
if conn == nil {
|
||||
dialed, err := netaddr.Dial(c.addr)
|
||||
if err != nil {
|
||||
return fmt.Errorf("%w: dial %s: %v", errWriteLost, c.addr, err)
|
||||
}
|
||||
c.conn = conn
|
||||
c.setConn(dialed)
|
||||
conn = dialed
|
||||
}
|
||||
if err := writeFrame(c.conn, Request{Method: m, Params: raw}); err != nil {
|
||||
if dl, ok := ctx.Deadline(); ok {
|
||||
_ = conn.SetDeadline(dl)
|
||||
} else {
|
||||
_ = conn.SetDeadline(time.Now().Add(defaultCallTimeout))
|
||||
}
|
||||
defer conn.SetDeadline(time.Time{})
|
||||
|
||||
// The watchdog and the end of the call race by construction: a cancellation
|
||||
// landing just as the reply arrives can close a conn this call is already
|
||||
// done with, and c.conn would still point at the closed socket. So a call
|
||||
// whose context ended does not leave the conn behind for the next one,
|
||||
// whichever of the two got there first.
|
||||
done := make(chan struct{})
|
||||
defer func() {
|
||||
close(done)
|
||||
if ctx.Err() != nil {
|
||||
c.drop()
|
||||
}
|
||||
}()
|
||||
go func() {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
_ = conn.Close()
|
||||
case <-done:
|
||||
}
|
||||
}()
|
||||
|
||||
if err := writeFrame(conn, Request{Method: m, Params: raw}); err != nil {
|
||||
c.drop()
|
||||
return fmt.Errorf("%w: %v", errWriteLost, err)
|
||||
}
|
||||
if err := readFrame(c.conn, resp); err != nil {
|
||||
if err := readFrame(conn, resp); err != nil {
|
||||
c.drop()
|
||||
return fmt.Errorf("%w: %v", errReadLost, err)
|
||||
}
|
||||
@@ -220,12 +288,26 @@ func (c *Client) roundtrip(m Method, raw json.RawMessage, resp *Response) error
|
||||
|
||||
// drop closes and forgets the current conn so the next call re-dials.
|
||||
func (c *Client) drop() {
|
||||
c.connMu.Lock()
|
||||
defer c.connMu.Unlock()
|
||||
if c.conn != nil {
|
||||
_ = c.conn.Close()
|
||||
c.conn = nil
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Client) currentConn() net.Conn {
|
||||
c.connMu.Lock()
|
||||
defer c.connMu.Unlock()
|
||||
return c.conn
|
||||
}
|
||||
|
||||
func (c *Client) setConn(conn net.Conn) {
|
||||
c.connMu.Lock()
|
||||
defer c.connMu.Unlock()
|
||||
c.conn = conn
|
||||
}
|
||||
|
||||
// hydrate rehydrates a wire RpcError into the matching package sentinel. The
|
||||
// code↔sentinel table is the only place the wire "knows" about errors; keep it
|
||||
// in sync with codeOf in wire.go.
|
||||
|
||||
+36
-8
@@ -17,9 +17,11 @@ import (
|
||||
// Server — the core side of the boundary. Listens on a unix domain socket,
|
||||
// accepts module connections, frames requests to a CoreAPI and responses back.
|
||||
// One Server per daemon process; concurrent connections are handled in their
|
||||
// own goroutine but share the single CoreAPI (and therefore the single store
|
||||
// writer — store is single-connection, SetMaxOpenConns(1), so serialization is
|
||||
// already guaranteed at the db; the Server adds no locking of its own).
|
||||
// own goroutine but share the single CoreAPI, and so the single store writer.
|
||||
// The store opens at SetMaxOpenConns(1), so serialisation is already guaranteed
|
||||
// at the database and the Server adds no locking of its own. That cap is an
|
||||
// invariant this comment depends on, measured and kept on 07-08-2026 (V-642,
|
||||
// docs/evals/2026-08-07-store-connection-cap.md).
|
||||
type Server struct {
|
||||
api atomic.Value // stores CoreAPI
|
||||
path string
|
||||
@@ -30,6 +32,14 @@ type Server struct {
|
||||
done chan struct{}
|
||||
accept sync.Mutex // guards wg.Add vs Close's wg.Wait sequence
|
||||
|
||||
// ctx — server-scoped, cancelled by Close, and the parent of every request
|
||||
// context. serveConn dispatched under context.Background() until V-638, so
|
||||
// a dispatch in flight during shutdown could not be told to stop and the
|
||||
// closeGrace below could only abandon it. Cancelling gives a handler that
|
||||
// respects its context the chance to return instead.
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
|
||||
// conns — every accepted connection still being served. Close needs these
|
||||
// because closing the listener does nothing to a connection already
|
||||
// accepted: serveConn is parked in readFrame waiting for a peer that may
|
||||
@@ -208,11 +218,14 @@ func Listen(path string, api CoreAPI) (*Server, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
s := &Server{
|
||||
path: path,
|
||||
addr: addr,
|
||||
ln: ln,
|
||||
done: make(chan struct{}),
|
||||
path: path,
|
||||
addr: addr,
|
||||
ln: ln,
|
||||
done: make(chan struct{}),
|
||||
ctx: ctx,
|
||||
cancel: cancel,
|
||||
}
|
||||
s.api.Store(api)
|
||||
return s, nil
|
||||
@@ -250,7 +263,10 @@ func (s *Server) Serve() error {
|
||||
|
||||
func (s *Server) serveConn(c net.Conn) {
|
||||
caller, callerOK := peerCaller(c)
|
||||
ctx := context.Background()
|
||||
// Derived from the server's, so Close cancels a dispatch in flight, and
|
||||
// cancelled when this conn ends so nothing a handler spawned outlives it.
|
||||
ctx, cancel := context.WithCancel(s.serverContext())
|
||||
defer cancel()
|
||||
if callerOK {
|
||||
ctx = WithCaller(ctx, caller)
|
||||
}
|
||||
@@ -274,6 +290,15 @@ func (s *Server) serveConn(c net.Conn) {
|
||||
}
|
||||
}
|
||||
|
||||
// serverContext is s.ctx, or Background for a Server built as a zero value
|
||||
// rather than by Listen (the wiring tests do that).
|
||||
func (s *Server) serverContext() context.Context {
|
||||
if s.ctx == nil {
|
||||
return context.Background()
|
||||
}
|
||||
return s.ctx
|
||||
}
|
||||
|
||||
func (s *Server) safeDispatch(ctx context.Context, req Request) (result json.RawMessage, err error) {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
@@ -720,6 +745,9 @@ func (s *Server) Close() error {
|
||||
default:
|
||||
close(s.done)
|
||||
}
|
||||
if s.cancel != nil {
|
||||
s.cancel()
|
||||
}
|
||||
err := s.ln.Close()
|
||||
// Closing the listener stops new connections; it does nothing to the ones
|
||||
// already accepted. Close those too, or every serveConn parked in readFrame
|
||||
|
||||
@@ -105,6 +105,12 @@ type Decision struct {
|
||||
Slots Slots
|
||||
Clarify bool // stage 3: below threshold — ask, don't guess
|
||||
|
||||
// Source — where the answer lives, for a query. The second half of the
|
||||
// route, and empty on every other intent. SourceUnknown means no decider
|
||||
// named one and the daemon walks its whole chain, which is what shipped
|
||||
// before this field existed. See source.go for why it is twelve values.
|
||||
Source Source
|
||||
|
||||
// Continued — this decision was rebuilt from the previous turn rather
|
||||
// than routed, because the utterance was an ellipsis ("а завтра?").
|
||||
// Handlers use it to know that Slots.Text is the PREVIOUS turn's topic
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
package router
|
||||
|
||||
// Source — where the answer to a query lives. It is the second half of a
|
||||
// routing decision and it used to be made outside the router entirely (V-655).
|
||||
//
|
||||
// The cascade sorted an utterance into one of seven intents with stage 0 rules,
|
||||
// the resident model and the classifier behind it, a fixture measuring it and
|
||||
// the decision trace recording it. Then IntentQuery handed the turn to
|
||||
// querySources in the daemon, a chain of twenty-two branches deciding by seed
|
||||
// similarity in a fixed order, with none of that. So the careful sorter did the
|
||||
// easy half and the sloppy one did the hard half: on 2026-08-07 weather claimed
|
||||
// "что такое TCP?" and answered "для какого города?", because weather read one
|
||||
// percent closer to the turn than the pile of leftover seeds did, and one
|
||||
// percent was enough. Search would have answered it and search was never asked.
|
||||
//
|
||||
// "query" is not a destination. It is a shrug. This is the field that says
|
||||
// where to look.
|
||||
//
|
||||
// # Why twelve and not twenty-two
|
||||
//
|
||||
// A destination is what a decider can plausibly name from the utterance alone,
|
||||
// not one entry per source. Three of the daemon's sources are successive passes
|
||||
// over his own words and a fourth reads the facts by key: which of them lands
|
||||
// the hit is an ordering detail inside the chain, and no utterance says. They
|
||||
// are SourceRecall together. The same goes for the metasearch, the offline
|
||||
// encyclopedia and a page he named by URL, which are SourceWorld.
|
||||
//
|
||||
// # Empty is a real value and it is the floor
|
||||
//
|
||||
// SourceUnknown means nobody decided. The daemon then walks the whole chain in
|
||||
// its original order, which is the behaviour that shipped before this field
|
||||
// existed. So the classifier arm sets nothing and costs nothing, and a box
|
||||
// whose model is down routes queries exactly as it did.
|
||||
type Source string
|
||||
|
||||
const (
|
||||
// SourceUnknown — no decider named a destination. Walk the chain.
|
||||
SourceUnknown Source = ""
|
||||
|
||||
// His own data.
|
||||
SourceRecall Source = "recall" // notes, facts and what he has said before
|
||||
SourceCalendar Source = "calendar" // events, and the only date-aware destination
|
||||
SourceTasks Source = "tasks" // the task list
|
||||
SourceList Source = "list" // the shopping and other named lists
|
||||
SourceMoney Source = "money" // the spending facts the poller writes
|
||||
|
||||
// The surroundings.
|
||||
SourceWeather Source = "weather" // the forecast for a place
|
||||
SourceHome Source = "home" // lights, devices, the house
|
||||
SourceNetwork Source = "network" // the LAN and what is on it
|
||||
SourceFeeds Source = "feeds" // the RSS she reads
|
||||
SourceAttention Source = "attention" // what Praxis says needs looking at
|
||||
|
||||
// Everything else.
|
||||
SourceSelf Source = "self" // a question about Maven herself
|
||||
SourceWorld Source = "world" // search, the ZIMs, a page he named
|
||||
)
|
||||
|
||||
// Sources — every destination a decider may name, in a fixed order so a prompt,
|
||||
// a grammar table and a test all read the same list. SourceUnknown is not a
|
||||
// member: it is the absence of a choice, not one of the choices.
|
||||
var Sources = []Source{
|
||||
SourceRecall, SourceCalendar, SourceTasks, SourceList, SourceMoney,
|
||||
SourceWeather, SourceHome, SourceNetwork, SourceFeeds, SourceAttention,
|
||||
SourceSelf, SourceWorld,
|
||||
}
|
||||
|
||||
// ValidSource reports whether s is one a decider may name. Anything else,
|
||||
// including a destination invented by a model, is dropped back to
|
||||
// SourceUnknown by the caller rather than trusted.
|
||||
func ValidSource(s Source) bool {
|
||||
for _, known := range Sources {
|
||||
if s == known {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
@@ -187,9 +187,11 @@ func SystemTimeDateGrammars() []Grammar {
|
||||
// written ("the clock/date system rule must not swallow it"); the daemon
|
||||
// disagreed with the fixture and the daemon was wrong.
|
||||
//
|
||||
// Routing, not answering. These set the intent and nothing else — which source
|
||||
// in the query chain claims the turn stays the chain's decision, and a
|
||||
// question with no date still falls through queryCalendar to recall.
|
||||
// Routing, not answering. Two of the five also name the calendar as the
|
||||
// destination (V-655), which narrows who may GUESS their way onto the turn and
|
||||
// claims nothing. Every source that looks something up still runs, in the order
|
||||
// it always did, so a question with no date still falls through queryCalendar
|
||||
// to recall.
|
||||
//
|
||||
// Deliberately not folded into SystemTimeDateGrammars: those exist to send
|
||||
// utterances TO system, these exist to keep utterances OUT of it, and one
|
||||
@@ -199,9 +201,15 @@ func AgendaQueryGrammars() []Grammar {
|
||||
{
|
||||
// An explicit calendar noun is unambiguous wherever it appears:
|
||||
// "что в календаре на завтра", "покажи расписание на среду".
|
||||
//
|
||||
// The one agenda rule that names its destination, because an
|
||||
// explicit calendar noun leaves nothing to weigh (V-655). The
|
||||
// possessive rules below deliberately do not: "что у меня в списке
|
||||
// покупок" matches agenda-query, and naming the calendar there
|
||||
// would take the list source off the turn.
|
||||
Name: "calendar-query",
|
||||
Pattern: regexp.MustCompile(`(?i)(календар|расписани|повестк)`),
|
||||
Build: agendaQueryBuild,
|
||||
Build: queryTo(SourceCalendar),
|
||||
},
|
||||
{
|
||||
// The agenda phrasing with no calendar noun. Anchored at the start
|
||||
@@ -251,9 +259,11 @@ func AgendaQueryGrammars() []Grammar {
|
||||
// "во сколько созвон". He is asking when something on his calendar
|
||||
// happens, and the noun is the only signal. Closed list, so "когда
|
||||
// битва при Ватерлоо" is still a world question.
|
||||
// Names the calendar (V-655): the noun list is closed and every
|
||||
// member of it is an event, so there is nothing else to weigh.
|
||||
Name: "event-time-query",
|
||||
Pattern: regexp.MustCompile(`(?i)^\s*(когда|во\s+сколько|в\s+котором\s+часу)\s+(будет\s+|у\s+нас\s+)?(планёрк|планерк|встреч|созвон|митинг|совещани|звонок|созвон|приём|прием|интервью|собеседовани|тренировк|урок|занятие|пара)[а-я]*(\s|[?!.]|$)`),
|
||||
Build: agendaQueryBuild,
|
||||
Build: queryTo(SourceCalendar),
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -340,6 +350,11 @@ func narrativeQueryBuild(m []string) (Decision, bool) {
|
||||
Intent: IntentQuery,
|
||||
Confidence: 1.0,
|
||||
Slots: Slots{Text: topic},
|
||||
// The world, because that is the shape this asks for and the rule has
|
||||
// already declined the two cases where it is not: entertainment, and
|
||||
// questions about her (V-655). His own notes are still read first — a
|
||||
// destination narrows who may guess and reorders nothing.
|
||||
Source: SourceWorld,
|
||||
}, true
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,86 @@
|
||||
package router
|
||||
|
||||
import "regexp"
|
||||
|
||||
// WorldQueryGrammars — stage-0 rules for the two question shapes that name the
|
||||
// world in their own words, and say so plainly enough that no scorer is needed
|
||||
// (V-655).
|
||||
//
|
||||
// They exist because of what happens when nothing deterministic claims these.
|
||||
// Measured on the box on 2026-08-07 (docs/evals/2026-08-07-week-of-usage.md,
|
||||
// section 4): "что такое TCP?" and "сколько будет 17 на 23?" were both answered
|
||||
// "для какого города?", and "кто такой Линус Торвальдс?" was answered "не знаю —
|
||||
// не нашла у тебя такой записи". None of those three is about him, about the
|
||||
// weather, or about anything on this box.
|
||||
//
|
||||
// The mechanism is the destination, not the answer. Naming SourceWorld does not
|
||||
// send the turn outside and does not skip a single source that looks something
|
||||
// up: his notes, his facts and the personal boundary all still run first, in the
|
||||
// order they always did. What it does is stop the sources that claim on seed
|
||||
// similarity from taking the turn on the way past. Weather cannot claim a
|
||||
// question about a protocol once the utterance has said which side it is on.
|
||||
//
|
||||
// Both patterns are spelled out here rather than drawn from internal/lexicon,
|
||||
// which is the same call the agenda rules made: these are interrogative FRAMES
|
||||
// of two words, not a closed class of single words, and the lexicon holds
|
||||
// classes. Nothing here is a stem pattern over open vocabulary — the variable
|
||||
// part of each rule is the topic, and the rule reads none of it.
|
||||
func WorldQueryGrammars() []Grammar {
|
||||
return []Grammar{
|
||||
{
|
||||
// "что такое X", "кто такой X". A request for what a thing or a
|
||||
// person IS, which his own data can answer and usually cannot.
|
||||
//
|
||||
// The topic is deliberately not captured into Slots.Text. Every
|
||||
// source below reads the utterance, "что такое TCP?" is already the
|
||||
// best query string for it, and the agenda rules make the same call
|
||||
// for the same reason.
|
||||
Name: "definition-query",
|
||||
Pattern: definitionQueryPattern,
|
||||
Build: queryTo(SourceWorld),
|
||||
},
|
||||
{
|
||||
// "сколько будет 17 на 23", "сколько будет 2+2". Arithmetic, which
|
||||
// the metasearch answers and no local source holds. The digits are
|
||||
// what make it arithmetic: "сколько будет гостей" names no number
|
||||
// and is a question about his evening.
|
||||
Name: "arithmetic-query",
|
||||
Pattern: arithmeticQueryPattern,
|
||||
Build: queryTo(SourceWorld),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// definitionQueryPattern — anchored at the start, because "напомни узнать что
|
||||
// такое TCP" is a reminder that happens to contain the frame.
|
||||
//
|
||||
// (\s|[?!.]|$) and not \b: Go's \b is ASCII-only and never fires after a
|
||||
// Cyrillic letter, so the ASCII form silently matches nothing. The agenda rules
|
||||
// carry the same note.
|
||||
var definitionQueryPattern = regexp.MustCompile(
|
||||
`(?i)^\s*(что\s+так(ое|ая)|кто\s+так(ой|ая|ие)|what\s+is|who\s+is)(\s|[?!.]|$)`)
|
||||
|
||||
// arithmeticQueryPattern — the ask, then a digit somewhere after it. Loose on
|
||||
// what sits between them on purpose: the operator is spoken half a dozen ways
|
||||
// ("на", "умножить на", "плюс", "+") and reading them is the calculator's job,
|
||||
// not this rule's. All this decides is which side of the boundary the turn is
|
||||
// on.
|
||||
var arithmeticQueryPattern = regexp.MustCompile(
|
||||
`(?i)^\s*(сколько\s+будет|посчитай|вычисли|how\s+much\s+is)\s.*\d`)
|
||||
|
||||
// queryTo builds a stage-0 query Decision that names where the answer lives.
|
||||
//
|
||||
// The utterance travels intact and no slot is filled, which is the same
|
||||
// contract agendaQueryBuild has: confidence 1.0 on the intent and the
|
||||
// destination, and every source below still decides for itself whether it has
|
||||
// an answer. Naming a destination narrows who may guess. It promises nothing.
|
||||
func queryTo(dest Source) func([]string) (Decision, bool) {
|
||||
return func([]string) (Decision, bool) {
|
||||
return Decision{
|
||||
Stage: 0,
|
||||
Intent: IntentQuery,
|
||||
Confidence: 1.0,
|
||||
Source: dest,
|
||||
}, true
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
package router
|
||||
|
||||
import "testing"
|
||||
|
||||
// The three utterances from the 2026-08-07 week on the box that no local source
|
||||
// could answer and three different local sources claimed anyway. Stage 0 has to
|
||||
// say which side of the boundary they are on, because by the time the chain is
|
||||
// walking, the only thing separating them from a weather forecast is a cosine.
|
||||
func TestAWorldQuestionNamesTheWorld(t *testing.T) {
|
||||
cases := []struct {
|
||||
utterance string
|
||||
rule string
|
||||
}{
|
||||
{"что такое TCP?", "definition-query"},
|
||||
{"кто такой Линус Торвальдс?", "definition-query"},
|
||||
{"что такая мембрана", "definition-query"},
|
||||
{"кто такая Ада Лавлейс?", "definition-query"},
|
||||
{"what is TCP?", "definition-query"},
|
||||
{"сколько будет 17 на 23?", "arithmetic-query"},
|
||||
{"посчитай 2+2", "arithmetic-query"},
|
||||
{"сколько будет 5 умножить на 6", "arithmetic-query"},
|
||||
}
|
||||
for _, c := range cases {
|
||||
dec, rule, ok := matchWorldQuery(c.utterance)
|
||||
if !ok {
|
||||
t.Errorf("%q: no world rule claimed it", c.utterance)
|
||||
continue
|
||||
}
|
||||
if rule != c.rule {
|
||||
t.Errorf("%q: claimed by %q, want %q", c.utterance, rule, c.rule)
|
||||
}
|
||||
if dec.Intent != IntentQuery {
|
||||
t.Errorf("%q: intent %q, want query", c.utterance, dec.Intent)
|
||||
}
|
||||
if dec.Source != SourceWorld {
|
||||
t.Errorf("%q: source %q, want %q", c.utterance, dec.Source, SourceWorld)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The frame has to be the whole opening or the rule is reading somebody else's
|
||||
// sentence. Every case here contains a world-question shape and is not one.
|
||||
func TestAWorldRuleDeclinesWhatIsNotItsShape(t *testing.T) {
|
||||
cases := []struct {
|
||||
utterance string
|
||||
why string
|
||||
}{
|
||||
{"напомни узнать что такое TCP", "a reminder that happens to quote the frame"},
|
||||
{"запиши что такое TCP", "a capture that happens to quote the frame"},
|
||||
{"сколько будет гостей", "an ask with no number is not arithmetic"},
|
||||
{"что у меня сегодня?", "his agenda, and the agenda rules own it"},
|
||||
{"кто там?", "not the frame"},
|
||||
{"посчитай расходы", "no number, so the money source keeps it"},
|
||||
}
|
||||
for _, c := range cases {
|
||||
if _, rule, ok := matchWorldQuery(c.utterance); ok {
|
||||
t.Errorf("%q: claimed by %q, want no claim — %s", c.utterance, rule, c.why)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The destination is advice about who may guess, never a filled slot. A rule
|
||||
// that quietly captured the topic would change what every source below reads.
|
||||
func TestNamingTheWorldFillsNoSlot(t *testing.T) {
|
||||
dec, _, ok := matchWorldQuery("что такое TCP?")
|
||||
if !ok {
|
||||
t.Fatal("definition-query did not claim it")
|
||||
}
|
||||
if dec.Slots.Text != "" || dec.Slots.HasTime || dec.Slots.HasFn || dec.Slots.HasKey {
|
||||
t.Errorf("slots = %+v, want none filled", dec.Slots)
|
||||
}
|
||||
if dec.Confidence != 1.0 {
|
||||
t.Errorf("confidence = %v, want 1.0 for a stage-0 match", dec.Confidence)
|
||||
}
|
||||
}
|
||||
|
||||
func matchWorldQuery(utterance string) (Decision, string, bool) {
|
||||
for _, g := range WorldQueryGrammars() {
|
||||
m := g.Pattern.FindStringSubmatch(utterance)
|
||||
if m == nil {
|
||||
continue
|
||||
}
|
||||
if dec, ok := g.Build(m); ok {
|
||||
return dec, g.Name, true
|
||||
}
|
||||
}
|
||||
return Decision{}, "", false
|
||||
}
|
||||
+37
-6
@@ -89,8 +89,8 @@ type Poller struct {
|
||||
ranker Ranker
|
||||
cfg Config
|
||||
nextDue map[string]time.Time
|
||||
seen map[string]map[string]bool // feed → item ID, for items with no date
|
||||
polled map[string]bool // feed → polled at least once in THIS process
|
||||
seen map[string]*seenIDs // feed → item IDs, for items with no date
|
||||
polled map[string]bool // feed → polled at least once in THIS process
|
||||
}
|
||||
|
||||
// NewPoller wires a poller. Returns nil when there is nothing to poll — a
|
||||
@@ -122,7 +122,7 @@ func NewPoller(feeds []FeedConfig, fetch Fetcher, notes Notes, marks Marks, embe
|
||||
feeds: valid, fetch: fetch, notes: notes, marks: marks,
|
||||
embed: embed, ranker: ranker, cfg: cfg,
|
||||
nextDue: map[string]time.Time{},
|
||||
seen: map[string]map[string]bool{},
|
||||
seen: map[string]*seenIDs{},
|
||||
polled: map[string]bool{},
|
||||
}
|
||||
}
|
||||
@@ -283,6 +283,38 @@ func (p *Poller) mark(ctx context.Context, feed string, now time.Time) (time.Tim
|
||||
return at, true
|
||||
}
|
||||
|
||||
// maxSeenPerFeed bounds the undated-item set. It has to stay comfortably above
|
||||
// any one feed's front page, or an item still listed there would fall out of the
|
||||
// set and be written a second time. A few hundred entries covers the largest
|
||||
// page anyone publishes, and the set only has to span one poll window plus the
|
||||
// resync guard, not all of history.
|
||||
const maxSeenPerFeed = 512
|
||||
|
||||
// seenIDs is a bounded insertion-ordered set. The map answers the lookup, the
|
||||
// slice remembers what to drop first, so an undated feed cannot grow the poller
|
||||
// for as long as mavend runs.
|
||||
type seenIDs struct {
|
||||
ids map[string]bool
|
||||
order []string
|
||||
}
|
||||
|
||||
// add records id and reports whether it was new.
|
||||
func (s *seenIDs) add(id string) bool {
|
||||
if s.ids == nil {
|
||||
s.ids = make(map[string]bool, maxSeenPerFeed)
|
||||
}
|
||||
if s.ids[id] {
|
||||
return false
|
||||
}
|
||||
s.ids[id] = true
|
||||
s.order = append(s.order, id)
|
||||
if len(s.order) > maxSeenPerFeed {
|
||||
delete(s.ids, s.order[0])
|
||||
s.order = s.order[1:]
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// fresh — two dedup rules, because feeds are inconsistent about dates. A dated
|
||||
// item must be newer than the mark; an undated one is kept once per process by
|
||||
// ID.
|
||||
@@ -308,12 +340,11 @@ func (p *Poller) fresh(f FeedConfig, it Item, mark, now time.Time, resync bool)
|
||||
id = it.Title
|
||||
}
|
||||
if p.seen[f.Name] == nil {
|
||||
p.seen[f.Name] = map[string]bool{}
|
||||
p.seen[f.Name] = &seenIDs{}
|
||||
}
|
||||
if p.seen[f.Name][id] {
|
||||
if !p.seen[f.Name].add(id) {
|
||||
return false
|
||||
}
|
||||
p.seen[f.Name][id] = true
|
||||
return !resync
|
||||
}
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@ package rss
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -208,3 +209,27 @@ func TestNoFeedsMeansNoPoller(t *testing.T) {
|
||||
t.Fatal("a feed with no name or url is not a configuration")
|
||||
}
|
||||
}
|
||||
|
||||
// An undated feed used to grow p.seen for as long as mavend ran. The set is
|
||||
// bounded now, and the bound must not cost the dedupe an item still on the
|
||||
// front page — only ids far older than any page fall out.
|
||||
func TestSeenIDsBounded(t *testing.T) {
|
||||
var s seenIDs
|
||||
for i := 0; i < maxSeenPerFeed*3; i++ {
|
||||
if !s.add(fmt.Sprintf("item-%d", i)) {
|
||||
t.Fatalf("item-%d read as already seen", i)
|
||||
}
|
||||
if len(s.ids) > maxSeenPerFeed || len(s.order) > maxSeenPerFeed {
|
||||
t.Fatalf("after %d inserts: ids=%d order=%d, cap is %d",
|
||||
i+1, len(s.ids), len(s.order), maxSeenPerFeed)
|
||||
}
|
||||
}
|
||||
// The newest insert is still deduped; the oldest was evicted.
|
||||
last := fmt.Sprintf("item-%d", maxSeenPerFeed*3-1)
|
||||
if s.add(last) {
|
||||
t.Fatalf("%s read as new, so the most recent id was dropped", last)
|
||||
}
|
||||
if !s.add("item-0") {
|
||||
t.Fatal("item-0 survived, so nothing was evicted")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,154 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// conncap_test.go measures whether a read queues behind a write at
|
||||
// SetMaxOpenConns(1), which is what openAt sets (V-642). It is a measurement
|
||||
// harness, not an assertion: the numbers it prints are the evidence, and the
|
||||
// decision to move the cap or leave it belongs in docs/evals.
|
||||
//
|
||||
// Run it with -v, and note that it is skipped under -short because it spends
|
||||
// seconds on purpose.
|
||||
|
||||
// openCapped opens a plaintext store at the given connection cap. In-package,
|
||||
// so it can reach the handle openAt caps at 1.
|
||||
func openCapped(t *testing.T, cap int) *Store {
|
||||
t.Helper()
|
||||
path := filepath.Join(t.TempDir(), "cap.db")
|
||||
db, err := openAt(context.Background(), path)
|
||||
if err != nil {
|
||||
t.Fatalf("openAt: %v", err)
|
||||
}
|
||||
db.SetMaxOpenConns(cap)
|
||||
s := &Store{db: db}
|
||||
t.Cleanup(func() { _ = s.Close() })
|
||||
return s
|
||||
}
|
||||
|
||||
func percentile(d []time.Duration, p float64) time.Duration {
|
||||
if len(d) == 0 {
|
||||
return 0
|
||||
}
|
||||
i := int(float64(len(d)-1) * p)
|
||||
return d[i]
|
||||
}
|
||||
|
||||
// seedFacts writes n facts so a read has rows to decode.
|
||||
func seedFacts(t *testing.T, s *Store, n int) {
|
||||
t.Helper()
|
||||
ctx := context.Background()
|
||||
now := time.Now().UTC()
|
||||
for i := 0; i < n; i++ {
|
||||
key := fmt.Sprintf("seed_%d", i)
|
||||
if _, err := s.SetValue(ctx, KindSelf, key, "tap:test",
|
||||
map[string]int{"ml": i}, now.Add(time.Duration(i)*time.Millisecond)); err != nil {
|
||||
t.Fatalf("seed %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// measureReadsUnderWrites reports read latency percentiles while a writer
|
||||
// writes at a fixed pace. The pace matters: an unpaced writer completes a
|
||||
// different number of writes at each cap, because at a higher cap it competes
|
||||
// with the readers for the write lock instead of taking turns on one
|
||||
// connection. Two runs that did different work cannot be compared.
|
||||
// It runs for a fixed wall-clock window rather than a fixed read count, so the
|
||||
// paced writer does the same work at every cap. Tying the window to a read
|
||||
// count made the faster configuration receive fewer writes.
|
||||
func measureReadsUnderWrites(t *testing.T, s *Store, window, pace time.Duration) []time.Duration {
|
||||
t.Helper()
|
||||
ctx := context.Background()
|
||||
|
||||
var stop atomic.Bool
|
||||
var writes atomic.Int64
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
now := time.Now().UTC()
|
||||
for i := 0; !stop.Load(); i++ {
|
||||
key := fmt.Sprintf("hot_%d", i%16)
|
||||
if _, err := s.SetValue(ctx, KindSelf, key, "tap:test",
|
||||
map[string]int{"n": i}, now.Add(time.Duration(i)*time.Millisecond)); err != nil {
|
||||
t.Errorf("write: %v", err)
|
||||
return
|
||||
}
|
||||
writes.Add(1)
|
||||
time.Sleep(pace)
|
||||
}
|
||||
}()
|
||||
|
||||
var lat []time.Duration
|
||||
deadline := time.Now().Add(window)
|
||||
for time.Now().Before(deadline) {
|
||||
start := time.Now()
|
||||
if _, err := s.RecentFacts(ctx, 50); err != nil {
|
||||
t.Fatalf("RecentFacts: %v", err)
|
||||
}
|
||||
lat = append(lat, time.Since(start))
|
||||
}
|
||||
stop.Store(true)
|
||||
wg.Wait()
|
||||
t.Logf("in %v: %d reads, %d writes", window, len(lat), writes.Load())
|
||||
|
||||
sort.Slice(lat, func(i, j int) bool { return lat[i] < lat[j] })
|
||||
return lat
|
||||
}
|
||||
|
||||
// TestConnCap_ReadLatencyUnderWrites is the V-642 measurement: read latency at
|
||||
// cap 1 against cap 4, same workload, same schema, same driver.
|
||||
func TestConnCap_ReadLatencyUnderWrites(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("measurement harness; runs for seconds")
|
||||
}
|
||||
for _, cap := range []int{1, 4} {
|
||||
t.Run(fmt.Sprintf("cap=%d", cap), func(t *testing.T) {
|
||||
s := openCapped(t, cap)
|
||||
seedFacts(t, s, 500)
|
||||
lat := measureReadsUnderWrites(t, s, 2*time.Second, 2*time.Millisecond)
|
||||
t.Logf("cap=%d reads=%d p50=%v p95=%v max=%v",
|
||||
cap, len(lat), percentile(lat, 0.50), percentile(lat, 0.95), lat[len(lat)-1])
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestConnCap_ReadBlocksBehindOpenSnapshot is the sharper claim: at cap 1 an
|
||||
// open read-only transaction holds the only connection, so an unrelated read
|
||||
// cannot proceed until it commits. This is why the store exposes no way to
|
||||
// begin one — `Store.DB` used to, and was deleted in V-642 with no caller. The
|
||||
// test stays as the reason, so re-adding that seam fails a measurement rather
|
||||
// than shipping a stall.
|
||||
func TestConnCap_ReadBlocksBehindOpenSnapshot(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("measurement harness; waits on a timeout")
|
||||
}
|
||||
for _, cap := range []int{1, 4} {
|
||||
t.Run(fmt.Sprintf("cap=%d", cap), func(t *testing.T) {
|
||||
s := openCapped(t, cap)
|
||||
seedFacts(t, s, 50)
|
||||
|
||||
tx, err := s.db.BeginTx(context.Background(), &sql.TxOptions{ReadOnly: true})
|
||||
if err != nil {
|
||||
t.Fatalf("BeginTx: %v", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||
defer cancel()
|
||||
start := time.Now()
|
||||
_, err = s.RecentFacts(ctx, 10)
|
||||
t.Logf("cap=%d read alongside an open snapshot: waited %v, err=%v",
|
||||
cap, time.Since(start).Round(time.Millisecond), err)
|
||||
})
|
||||
}
|
||||
}
|
||||
+135
-11
@@ -68,6 +68,13 @@ func (m *MemoryStore) Insert(ctx context.Context, id string, vec []float32, meta
|
||||
// Rows under memory.NonRecallPrefix are excluded in SQL. They are speaker
|
||||
// voiceprints sharing this table, and note recall must not rank them; see that
|
||||
// constant for why the previous arrangement only appeared to do this.
|
||||
//
|
||||
// Every row is still scored, because a full scan is what picks the winners.
|
||||
// What the scan does NOT do is pay for a row it is about to discard: the score
|
||||
// is read straight off the stored bytes without materializing a []float32, and
|
||||
// the meta blob is copied and unmarshalled only for a row that has entered the
|
||||
// topK. Losers cost one dot product and nothing else. Ranking is unchanged —
|
||||
// same scores, same order, same ties.
|
||||
func (m *MemoryStore) Search(ctx context.Context, vec []float32, topK int) ([]memory.Result, error) {
|
||||
if topK <= 0 {
|
||||
topK = 10
|
||||
@@ -80,30 +87,127 @@ func (m *MemoryStore) Search(ctx context.Context, vec []float32, topK int) ([]me
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var out []memory.Result
|
||||
// sql.RawBytes hands us the driver's own buffer, valid only until the next
|
||||
// Next(). Nothing here outlives the row except what topK.offer copies on a
|
||||
// survivor, so the three columns cost no allocation per row.
|
||||
var id, blob, metaJSON sql.RawBytes
|
||||
top := newTopK(topK)
|
||||
for rows.Next() {
|
||||
var id, metaJSON string
|
||||
var blob []byte
|
||||
if err := rows.Scan(&id, &blob, &metaJSON); err != nil {
|
||||
return nil, fmt.Errorf("memory: row: %w", err)
|
||||
}
|
||||
meta := map[string]string{}
|
||||
if err := json.Unmarshal([]byte(metaJSON), &meta); err != nil {
|
||||
return nil, fmt.Errorf("memory: unmarshal meta for %q: %w", id, err)
|
||||
}
|
||||
out = append(out, memory.Result{ID: id, Score: dot(vec, decodeVec(blob)), Meta: meta})
|
||||
top.offer(dotBlob(vec, blob), id, metaJSON)
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, fmt.Errorf("memory: rows: %w", err)
|
||||
}
|
||||
|
||||
sort.Slice(out, func(i, j int) bool { return out[i].Score > out[j].Score })
|
||||
if topK < len(out) {
|
||||
out = out[:topK]
|
||||
survivors := top.sorted()
|
||||
out := make([]memory.Result, 0, len(survivors))
|
||||
for _, c := range survivors {
|
||||
meta := map[string]string{}
|
||||
if err := json.Unmarshal(c.meta, &meta); err != nil {
|
||||
return nil, fmt.Errorf("memory: unmarshal meta for %q: %w", c.id, err)
|
||||
}
|
||||
out = append(out, memory.Result{ID: c.id, Score: c.score, Meta: meta})
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// candidate is one row that is currently in the topK: its score, its id, and
|
||||
// its meta blob copied out of the driver's buffer. The copy is the price of
|
||||
// surviving, and only survivors pay it.
|
||||
type candidate struct {
|
||||
score float64
|
||||
id string
|
||||
meta []byte
|
||||
}
|
||||
|
||||
// topK keeps the k highest-scoring candidates seen so far as a min-heap, so the
|
||||
// weakest survivor is always heap[0] and one comparison decides whether a new
|
||||
// row is worth copying. k is 10 in practice, so the heap is tiny and the whole
|
||||
// structure fits in cache.
|
||||
//
|
||||
// It is a plain slice with hand-written sift operations rather than
|
||||
// container/heap, because that interface boxes every element into an `any` on
|
||||
// Push and costs an allocation per surviving row.
|
||||
type topK struct {
|
||||
k int
|
||||
heap []candidate
|
||||
}
|
||||
|
||||
func newTopK(k int) *topK {
|
||||
return &topK{k: k, heap: make([]candidate, 0, k)}
|
||||
}
|
||||
|
||||
// offer admits a row if it beats the weakest survivor, or if the heap is not
|
||||
// full yet. id and meta are the driver's buffers and are copied here, never
|
||||
// retained.
|
||||
//
|
||||
// A row that only ties the weakest survivor does not displace it, so among
|
||||
// equal scores the earliest k rows are kept. The full sort this replaced used
|
||||
// sort.Slice, which is not stable, so it broke such a tie arbitrarily. That is
|
||||
// the ONE observable difference between the two, and it is deliberate:
|
||||
// deterministic beats arbitrary.
|
||||
//
|
||||
// It is not academic. Under the real embedder an exact tie means duplicate
|
||||
// vectors and nothing in the recall eval moved (V-643). Under the hash
|
||||
// embedder the eval's deterministic floor uses, ties are everywhere — it is
|
||||
// bag-of-words, so every note sharing no word with the query scores exactly 0
|
||||
// — and recall@3 on that run moved 74.1% to 81.5% purely because the zeros now
|
||||
// come out in a fixed order. Neither number measures retrieval. recall@1 and
|
||||
// false recall, which the eval actually asserts, are unchanged on both runs.
|
||||
func (t *topK) offer(score float64, id, meta []byte) {
|
||||
if t.k == 0 {
|
||||
return
|
||||
}
|
||||
if len(t.heap) < t.k {
|
||||
t.heap = append(t.heap, candidate{score: score, id: string(id), meta: append([]byte(nil), meta...)})
|
||||
t.up(len(t.heap) - 1)
|
||||
return
|
||||
}
|
||||
if score <= t.heap[0].score {
|
||||
return
|
||||
}
|
||||
t.heap[0] = candidate{score: score, id: string(id), meta: append([]byte(nil), meta...)}
|
||||
t.down(0)
|
||||
}
|
||||
|
||||
func (t *topK) up(i int) {
|
||||
for i > 0 {
|
||||
parent := (i - 1) / 2
|
||||
if t.heap[parent].score <= t.heap[i].score {
|
||||
return
|
||||
}
|
||||
t.heap[parent], t.heap[i] = t.heap[i], t.heap[parent]
|
||||
i = parent
|
||||
}
|
||||
}
|
||||
|
||||
func (t *topK) down(i int) {
|
||||
for {
|
||||
l, r, small := 2*i+1, 2*i+2, i
|
||||
if l < len(t.heap) && t.heap[l].score < t.heap[small].score {
|
||||
small = l
|
||||
}
|
||||
if r < len(t.heap) && t.heap[r].score < t.heap[small].score {
|
||||
small = r
|
||||
}
|
||||
if small == i {
|
||||
return
|
||||
}
|
||||
t.heap[small], t.heap[i] = t.heap[i], t.heap[small]
|
||||
i = small
|
||||
}
|
||||
}
|
||||
|
||||
// sorted drains the heap into descending score order — what Search returns.
|
||||
func (t *topK) sorted() []candidate {
|
||||
out := t.heap
|
||||
sort.Slice(out, func(i, j int) bool { return out[i].score > out[j].score })
|
||||
return out
|
||||
}
|
||||
|
||||
// ByPrefix returns every row whose id starts with prefix, vectors included.
|
||||
//
|
||||
// This is not a similarity query and deliberately does not score anything:
|
||||
@@ -240,6 +344,26 @@ func decodeVec(b []byte) []float32 {
|
||||
return v
|
||||
}
|
||||
|
||||
// dotBlob is dot against a vector still in its stored encoding, so scoring a
|
||||
// row the query is about to discard does not allocate the []float32 that
|
||||
// decodeVec would build. Same arithmetic, same order of operations, so it
|
||||
// returns bit-identical scores to dot(a, decodeVec(b)).
|
||||
//
|
||||
// A blob whose length isn't a multiple of 4 is truncated to the whole-element
|
||||
// prefix, matching decodeVec, and a length mismatch is 0, matching dot.
|
||||
func dotBlob(a []float32, b []byte) float64 {
|
||||
n := len(b) / 4
|
||||
if len(a) != n || n == 0 {
|
||||
return 0
|
||||
}
|
||||
var sum float64
|
||||
for i := 0; i < n; i++ {
|
||||
f := math.Float32frombits(binary.LittleEndian.Uint32(b[4*i:]))
|
||||
sum += float64(a[i]) * float64(f)
|
||||
}
|
||||
return sum
|
||||
}
|
||||
|
||||
// dot is the cosine similarity for L2-normalized vectors (mismatched lengths ⇒
|
||||
// 0, matching internal/memory's cosine).
|
||||
func dot(a, b []float32) float64 {
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math"
|
||||
"math/rand"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// benchDim is the resident embedder's width (multilingual-e5-small, 384), so
|
||||
// the per-row decode cost the benchmark measures is the real one.
|
||||
const benchDim = 384
|
||||
|
||||
// seedMemVectors fills a fresh store with n L2-normalized rows carrying a meta
|
||||
// blob the size recall actually stores — the note text plus its type — because
|
||||
// the cost this benchmark exists to measure is unmarshalling that blob for
|
||||
// every row when only topK survivors need it.
|
||||
func seedMemVectors(tb testing.TB, n int) *MemoryStore {
|
||||
tb.Helper()
|
||||
path := filepath.Join(tb.TempDir(), "mem_bench.db")
|
||||
st, err := Open(context.Background(), path)
|
||||
if err != nil {
|
||||
tb.Fatalf("Open: %v", err)
|
||||
}
|
||||
tb.Cleanup(func() { _ = st.Close() })
|
||||
m := st.VectorMemory()
|
||||
|
||||
rng := rand.New(rand.NewSource(1))
|
||||
ctx := context.Background()
|
||||
for i := 0; i < n; i++ {
|
||||
if err := m.Insert(ctx, fmt.Sprintf("note:%d", i), randUnitVec(rng, benchDim), map[string]string{
|
||||
"type": "note",
|
||||
"text": fmt.Sprintf("заметка номер %d о том, что надо не забыть сделать на неделе", i),
|
||||
}); err != nil {
|
||||
tb.Fatalf("Insert %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
return m
|
||||
}
|
||||
|
||||
func randUnitVec(rng *rand.Rand, dim int) []float32 {
|
||||
v := make([]float32, dim)
|
||||
var norm float64
|
||||
for i := range v {
|
||||
f := rng.NormFloat64()
|
||||
v[i] = float32(f)
|
||||
norm += f * f
|
||||
}
|
||||
norm = math.Sqrt(norm)
|
||||
for i := range v {
|
||||
v[i] = float32(float64(v[i]) / norm)
|
||||
}
|
||||
return v
|
||||
}
|
||||
|
||||
// BenchmarkMemoryStoreSearch measures one recall query against a store of n
|
||||
// rows. Row counts bracket the documented scale: 1000 is a plausible today,
|
||||
// 10000 is the "thousands, not millions" ceiling the type doc claims a full
|
||||
// scan is fine at.
|
||||
func BenchmarkMemoryStoreSearch(b *testing.B) {
|
||||
for _, n := range []int{1000, 10000} {
|
||||
b.Run(fmt.Sprintf("rows=%d", n), func(b *testing.B) {
|
||||
m := seedMemVectors(b, n)
|
||||
q := randUnitVec(rand.New(rand.NewSource(2)), benchDim)
|
||||
ctx := context.Background()
|
||||
b.ReportAllocs()
|
||||
b.ResetTimer()
|
||||
for i := 0; i < b.N; i++ {
|
||||
if _, err := m.Search(ctx, q, 10); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"sort"
|
||||
"testing"
|
||||
|
||||
"github.com/kami/maven/internal/memory"
|
||||
)
|
||||
|
||||
// naiveSearch is the implementation Search replaced: score every row into a
|
||||
// slice, sort the whole slice, truncate. It stays in the test file as the
|
||||
// reference the bounded-heap version is judged against, because "recall must
|
||||
// not change" is a claim about output, not about the code that produces it.
|
||||
func naiveSearch(t *testing.T, m *MemoryStore, vec []float32, topK int) []memory.Result {
|
||||
t.Helper()
|
||||
rows, err := m.db.QueryContext(context.Background(),
|
||||
`SELECT id, vec FROM memory_vectors WHERE id NOT LIKE ? ESCAPE '\'`,
|
||||
escapeLike(memory.NonRecallPrefix)+"%")
|
||||
if err != nil {
|
||||
t.Fatalf("naive scan: %v", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
var out []memory.Result
|
||||
for rows.Next() {
|
||||
var id string
|
||||
var blob []byte
|
||||
if err := rows.Scan(&id, &blob); err != nil {
|
||||
t.Fatalf("naive row: %v", err)
|
||||
}
|
||||
out = append(out, memory.Result{ID: id, Score: dot(vec, decodeVec(blob))})
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
t.Fatalf("naive rows: %v", err)
|
||||
}
|
||||
sort.Slice(out, func(i, j int) bool { return out[i].Score > out[j].Score })
|
||||
if topK < len(out) {
|
||||
out = out[:topK]
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// TestMemoryStoreSearchMatchesNaive is the constraint on V-643: the bounded
|
||||
// heap must return exactly what a full scan and sort returned. Distinct random
|
||||
// vectors, so no two scores tie and the ranking is total — a mismatch here is
|
||||
// arithmetic or heap logic, not a tie-break difference.
|
||||
func TestMemoryStoreSearchMatchesNaive(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
m := newMemTestStore(t).VectorMemory()
|
||||
|
||||
rng := rand.New(rand.NewSource(7))
|
||||
const rows, dim = 500, 64
|
||||
for i := 0; i < rows; i++ {
|
||||
if err := m.Insert(ctx, fmt.Sprintf("n%d", i), randUnitVec(rng, dim), map[string]string{
|
||||
"text": fmt.Sprintf("note %d", i),
|
||||
}); err != nil {
|
||||
t.Fatalf("Insert %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
|
||||
for _, topK := range []int{1, 3, 10, 50, rows, rows + 100} {
|
||||
q := randUnitVec(rng, dim)
|
||||
got, err := m.Search(ctx, q, topK)
|
||||
if err != nil {
|
||||
t.Fatalf("Search topK=%d: %v", topK, err)
|
||||
}
|
||||
want := naiveSearch(t, m, q, topK)
|
||||
if len(got) != len(want) {
|
||||
t.Fatalf("topK=%d: got %d results, naive returned %d", topK, len(got), len(want))
|
||||
}
|
||||
for i := range want {
|
||||
if got[i].ID != want[i].ID {
|
||||
t.Errorf("topK=%d rank %d: got %q, naive says %q", topK, i, got[i].ID, want[i].ID)
|
||||
}
|
||||
if got[i].Score != want[i].Score {
|
||||
t.Errorf("topK=%d rank %d (%s): score %v, naive says %v",
|
||||
topK, i, got[i].ID, got[i].Score, want[i].Score)
|
||||
}
|
||||
}
|
||||
if len(got) > 0 && got[0].Meta["text"] == "" {
|
||||
t.Errorf("topK=%d: survivor %s has no meta — it was never unmarshalled", topK, got[0].ID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestDotBlobMatchesDot pins the claim in dotBlob's doc comment: reading the
|
||||
// vector out of its stored bytes is bit-identical to decoding it first. Scores
|
||||
// feed a gate with a 0.008 margin, so "close enough" is not the bar.
|
||||
func TestDotBlobMatchesDot(t *testing.T) {
|
||||
rng := rand.New(rand.NewSource(11))
|
||||
for i := 0; i < 200; i++ {
|
||||
a := randUnitVec(rng, 384)
|
||||
b := randUnitVec(rng, 384)
|
||||
if got, want := dotBlob(a, encodeVec(b)), dot(a, b); got != want {
|
||||
t.Fatalf("dotBlob = %v, dot = %v", got, want)
|
||||
}
|
||||
}
|
||||
// Length mismatch is 0 in both, and so is an empty vector.
|
||||
if got := dotBlob([]float32{1, 0}, encodeVec([]float32{1, 0, 0})); got != 0 {
|
||||
t.Errorf("mismatched lengths scored %v, want 0", got)
|
||||
}
|
||||
if got := dotBlob(nil, nil); got != 0 {
|
||||
t.Errorf("empty scored %v, want 0", got)
|
||||
}
|
||||
}
|
||||
+14
-8
@@ -95,7 +95,20 @@ func openAt(ctx context.Context, path string) (*sql.DB, error) {
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("open %s: %w", path, err)
|
||||
}
|
||||
// single writer expected; the daemon is the only process touching the db.
|
||||
// One connection, so every statement is serialised at the database and no
|
||||
// caller above needs a lock of its own. internal/ipc's Server relies on
|
||||
// exactly this, which is why the cap is an invariant rather than a tuning
|
||||
// knob: raising it moves the serialisation guarantee somewhere it is not
|
||||
// written down.
|
||||
//
|
||||
// Measured on 07-08-2026 (V-642, docs/evals/2026-08-07-store-connection-cap.md).
|
||||
// WAL exists to let readers run beside one writer, and the cap gives that
|
||||
// up, but reads do not queue: p50 594µs against 525µs at a cap of four,
|
||||
// while write throughput more than halves. The one thing the cap cannot
|
||||
// survive is a long-lived transaction, which holds the only connection and
|
||||
// stalls every read for its lifetime. So the store begins none, and
|
||||
// TestConnCap_ReadBlocksBehindOpenSnapshot is the standing measurement of
|
||||
// what re-adding one would cost.
|
||||
db.SetMaxOpenConns(1)
|
||||
if _, err := db.ExecContext(ctx, schemaSQL); err != nil {
|
||||
if closeErr := db.Close(); closeErr != nil {
|
||||
@@ -133,13 +146,6 @@ func (s *Store) Close() error {
|
||||
return s.enc.closeAndSeal(s.db)
|
||||
}
|
||||
|
||||
// DB exposes the underlying handle for internal read-only snapshots.
|
||||
// Used by the loop to take a consistent read under a single transaction.
|
||||
// Modules never receive this handle — core mediates.
|
||||
func (s *Store) DB(ctx context.Context) (*sql.Tx, error) {
|
||||
return s.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true})
|
||||
}
|
||||
|
||||
var (
|
||||
// ErrNoFact — no non-voided row exists for this key.
|
||||
ErrNoFact = errors.New("store: no fact for key")
|
||||
|
||||
@@ -328,7 +328,7 @@ func TestExecEmptyCmdRefuses(t *testing.T) {
|
||||
func TestExecRefusesATargetTheSystemCannotHave(t *testing.T) {
|
||||
api := fakeAPI{tools: map[string]ipc.Tool{
|
||||
"restart": {Name: "restart", Cmd: []string{"systemctl", "restart"}, Status: "enabled"},
|
||||
"drop": {Name: "drop", Cmd: []string{"dropdb"}, Destructive: true, Status: "enabled"},
|
||||
"drop": {Name: "drop", Cmd: []string{"dropdb"}, Destructive: true, Status: "enabled"},
|
||||
}}
|
||||
ran := false
|
||||
e := NewExecutor(api, 0)
|
||||
|
||||
@@ -26,6 +26,8 @@
|
||||
package voice
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/kami/maven/internal/phraser"
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
@@ -38,8 +40,10 @@ import (
|
||||
// decision's Intent + Slots + Clarify. The Intent largely names the reply
|
||||
// shape (act/reminder/fact/note/query/clarify); the Slots carry the
|
||||
// specifics that personalise it ("got it: water at 14:00").
|
||||
// The context is the turn's, and it is the only bound an LLM-backed impl has
|
||||
// besides the phraser timeout (V-638). A floor impl ignores it.
|
||||
type Replier interface {
|
||||
Reply(d router.Decision) string
|
||||
Reply(ctx context.Context, d router.Decision) string
|
||||
}
|
||||
|
||||
// StubReplier — the deterministic, no-model floor. Canned per intent;
|
||||
@@ -54,7 +58,8 @@ func NewStubReplier() *StubReplier { return &StubReplier{} }
|
||||
|
||||
// Reply dispatches on Intent + Clarify. Each branch is short; the LLM impl
|
||||
// will replace this with prompted text and the same dispatch shape.
|
||||
func (s *StubReplier) Reply(d router.Decision) string {
|
||||
// It makes no model call, so the context is unused.
|
||||
func (s *StubReplier) Reply(_ context.Context, d router.Decision) string {
|
||||
if d.Clarify {
|
||||
return "не совсем поняла — можешь переформулировать?"
|
||||
}
|
||||
|
||||
@@ -295,6 +295,27 @@ func (f *Fetcher) checkURL(u *url.URL) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// pruneHostsAbove is when pruneLocked bothers to walk the map. Below it the
|
||||
// walk costs more than the entries do, and `crawl.on_demand` means the host set
|
||||
// is whatever he names out loud, so it grows slowly.
|
||||
const pruneHostsAbove = 64
|
||||
|
||||
// pruneLocked drops hosts whose last dial is further back than HostInterval.
|
||||
// Such an entry cannot delay anything — waitTurn would let the next request
|
||||
// through immediately — so keeping it only holds memory for the life of the
|
||||
// process. Caller holds f.mu.
|
||||
func (f *Fetcher) pruneLocked(now time.Time) {
|
||||
if len(f.last) <= pruneHostsAbove {
|
||||
return
|
||||
}
|
||||
cutoff := now.Add(-f.cfg.HostInterval)
|
||||
for h, at := range f.last {
|
||||
if at.Before(cutoff) {
|
||||
delete(f.last, h)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// waitTurn blocks until this host's rate-limit interval has elapsed. It holds
|
||||
// no lock while sleeping, so two hosts never wait on each other.
|
||||
func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
|
||||
@@ -304,6 +325,7 @@ func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
|
||||
earliest := f.last[host].Add(f.cfg.HostInterval)
|
||||
if !now.Before(earliest) {
|
||||
f.last[host] = now
|
||||
f.pruneLocked(now)
|
||||
f.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
@@ -318,3 +319,37 @@ func TestPostObeysDenylist(t *testing.T) {
|
||||
t.Fatalf("error = %v, want ErrBlocked", err)
|
||||
}
|
||||
}
|
||||
|
||||
// f.last used to hold one entry per host ever dialed, for the life of the
|
||||
// process. A host whose last dial is older than HostInterval cannot delay
|
||||
// anything, so it is dropped once the map is worth walking.
|
||||
func TestHostRateMapIsPruned(t *testing.T) {
|
||||
f := New(Config{HostInterval: time.Minute, AllowPrivate: true})
|
||||
stale := time.Now().Add(-time.Hour)
|
||||
for i := 0; i < pruneHostsAbove*2; i++ {
|
||||
f.last[fmt.Sprintf("h%d.example", i)] = stale
|
||||
}
|
||||
|
||||
// One real turn is what triggers the sweep.
|
||||
if err := f.waitTurn(context.Background(), "fresh.example"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(f.last) != 1 {
|
||||
t.Fatalf("len(f.last) = %d after the sweep, want 1 (only the host just dialed)", len(f.last))
|
||||
}
|
||||
if _, ok := f.last["fresh.example"]; !ok {
|
||||
t.Fatal("the host just dialed was pruned, so its own rate limit is lost")
|
||||
}
|
||||
|
||||
// A host inside the interval is kept: pruning must not hand out a free turn.
|
||||
f.last["recent.example"] = time.Now()
|
||||
for i := 0; i < pruneHostsAbove*2; i++ {
|
||||
f.last[fmt.Sprintf("g%d.example", i)] = stale
|
||||
}
|
||||
if err := f.waitTurn(context.Background(), "other.example"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, ok := f.last["recent.example"]; !ok {
|
||||
t.Fatal("a host dialed inside HostInterval was pruned")
|
||||
}
|
||||
}
|
||||
|
||||
Executable
+143
@@ -0,0 +1,143 @@
|
||||
#!/usr/bin/env bash
|
||||
# audit.sh — the repo inventory, in one command.
|
||||
#
|
||||
# Every section here was reconstructed by hand, from scratch, in session after
|
||||
# session: 93 grep sweeps across the four longest ones before a single edit was
|
||||
# made. The answers move slowly and the sweeps did not, so they are written down
|
||||
# once here instead.
|
||||
#
|
||||
# It prints and never writes. A committed inventory file goes stale silently and
|
||||
# then lies; a report you regenerate cannot.
|
||||
#
|
||||
# Read-only. Safe to run at any point, including mid-conflict.
|
||||
#
|
||||
# make audit # everything
|
||||
# make audit SECTION=todo # one section: loc, todo, stubs, docs, tests, gaps
|
||||
|
||||
set -uo pipefail
|
||||
cd "$(dirname "$0")/.." || exit 1
|
||||
|
||||
SECTION="${SECTION:-all}"
|
||||
want() { [ "$SECTION" = all ] || [ "$SECTION" = "$1" ]; }
|
||||
rule() { printf '\n=== %s %s\n' "$1" "$(printf '%.0s=' $(seq 1 $((66 - ${#1}))))"; }
|
||||
|
||||
# git grep over tracked files only. deps/ and models/ are gitignored and huge,
|
||||
# and a plain grep -r walks into both.
|
||||
g() { git grep -nI "$@" 2>/dev/null; }
|
||||
|
||||
printf 'Maven repo inventory @ %s (%s)\n' \
|
||||
"$(git rev-parse --short HEAD 2>/dev/null || echo '?')" \
|
||||
"$(git log -1 --format=%cs 2>/dev/null || echo '?')"
|
||||
|
||||
# --- loc -------------------------------------------------------------------
|
||||
# Non-test Go lines per package. Size is the cheapest proxy for "where does the
|
||||
# complexity actually sit", and it is the first thing every audit asked for.
|
||||
if want loc; then
|
||||
rule "PACKAGES BY LOC (non-test)"
|
||||
for d in $(find ./cmd ./internal ./pkg -maxdepth 2 -type d 2>/dev/null | sort); do
|
||||
files=$(find "$d" -maxdepth 1 -name '*.go' ! -name '*_test.go' 2>/dev/null)
|
||||
[ -z "$files" ] && continue
|
||||
n=$(printf '%s\n' "$files" | wc -l)
|
||||
l=$(printf '%s\0' $files | xargs -0 cat 2>/dev/null | wc -l)
|
||||
printf '%7d %3d files %s\n' "$l" "$n" "$d"
|
||||
done | sort -rn
|
||||
fi
|
||||
|
||||
# --- todo ------------------------------------------------------------------
|
||||
if want todo; then
|
||||
rule "TODO / FIXME / XXX / HACK / BUG (non-test)"
|
||||
g -E '(^|[^a-zA-Z])(TODO|FIXME|XXX|HACK|BUG:)' -- 'cmd/**/*.go' 'internal/**/*.go' 'pkg/**/*.go' \
|
||||
| grep -v '_test\.go:' | sed 's/^/ /' || echo " none"
|
||||
fi
|
||||
|
||||
# --- stubs -----------------------------------------------------------------
|
||||
# Code only, never doc comments. "not wired" is this repo's design vocabulary
|
||||
# for a nil dependency and appears in ~30 comments that describe working code,
|
||||
# so searching prose here reports the architecture back as a gap. Likewise
|
||||
# internal/ipc/unimplemented.go is skipped whole: the file IS the deliberate
|
||||
# Unimplemented*Server pattern, not 60 missing methods. "placeholder" is not a
|
||||
# term here either -- it names real identifiers (SQL placeholders,
|
||||
# Deck.RequirePlaceholder, tokenPlaceholder) and matched 16 working lines.
|
||||
if want stubs; then
|
||||
rule "STUBS / NOT IMPLEMENTED (code, non-test)"
|
||||
g -iE 'not (yet )?implemented|unimplemented|пока не умею|panic\("TODO' \
|
||||
-- 'cmd/**/*.go' 'internal/**/*.go' 'pkg/**/*.go' \
|
||||
| grep -v '_test\.go:' \
|
||||
| grep -v '^internal/ipc/unimplemented\.go:' \
|
||||
| grep -vE '^[^:]+:[0-9]+:[[:space:]]*//' \
|
||||
| sed 's/^/ /' || echo " none"
|
||||
|
||||
# CoreAPI stub parity. A stub missing from unimplemented.go breaks the build
|
||||
# via `var _ CoreAPI = UnimplementedCoreAPI{}`. A stub left behind after its
|
||||
# method leaves an interface compiles forever and is caught by nothing, so it
|
||||
# is counted here until V-652 turns it into a test.
|
||||
printf '\n CoreAPI stub parity:\n'
|
||||
python3 - <<'PY' 2>/dev/null | sed 's/^/ /' || echo " (skipped: python3 unavailable)"
|
||||
import re
|
||||
src = open('internal/ipc/coreapi.go').read()
|
||||
decl = set()
|
||||
for m in re.finditer(r'type (\w+API) interface \{(.*?)\n\}', src, re.S):
|
||||
decl |= set(re.findall(r'^\t([A-Z]\w*)\(', m.group(2), re.M))
|
||||
stub = set(re.findall(r'func \(UnimplementedCoreAPI\) (\w+)\(',
|
||||
open('internal/ipc/unimplemented.go').read()))
|
||||
print(f"{len(decl)} declared, {len(stub)} stubbed")
|
||||
for name in sorted(stub - decl):
|
||||
print(f"STALE {name} (stubbed, on no interface)")
|
||||
for name in sorted(decl - stub):
|
||||
print(f"MISSING {name} (declared, no stub)")
|
||||
PY
|
||||
fi
|
||||
|
||||
# --- docs ------------------------------------------------------------------
|
||||
# Living tier only. docs/evals/ are dated measurements that are never edited
|
||||
# after the day, and docs/archive/ is dead by definition, so neither can be
|
||||
# stale. Age is against the recorded date, not against HEAD: every commit moves
|
||||
# HEAD, so a sha comparison would mark the whole tier stale every day.
|
||||
if want docs; then
|
||||
rule "LIVING DOCS — Last verified"
|
||||
today=$(date +%s)
|
||||
for f in docs/*.md; do
|
||||
[ -e "$f" ] || continue
|
||||
line=$(grep -m1 -o 'Last verified: *[0-9-]\{8,10\}[^ ]*\( *@ *[0-9a-f]\{7,\}\)\?' "$f" 2>/dev/null)
|
||||
if [ -z "$line" ]; then
|
||||
printf ' %-44s %s\n' "$(basename "$f")" "MISSING"
|
||||
continue
|
||||
fi
|
||||
d=$(printf '%s' "$line" | grep -o '[0-9]\{4\}-[0-9]\{2\}-[0-9]\{2\}' | head -1)
|
||||
age=""
|
||||
if [ -n "$d" ] && when=$(date -d "$d" +%s 2>/dev/null); then
|
||||
days=$(( (today - when) / 86400 ))
|
||||
age="${days}d"
|
||||
[ "$days" -gt 30 ] && age="${days}d <-- STALE"
|
||||
fi
|
||||
printf ' %-44s %-34s %s\n' "$(basename "$f")" "${line#Last verified: }" "$age"
|
||||
done
|
||||
fi
|
||||
|
||||
# --- tests -----------------------------------------------------------------
|
||||
if want tests; then
|
||||
rule "TEST SHAPE"
|
||||
printf ' benchmarks : %s\n' "$(g -c 'func Benchmark' -- '**/*_test.go' | awk -F: '{s+=$2} END{print s+0}')"
|
||||
printf ' fuzz : %s\n' "$(g -c 'func Fuzz' -- '**/*_test.go' | awk -F: '{s+=$2} END{print s+0}')"
|
||||
printf ' table t.Run: %s\n' "$(g -c 't.Run(' -- '**/*_test.go' | awk -F: '{s+=$2} END{print s+0}')"
|
||||
printf ' golden files: %s\n' "$(g -l 'golden' -- '**/*_test.go' | wc -l)"
|
||||
fi
|
||||
|
||||
# --- gaps ------------------------------------------------------------------
|
||||
# A package with production code and no test file at all. Not a verdict — some
|
||||
# are pure wiring — but it is the list worth looking at before adding more.
|
||||
if want gaps; then
|
||||
rule "PACKAGES WITH NO TEST FILE"
|
||||
found=0
|
||||
for d in $(find ./cmd ./internal ./pkg -maxdepth 2 -type d 2>/dev/null | sort); do
|
||||
ls "$d"/*.go >/dev/null 2>&1 || continue
|
||||
ls "$d"/*_test.go >/dev/null 2>&1 && continue
|
||||
n=$(ls "$d"/*.go 2>/dev/null | wc -l)
|
||||
l=$(cat "$d"/*.go 2>/dev/null | wc -l)
|
||||
printf ' %6d lines %2d files %s\n' "$l" "$n" "$d"
|
||||
found=1
|
||||
done
|
||||
[ "$found" = 0 ] && echo " none"
|
||||
fi
|
||||
|
||||
exit 0
|
||||
Reference in New Issue
Block a user