Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 2e6b274bff | |||
| 60e64dd83c | |||
| cbd8077d2c | |||
| 123b9aa961 | |||
| 1607009215 |
@@ -53,32 +53,21 @@ 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`:
|
and libs wired through the Makefile — **do not** call `go build` on them bare, use `make`:
|
||||||
|
|
||||||
```sh
|
```sh
|
||||||
make build # all 11 binaries
|
make build # all 9 binaries
|
||||||
make build-web # single daemon (pure-Go ones: web/waked/poll/caldav build without CGO)
|
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
|
make test # go test -race across ./internal/... ./cmd/... with CGO env set
|
||||||
```
|
```
|
||||||
|
|
||||||
Run one package or one test with `make t`. **Do not hand-write the CGO preamble.**
|
Run a single test (must carry the CGO env for packages that touch STT/TTS/voice):
|
||||||
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
|
```sh
|
||||||
make t PKG=./internal/router/
|
CGO_CFLAGS="-I$(pwd)/deps/include -I$(pwd)/deps/whisper.cpp/ggml/include" \
|
||||||
make t PKG=./cmd/mavend/ RUN=TestSimulator
|
CGO_LDFLAGS="-L$(pwd)/deps/lib -Wl,-rpath,$(pwd)/deps/lib" \
|
||||||
make t PKG=./internal/router/eval/ RUN='TestONNX' V=1 # V=1 for -v, RACE=0 to drop -race
|
LD_LIBRARY_PATH="$(pwd)/deps/lib" \
|
||||||
|
deps/go/go/bin/go test -run TestName ./internal/router/
|
||||||
```
|
```
|
||||||
|
|
||||||
`t` carries `-race`, so a green `make t` cannot turn red under `make test`. It carries
|
Pure-Go packages (`router`, `memory`, `mavweb`, …) run under a plain `go test ./pkg/`.
|
||||||
`-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/`)
|
## The daemons (`cmd/`)
|
||||||
|
|
||||||
@@ -93,37 +82,13 @@ but `make t` works everywhere and is one thing to remember.
|
|||||||
| `mavpoll` | Environment poller: netdata alarms, uptime-kuma, zenmoney, wireguard presence. Writes facts, sends nothing. Telegram is `internal/delivery/telegramsink`, not this. |
|
| `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. |
|
| `mavcaldav` | CalDAV calendar sync. |
|
||||||
| `mavmaild` | Mail reader (IMAP, read-only). Holds the IMAP password; core never sees it. |
|
| `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
|
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
|
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.
|
`deploy/telegram.env`) sets socket paths, model paths, and the phraser/embedder blocks.
|
||||||
|
|
||||||
**`docker-compose.yml` runs five: `mavend`, `mavsttd`, `mavttsd`, `mavweb`, `mavpoll`.**
|
**Seven of the nine run on homesrv. `mavwaked` and `mavenclient` do not, and that is the
|
||||||
Count against compose, not against the table. Four of the nine daemons are absent, and each
|
decision, not an oversight** (Vikunja #463, `docs/plans/17-where-the-voice-loop-runs.md`).
|
||||||
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
|
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.
|
daemon there listens to nobody. They belong on a client machine where the owner is standing.
|
||||||
|
|
||||||
@@ -357,53 +322,6 @@ Adding a rung to the ladder
|
|||||||
in `runTurn` means adding its name to `preRouteLadder` in
|
in `runTurn` means adding its name to `preRouteLadder` in
|
||||||
`cmd/mavend/decisiontrace.go`, or that rung is silently missing from the record.
|
`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
|
## LLM output contract
|
||||||
|
|
||||||
All phrasing paths emit `{"response":"...","mood":"..."}`, with fallback to plain text when
|
All phrasing paths emit `{"response":"...","mood":"..."}`, with fallback to plain text when
|
||||||
@@ -518,11 +436,8 @@ 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")
|
ToolSearch("select:mcp__vikunja__list_tasks,mcp__vikunja__get_task_details,mcp__vikunja__create_task,mcp__vikunja__update_task")
|
||||||
```
|
```
|
||||||
|
|
||||||
**Close a finished task with `done: true` and nothing else** (owner's call, 07-08-2026).
|
`update_task` carrying a `description` resets `done` to false, so closing a task with a
|
||||||
Do not write a completion summary into the description on the way out. It is lost anyway,
|
write-up takes two calls: the description, then `done: true`.
|
||||||
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
|
## 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_MODEL := $(shell pwd)/models/tts/ru_RU-irina-medium.onnx
|
||||||
PIPER_ESPEAK := $(shell pwd)/deps/piper/espeak-ng-data
|
PIPER_ESPEAK := $(shell pwd)/deps/piper/espeak-ng-data
|
||||||
|
|
||||||
.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
|
.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
|
||||||
|
|
||||||
all: build
|
all: build
|
||||||
|
|
||||||
@@ -128,35 +128,6 @@ test: fmt-check vet
|
|||||||
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
|
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
|
||||||
$(GO) test -race -coverprofile=coverage.out ./internal/... ./cmd/...
|
$(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).
|
# 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
|
# 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
|
# prod-representative baseline at the vendored runtime; override it or set it
|
||||||
@@ -218,16 +189,6 @@ eval-models:
|
|||||||
# scores the fixtures against ggml-small and self-skips when the model is
|
# scores the fixtures against ggml-small and self-skips when the model is
|
||||||
# absent, and TestGoldenFixturesAreCanonical, which checks the committed audio
|
# absent, and TestGoldenFixturesAreCanonical, which checks the committed audio
|
||||||
# and the manifest with no model at all.
|
# 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:
|
stt-fixtures:
|
||||||
./scripts/gen-stt-fixtures.sh
|
./scripts/gen-stt-fixtures.sh
|
||||||
|
|
||||||
|
|||||||
+8
-35
@@ -50,10 +50,10 @@ func run(args []string) error {
|
|||||||
socket := fs.String("socket", "", "core IPC socket path (required)")
|
socket := fs.String("socket", "", "core IPC socket path (required)")
|
||||||
url := fs.String("url", "", "CalDAV calendar URL, e.g. http://localhost:5232/kami/personal (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)")
|
user := fs.String("user", "", "CalDAV basic-auth username (required)")
|
||||||
passFile := fs.String("pass-file", "", "file holding the CalDAV basic-auth password (required — never passed as a flag value)")
|
pass := fs.String("pass", "", "CalDAV basic-auth password (required)")
|
||||||
renderURL := fs.String("render-url", "", "CalDAV collection maven publishes her own reminders to; empty disables rendering")
|
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)")
|
renderUser := fs.String("render-user", "", "basic-auth username for -render-url (defaults to -user)")
|
||||||
renderPassFile := fs.String("render-pass-file", "", "file holding the password for -render-url (defaults to -pass-file)")
|
renderPass := fs.String("render-pass", "", "basic-auth password for -render-url (defaults to -pass)")
|
||||||
renderDur := fs.Duration("render-duration", calendar.DefaultReminderDuration, "how long a rendered reminder occupies")
|
renderDur := fs.Duration("render-duration", calendar.DefaultReminderDuration, "how long a rendered reminder occupies")
|
||||||
interval := fs.Duration("interval", 5*time.Minute, "poll cadence")
|
interval := fs.Duration("interval", 5*time.Minute, "poll cadence")
|
||||||
timeout := fs.Duration("timeout", 10*time.Second, "per-request HTTP timeout")
|
timeout := fs.Duration("timeout", 10*time.Second, "per-request HTTP timeout")
|
||||||
@@ -63,22 +63,13 @@ func run(args []string) error {
|
|||||||
if *socket == "" {
|
if *socket == "" {
|
||||||
return fmt.Errorf("-socket is required")
|
return fmt.Errorf("-socket is required")
|
||||||
}
|
}
|
||||||
if *url == "" || *user == "" || *passFile == "" {
|
if *url == "" || *user == "" || *pass == "" {
|
||||||
return fmt.Errorf("-url, -user, -pass-file are required")
|
return fmt.Errorf("-url, -user, -pass are required")
|
||||||
}
|
}
|
||||||
if err := checkRenderTarget([]string{*url}, *renderURL); err != nil {
|
if err := checkRenderTarget([]string{*url}, *renderURL); err != nil {
|
||||||
return err
|
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)
|
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||||
defer stop()
|
defer stop()
|
||||||
|
|
||||||
@@ -94,20 +85,17 @@ func run(args []string) error {
|
|||||||
http: hc,
|
http: hc,
|
||||||
url: strings.TrimRight(*url, "/"),
|
url: strings.TrimRight(*url, "/"),
|
||||||
user: *user,
|
user: *user,
|
||||||
pass: pass,
|
pass: *pass,
|
||||||
}
|
}
|
||||||
|
|
||||||
var rend *renderer
|
var rend *renderer
|
||||||
if *renderURL != "" {
|
if *renderURL != "" {
|
||||||
ru, rp := *renderUser, pass
|
ru, rp := *renderUser, *renderPass
|
||||||
if ru == "" {
|
if ru == "" {
|
||||||
ru = *user
|
ru = *user
|
||||||
}
|
}
|
||||||
if *renderPassFile != "" {
|
if rp == "" {
|
||||||
rp, err = readSecret(*renderPassFile)
|
rp = *pass
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
rend = newRenderer(core, hc, *renderURL, ru, rp, *renderDur)
|
rend = newRenderer(core, hc, *renderURL, ru, rp, *renderDur)
|
||||||
log.Printf("mavcaldav: rendering reminders to %s", *renderURL)
|
log.Printf("mavcaldav: rendering reminders to %s", *renderURL)
|
||||||
@@ -143,21 +131,6 @@ func run(args []string) error {
|
|||||||
// It takes the whole read set, not one URL. The guarantee in the package
|
// 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
|
// comment is about every calendar maven reads, and a second read target added
|
||||||
// later must not quietly fall outside the check.
|
// 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 {
|
func checkRenderTarget(readURLs []string, renderURL string) error {
|
||||||
if renderURL == "" {
|
if renderURL == "" {
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -5,38 +5,12 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"os"
|
|
||||||
"path/filepath"
|
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/kami/maven/internal/ipc"
|
"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 {
|
type fakeCore struct {
|
||||||
ipc.UnimplementedCoreAPI
|
ipc.UnimplementedCoreAPI
|
||||||
facts map[string]ipc.Fact // composite key "key|source" → Fact
|
facts map[string]ipc.Fact // composite key "key|source" → Fact
|
||||||
|
|||||||
+23
-85
@@ -59,26 +59,6 @@ type querySource struct {
|
|||||||
// sources search text with no notion of a day. When one of them grows a
|
// sources search text with no notion of a day. When one of them grows a
|
||||||
// date parameter, flip its flag here.
|
// date parameter, flip its flag here.
|
||||||
dateAware bool
|
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
|
// querySources is the ordered chain actionQuery walks; first source to claim
|
||||||
@@ -87,85 +67,85 @@ type querySource struct {
|
|||||||
// gate was never the bug. Adding a source (Kiwix, RSS, crawler, email) is one
|
// 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.
|
// line here plus its method; where you put the line is the whole decision.
|
||||||
var querySources = []querySource{
|
var querySources = []querySource{
|
||||||
{name: "fact-by-key", answer: (*reactiveHandler).queryFactByKey, dest: router.SourceRecall},
|
{name: "fact-by-key", answer: (*reactiveHandler).queryFactByKey},
|
||||||
// Before "calendar" on purpose: both match "…на сегодня", and the plan is
|
// Before "calendar" on purpose: both match "…на сегодня", and the plan is
|
||||||
// the more specific ask (its matcher requires a plan word), so the calendar
|
// the more specific ask (its matcher requires a plan word), so the calendar
|
||||||
// listing would otherwise swallow it.
|
// listing would otherwise swallow it.
|
||||||
{name: "day-plan", answer: (*reactiveHandler).queryDayPlan, dest: router.SourceCalendar},
|
{name: "day-plan", answer: (*reactiveHandler).queryDayPlan},
|
||||||
// Also before "calendar": "что я обычно делаю по средам?" names a weekday,
|
// Also before "calendar": "что я обычно делаю по средам?" names a weekday,
|
||||||
// and the habit question is the more specific one. Its matcher requires a
|
// and the habit question is the more specific one. Its matcher requires a
|
||||||
// habit marker ("обычно", "каждый", …), so a question about this coming
|
// habit marker ("обычно", "каждый", …), so a question about this coming
|
||||||
// Wednesday still reaches the calendar.
|
// Wednesday still reaches the calendar.
|
||||||
{name: "habits", answer: (*reactiveHandler).queryHabits, dest: router.SourceCalendar},
|
{name: "habits", answer: (*reactiveHandler).queryHabits},
|
||||||
// Before "calendar" and before the recall sources: "что мне нужно
|
// Before "calendar" and before the recall sources: "что мне нужно
|
||||||
// сделать?" is a question about the task list, and the notes pass would
|
// сделать?" is a question about the task list, and the notes pass would
|
||||||
// otherwise answer it with whatever note happens to be nearest. Its
|
// otherwise answer it with whatever note happens to be nearest. Its
|
||||||
// matcher requires a task noun or an explicit "что … сделать", so a
|
// matcher requires a task noun or an explicit "что … сделать", so a
|
||||||
// date-bearing question still reaches the calendar.
|
// date-bearing question still reaches the calendar.
|
||||||
{name: "tasks", answer: (*reactiveHandler).queryTasks, dest: router.SourceTasks},
|
{name: "tasks", answer: (*reactiveHandler).queryTasks},
|
||||||
// Next to "tasks" and for the same reason: "что требует внимания?" is a
|
// Next to "tasks" and for the same reason: "что требует внимания?" is a
|
||||||
// question about the operational state Praxis holds, and it used to fall
|
// question about the operational state Praxis holds, and it used to fall
|
||||||
// through every source to the web search (Vikunja #475). Its matcher needs
|
// through every source to the web search (Vikunja #475). Its matcher needs
|
||||||
// an attention marker, and it falls through when Praxis is not configured.
|
// an attention marker, and it falls through when Praxis is not configured.
|
||||||
{name: "attention", answer: (*reactiveHandler).queryAttention, dest: router.SourceAttention, guesses: true},
|
{name: "attention", answer: (*reactiveHandler).queryAttention},
|
||||||
// Next to "tasks" and for the same reason: "что мне купить?" is a question
|
// Next to "tasks" and for the same reason: "что мне купить?" is a question
|
||||||
// about the shopping list, and the recall pass would otherwise answer it
|
// 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
|
// from an old note about the shop. Its matcher needs an explicit list
|
||||||
// marker, so "надо бы съездить в магазин" is untouched.
|
// marker, so "надо бы съездить в магазин" is untouched.
|
||||||
{name: "list", answer: (*reactiveHandler).queryList, dest: router.SourceList, guesses: true},
|
{name: "list", answer: (*reactiveHandler).queryList},
|
||||||
// Before the recall sources too: "сколько я потратил?" is a question about
|
// Before the recall sources too: "сколько я потратил?" is a question about
|
||||||
// the money facts the poller wrote, and the notes pass would otherwise
|
// 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
|
// answer it from whatever he once said about spending. Its matcher needs a
|
||||||
// money noun plus an actual ask, so "я потратил весь день" is untouched.
|
// money noun plus an actual ask, so "я потратил весь день" is untouched.
|
||||||
{name: "money", answer: (*reactiveHandler).queryMoney, dest: router.SourceMoney},
|
{name: "money", answer: (*reactiveHandler).queryMoney},
|
||||||
// Also above the recall sources: "что я тебе говорил?" is a question about
|
// Also above the recall sources: "что я тебе говорил?" is a question about
|
||||||
// the facts he tapped in, and the notes pass would answer it with whatever
|
// 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
|
// note is nearest (Vikunja #456). Its matcher needs both halves of a
|
||||||
// history phrase and bails out when he names a topic, so "что я говорил
|
// history phrase and bails out when he names a topic, so "что я говорил
|
||||||
// про сервер" is still recall.
|
// про сервер" is still recall.
|
||||||
{name: "history", answer: (*reactiveHandler).queryHistory, dest: router.SourceRecall},
|
{name: "history", answer: (*reactiveHandler).queryHistory},
|
||||||
// Before the recall sources and before general knowledge: "что нового?" is
|
// Before the recall sources and before general knowledge: "что нового?" is
|
||||||
// a question about the feeds she reads, and general knowledge would answer
|
// 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
|
// it by inventing news. Its matcher needs a feed noun plus an ask, so
|
||||||
// "у меня новая лента в инстаграме" is untouched.
|
// "у меня новая лента в инстаграме" is untouched.
|
||||||
{name: "feeds", answer: (*reactiveHandler).queryFeeds, dest: router.SourceFeeds, guesses: true},
|
{name: "feeds", answer: (*reactiveHandler).queryFeeds},
|
||||||
// Before "calendar" and before the recall sources: "что включено дома?" is
|
// Before "calendar" and before the recall sources: "что включено дома?" is
|
||||||
// a question about the house, and the notes pass would otherwise answer it
|
// 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
|
// 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
|
// marker plus an ask plus a device word, and it bails out on weather
|
||||||
// wording, so "какая температура на улице?" still reaches the weather
|
// wording, so "какая температура на улице?" still reaches the weather
|
||||||
// source.
|
// source.
|
||||||
{name: "home", answer: (*reactiveHandler).queryHome, dest: router.SourceHome, guesses: true},
|
{name: "home", answer: (*reactiveHandler).queryHome},
|
||||||
// Next to "home" and for the same reason: "какие устройства в сети?" is a
|
// Next to "home" and for the same reason: "какие устройства в сети?" is a
|
||||||
// question about the LAN, and the recall pass would otherwise answer it
|
// 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
|
// from an old note about the router. Its matcher needs a network word plus
|
||||||
// an ask plus a device noun, so "интернет не работает" is untouched.
|
// an ask plus a device noun, so "интернет не работает" is untouched.
|
||||||
{name: "network", answer: (*reactiveHandler).queryNetwork, dest: router.SourceNetwork, guesses: true},
|
{name: "network", answer: (*reactiveHandler).queryNetwork},
|
||||||
{name: "calendar", answer: (*reactiveHandler).queryCalendar, dateAware: true, dest: router.SourceCalendar},
|
{name: "calendar", answer: (*reactiveHandler).queryCalendar, dateAware: true},
|
||||||
{name: "weather", answer: (*reactiveHandler).queryWeather, dest: router.SourceWeather, guesses: true},
|
{name: "weather", answer: (*reactiveHandler).queryWeather},
|
||||||
// A question about her, above the three sources that search his own data
|
// A question about her, above the three sources that search his own data
|
||||||
// (Vikunja #555). It has no answer anywhere else: below the boundary
|
// (Vikunja #555). It has no answer anywhere else: below the boundary
|
||||||
// SearXNG answers about somebody else's assistant, and above it his notes
|
// SearXNG answers about somebody else's assistant, and above it his notes
|
||||||
// answer by proximity — "кто ты" came back from a note of his, measured on
|
// 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.
|
// the box, because the recall index has no idea the subject is her.
|
||||||
{name: "self", answer: (*reactiveHandler).querySelf, dest: router.SourceSelf, guesses: true},
|
{name: "self", answer: (*reactiveHandler).querySelf},
|
||||||
{name: "embed", answer: (*reactiveHandler).queryEmbed, dest: router.SourceRecall},
|
{name: "embed", answer: (*reactiveHandler).queryEmbed},
|
||||||
{name: "memory", answer: (*reactiveHandler).queryMemory, dest: router.SourceRecall},
|
{name: "memory", answer: (*reactiveHandler).queryMemory},
|
||||||
{name: "notes", answer: (*reactiveHandler).queryNotes, dest: router.SourceRecall},
|
{name: "notes", answer: (*reactiveHandler).queryNotes},
|
||||||
// THE BOUNDARY. Everything above answers from his own data; everything
|
// THE BOUNDARY. Everything above answers from his own data; everything
|
||||||
// below answers from the world's. A question about him that got this far
|
// 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
|
// 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.
|
// stops the walk rather than let the encyclopedia and the model guess.
|
||||||
{name: "personal", answer: (*reactiveHandler).queryPersonal, dest: router.SourceRecall, guesses: true},
|
{name: "personal", answer: (*reactiveHandler).queryPersonal},
|
||||||
// The world, read live. Owner's ruling of 2026-08-02: a metasearch hit beats
|
// 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
|
// 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
|
// stake by this point — the boundary above already stopped every question
|
||||||
// about him, and only the query string leaves the box.
|
// about him, and only the query string leaves the box.
|
||||||
{name: "search", answer: (*reactiveHandler).querySearch, dest: router.SourceWorld},
|
{name: "search", answer: (*reactiveHandler).querySearch},
|
||||||
// The offline encyclopedia, now the fallback for when the line is down or
|
// 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
|
// 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.
|
// is that it no longer gets first refusal on a world question.
|
||||||
{name: "kiwix", answer: (*reactiveHandler).queryKiwix, dest: router.SourceWorld},
|
{name: "kiwix", answer: (*reactiveHandler).queryKiwix},
|
||||||
// LAST before the model answers from memory, and that position is the whole
|
// 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,
|
// 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
|
// and only then a page he named. The model does NOT come first: it
|
||||||
@@ -173,43 +153,8 @@ var querySources = []querySource{
|
|||||||
// a 1.7B guessing at a page it cannot read is how contents get invented.
|
// 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
|
// This source only claims a turn where he named a URL, so it never competes
|
||||||
// with a local answer.
|
// with a local answer.
|
||||||
{name: "web", answer: (*reactiveHandler).queryWeb, dest: router.SourceWorld},
|
{name: "web", answer: (*reactiveHandler).queryWeb},
|
||||||
{name: "general-knowledge", answer: (*reactiveHandler).queryGeneral, dest: router.SourceWorld},
|
{name: "general-knowledge", answer: (*reactiveHandler).queryGeneral},
|
||||||
}
|
|
||||||
|
|
||||||
// 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 {
|
func (h *reactiveHandler) actionQuery(ctx context.Context, dec router.Decision) string {
|
||||||
@@ -219,14 +164,7 @@ func (h *reactiveHandler) actionQuery(ctx context.Context, dec router.Decision)
|
|||||||
// (V-564). Finish names everyone below the winner.
|
// (V-564). Finish names everyone below the winner.
|
||||||
decision.Expect(ctx, decision.StageQuery, querySourceNames())
|
decision.Expect(ctx, decision.StageQuery, querySourceNames())
|
||||||
rec := decision.From(ctx)
|
rec := decision.From(ctx)
|
||||||
walk, skipped := queryWalk(dec.Source)
|
for _, src := range querySources {
|
||||||
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 {
|
if dec.Continued && !src.dateAware {
|
||||||
rec.Note(decision.Claim{
|
rec.Note(decision.Claim{
|
||||||
Stage: decision.StageQuery, Claimant: src.name, Outcome: decision.NeverAsked,
|
Stage: decision.StageQuery, Claimant: src.name, Outcome: decision.NeverAsked,
|
||||||
|
|||||||
@@ -1,120 +0,0 @@
|
|||||||
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())
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,96 +0,0 @@
|
|||||||
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+22
-42
@@ -6,6 +6,7 @@ import (
|
|||||||
"math/rand"
|
"math/rand"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
"unicode"
|
||||||
"unicode/utf8"
|
"unicode/utf8"
|
||||||
|
|
||||||
"github.com/kami/maven/internal/dialogue"
|
"github.com/kami/maven/internal/dialogue"
|
||||||
@@ -161,13 +162,12 @@ func withNotice(notice, reply string) string {
|
|||||||
// напоминание?" — answer first, then the open question. A question in front of
|
// напоминание?" — answer first, then the open question. A question in front of
|
||||||
// its own answer would read as ignoring what he asked.
|
// its own answer would read as ignoring what he asked.
|
||||||
//
|
//
|
||||||
// Two sentences, not one (V-654). This used to fold the answer's full stop into
|
// A statement's full stop is folded into a comma, so the two acts read as one
|
||||||
// a comma, on the strength of the owner having written it that way once. Spliced
|
// sentence — that is the owner's own punctuation, "в Риме сейчас ..., на какое
|
||||||
// onto a real answer it reads as one run-on thought — "вот что я нашла: вайфай
|
// время поставить напоминание?". An answer that is ITSELF a question keeps its
|
||||||
// пароль лежит в ящике стола, на какое время поставить напоминание?" — and the
|
// mark and the resume starts a new sentence: she sometimes answers a side query
|
||||||
// question disappears into the tail of a sentence about something else. A reply
|
// by asking him to say it again, and "переформулировать?, на какое время" folds
|
||||||
// with no terminator of its own is given one, so the join never depends on how
|
// two questions into one unreadable line.
|
||||||
// the phraser chose to end.
|
|
||||||
//
|
//
|
||||||
// A resume with no answer in front of it is just the question.
|
// A resume with no answer in front of it is just the question.
|
||||||
func withResumed(reply, resumed string) string {
|
func withResumed(reply, resumed string) string {
|
||||||
@@ -178,17 +178,23 @@ func withResumed(reply, resumed string) string {
|
|||||||
if reply == "" {
|
if reply == "" {
|
||||||
return resumed
|
return resumed
|
||||||
}
|
}
|
||||||
if !endsSentence(reply) {
|
if strings.HasSuffix(reply, "?") {
|
||||||
reply += "."
|
return reply + " " + resumed
|
||||||
}
|
}
|
||||||
return reply + " " + resumed
|
if trimmed := strings.TrimRight(reply, ".!"); trimmed != "" {
|
||||||
|
reply = trimmed
|
||||||
|
}
|
||||||
|
return reply + ", " + lowerFirst(resumed)
|
||||||
}
|
}
|
||||||
|
|
||||||
// endsSentence reports whether s already closes itself. The ellipsis counts: a
|
// lowerFirst lowercases the opening rune, so a deck line written as a standalone
|
||||||
// trailing "…" is a deliberate end, and a full stop after it reads as a typo.
|
// sentence reads as the second half of one. Only the first rune: "На какое
|
||||||
func endsSentence(s string) bool {
|
// время" must become "на какое время" and nothing else in it may move.
|
||||||
r, _ := utf8.DecodeLastRuneInString(s)
|
func lowerFirst(s string) string {
|
||||||
return strings.ContainsRune(".!?…", r)
|
for i, r := range s {
|
||||||
|
return string(unicode.ToLower(r)) + s[i+utf8.RuneLen(r):]
|
||||||
|
}
|
||||||
|
return s
|
||||||
}
|
}
|
||||||
|
|
||||||
// missingFor returns the slots a decision still needs, most important first.
|
// missingFor returns the slots a decision still needs, most important first.
|
||||||
@@ -375,14 +381,6 @@ func (h *reactiveHandler) resolveClarifyAnswer(ctx context.Context, text string)
|
|||||||
return "", false
|
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))
|
merged := q.Answer(text, toDialogueSlots(answer))
|
||||||
// Fold a newly answered subject into the raw utterance. Downstream actions
|
// Fold a newly answered subject into the raw utterance. Downstream actions
|
||||||
// phrase from Utterance, not from the text slot — actionReminder stores it
|
// phrase from Utterance, not from the text slot — actionReminder stores it
|
||||||
@@ -469,14 +467,6 @@ func (h *reactiveHandler) noteDropped(ctx context.Context) {
|
|||||||
//
|
//
|
||||||
// A slot with no resumed wording (clarifyResumedFor says so) resumes nothing and
|
// 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.
|
// 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) {
|
func (h *reactiveHandler) noteSuspended(ctx context.Context, q *dialogue.PendingQuestion) {
|
||||||
rt := turnRouteFrom(ctx)
|
rt := turnRouteFrom(ctx)
|
||||||
if rt == nil || len(q.Missing) == 0 {
|
if rt == nil || len(q.Missing) == 0 {
|
||||||
@@ -486,18 +476,11 @@ func (h *reactiveHandler) noteSuspended(ctx context.Context, q *dialogue.Pending
|
|||||||
if !ok {
|
if !ok {
|
||||||
return
|
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()
|
q.Asked = h.now()
|
||||||
h.clarifyStore.Put(dialogueIDOf(ctx), q)
|
h.clarifyStore.Put(dialogueIDOf(ctx), q)
|
||||||
rt.resume = question
|
rt.resume = question
|
||||||
rt.suspended = true
|
rt.suspended = true
|
||||||
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)
|
log.Printf("voice: clarify — is its own request; suspending the question about %s and resuming it in the same reply", q.Missing[0])
|
||||||
}
|
}
|
||||||
|
|
||||||
// foldAnswerIntoUtterance appends an answered subject to the original words,
|
// foldAnswerIntoUtterance appends an answered subject to the original words,
|
||||||
@@ -537,9 +520,6 @@ func (h *reactiveHandler) askRemainingGap(ctx context.Context, q *dialogue.Pendi
|
|||||||
if !ok || !q.CanAsk() {
|
if !ok || !q.CanAsk() {
|
||||||
return "", false
|
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{
|
h.clarifyStore.Put(dialogueIDOf(ctx), &dialogue.PendingQuestion{
|
||||||
Intent: q.Intent,
|
Intent: q.Intent,
|
||||||
Slots: merged,
|
Slots: merged,
|
||||||
|
|||||||
@@ -740,53 +740,3 @@ 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
|
nextTry map[int64]time.Time // fact id → earliest retry
|
||||||
}
|
}
|
||||||
|
|
||||||
// enrichmentScanLimit bounds how deep a single tick walks the pending queue
|
// enrichmentScanLimit bounds how deep a single tick (or status report) walks
|
||||||
// looking for facts whose backoff has elapsed. The queue is
|
// 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
|
// 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
|
// whether or not they are eligible, and one permanently failing fact stalls
|
||||||
// every younger one behind it.
|
// 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
|
// has been down all day must be visible as a backlog, not as facts that
|
||||||
// silently never got tagged.
|
// silently never got tagged.
|
||||||
//
|
//
|
||||||
// All three numbers describe the same set of rows, whatever is still pending
|
// All three numbers describe the same set of rows, the first
|
||||||
// out of the first enrichmentScanLimit facts. Counting Pending over a thousand rows
|
// enrichmentScanLimit pending facts. Counting Pending over a thousand rows
|
||||||
// while counting InBackoff over the twenty that reached the head of a batch
|
// while counting InBackoff over the twenty that reached the head of a batch
|
||||||
// described two different populations under one struct.
|
// described two different populations under one struct.
|
||||||
type enrichmentStatus struct {
|
type enrichmentStatus struct {
|
||||||
@@ -86,22 +86,13 @@ type enrichmentStatus struct {
|
|||||||
Scanned int // rows the other three counts were taken over
|
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 {
|
func (w *factEnrichmentWorker) status(ctx context.Context) enrichmentStatus {
|
||||||
|
var st enrichmentStatus
|
||||||
pending, err := w.store.PendingFactResolutions(ctx, enrichmentScanLimit)
|
pending, err := w.store.PendingFactResolutions(ctx, enrichmentScanLimit)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("factenrichment: status: %v", err)
|
log.Printf("factenrichment: status: %v", err)
|
||||||
return enrichmentStatus{}
|
return st
|
||||||
}
|
}
|
||||||
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.Pending = len(pending)
|
||||||
st.Scanned = len(pending)
|
st.Scanned = len(pending)
|
||||||
w.mu.Lock()
|
w.mu.Lock()
|
||||||
@@ -153,24 +144,17 @@ func (w *factEnrichmentWorker) tick(ctx context.Context) {
|
|||||||
}
|
}
|
||||||
w.forgetDeparted(pending)
|
w.forgetDeparted(pending)
|
||||||
skipped, failed, attempted := 0, 0, 0
|
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 {
|
for _, f := range pending {
|
||||||
if attempted >= w.batch {
|
if attempted >= w.batch {
|
||||||
remaining = append(remaining, f)
|
break
|
||||||
continue
|
|
||||||
}
|
}
|
||||||
if !w.due(f.ID) {
|
if !w.due(f.ID) {
|
||||||
skipped++
|
skipped++
|
||||||
remaining = append(remaining, f)
|
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
attempted++
|
attempted++
|
||||||
if !w.resolveOne(ctx, f) {
|
if !w.resolveOne(ctx, f) {
|
||||||
failed++
|
failed++
|
||||||
remaining = append(remaining, f)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if failed > 0 {
|
if failed > 0 {
|
||||||
@@ -180,7 +164,7 @@ func (w *factEnrichmentWorker) tick(ctx context.Context) {
|
|||||||
// Report the backlog every tick, not only when something failed: the
|
// Report the backlog every tick, not only when something failed: the
|
||||||
// stalled state worth seeing is the one where nothing failed because
|
// stalled state worth seeing is the one where nothing failed because
|
||||||
// nothing was attempted.
|
// nothing was attempted.
|
||||||
if st := w.statusOf(remaining); st.Pending > 0 {
|
if st := w.status(ctx); st.Pending > 0 {
|
||||||
log.Printf("factenrichment: %d facts pending entity resolution, %d in backoff, worst attempt %d (scanned %d)",
|
log.Printf("factenrichment: %d facts pending entity resolution, %d in backoff, worst attempt %d (scanned %d)",
|
||||||
st.Pending, st.InBackoff, st.MaxAttempts, st.Scanned)
|
st.Pending, st.InBackoff, st.MaxAttempts, st.Scanned)
|
||||||
}
|
}
|
||||||
|
|||||||
+112
-24
@@ -252,23 +252,6 @@ func run(args []string) error {
|
|||||||
// envelope per successful intake write.
|
// envelope per successful intake write.
|
||||||
coreFor := func() ipc.CoreAPI { return newIntakeAPI(ipc.NewStoreAPI(st), evBus, time.Now) }
|
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 {
|
if !locked {
|
||||||
rules = wireRules(cfg)
|
rules = wireRules(cfg)
|
||||||
gatherer = wireGatherer(st, cfg, rules)
|
gatherer = wireGatherer(st, cfg, rules)
|
||||||
@@ -301,7 +284,26 @@ func run(args []string) error {
|
|||||||
feedWkr = newFeedWorker(coreFor(), embedderOf(voiceW), cfg)
|
feedWkr = newFeedWorker(coreFor(), embedderOf(voiceW), cfg)
|
||||||
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
|
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
|
||||||
|
|
||||||
coreAPI = newDaemonAPI(depsNow())
|
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
|
||||||
|
}
|
||||||
} else {
|
} else {
|
||||||
// locked mode: no real store yet, so there's no meaningful CoreAPI to
|
// 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
|
// serve. srv.Check below is the actual guard — every CoreAPI call is
|
||||||
@@ -495,7 +497,19 @@ func run(args []string) error {
|
|||||||
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
|
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
|
||||||
|
|
||||||
// Swap the CoreAPI from the locked placeholder to the real store adapter.
|
// Swap the CoreAPI from the locked placeholder to the real store adapter.
|
||||||
newAPI := newDaemonAPI(depsNow())
|
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)
|
||||||
|
}
|
||||||
srv.SetAPI(newAPI)
|
srv.SetAPI(newAPI)
|
||||||
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
|
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
|
||||||
wireMailIntake(srv, st, phr, cfg, evBus)
|
wireMailIntake(srv, st, phr, cfg, evBus)
|
||||||
@@ -510,10 +524,59 @@ func run(args []string) error {
|
|||||||
// block, so no wire path takes a voiceprint on a default box.
|
// block, so no wire path takes a voiceprint on a default box.
|
||||||
wireSpeaker(srv, st, cfg)
|
wireSpeaker(srv, st, cfg)
|
||||||
|
|
||||||
// The voice server and every background worker, on the outer wg
|
// Start voice server.
|
||||||
// so shutdown waits for them. This used to be nine bare
|
if voiceW != nil {
|
||||||
// `go func()` calls and a shadowed WaitGroup (V-639).
|
var wg sync.WaitGroup
|
||||||
startBackground(ctx, &wg, depsNow())
|
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)
|
||||||
|
}
|
||||||
|
|
||||||
dl.unlock(st)
|
dl.unlock(st)
|
||||||
log.Printf("mavend: unlocked via passkey assertion")
|
log.Printf("mavend: unlocked via passkey assertion")
|
||||||
@@ -528,8 +591,33 @@ func run(args []string) error {
|
|||||||
})
|
})
|
||||||
log.Printf("mavend: ipc listening on %s", srv.Path())
|
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 {
|
if !locked {
|
||||||
startBackground(ctx, &wg, depsNow())
|
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) })
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
<-ctx.Done()
|
<-ctx.Done()
|
||||||
|
|||||||
@@ -1,115 +0,0 @@
|
|||||||
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
|
|
||||||
}
|
|
||||||
@@ -279,97 +279,3 @@ func TestTheTurnIsRoutedOnce(t *testing.T) {
|
|||||||
t.Fatalf("the pipeline routed again and got something else: %+v vs %+v", second, first)
|
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -401,10 +401,6 @@ func buildRouter(emb router.Embedder, acts router.ActMatcher, threshold float64,
|
|||||||
grammars = append(grammars, router.AgendaQueryGrammars()...)
|
grammars = append(grammars, router.AgendaQueryGrammars()...)
|
||||||
// Same reason as the agenda rules, for the feeds: "что нового в лентах?"
|
// Same reason as the agenda rules, for the feeds: "что нового в лентах?"
|
||||||
// routed system and answered "пока не умею" (Vikunja #474).
|
// 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())
|
grammars = append(grammars, router.FeedQueryGrammar())
|
||||||
// The list side of the same exposure: a phrasing with no possessive in it
|
// The list side of the same exposure: a phrasing with no possessive in it
|
||||||
// ("список дел") routed system and never reached queryTasks (Vikunja #467).
|
// ("список дел") routed system and never reached queryTasks (Vikunja #467).
|
||||||
|
|||||||
+1
-28
@@ -25,25 +25,6 @@
|
|||||||
"llm_nudges": false
|
"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": {
|
"telegram": {
|
||||||
"bot_token": "${TELEGRAM_BOT_TOKEN}",
|
"bot_token": "${TELEGRAM_BOT_TOKEN}",
|
||||||
"chat_id": "${TELEGRAM_CHAT_ID}",
|
"chat_id": "${TELEGRAM_CHAT_ID}",
|
||||||
@@ -56,15 +37,7 @@
|
|||||||
"This needs a matching ufw rule or the container's SYN is dropped:",
|
"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"
|
" 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": [
|
"//workstation": [
|
||||||
|
|||||||
@@ -1,12 +1,5 @@
|
|||||||
# Secrets for mavend's away-channel reaches. The file is still called
|
# Telegram bot token and chat ID for mavend's away-channel reach.
|
||||||
# telegram.env because compose names it that; it holds both reaches now.
|
|
||||||
# Copy this file to deploy/telegram.env and fill in real values.
|
# Copy this file to deploy/telegram.env and fill in real values.
|
||||||
# deploy/telegram.env is gitignored — never commit the real secrets.
|
# deploy/telegram.env is gitignored — never commit the real secrets.
|
||||||
TELEGRAM_BOT_TOKEN=
|
TELEGRAM_BOT_TOKEN=
|
||||||
TELEGRAM_CHAT_ID=
|
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,44 +157,6 @@ services:
|
|||||||
# - maildata:/var/lib/mavmaild
|
# - maildata:/var/lib/mavmaild
|
||||||
# - ./deploy/imap.password:/run/secrets/imap.password:ro
|
# - ./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:
|
volumes:
|
||||||
dbdata:
|
dbdata:
|
||||||
sockets:
|
sockets:
|
||||||
|
|||||||
+1
-34
@@ -1,6 +1,6 @@
|
|||||||
# Maven — Design
|
# Maven — Design
|
||||||
|
|
||||||
*Last verified: 2026-08-07 @ beb093a. Living doc: correct it in place, do not append.*
|
*Last verified: 2026-08-02 @ 7079a24. Living doc: correct it in place, do not append.*
|
||||||
|
|
||||||
> Folded 2026-07-30 from `SPEC.md` (north star, 2026-07-03), `maven.md`
|
> Folded 2026-07-30 from `SPEC.md` (north star, 2026-07-03), `maven.md`
|
||||||
> (consolidated decisions, 2026-06-30) and `ROADMAP.md` (execution plan,
|
> (consolidated decisions, 2026-06-30) and `ROADMAP.md` (execution plan,
|
||||||
@@ -282,39 +282,6 @@ Three reasons, in the order they settle it:
|
|||||||
So the notice stays what it is: the in-process TTL case, where she really did
|
So the notice stays what it is: the in-process TTL case, where she really did
|
||||||
wait and really did let go.
|
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
|
### save-where — the two-memory routing axis
|
||||||
|
|
||||||
One discriminator: **does the loop evaluate a predicate against it?**
|
One discriminator: **does the loop evaluate a predicate against it?**
|
||||||
|
|||||||
@@ -1,69 +0,0 @@
|
|||||||
# 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.
|
|
||||||
@@ -1,321 +0,0 @@
|
|||||||
# 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]
|
|
||||||
|
|
||||||
```
|
|
||||||
|
|
||||||
@@ -1,195 +0,0 @@
|
|||||||
# 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.
|
|
||||||
+2
-12
@@ -1,6 +1,6 @@
|
|||||||
# Start Commands
|
# Start Commands
|
||||||
|
|
||||||
*Last verified: 2026-08-07 @ a4630b9. Living doc: correct it in place, do not append.*
|
*Last verified: 2026-08-02 @ 7079a24. 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`.
|
All commands assume `ROOT=/home/kami/apps/Maven` and the local Go toolchain at `$ROOT/deps/go/go/bin/go`.
|
||||||
|
|
||||||
@@ -44,8 +44,7 @@ Config path: `~/.config/maven/mavend.json`. Full example with all options.
|
|||||||
"repeat_interval": "5m",
|
"repeat_interval": "5m",
|
||||||
"ntfy": {
|
"ntfy": {
|
||||||
"base_url": "https://ntfy.kvmx.ru",
|
"base_url": "https://ntfy.kvmx.ru",
|
||||||
"topic": "maven",
|
"topic": "maven"
|
||||||
"token": "${NTFY_TOKEN}"
|
|
||||||
},
|
},
|
||||||
"phraser": {
|
"phraser": {
|
||||||
"model_path": "/mnt/hdd1/llms/Qwen3-Maven-1.7B-Q8_0.gguf",
|
"model_path": "/mnt/hdd1/llms/Qwen3-Maven-1.7B-Q8_0.gguf",
|
||||||
@@ -67,15 +66,6 @@ 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.
|
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)
|
## mavsttd — STT worker (optional, remote whisper.cpp)
|
||||||
|
|
||||||
Requires `LD_LIBRARY_PATH` to include deps/lib (for libwhisper.so, libggml-vulkan.so).
|
Requires `LD_LIBRARY_PATH` to include deps/lib (for libwhisper.so, libggml-vulkan.so).
|
||||||
|
|||||||
@@ -1,26 +1,9 @@
|
|||||||
# The two boot paths have drifted
|
# The two boot paths have drifted
|
||||||
|
|
||||||
Last verified: 06-08-2026 @ 69d0f5e
|
Last verified: 06-08-2026 @ 06c1cf2
|
||||||
|
|
||||||
V-639. Reads with `docs/operations.md`.
|
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
|
## What is wrong
|
||||||
|
|
||||||
`run()` in `cmd/mavend/main.go` brings the daemon up two ways. A box with a key in the
|
`run()` in `cmd/mavend/main.go` brings the daemon up two ways. A box with a key in the
|
||||||
|
|||||||
@@ -456,30 +456,9 @@ func (c *Config) validate() error {
|
|||||||
if err := c.validateCapture(); err != nil {
|
if err := c.validateCapture(); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := c.validateTelegram(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
return nil
|
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
|
// DBEncryptionKey resolves the at-rest encryption key: DBKeyEnv (if set) wins
|
||||||
// over DBKeyB64. Returns (nil, nil) when neither is set — the caller then opens
|
// 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,
|
// a plaintext store. A configured-but-invalid key is an error (fail closed,
|
||||||
|
|||||||
@@ -466,29 +466,3 @@ func TestNormaliseKeepsExplicitWorkstationHealth(t *testing.T) {
|
|||||||
t.Errorf("Health = %q, want %q", got, want)
|
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,17 +45,4 @@ func TestDeployConfigLoads(t *testing.T) {
|
|||||||
if cfg.Voice.RouterThreshold <= 0 {
|
if cfg.Voice.RouterThreshold <= 0 {
|
||||||
t.Error("router threshold did not get its default")
|
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,22 +94,20 @@ func openTestStore(t *testing.T) *store.Store {
|
|||||||
|
|
||||||
// attemptStatus reads one attempt row back. Returns ok=false when the row is
|
// attemptStatus reads one attempt row back. Returns ok=false when the row is
|
||||||
// gone, which would itself be a broken promise (a dropped attempt).
|
// 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) {
|
func attemptStatus(t *testing.T, st *store.Store, id int64) (status string, completed bool, ok bool) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
attempts, err := st.ListDeliveryAttempts(context.Background(), "", 200)
|
tx, err := st.DB(context.Background())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("ListDeliveryAttempts: %v", err)
|
t.Fatalf("read tx: %v", err)
|
||||||
}
|
}
|
||||||
for _, a := range attempts {
|
defer func() { _ = tx.Rollback() }()
|
||||||
if a.ID == id {
|
var completedTS *int64
|
||||||
return a.Status, a.HasComplete, true
|
err = tx.QueryRowContext(context.Background(),
|
||||||
}
|
`SELECT status, completed_ts FROM delivery_attempts WHERE id = ?`, id).Scan(&status, &completedTS)
|
||||||
|
if err != nil {
|
||||||
|
return "", false, false
|
||||||
}
|
}
|
||||||
return "", false, false
|
return status, completedTS != nil, true
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestCrashBetweenBeginAndCompleteBecomesUnknown — simulate the crash window:
|
// TestCrashBetweenBeginAndCompleteBecomesUnknown — simulate the crash window:
|
||||||
|
|||||||
@@ -7,17 +7,11 @@
|
|||||||
// the relay). the dispatcher already strips detail off away sendables; the
|
// 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.
|
// sink uses the same helper so it can't leak the body on its own either.
|
||||||
//
|
//
|
||||||
// ntfy is a self-hosted server with deny-all auth — ntfy.kvmx.ru as of
|
// ntfy runs locally (docker, 127.0.0.1:8085, deny-all auth). maven publishes
|
||||||
// 07-08-2026, reached directly, not through the socks relay telegram needs.
|
// with a dedicated user (write-only to maven-* topics) — the credential is a
|
||||||
// maven publishes with a write-only token scoped to its own topic; the
|
// delivery-config secret, not a db key; a popped ntfy sink can push spam to
|
||||||
// credential is a delivery-config secret, not a db key. a popped ntfy sink
|
// your phone, nothing else. matches the module key-isolation invariant: the
|
||||||
// can push spam to that one topic, nothing else — it cannot read the topic
|
// sink never holds the sqlcipher key.
|
||||||
// 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
|
package ntfysink
|
||||||
|
|
||||||
import (
|
import (
|
||||||
@@ -37,29 +31,11 @@ import (
|
|||||||
// the credential lives in the daemon's config (or a systemd credential),
|
// the credential lives in the daemon's config (or a systemd credential),
|
||||||
// never in the binary.
|
// never in the binary.
|
||||||
type Config struct {
|
type Config struct {
|
||||||
// BaseURL — the ntfy server, no trailing path. Required.
|
BaseURL string // e.g. http://127.0.0.1:8085 (no trailing path)
|
||||||
BaseURL string `json:"base_url"`
|
Topic string // e.g. maven (all maven notifications land here)
|
||||||
|
Username string // basic auth; empty = anonymous (won't work with deny-all)
|
||||||
// Topic — where maven publishes. Required. All maven notifications land
|
Password string // basic auth
|
||||||
// on this one topic; severity rides the Priority header, not the topic.
|
Timeout time.Duration // per-request; 0 = DefaultTimeout
|
||||||
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
|
const DefaultTimeout = 10 * time.Second
|
||||||
@@ -83,12 +59,6 @@ func New(cfg Config) (*Sink, error) {
|
|||||||
if cfg.Topic == "" {
|
if cfg.Topic == "" {
|
||||||
return nil, fmt.Errorf("ntfysink: Topic is required")
|
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
|
to := cfg.Timeout
|
||||||
if to == 0 {
|
if to == 0 {
|
||||||
to = DefaultTimeout
|
to = DefaultTimeout
|
||||||
@@ -114,9 +84,7 @@ func (s *Sink) Send(ctx context.Context, d delivery.Sendable) error {
|
|||||||
}
|
}
|
||||||
req.Header.Set("Title", "maven")
|
req.Header.Set("Title", "maven")
|
||||||
req.Header.Set("Priority", priorityFor(d).String())
|
req.Header.Set("Priority", priorityFor(d).String())
|
||||||
if s.cfg.Token != "" {
|
if s.cfg.Username != "" {
|
||||||
req.Header.Set("Authorization", "Bearer "+s.cfg.Token)
|
|
||||||
} else if s.cfg.Username != "" {
|
|
||||||
req.SetBasicAuth(s.cfg.Username, s.cfg.Password)
|
req.SetBasicAuth(s.cfg.Username, s.cfg.Password)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -224,37 +224,6 @@ 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) {
|
func TestSendTitleIsMaven(t *testing.T) {
|
||||||
rs := newRecordingServer(t, 200, "")
|
rs := newRecordingServer(t, 200, "")
|
||||||
srv := httptest.NewServer(rs.handler())
|
srv := httptest.NewServer(rs.handler())
|
||||||
|
|||||||
@@ -49,24 +49,6 @@ type Poller struct {
|
|||||||
offset int64
|
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
|
// 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
|
// 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.
|
// required; correct may be nil, and then the reply carries no buttons.
|
||||||
@@ -77,8 +59,12 @@ func NewPoller(s *Sink, turn Turn, correct Correct) (*Poller, error) {
|
|||||||
if turn == nil {
|
if turn == nil {
|
||||||
return nil, errors.New("telegramsink: intake needs a turn handler")
|
return nil, errors.New("telegramsink: intake needs a turn handler")
|
||||||
}
|
}
|
||||||
if err := ValidateIntakeChatID(s.cfg.ChatID); err != nil {
|
// The push half accepts @channelusername as a destination. The intake half
|
||||||
return nil, err
|
// 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.
|
||||||
|
// Refusing here is the difference between a boot error and a dead reach.
|
||||||
|
if strings.HasPrefix(strings.TrimSpace(s.cfg.ChatID), "@") {
|
||||||
|
return nil, fmt.Errorf("telegramsink: intake needs the numeric chat id, not %s", s.cfg.ChatID)
|
||||||
}
|
}
|
||||||
// The sink's transport already carries the relay. Only the timeout differs,
|
// The sink's transport already carries the relay. Only the timeout differs,
|
||||||
// and it has to clear the long poll.
|
// and it has to clear the long poll.
|
||||||
@@ -178,10 +164,9 @@ func (p *Poller) onMessage(ctx context.Context, m *message) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// onCallback handles a tap on a correction button. Every path from the owner
|
// onCallback handles a tap on a correction button. Every path answers the
|
||||||
// answers the callback: telegram spins a clock on the button until it is
|
// callback: telegram spins a clock on the button until it is answered, and an
|
||||||
// answered, and an unanswered tap reads as a gesture that was dropped. A tap
|
// unanswered tap reads as a gesture that was dropped.
|
||||||
// from anyone else gets silence, the same as a message from a stranger.
|
|
||||||
func (p *Poller) onCallback(ctx context.Context, cb *callbackQuery) {
|
func (p *Poller) onCallback(ctx context.Context, cb *callbackQuery) {
|
||||||
if !p.fromOwner(cb.Message.Chat.idString()) {
|
if !p.fromOwner(cb.Message.Chat.idString()) {
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -199,18 +199,3 @@ func TestNewPollerNeedsATurn(t *testing.T) {
|
|||||||
t.Error("built a poller with no sink to answer through")
|
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -41,33 +41,8 @@ type PendingQuestion struct {
|
|||||||
Attempts int // questions already asked
|
Attempts int // questions already asked
|
||||||
// MaxAttempts caps Attempts. 0 ⇒ DefaultMaxAttempts.
|
// MaxAttempts caps Attempts. 0 ⇒ DefaultMaxAttempts.
|
||||||
MaxAttempts int
|
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
|
// 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
|
// (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
|
// copy of the truth, so a caller that fills them the old way cannot end up with
|
||||||
|
|||||||
@@ -168,25 +168,3 @@ func TestServerCloseCancelsADispatchInFlight(t *testing.T) {
|
|||||||
t.Fatal("Close did not cancel the dispatch")
|
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
+6
-20
@@ -117,19 +117,15 @@ func Dial(path string) (*Client, error) {
|
|||||||
// Close closes the connection out from under a call in flight, on purpose: a
|
// 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
|
// 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.
|
// 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 {
|
func (c *Client) Close() error {
|
||||||
c.connMu.Lock()
|
c.connMu.Lock()
|
||||||
conn := c.conn
|
defer c.connMu.Unlock()
|
||||||
c.conn = nil
|
if c.conn == nil {
|
||||||
c.connMu.Unlock()
|
|
||||||
if conn == nil {
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
return conn.Close()
|
err := c.conn.Close()
|
||||||
|
c.conn = nil
|
||||||
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
// DialWait is Dial with patience: it retries with capped backoff until the
|
// DialWait is Dial with patience: it retries with capped backoff until the
|
||||||
@@ -255,18 +251,8 @@ func (c *Client) roundtrip(ctx context.Context, m Method, raw json.RawMessage, r
|
|||||||
}
|
}
|
||||||
defer conn.SetDeadline(time.Time{})
|
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{})
|
done := make(chan struct{})
|
||||||
defer func() {
|
defer close(done)
|
||||||
close(done)
|
|
||||||
if ctx.Err() != nil {
|
|
||||||
c.drop()
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
go func() {
|
go func() {
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
|
|||||||
@@ -17,11 +17,9 @@ import (
|
|||||||
// Server — the core side of the boundary. Listens on a unix domain socket,
|
// Server — the core side of the boundary. Listens on a unix domain socket,
|
||||||
// accepts module connections, frames requests to a CoreAPI and responses back.
|
// accepts module connections, frames requests to a CoreAPI and responses back.
|
||||||
// One Server per daemon process; concurrent connections are handled in their
|
// One Server per daemon process; concurrent connections are handled in their
|
||||||
// own goroutine but share the single CoreAPI, and so the single store writer.
|
// own goroutine but share the single CoreAPI (and therefore the single store
|
||||||
// The store opens at SetMaxOpenConns(1), so serialisation is already guaranteed
|
// writer — store is single-connection, SetMaxOpenConns(1), so serialization is
|
||||||
// at the database and the Server adds no locking of its own. That cap is an
|
// already guaranteed at the db; the Server adds no locking of its own).
|
||||||
// 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 {
|
type Server struct {
|
||||||
api atomic.Value // stores CoreAPI
|
api atomic.Value // stores CoreAPI
|
||||||
path string
|
path string
|
||||||
|
|||||||
@@ -105,12 +105,6 @@ type Decision struct {
|
|||||||
Slots Slots
|
Slots Slots
|
||||||
Clarify bool // stage 3: below threshold — ask, don't guess
|
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
|
// Continued — this decision was rebuilt from the previous turn rather
|
||||||
// than routed, because the utterance was an ellipsis ("а завтра?").
|
// than routed, because the utterance was an ellipsis ("а завтра?").
|
||||||
// Handlers use it to know that Slots.Text is the PREVIOUS turn's topic
|
// Handlers use it to know that Slots.Text is the PREVIOUS turn's topic
|
||||||
|
|||||||
@@ -1,78 +0,0 @@
|
|||||||
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,11 +187,9 @@ func SystemTimeDateGrammars() []Grammar {
|
|||||||
// written ("the clock/date system rule must not swallow it"); the daemon
|
// written ("the clock/date system rule must not swallow it"); the daemon
|
||||||
// disagreed with the fixture and the daemon was wrong.
|
// disagreed with the fixture and the daemon was wrong.
|
||||||
//
|
//
|
||||||
// Routing, not answering. Two of the five also name the calendar as the
|
// Routing, not answering. These set the intent and nothing else — which source
|
||||||
// destination (V-655), which narrows who may GUESS their way onto the turn and
|
// in the query chain claims the turn stays the chain's decision, and a
|
||||||
// claims nothing. Every source that looks something up still runs, in the order
|
// question with no date still falls through queryCalendar to recall.
|
||||||
// it always did, so a question with no date still falls through queryCalendar
|
|
||||||
// to recall.
|
|
||||||
//
|
//
|
||||||
// Deliberately not folded into SystemTimeDateGrammars: those exist to send
|
// Deliberately not folded into SystemTimeDateGrammars: those exist to send
|
||||||
// utterances TO system, these exist to keep utterances OUT of it, and one
|
// utterances TO system, these exist to keep utterances OUT of it, and one
|
||||||
@@ -201,15 +199,9 @@ func AgendaQueryGrammars() []Grammar {
|
|||||||
{
|
{
|
||||||
// An explicit calendar noun is unambiguous wherever it appears:
|
// 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",
|
Name: "calendar-query",
|
||||||
Pattern: regexp.MustCompile(`(?i)(календар|расписани|повестк)`),
|
Pattern: regexp.MustCompile(`(?i)(календар|расписани|повестк)`),
|
||||||
Build: queryTo(SourceCalendar),
|
Build: agendaQueryBuild,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
// The agenda phrasing with no calendar noun. Anchored at the start
|
// The agenda phrasing with no calendar noun. Anchored at the start
|
||||||
@@ -259,11 +251,9 @@ func AgendaQueryGrammars() []Grammar {
|
|||||||
// "во сколько созвон". He is asking when something on his calendar
|
// "во сколько созвон". He is asking when something on his calendar
|
||||||
// happens, and the noun is the only signal. Closed list, so "когда
|
// happens, and the noun is the only signal. Closed list, so "когда
|
||||||
// битва при Ватерлоо" is still a world question.
|
// битва при Ватерлоо" 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",
|
Name: "event-time-query",
|
||||||
Pattern: regexp.MustCompile(`(?i)^\s*(когда|во\s+сколько|в\s+котором\s+часу)\s+(будет\s+|у\s+нас\s+)?(планёрк|планерк|встреч|созвон|митинг|совещани|звонок|созвон|приём|прием|интервью|собеседовани|тренировк|урок|занятие|пара)[а-я]*(\s|[?!.]|$)`),
|
Pattern: regexp.MustCompile(`(?i)^\s*(когда|во\s+сколько|в\s+котором\s+часу)\s+(будет\s+|у\s+нас\s+)?(планёрк|планерк|встреч|созвон|митинг|совещани|звонок|созвон|приём|прием|интервью|собеседовани|тренировк|урок|занятие|пара)[а-я]*(\s|[?!.]|$)`),
|
||||||
Build: queryTo(SourceCalendar),
|
Build: agendaQueryBuild,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -350,11 +340,6 @@ func narrativeQueryBuild(m []string) (Decision, bool) {
|
|||||||
Intent: IntentQuery,
|
Intent: IntentQuery,
|
||||||
Confidence: 1.0,
|
Confidence: 1.0,
|
||||||
Slots: Slots{Text: topic},
|
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
|
}, true
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,86 +0,0 @@
|
|||||||
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
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,88 +0,0 @@
|
|||||||
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
|
|
||||||
}
|
|
||||||
+6
-37
@@ -89,8 +89,8 @@ type Poller struct {
|
|||||||
ranker Ranker
|
ranker Ranker
|
||||||
cfg Config
|
cfg Config
|
||||||
nextDue map[string]time.Time
|
nextDue map[string]time.Time
|
||||||
seen map[string]*seenIDs // feed → item IDs, for items with no date
|
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
|
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
|
// 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,
|
feeds: valid, fetch: fetch, notes: notes, marks: marks,
|
||||||
embed: embed, ranker: ranker, cfg: cfg,
|
embed: embed, ranker: ranker, cfg: cfg,
|
||||||
nextDue: map[string]time.Time{},
|
nextDue: map[string]time.Time{},
|
||||||
seen: map[string]*seenIDs{},
|
seen: map[string]map[string]bool{},
|
||||||
polled: map[string]bool{},
|
polled: map[string]bool{},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -283,38 +283,6 @@ func (p *Poller) mark(ctx context.Context, feed string, now time.Time) (time.Tim
|
|||||||
return at, true
|
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
|
// 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
|
// item must be newer than the mark; an undated one is kept once per process by
|
||||||
// ID.
|
// ID.
|
||||||
@@ -340,11 +308,12 @@ func (p *Poller) fresh(f FeedConfig, it Item, mark, now time.Time, resync bool)
|
|||||||
id = it.Title
|
id = it.Title
|
||||||
}
|
}
|
||||||
if p.seen[f.Name] == nil {
|
if p.seen[f.Name] == nil {
|
||||||
p.seen[f.Name] = &seenIDs{}
|
p.seen[f.Name] = map[string]bool{}
|
||||||
}
|
}
|
||||||
if !p.seen[f.Name].add(id) {
|
if p.seen[f.Name][id] {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
p.seen[f.Name][id] = true
|
||||||
return !resync
|
return !resync
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -3,7 +3,6 @@ package rss
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -209,27 +208,3 @@ func TestNoFeedsMeansNoPoller(t *testing.T) {
|
|||||||
t.Fatal("a feed with no name or url is not a configuration")
|
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")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -1,154 +0,0 @@
|
|||||||
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)
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+11
-135
@@ -68,13 +68,6 @@ func (m *MemoryStore) Insert(ctx context.Context, id string, vec []float32, meta
|
|||||||
// Rows under memory.NonRecallPrefix are excluded in SQL. They are speaker
|
// Rows under memory.NonRecallPrefix are excluded in SQL. They are speaker
|
||||||
// voiceprints sharing this table, and note recall must not rank them; see that
|
// voiceprints sharing this table, and note recall must not rank them; see that
|
||||||
// constant for why the previous arrangement only appeared to do this.
|
// 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) {
|
func (m *MemoryStore) Search(ctx context.Context, vec []float32, topK int) ([]memory.Result, error) {
|
||||||
if topK <= 0 {
|
if topK <= 0 {
|
||||||
topK = 10
|
topK = 10
|
||||||
@@ -87,127 +80,30 @@ func (m *MemoryStore) Search(ctx context.Context, vec []float32, topK int) ([]me
|
|||||||
}
|
}
|
||||||
defer rows.Close()
|
defer rows.Close()
|
||||||
|
|
||||||
// sql.RawBytes hands us the driver's own buffer, valid only until the next
|
var out []memory.Result
|
||||||
// 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() {
|
for rows.Next() {
|
||||||
|
var id, metaJSON string
|
||||||
|
var blob []byte
|
||||||
if err := rows.Scan(&id, &blob, &metaJSON); err != nil {
|
if err := rows.Scan(&id, &blob, &metaJSON); err != nil {
|
||||||
return nil, fmt.Errorf("memory: row: %w", err)
|
return nil, fmt.Errorf("memory: row: %w", err)
|
||||||
}
|
}
|
||||||
top.offer(dotBlob(vec, blob), id, metaJSON)
|
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})
|
||||||
}
|
}
|
||||||
if err := rows.Err(); err != nil {
|
if err := rows.Err(); err != nil {
|
||||||
return nil, fmt.Errorf("memory: rows: %w", err)
|
return nil, fmt.Errorf("memory: rows: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
survivors := top.sorted()
|
sort.Slice(out, func(i, j int) bool { return out[i].Score > out[j].Score })
|
||||||
out := make([]memory.Result, 0, len(survivors))
|
if topK < len(out) {
|
||||||
for _, c := range survivors {
|
out = out[:topK]
|
||||||
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
|
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.
|
// ByPrefix returns every row whose id starts with prefix, vectors included.
|
||||||
//
|
//
|
||||||
// This is not a similarity query and deliberately does not score anything:
|
// This is not a similarity query and deliberately does not score anything:
|
||||||
@@ -344,26 +240,6 @@ func decodeVec(b []byte) []float32 {
|
|||||||
return v
|
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 ⇒
|
// dot is the cosine similarity for L2-normalized vectors (mismatched lengths ⇒
|
||||||
// 0, matching internal/memory's cosine).
|
// 0, matching internal/memory's cosine).
|
||||||
func dot(a, b []float32) float64 {
|
func dot(a, b []float32) float64 {
|
||||||
|
|||||||
@@ -1,77 +0,0 @@
|
|||||||
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,107 +0,0 @@
|
|||||||
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+8
-14
@@ -95,20 +95,7 @@ func openAt(ctx context.Context, path string) (*sql.DB, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("open %s: %w", path, err)
|
return nil, fmt.Errorf("open %s: %w", path, err)
|
||||||
}
|
}
|
||||||
// One connection, so every statement is serialised at the database and no
|
// single writer expected; the daemon is the only process touching the db.
|
||||||
// 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)
|
db.SetMaxOpenConns(1)
|
||||||
if _, err := db.ExecContext(ctx, schemaSQL); err != nil {
|
if _, err := db.ExecContext(ctx, schemaSQL); err != nil {
|
||||||
if closeErr := db.Close(); closeErr != nil {
|
if closeErr := db.Close(); closeErr != nil {
|
||||||
@@ -146,6 +133,13 @@ func (s *Store) Close() error {
|
|||||||
return s.enc.closeAndSeal(s.db)
|
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 (
|
var (
|
||||||
// ErrNoFact — no non-voided row exists for this key.
|
// ErrNoFact — no non-voided row exists for this key.
|
||||||
ErrNoFact = errors.New("store: no fact for key")
|
ErrNoFact = errors.New("store: no fact for key")
|
||||||
|
|||||||
@@ -295,27 +295,6 @@ func (f *Fetcher) checkURL(u *url.URL) error {
|
|||||||
return nil
|
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
|
// 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.
|
// no lock while sleeping, so two hosts never wait on each other.
|
||||||
func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
|
func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
|
||||||
@@ -325,7 +304,6 @@ func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
|
|||||||
earliest := f.last[host].Add(f.cfg.HostInterval)
|
earliest := f.last[host].Add(f.cfg.HostInterval)
|
||||||
if !now.Before(earliest) {
|
if !now.Before(earliest) {
|
||||||
f.last[host] = now
|
f.last[host] = now
|
||||||
f.pruneLocked(now)
|
|
||||||
f.mu.Unlock()
|
f.mu.Unlock()
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,7 +4,6 @@ import (
|
|||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
|
||||||
"io"
|
"io"
|
||||||
"net"
|
"net"
|
||||||
"net/http"
|
"net/http"
|
||||||
@@ -319,37 +318,3 @@ func TestPostObeysDenylist(t *testing.T) {
|
|||||||
t.Fatalf("error = %v, want ErrBlocked", err)
|
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")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -1,143 +0,0 @@
|
|||||||
#!/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