Compare commits

...

18 Commits

Author SHA1 Message Date
claude 39d44bb384 Close a Vikunja task with done, and nothing else (V-641)
Owner's call, 07-08-2026. A completion summary written into the
description on the way out is lost anyway, and the durable record is the
commit messages and the merged PR.

Written during the V-641 session and left uncommitted; it rides this
branch rather than being dropped.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01YMNNEkYx1mZFtHNrFk7uqb
2026-08-07 01:48:19 +04:00
claude 65ee0f9c61 Score every row, pay for only the ten that survive (V-643)
Search decoded the vector blob into a []float32 and JSON-unmarshalled the
meta map for every row, then sorted all N and threw away everything past
topK. Meta only ever matters for a survivor, and the sort answered a
question a bounded heap answers cheaper.

The scan still visits every row — that is what picks the winners. What it
no longer does is allocate for a row it is about to discard. dotBlob reads
the vector out of its stored bytes, so scoring costs nothing; a row is
copied and its meta unmarshalled only once it has entered the topK.

At 10000 rows and topK 10: 70.6ms to 26.8ms, 58MB to 17.5MB, 240k allocs
to 60k.

Recall is unchanged where it is measured. recall+onnx scores 22/32 with
recall@1 70.4% and recall@3 85.2%, identical to before.
TestMemoryStoreSearchMatchesNaive pins the ranking against the full-sort
implementation it replaced, and TestDotBlobMatchesDot pins bit-identical
scores, which the 0.008 gate margin demands.

One behaviour did move: ties. sort.Slice is not stable, so equal scores
were ordered arbitrarily; the heap now keeps the earliest. Under the real
embedder an exact tie is a duplicate vector and nothing moved. Under the
hash embedder the eval's floor uses, everything ties at 0 and that run's
recall@3 went 74.1% to 81.5% — a number that measures tie order, not
retrieval. recall@1 and false recall, the two the eval asserts, are
unchanged on both runs.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01YMNNEkYx1mZFtHNrFk7uqb
2026-08-07 01:48:05 +04:00
claude 76938e206d Put a number on the recall scan before changing it (V-643)
MemoryStore.Search is on the per-turn recall path and had no benchmark, so
any claim about its cost was an argument rather than a measurement.

Seeds a store with rows the shape recall actually stores — 384-wide
vectors, the resident embedder's width, and a meta blob carrying the note
text — at 1000 and 10000 rows. 10000 is the ceiling the type doc claims a
full scan is fine at.

Measured as it stands: 5.3ms and 24k allocs at 1000 rows, 70.6ms and 240k
allocs at 10000.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01YMNNEkYx1mZFtHNrFk7uqb
2026-08-07 01:48:05 +04:00
claude 0b3d81ecbf Merge pull request 'Two maps grow for the process lifetime with no eviction' (#194) from task/641-two-maps-grow-for-the-process-lifetime-w into master 2026-08-06 23:33:23 +02:00
claude 4be6852b94 Drop host rate-limit entries that can no longer delay anything (V-641)
webfetch.Fetcher.last held one entry per distinct host the crawler ever
dialed, never pruned. Bounded in practice by how many hosts get crawled, but
crawl.on_demand is true in deploy, so the host set is whatever he names out
loud.

An entry older than HostInterval cannot delay a request — waitTurn would let
the next one straight through — so it is dropped. The sweep runs on write and
only once the map passes 64 entries, below which walking it costs more than
the entries do.

Rate limiting is unchanged: a host dialed inside the interval is kept, which
the test asserts, because pruning one would hand out a free turn.
2026-08-07 01:32:26 +04:00
claude f7b76c572f Bound the undated-item set per feed (V-641)
rss.Poller.seen held every undated item ever seen, one entry per id, for as
long as mavend ran. fresh() added and nothing removed. A feed that ships items
with no <pubDate> grew it forever.

seenIDs is the same set with a bound: the map answers the lookup, a slice
remembers insertion order, and the oldest id falls out past 512. The cap has
to stay above any one feed's front page or an item still listed there would be
written a second time, and a few hundred covers the largest page anyone
publishes. The set only ever had to span one poll window plus the resync
guard, not all of history.

Dedupe behaviour is unchanged. The comment at fresh() explains why the set
does not survive a restart; it never bounded it within one run.
2026-08-07 01:32:15 +04:00
claude 05ddc5c92e Merge pull request 'mavcaldav is built, documented as running, and deployed nowhere' (#193) from task/644-mavcaldav-is-built-documented-as-running into master 2026-08-06 23:22:18 +02:00
claude b55e68f98d Say in compose that the calendar is off, and why (V-644)
mavcaldav was built, in `make build`, listed in CLAUDE.md's daemon table, and
deployed nowhere. Not commented out the way mavmaild is, which at least
records the decision and the enable steps. Built and mentioned nowhere is the
worst of the three states, so this writes the decision down.

The box has no CalDAV account, so the block stays commented. It names what the
absence costs, because both costs are invisible from the daemon table. Agenda
questions route correctly and answer from nothing: stage 0 sends "что у меня
сегодня" to IntentQuery (V-498) and the calendar query source then reads facts
nobody writes. And loop.State.CalendarBusy is fed by those same facts, so the
gate's "do not nag mid-meeting" is permanently false.

CLAUDE.md said the absence was an oversight. It is a decision now.
2026-08-07 01:19:44 +04:00
claude beaa24754c Read the CalDAV password from a file, not from argv (V-644)
mavcaldav took -pass and -render-pass as flag values, so enabling it would
have put his calendar password in `ps` inside the container, in the compose
file, and in shell history. mavpoll and mavmaild both read their secret from
a file for exactly that reason.

readSecret reads once at start, trims, and refuses an empty or missing file.
An empty file is a deployment mistake, not a password, and basic auth would
otherwise send "" and collect a 401 every poll. A rotated password means a
restart, which is cheaper than re-reading the credential every five minutes.

Nothing called the old flags: no compose service, no systemd unit, no test.
So they are replaced rather than kept beside the new ones.
2026-08-07 01:19:33 +04:00
kami aed8cac439 Merge pull request 'The store caps sqlite at one connection under WAL, so every read queues behind every write' (#192) from task/642-the-store-caps-sqlite-at-one-connection into master 2026-08-06 23:04:39 +02:00
claude af4eeceb6a Keep the store's one connection, delete the seam it cannot survive (V-642)
`SetMaxOpenConns(1)` under WAL gives up concurrent reads, and the task
asked whether that costs anything. Measured over a fixed two-second
window, a paced writer against a read loop, three runs per cap:
reads do not queue. Four connections buy 70µs at p50 on a turn that
spends 1.19s in the resident model, and write throughput more than
halves. A 19ms worst case also cannot be the source of the 2.7s router
figure, so that line of enquiry is closed.

What the cap cannot survive is a long-lived transaction. It holds the
only connection, so a second read never completes: two seconds and
`context deadline exceeded`, against 1ms at a cap of four.

`Store.DB` handed out exactly that transaction. It had been there since
the initial commit with no production caller, and its comment described
a loop that never materialised. Its one user was a test helper reading
`delivery_attempts` by raw SQL, which `ListDeliveryAttempts` has covered
since V-390. So the cap stays and the seam goes, and the hazard is gone
by construction rather than by documentation.

`internal/store/conncap_test.go` stays as the standing measurement,
skipped under -short. The comment at the cap and the one in
`internal/ipc/server.go` that leans on it now state the invariant and
cite the numbers.

Measurement: docs/evals/2026-08-07-store-connection-cap.md

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-07 01:01:27 +04:00
kami 7b507dec94 Merge pull request 'factEnrichmentWorker walks the pending queue twice per tick to write one log line' (#191) from task/647-factenrichmentworker-walks-the-pending-q into master 2026-08-06 22:33:54 +02:00
claude 2c0334c4fe Count the enrichment backlog without a second query (V-647)
`tick` read `PendingFactResolutions` at the scan limit, then `status`
read it again with the same limit for one log line. Up to 2000 rows per
tick on a database that serialises reads, to say how long the queue is.

`statusOf` counts over a batch the caller already holds, and the tick
passes it the batch it just read. A resolved fact leaves the queue, so
the loop collects what is still pending rather than reporting the
pre-tick count. `status(ctx)` stays as the querying form, for a caller
outside the tick with no batch in hand.

No behaviour change: the three counts still describe one row set, and
the same facts are attempted per tick.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-07 00:32:54 +04:00
kami 92cbdbfdd3 Merge pull request 'V-637 follow-up: telegram intake has no deploy switch, and the chat-id check cannot fail a boot' (#190) from task/646-v-637-follow-up-telegram-intake-has-no-d into master 2026-08-06 22:18:25 +02:00
claude e78b2d8992 the daemon table, against make build and compose (V-648)
The table listed nine binaries. make build builds eleven, and mavseal and
labelgen exist without targets. The running count said seven on homesrv;
docker-compose.yml runs five.

Adds mavgpud, mavupdate, mavseal and labelgen, and names why each absent daemon
is absent: mavmaild has no mail account, mavwaked and mavenclient belong on
workpc, and mavcaldav is an oversight (V-644).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-07 00:14:42 +04:00
claude 9d58922462 Refuse a telegram intake chat id the poller cannot match (V-646)
The push half accepts an @channelusername and the intake half cannot: an
inbound update names its chat by number, so an @-name matches nothing. The
check lived in NewPoller, which wireTelegramIntake logs and returns from, so a
box configured that way booted clean with a dead intake half and a working push
half. Nothing looked broken from the chat.

ValidateIntakeChatID moves the rule where config validation can reach it, the
same shape validateNetScan uses. It is stricter than the old prefix test: any
non-digit is refused, not just a leading @. An empty token or chat id still
means telegram is not wired, because an unset ${TELEGRAM_*} expands to empty
and that must not fail a box with no bot.

deploy/mavend.json turns intake on. The chat id on this box is numeric.

The onCallback comment claimed every path answers the callback. The fromOwner
early return does not, and silence toward a stranger is correct, so the comment
was what was wrong.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-07 00:14:42 +04:00
claude b5ac48c126 One boot path for the workers and the API (#189) 2026-08-06 21:54:13 +02:00
claude 69d0f5ee78 No deadline survives the turn path, from mavweb down to llama-server (#188)
Co-authored-by: claude <no-reply@agents.claude.kvmx.ru>
Co-committed-by: claude <no-reply@agents.claude.kvmx.ru>
2026-08-06 21:11:42 +02:00
35 changed files with 1562 additions and 213 deletions
+32 -5
View File
@@ -53,7 +53,7 @@ CGO daemons (`mavend`, `mavsttd`, `mavttsd`, `mavenclient`) need the vendored to
and libs wired through the Makefile — **do not** call `go build` on them bare, use `make`:
```sh
make build # all 9 binaries
make build # all 11 binaries
make build-web # single daemon (pure-Go ones: web/waked/poll/caldav build without CGO)
make test # go test -race across ./internal/... ./cmd/... with CGO env set
```
@@ -82,13 +82,37 @@ Pure-Go packages (`router`, `memory`, `mavweb`, …) run under a plain `go test
| `mavpoll` | Environment poller: netdata alarms, uptime-kuma, zenmoney, wireguard presence. Writes facts, sends nothing. Telegram is `internal/delivery/telegramsink`, not this. |
| `mavcaldav` | CalDAV calendar sync. |
| `mavmaild` | Mail reader (IMAP, read-only). Holds the IMAP password; core never sees it. |
| `mavgpud` | GPU supervisor. **Runs on workpc, not homesrv** — own unit, `deploy/mavgpud.service`. Keeps llama-server loaded while the card is free (V-488). Maven never asks it for anything, it reads `/health` through `llm.Pair`. |
| `mavupdate` | Not a daemon. Operator CLI a human runs on the box to deploy a new build. |
Two more binaries have no Makefile target and are built with `go run` or `go build` when
they are needed. Neither is deployed.
| Binary | Role |
|---|---|
| `mavseal` | Recovery tool. Encrypts a live tmpfs working copy back to the ciphertext file when mavend was killed before `defer st.Close()` sealed it. |
| `labelgen` | Runs the stage 0 grammars over utterances and prints JSONL, the training data for the routing heads (V-546). |
Daemons are wired socket-to-socket, not linked. `internal/ipc` is the client/server wire
protocol; the config in `deploy/mavend.json` (with `${VAR}` env expansion from gitignored
`deploy/telegram.env`) sets socket paths, model paths, and the phraser/embedder blocks.
**Seven of the nine run on homesrv. `mavwaked` and `mavenclient` do not, and that is the
decision, not an oversight** (Vikunja #463, `docs/plans/17-where-the-voice-loop-runs.md`).
**`docker-compose.yml` runs five: `mavend`, `mavsttd`, `mavttsd`, `mavweb`, `mavpoll`.**
Count against compose, not against the table. Four of the nine daemons are absent, and each
absence has a different reason.
`mavmaild` and `mavcaldav` are commented out in compose, each with the reason written
beside it: the first needs a mail account, the second a CalDAV account, and this box has
neither. `mavcaldav` used to appear nowhere at all, which was an oversight; it became a
recorded decision on 07-08-2026 (V-644). Two things ride on that absence and the block
names them. Agenda questions route to `IntentQuery` at stage 0 (V-498) and the `calendar`
query source then reads a table nobody writes. And `loop.State.CalendarBusy` is fed by the
same facts, so the gate's "do not nag mid-meeting" is permanently false. Its password is
read from a file (`-pass-file`, and `-render-pass-file` for the render collection), never
taken as a flag value, which is the rule `mavpoll` and `mavmaild` follow too.
**`mavwaked` and `mavenclient` are absent by decision, not oversight** (Vikunja #463,
`docs/plans/17-where-the-voice-loop-runs.md`).
homesrv has a microphone — it is a laptop — but it is in the wrong room, so a wake-word
daemon there listens to nobody. They belong on a client machine where the owner is standing.
@@ -436,8 +460,11 @@ start of a session rather than one lookup per first use:
ToolSearch("select:mcp__vikunja__list_tasks,mcp__vikunja__get_task_details,mcp__vikunja__create_task,mcp__vikunja__update_task")
```
`update_task` carrying a `description` resets `done` to false, so closing a task with a
write-up takes two calls: the description, then `done: true`.
**Close a finished task with `done: true` and nothing else** (owner's call, 07-08-2026).
Do not write a completion summary into the description on the way out. It is lost anyway,
and the durable record is the commit messages and the merged PR. Note that `update_task`
carrying a `description` resets `done` to false, which is why a write-up ever took two
calls.
## Session workflow
+35 -8
View File
@@ -50,10 +50,10 @@ func run(args []string) error {
socket := fs.String("socket", "", "core IPC socket path (required)")
url := fs.String("url", "", "CalDAV calendar URL, e.g. http://localhost:5232/kami/personal (required)")
user := fs.String("user", "", "CalDAV basic-auth username (required)")
pass := fs.String("pass", "", "CalDAV basic-auth password (required)")
passFile := fs.String("pass-file", "", "file holding the CalDAV basic-auth password (required — never passed as a flag value)")
renderURL := fs.String("render-url", "", "CalDAV collection maven publishes her own reminders to; empty disables rendering")
renderUser := fs.String("render-user", "", "basic-auth username for -render-url (defaults to -user)")
renderPass := fs.String("render-pass", "", "basic-auth password for -render-url (defaults to -pass)")
renderPassFile := fs.String("render-pass-file", "", "file holding the password for -render-url (defaults to -pass-file)")
renderDur := fs.Duration("render-duration", calendar.DefaultReminderDuration, "how long a rendered reminder occupies")
interval := fs.Duration("interval", 5*time.Minute, "poll cadence")
timeout := fs.Duration("timeout", 10*time.Second, "per-request HTTP timeout")
@@ -63,13 +63,22 @@ func run(args []string) error {
if *socket == "" {
return fmt.Errorf("-socket is required")
}
if *url == "" || *user == "" || *pass == "" {
return fmt.Errorf("-url, -user, -pass are required")
if *url == "" || *user == "" || *passFile == "" {
return fmt.Errorf("-url, -user, -pass-file are required")
}
if err := checkRenderTarget([]string{*url}, *renderURL); err != nil {
return err
}
// The password is read from a file, never taken as a flag value: an argv
// secret is visible in `ps` to every user on the box and lands in the compose
// file and the shell history. Same rule mavmaild and mavpoll follow. Read
// once at start, so a rotated password means a restart.
pass, err := readSecret(*passFile)
if err != nil {
return err
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
@@ -85,17 +94,20 @@ func run(args []string) error {
http: hc,
url: strings.TrimRight(*url, "/"),
user: *user,
pass: *pass,
pass: pass,
}
var rend *renderer
if *renderURL != "" {
ru, rp := *renderUser, *renderPass
ru, rp := *renderUser, pass
if ru == "" {
ru = *user
}
if rp == "" {
rp = *pass
if *renderPassFile != "" {
rp, err = readSecret(*renderPassFile)
if err != nil {
return err
}
}
rend = newRenderer(core, hc, *renderURL, ru, rp, *renderDur)
log.Printf("mavcaldav: rendering reminders to %s", *renderURL)
@@ -131,6 +143,21 @@ func run(args []string) error {
// It takes the whole read set, not one URL. The guarantee in the package
// comment is about every calendar maven reads, and a second read target added
// later must not quietly fall outside the check.
// readSecret reads one credential from a file and refuses an empty one. An
// empty file is a deployment mistake, not a password, and CalDAV basic auth
// would send it and get a 401 every poll.
func readSecret(path string) (string, error) {
raw, err := os.ReadFile(path)
if err != nil {
return "", fmt.Errorf("read password file: %w", err)
}
secret := strings.TrimSpace(string(raw))
if secret == "" {
return "", fmt.Errorf("password file %s is empty", path)
}
return secret, nil
}
func checkRenderTarget(readURLs []string, renderURL string) error {
if renderURL == "" {
return nil
+26
View File
@@ -5,12 +5,38 @@ import (
"fmt"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"testing"
"time"
"github.com/kami/maven/internal/ipc"
)
// The password comes from a file so it never reaches argv. An empty or missing
// file must fail at start rather than authenticate as "" against his calendar.
func TestReadSecret(t *testing.T) {
dir := t.TempDir()
good := filepath.Join(dir, "ok")
if err := os.WriteFile(good, []byte(" hunter2\n"), 0o600); err != nil {
t.Fatal(err)
}
if got, err := readSecret(good); err != nil || got != "hunter2" {
t.Fatalf("readSecret(good) = %q, %v; want \"hunter2\", nil", got, err)
}
empty := filepath.Join(dir, "empty")
if err := os.WriteFile(empty, []byte("\n \n"), 0o600); err != nil {
t.Fatal(err)
}
if _, err := readSecret(empty); err == nil {
t.Fatal("readSecret(empty) = nil error, want refusal")
}
if _, err := readSecret(filepath.Join(dir, "absent")); err == nil {
t.Fatal("readSecret(absent) = nil error, want refusal")
}
}
type fakeCore struct {
ipc.UnimplementedCoreAPI
facts map[string]ipc.Fact // composite key "key|source" → Fact
+120
View File
@@ -0,0 +1,120 @@
package main
import (
"context"
"errors"
"log"
"net"
"sync"
"time"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/store"
)
// The two boot paths meet here. run() wires the daemon twice: once at boot
// when a key is in the environment, and once inside UnlockFn after a passkey
// assertion, minutes or days later. Listing the same wiring in both places is
// what let them drift — seven workers started untracked on the unlock path and
// two daemonAPI fields were never set there, silently, for as long as anyone
// had been cold-starting (V-639).
//
// So both paths call newDaemonAPI and startBackground and nothing else. A
// field or a worker added later reaches both paths or neither.
// bootDeps is everything the two constructors below read. It is filled from
// the same variables on both paths, by depsNow in run().
type bootDeps struct {
coreFor func() ipc.CoreAPI
tl *tickLoop
evBus *event.Bus
voiceW *voiceWiring
st *store.Store
factWorker *factEnrichmentWorker
evalWorker *memoryEvalWorker // nil ⇒ memory evaluation off (the default)
feedWkr *feedWorker // nil ⇒ no feed is read (the default)
crawlWkr *crawlWorker // nil ⇒ no page is watched (the default)
}
// newDaemonAPI builds the real CoreAPI, with every field set. The unlock path
// used to leave nexus and getMCPServers nil, so after a cold start
// ResolveEntity refused with a nexus block configured and /tools rendered
// "not configured" with an mcp block configured. Empty is a wrong answer
// there, not a degraded one.
func newDaemonAPI(d bootDeps) *daemonAPI {
api := &daemonAPI{
CoreAPI: d.coreFor(),
getTrace: d.tl.trace,
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return d.tl.morningStatus(ctx, time.Now()) },
getDayPlan: func(ctx context.Context) ipc.DayPlan { return d.tl.dayPlan(ctx, time.Now()) },
getEvents: intakeEventsFn(d.evBus),
getDecisions: turnDecisionsFn(d.voiceW),
seedStore: seedStoreIfAllowed(d.st),
nexus: nexusOf(d.voiceW),
}
if d.voiceW != nil && d.voiceW.handler != nil {
api.chatFn = d.voiceW.handler.handleText
// And the reverse: the handler was wired with the bare store adapter,
// which cannot serve the day plan. See upgradeAPI.
d.voiceW.handler.upgradeAPI(api)
}
if d.voiceW != nil && d.voiceW.mcp != nil {
api.getMCPServers = d.voiceW.mcp.status
}
return api
}
// namedWorker is one long-running goroutine. The name exists so the set is
// assertable from a test and readable in a log; nothing dispatches on it.
type namedWorker struct {
name string
run func(ctx context.Context)
}
// backgroundWorkers lists what this deployment runs. It is pure — it starts
// nothing — so a test can compare the set the two paths would start without
// standing a daemon up.
func backgroundWorkers(d bootDeps) []namedWorker {
var ws []namedWorker
if d.voiceW != nil && d.voiceW.server != nil {
ws = append(ws, namedWorker{"voice", func(context.Context) {
if err := d.voiceW.server.Serve(); err != nil && !errors.Is(err, net.ErrClosed) {
log.Printf("voice serve: %v", err)
}
}})
}
ws = append(ws,
namedWorker{"tick", d.tl.run},
namedWorker{"fact-enrichment", d.factWorker.run},
)
if d.evalWorker != nil {
ws = append(ws, namedWorker{"memory-eval", d.evalWorker.run})
}
if d.feedWkr != nil {
ws = append(ws, namedWorker{"feed", d.feedWkr.run})
}
if d.crawlWkr != nil {
ws = append(ws, namedWorker{"crawl", d.crawlWkr.run})
}
if d.voiceW != nil && d.voiceW.mcp != nil {
ws = append(ws, namedWorker{"mcp", d.voiceW.mcp.run})
}
if d.voiceW != nil && d.voiceW.home != nil {
ws = append(ws, namedWorker{"home", d.voiceW.home.run})
}
return ws
}
// startBackground starts every worker through goWorker, so waitWorkers can
// wait for it at shutdown. A worker started as a bare `go func()` is the
// shutdown bug documented at the end of run(): run() never returns, the
// deferred Close never seals the database, and the ciphertext goes stale.
func startBackground(ctx context.Context, wg *sync.WaitGroup, d bootDeps) {
for _, w := range backgroundWorkers(d) {
goWorker(wg, func() { w.run(ctx) })
}
if d.voiceW != nil && d.voiceW.server != nil {
log.Printf("mavend: voice listening on %s", d.voiceW.server.Addr())
}
}
+96
View File
@@ -0,0 +1,96 @@
package main
import (
"reflect"
"testing"
"github.com/kami/maven/internal/decision"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/store"
"github.com/kami/maven/internal/voice"
)
// fullDeps — a deployment with every optional piece present. Nothing here is
// run: newDaemonAPI takes method values and backgroundWorkers is pure, so
// zero-value wirings are enough to say what WOULD be started.
func fullDeps() bootDeps {
h := &reactiveHandler{
ecosystem: &ecosystemWiring{nexus: &nexusClient{}},
decisions: decision.NewRing(),
}
return bootDeps{
coreFor: func() ipc.CoreAPI { return ipc.UnimplementedCoreAPI{} },
tl: &tickLoop{},
evBus: event.NewBus(4),
st: &store.Store{},
factWorker: &factEnrichmentWorker{},
evalWorker: &memoryEvalWorker{},
feedWkr: &feedWorker{},
crawlWkr: &crawlWorker{},
voiceW: &voiceWiring{
server: &voice.Server{},
handler: h,
mcp: &mcpWiring{},
home: &homeWiring{},
},
}
}
// The unlock path used to build its own daemonAPI literal and leave nexus and
// getMCPServers nil (V-639). Both paths call newDaemonAPI now, so the drift
// that can still happen is a field added to the struct and not to the
// constructor. This catches that one, by name.
func TestNewDaemonAPISetsEveryField(t *testing.T) {
prev := allowSeedOnStart
allowSeedOnStart = true
defer func() { allowSeedOnStart = prev }()
api := newDaemonAPI(fullDeps())
v := reflect.ValueOf(*api)
for i := range v.NumField() {
if v.Field(i).IsZero() {
t.Errorf("newDaemonAPI left %s unset — a fully wired deployment must fill every field", v.Type().Field(i).Name)
}
}
}
// The handler is wired with the bare store adapter and cannot serve the day
// plan until upgradeAPI hands it the real one. The unlocked path did that and
// the unlock path did it too; keep it a property of the constructor.
func TestNewDaemonAPIUpgradesTheHandler(t *testing.T) {
d := fullDeps()
api := newDaemonAPI(d)
if d.voiceW.handler.api != ipc.CoreAPI(api) {
t.Fatal("newDaemonAPI did not hand the handler the API it built")
}
}
// Every worker the daemon runs goes through startBackground, so shutdown can
// wait for it. The unlock path used to start seven of these as bare
// `go func()` under a shadowed WaitGroup.
func TestBackgroundWorkersFullSet(t *testing.T) {
want := []string{"voice", "tick", "fact-enrichment", "memory-eval", "feed", "crawl", "mcp", "home"}
var got []string
for _, w := range backgroundWorkers(fullDeps()) {
got = append(got, w.name)
}
if !reflect.DeepEqual(got, want) {
t.Errorf("workers = %v, want %v", got, want)
}
}
// A default box configures none of the optional blocks. Two workers always run
// and the rest stay dark, rather than a nil run being scheduled.
func TestBackgroundWorkersFloor(t *testing.T) {
d := fullDeps()
d.evalWorker, d.feedWkr, d.crawlWkr, d.voiceW = nil, nil, nil, nil
want := []string{"tick", "fact-enrichment"}
var got []string
for _, w := range backgroundWorkers(d) {
got = append(got, w.name)
}
if !reflect.DeepEqual(got, want) {
t.Errorf("workers = %v, want %v", got, want)
}
}
+1 -1
View File
@@ -571,7 +571,7 @@ func (h *reactiveHandler) finishClarified(ctx context.Context, dec router.Decisi
}
reply := h.applyAction(ctx, dec)
if reply == "" {
reply = h.replier.Reply(dec)
reply = h.replier.Reply(ctx, dec)
}
if reply == "" {
// Belt: an empty reply here would be a silent drop.
+24 -8
View File
@@ -37,8 +37,8 @@ type factEnrichmentWorker struct {
nextTry map[int64]time.Time // fact id → earliest retry
}
// enrichmentScanLimit bounds how deep a single tick (or status report) walks
// the pending queue looking for facts whose backoff has elapsed. The queue is
// enrichmentScanLimit bounds how deep a single tick walks the pending queue
// looking for facts whose backoff has elapsed. The queue is
// ordered by id, so without a scan the oldest facts hold every batch slot
// whether or not they are eligible, and one permanently failing fact stalls
// every younger one behind it.
@@ -75,8 +75,8 @@ func newFactEnrichmentWorker(st *store.Store, eco *ecosystemWiring, interval tim
// has been down all day must be visible as a backlog, not as facts that
// silently never got tagged.
//
// All three numbers describe the same set of rows, the first
// enrichmentScanLimit pending facts. Counting Pending over a thousand rows
// All three numbers describe the same set of rows, whatever is still pending
// out of the first enrichmentScanLimit facts. Counting Pending over a thousand rows
// while counting InBackoff over the twenty that reached the head of a batch
// described two different populations under one struct.
type enrichmentStatus struct {
@@ -86,13 +86,22 @@ type enrichmentStatus struct {
Scanned int // rows the other three counts were taken over
}
// status reads the queue and counts over it. For a caller with no batch in
// hand — anything asking the worker how it is doing from outside the tick.
func (w *factEnrichmentWorker) status(ctx context.Context) enrichmentStatus {
var st enrichmentStatus
pending, err := w.store.PendingFactResolutions(ctx, enrichmentScanLimit)
if err != nil {
log.Printf("factenrichment: status: %v", err)
return st
return enrichmentStatus{}
}
return w.statusOf(pending)
}
// statusOf counts over a batch the caller already has. The batch is the query
// the tick already ran, so reporting the backlog costs no second read of the
// scan limit — up to a thousand rows, on a database that serialises them.
func (w *factEnrichmentWorker) statusOf(pending []store.Fact) enrichmentStatus {
var st enrichmentStatus
st.Pending = len(pending)
st.Scanned = len(pending)
w.mu.Lock()
@@ -144,17 +153,24 @@ func (w *factEnrichmentWorker) tick(ctx context.Context) {
}
w.forgetDeparted(pending)
skipped, failed, attempted := 0, 0, 0
// A resolved fact leaves the pending queue, so the batch in hand overstates
// the backlog by however many succeeded. Drop them here rather than
// re-reading the queue to find out.
remaining := make([]store.Fact, 0, len(pending))
for _, f := range pending {
if attempted >= w.batch {
break
remaining = append(remaining, f)
continue
}
if !w.due(f.ID) {
skipped++
remaining = append(remaining, f)
continue
}
attempted++
if !w.resolveOne(ctx, f) {
failed++
remaining = append(remaining, f)
}
}
if failed > 0 {
@@ -164,7 +180,7 @@ func (w *factEnrichmentWorker) tick(ctx context.Context) {
// Report the backlog every tick, not only when something failed: the
// stalled state worth seeing is the one where nothing failed because
// nothing was attempted.
if st := w.status(ctx); st.Pending > 0 {
if st := w.statusOf(remaining); st.Pending > 0 {
log.Printf("factenrichment: %d facts pending entity resolution, %d in backoff, worst attempt %d (scanned %d)",
st.Pending, st.InBackoff, st.MaxAttempts, st.Scanned)
}
+24 -112
View File
@@ -252,6 +252,23 @@ func run(args []string) error {
// envelope per successful intake write.
coreFor := func() ipc.CoreAPI { return newIntakeAPI(ipc.NewStoreAPI(st), evBus, time.Now) }
// depsNow reads whatever the current path has wired. Both boot paths build
// the CoreAPI and start the workers from this one value, so neither can
// hold a field the other misses. See cmd/mavend/boot.go.
depsNow := func() bootDeps {
return bootDeps{
coreFor: coreFor,
tl: tl,
evBus: evBus,
voiceW: voiceW,
st: st,
factWorker: factWorker,
evalWorker: evalWorker,
feedWkr: feedWkr,
crawlWkr: crawlWkr,
}
}
if !locked {
rules = wireRules(cfg)
gatherer = wireGatherer(st, cfg, rules)
@@ -284,26 +301,7 @@ func run(args []string) error {
feedWkr = newFeedWorker(coreFor(), embedderOf(voiceW), cfg)
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
coreAPI = &daemonAPI{
CoreAPI: coreFor(),
getTrace: tl.trace,
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
getDayPlan: func(ctx context.Context) ipc.DayPlan { return tl.dayPlan(ctx, time.Now()) },
getEvents: intakeEventsFn(evBus),
getDecisions: turnDecisionsFn(voiceW),
seedStore: seedStoreIfAllowed(st),
nexus: nexusOf(voiceW),
}
if voiceW != nil && voiceW.handler != nil {
api := coreAPI.(*daemonAPI)
api.chatFn = voiceW.handler.handleText
// And the reverse: the handler was wired with the bare store
// adapter, which cannot serve the day plan. See upgradeAPI.
voiceW.handler.upgradeAPI(api)
}
if voiceW != nil && voiceW.mcp != nil {
coreAPI.(*daemonAPI).getMCPServers = voiceW.mcp.status
}
coreAPI = newDaemonAPI(depsNow())
} else {
// locked mode: no real store yet, so there's no meaningful CoreAPI to
// serve. srv.Check below is the actual guard — every CoreAPI call is
@@ -497,19 +495,7 @@ func run(args []string) error {
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
// Swap the CoreAPI from the locked placeholder to the real store adapter.
newAPI := &daemonAPI{
CoreAPI: coreFor(),
getTrace: tl.trace,
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
getDayPlan: func(ctx context.Context) ipc.DayPlan { return tl.dayPlan(ctx, time.Now()) },
getEvents: intakeEventsFn(evBus),
getDecisions: turnDecisionsFn(voiceW),
seedStore: seedStoreIfAllowed(st),
}
if voiceW != nil && voiceW.handler != nil {
newAPI.chatFn = voiceW.handler.handleText
voiceW.handler.upgradeAPI(newAPI)
}
newAPI := newDaemonAPI(depsNow())
srv.SetAPI(newAPI)
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
wireMailIntake(srv, st, phr, cfg, evBus)
@@ -524,59 +510,10 @@ func run(args []string) error {
// block, so no wire path takes a voiceprint on a default box.
wireSpeaker(srv, st, cfg)
// Start voice server.
if voiceW != nil {
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
if err := voiceW.server.Serve(); err != nil && !errors.Is(err, net.ErrClosed) {
log.Printf("voice serve: %v", err)
}
}()
log.Printf("mavend: voice listening on %s", voiceW.server.Addr())
}
// Start tick loop.
go func() {
tl.run(ctx)
}()
// Start fact-entity enrichment worker.
go func() {
factWorker.run(ctx)
}()
// Start background memory evaluation (nil unless configured).
if evalWorker != nil {
go func() {
evalWorker.run(ctx)
}()
}
// Start feed reading (nil unless configured).
if feedWkr != nil {
go func() {
feedWkr.run(ctx)
}()
}
// Start the watched-page crawls (nil unless configured).
if crawlWkr != nil {
go func() {
crawlWkr.run(ctx)
}()
}
// Keep MCP connections alive (nil unless configured).
if voiceW != nil && voiceW.mcp != nil {
go voiceW.mcp.run(ctx)
}
// Re-enumerate the house for new devices (nil unless configured).
if voiceW != nil && voiceW.home != nil {
go voiceW.home.run(ctx)
}
// The voice server and every background worker, on the outer wg
// so shutdown waits for them. This used to be nine bare
// `go func()` calls and a shadowed WaitGroup (V-639).
startBackground(ctx, &wg, depsNow())
dl.unlock(st)
log.Printf("mavend: unlocked via passkey assertion")
@@ -591,33 +528,8 @@ func run(args []string) error {
})
log.Printf("mavend: ipc listening on %s", srv.Path())
if !locked && voiceW != nil {
goWorker(&wg, func() {
if err := voiceW.server.Serve(); err != nil && !errors.Is(err, net.ErrClosed) {
log.Printf("voice serve: %v", err)
}
})
log.Printf("mavend: voice listening on %s", voiceW.server.Addr())
}
if !locked {
goWorker(&wg, func() { tl.run(ctx) })
goWorker(&wg, func() { factWorker.run(ctx) })
if evalWorker != nil {
goWorker(&wg, func() { evalWorker.run(ctx) })
}
if feedWkr != nil {
goWorker(&wg, func() { feedWkr.run(ctx) })
}
if crawlWkr != nil {
goWorker(&wg, func() { crawlWkr.run(ctx) })
}
if voiceW != nil && voiceW.mcp != nil {
goWorker(&wg, func() { voiceW.mcp.run(ctx) })
}
if voiceW != nil && voiceW.home != nil {
goWorker(&wg, func() { voiceW.home.run(ctx) })
}
startBackground(ctx, &wg, depsNow())
}
<-ctx.Done()
+4 -4
View File
@@ -22,7 +22,7 @@ func newLLMReplier(c phraser.Completer, block func() string) *llmReplier {
// Reply never fails: a clarify, a model error and an unusable generation all
// answer from the stub, which is what keeps a turn from breaking on the model.
func (r *llmReplier) Reply(d router.Decision) string {
func (r *llmReplier) Reply(ctx context.Context, d router.Decision) string {
if d.Clarify {
// The deck, not the stub's single sentence: a clarify she cannot turn
// into a question is the line he hears most often when she misses him,
@@ -39,14 +39,14 @@ func (r *llmReplier) Reply(d router.Decision) string {
// что ты выпел стакан воды" for "я выпил воды".
return phraser.FactAck(d.Utterance)
}
out, err := r.p.PhraseReply(context.Background(), d)
out, err := r.p.PhraseReply(ctx, d)
if err != nil || out == "" {
return r.stub.Reply(d)
return r.stub.Reply(ctx, d)
}
// The persona checks, on the live path (personaguard.go). A reply that
// leaks reasoning or calls him "вы" is worse than a flat one.
if _, ok := guardSpoken("reply", out); !ok {
return r.stub.Reply(d)
return r.stub.Reply(ctx, d)
}
return out
}
+5 -5
View File
@@ -22,7 +22,7 @@ func (s stubCompleter) Complete(_ context.Context, _ llm.Req) (string, error) {
func TestLLMReplierPassesTheModelReplyThrough(t *testing.T) {
r := newLLMReplier(stubCompleter{out: `{"response":"записала, кофе закончился","mood":"neutral"}`}, nil)
got := r.Reply(router.Decision{Intent: router.IntentNote, Slots: router.Slots{Text: "кофе закончился"}})
got := r.Reply(context.Background(), router.Decision{Intent: router.IntentNote, Slots: router.Slots{Text: "кофе закончился"}})
if got != "записала, кофе закончился" {
t.Errorf("got %q, want %q", got, "записала, кофе закончился")
}
@@ -42,7 +42,7 @@ func TestLLMReplierFallsBackToStubOnEmpty(t *testing.T) {
// the clarify deck rather than the stub's single sentence.
func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
r := newLLMReplier(stubCompleter{out: "я всё поняла"}, nil)
got := r.Reply(router.Decision{Clarify: true, Utterance: "мгм"})
got := r.Reply(context.Background(), router.Decision{Clarify: true, Utterance: "мгм"})
if got == "я всё поняла" {
t.Fatal("a clarify must not be phrased by the model")
}
@@ -50,7 +50,7 @@ func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
t.Errorf("on clarify: got %q, want %q", got, want)
}
// Two different misses do not sound identical.
if same := r.Reply(router.Decision{Clarify: true, Utterance: "а"}); same == got {
if same := r.Reply(context.Background(), router.Decision{Clarify: true, Utterance: "а"}); same == got {
t.Log("two utterances hashed to the same line, which is allowed but should be rare")
}
}
@@ -60,14 +60,14 @@ func TestLLMReplierClarifyReadsTheDeck(t *testing.T) {
// produce, which is the same claim without pinning one wording.
func assertAck(t *testing.T, r *llmReplier, d router.Decision, key, what string) {
t.Helper()
if got := r.Reply(d); !phraser.IsAck(key, nil, got) {
if got := r.Reply(context.Background(), d); !phraser.IsAck(key, nil, got) {
t.Errorf("on %s: got %q, want a %q line", what, got, key)
}
}
func assertStub(t *testing.T, r *llmReplier, d router.Decision, what string) {
t.Helper()
got, want := r.Reply(d), voice.NewStubReplier().Reply(d)
got, want := r.Reply(context.Background(), d), voice.NewStubReplier().Reply(context.Background(), d)
if got != want {
t.Errorf("on %s: got %q, want stub %q", what, got, want)
}
+1 -1
View File
@@ -458,7 +458,7 @@ func (h *reactiveHandler) runTurn(ctx context.Context, text string, src turnSour
// 9. replier — phrase the reply across the router decision.
if replyText == "" {
replyText = h.replier.Reply(dec)
replyText = h.replier.Reply(ctx, dec)
}
return withNotice(expiredNotice, replyText)
}
+19 -1
View File
@@ -58,6 +58,12 @@ func main() {
// mutex, so sharing the connection would freeze every other page for the
// length of the load. See handleModels.
var swapConn modelController
// turnConn — a third connection, for POST /api/chat and nothing else, for
// the same reason /models has one (V-638). A chat turn routes, phrases and
// may act, bounded only by phraser.timeout at 60s, and every other handler
// on this server queues behind it on the shared client's one mutex. Nil ⇒
// chat shares the main connection, which is how it behaved before.
var turnConn ipc.CoreAPI
if *coreSock != "" {
c, err := ipc.DialWait(*coreSock, 60*time.Second)
if err != nil {
@@ -71,6 +77,12 @@ func main() {
defer sc.Close()
swapConn = sc
}
if tc, err := ipc.Dial(*coreSock); err != nil {
log.Printf("chat: third core connection failed (%v) — /api/chat will share the main one and a turn will block the other pages", err)
} else {
defer tc.Close()
turnConn = tc
}
}
// stepUpSession stays nil unless the passkey endpoints are wired below — it
@@ -208,7 +220,13 @@ func main() {
// decides how every utterance is routed and how every reply is worded.
mux.HandleFunc("/tools", gatedPage(handleTools))
mux.HandleFunc("/routines", gatedPage(handleRoutines))
mux.HandleFunc("/api/chat", gatedPage(handleChatAPI))
mux.HandleFunc("/api/chat", func(w http.ResponseWriter, r *http.Request) {
c := turnConn
if c == nil {
c = core
}
handleChatAPI(w, r, c, stepUpSession, *requireStepUp)
})
mux.HandleFunc("/api/revert", gatedPage(handleRevert))
mux.HandleFunc("/api/correct", gatedPage(handleCorrectAPI))
mux.HandleFunc("/models", func(w http.ResponseWriter, r *http.Request) {
+9 -1
View File
@@ -37,7 +37,15 @@
"This needs a matching ufw rule or the container's SYN is dropped:",
" ufw allow from 192.168.240.0/20 to any port 10808 proto tcp"
],
"proxy": "socks5://192.168.240.1:10808"
"proxy": "socks5://192.168.240.1:10808",
"//intake": [
"Read the chat as well as write to it (V-637). The poller long-polls",
"getUpdates through the same relay and accepts chat_id as the only",
"sender. Deleting this key turns inbound off again.",
"chat_id must be numeric here or the daemon refuses to start: an inbound",
"update names its chat by number, so an @-name would match nothing."
],
"intake": true
},
"//workstation": [
+38
View File
@@ -157,6 +157,44 @@ services:
# - maildata:/var/lib/mavmaild
# - ./deploy/imap.password:/run/secrets/imap.password:ro
# The calendar reader (Vikunja #644) is OFF and commented out: it needs a
# CalDAV account, and there is none on this box. It was built, listed in
# `make build`, and deployed nowhere, which is the worst of the three states —
# this block records the decision instead.
#
# What its absence costs, so the cost is visible from here:
# - Agenda questions route correctly and answer from nothing. Stage 0 sends
# "что у меня сегодня" to IntentQuery (V-498) and the `calendar` query
# source reads facts(kind=env, source=caldav:*) that nobody writes.
# - The nudge gate loses a suppressor. loop.State.CalendarBusy is fed by
# those same facts, so "do not nag mid-meeting" is permanently false.
#
# Core never sees the CalDAV password: the reader polls the collection itself
# and hands core one fact per event over WriteFact. Nothing here can create a
# reminder, so a misread event cannot fire.
#
# The password is read from a FILE, so it never appears in `ps`, in this file,
# or in shell history — the same rule mavpoll and mavmaild follow.
#
# To enable: write the password to deploy/caldav.password (0600, gitignored),
# point -url at the collection, and uncomment this service. No mavend.json
# block is needed — events arrive over IPC as facts. -render-url is optional
# and OFF here: it publishes Maven's own reminders back as events, and it must
# not name the collection -url reads, or the poller reads its own writes back
# in (checkRenderTarget refuses that). It takes -render-pass-file, and falls
# back to this password when that is not given.
# mavcaldav:
# <<: *image
# command: ["mavcaldav", "-socket", "/run/maven/mavend.sock",
# "-url", "http://localhost:5232/kami/personal",
# "-user", "kami",
# "-pass-file", "/run/secrets/caldav.password",
# "-interval", "5m"]
# depends_on: [mavend]
# volumes:
# - sockets:/run/maven
# - ./deploy/caldav.password:/run/secrets/caldav.password:ro
volumes:
dbdata:
sockets:
@@ -0,0 +1,69 @@
# Does one sqlite connection make reads queue? No (V-642)
Measured 07-08-2026 at `7b507de`, on homesrv. The harness is
`internal/store/conncap_test.go`. It stays in the repo, because this claim gets
re-argued and the numbers should be re-runnable rather than quoted.
`internal/store/store.go` opens the database with `SetMaxOpenConns(1)`, while
`schema.sql` sets `journal_mode=WAL`. WAL exists to let readers run beside one
writer, so the cap gives up the thing the journal mode was chosen for. The
question was whether that costs anything.
## What was measured
A fixed two-second window. One writer calling `SetValue` paced at 2ms, and a
reader loop calling `RecentFacts(50)` over 500 seeded rows as fast as it can.
Same schema, same modernc driver, same machine, three runs per cap.
The window is wall-clock rather than a read count on purpose. A first version ran
a fixed 300 reads. That finished sooner at the higher cap, so it received fewer
writes, and two runs that did different work cannot be compared.
| cap | reads | writes | p50 | p95 | max |
|---|---|---|---|---|---|
| 1 | ~3050 | ~760 | 594µs | 900µs | 16-19ms |
| 4 | ~3600 | ~340 | 525µs | 710µs | 1-2ms |
## What it says
**Reads do not queue behind writes.** Four connections buy about 70µs at p50. A
turn spends 1.19s in the resident model. The tail does improve, from 19ms to 2ms,
and 19ms is still not a figure anyone notices in a spoken reply.
**Write throughput more than halves at the higher cap**, 760 writes against 340.
inference, not measured directly: at one connection the reader and the writer take
turns with no lock contention. At four the writer contends for the WAL write lock
with a live reader. Whatever the mechanism, the trade runs the opposite way from
the one the task expected.
**The cap was not the source of the 2.7s router figure.** CLAUDE.md records that
figure as contention rather than the model. This task was a candidate for where
that contention came from. A 19ms worst case cannot produce it. That line of
enquiry is closed.
**One transaction is what the cap cannot survive.** With a read-only transaction
open, a second read at cap 1 never completes. The harness gave it two seconds and
got `context deadline exceeded`. The same read at cap 4 took 1ms. The transaction
holds the only connection, so this is not a slow read, it is a stalled database.
## What was done
The cap stays at 1. The reason is now written where the cap is set, rather than
inferred from a four-word comment.
`Store.DB` was deleted. It handed out exactly the read-only transaction measured
above. It had been there since the initial commit with no production caller, and
its doc comment described a loop that never materialised. Its one user was a test
helper reading `delivery_attempts` by raw SQL. `ListDeliveryAttempts` has covered
that since V-390, and the helper now goes through the reader.
So the hazard is gone by construction, not by documentation.
`TestConnCap_ReadBlocksBehindOpenSnapshot` is the standing measurement of what
re-adding the seam would cost.
## Not answered
Whether reads queue on the deployed box under real load, as opposed to a
synthetic loop. The harness writes and reads one table. Digestion reads four and
embeds while it does. The finding that closes this task is the transaction stall,
which is structural and does not depend on load.
@@ -1,6 +1,14 @@
# No deadline on the turn path
Last verified: 06-08-2026 @ 06c1cf2
Last verified: 06-08-2026 @ 60e64dd
**All four steps landed on 06-08-2026.** What follows describes the defect as it was and
the work as it was planned. Two things came out differently. `Client.Close` read the conn
field with no lock while `roundtrip` re-dialed and dropped it. `-race` caught that on the
new cancellation test. So the conn field now has a mutex of its own, held only across a
read or an assignment. And `/api/ptt` needed nothing: it proxies to the voice port and never
touches the shared client, so only `/api/chat` got the extra connection. The pool inside
`ipc.Client` is still unbuilt and still waiting on a second module measured queueing.
V-638. Sibling of V-607, which is the same class of bug in `internal/worker`.
Reads with `docs/offload.md` and `docs/protocol.md`.
+18 -1
View File
@@ -1,9 +1,26 @@
# The two boot paths have drifted
Last verified: 06-08-2026 @ 06c1cf2
Last verified: 06-08-2026 @ 69d0f5e
V-639. Reads with `docs/operations.md`.
## What landed
`cmd/mavend/boot.go`. `newDaemonAPI(deps)` builds the CoreAPI with every field
set, and `startBackground(ctx, &wg, deps)` starts the voice server and every
worker through `goWorker`. `backgroundWorkers(deps)` is the pure list behind it,
so a test can compare the set without standing a daemon up. Both paths in
`run()` now read `coreAPI = newDaemonAPI(depsNow())` and one
`startBackground(...)`, where `depsNow` reads whatever the current path wired.
The shadowed `wg` is gone. Four tests in `cmd/mavend/boot_test.go`. Every
`daemonAPI` field is set on a fully wired deployment. The handler gets the API
it was built with. The worker set is asserted by name, at the full set and at
the floor.
Still by hand: unlock a locked box by passkey, ask something that needs Nexus,
and check `/tools` lists the MCP servers.
## What is wrong
`run()` in `cmd/mavend/main.go` brings the daemon up two ways. A box with a key in the
+21
View File
@@ -456,9 +456,30 @@ func (c *Config) validate() error {
if err := c.validateCapture(); err != nil {
return err
}
if err := c.validateTelegram(); err != nil {
return err
}
return nil
}
// validateTelegram refuses an intake half that cannot read the chat it is
// pointed at. The push half accepts an @channelusername and the intake half
// does not, so a box configured with both boots clean, keeps pushing, and
// answers nothing — the failure is invisible from the chat. Same shape as
// validateNetScan: fail the config rather than the turn.
func (c *Config) validateTelegram() error {
if c.Telegram == nil || !c.Telegram.Intake {
return nil
}
// An unset ${TELEGRAM_*} expands to empty, and the daemon already reads an
// empty token or chat id as telegram not being wired at all. Validating a
// block that wires nothing would fail a box that merely has no bot.
if c.Telegram.BotToken == "" || c.Telegram.ChatID == "" {
return nil
}
return telegramsink.ValidateIntakeChatID(c.Telegram.ChatID)
}
// DBEncryptionKey resolves the at-rest encryption key: DBKeyEnv (if set) wins
// over DBKeyB64. Returns (nil, nil) when neither is set — the caller then opens
// a plaintext store. A configured-but-invalid key is an error (fail closed,
+26
View File
@@ -466,3 +466,29 @@ func TestNormaliseKeepsExplicitWorkstationHealth(t *testing.T) {
t.Errorf("Health = %q, want %q", got, want)
}
}
func TestTelegramIntakeRefusesNamedChat(t *testing.T) {
// The push half accepts an @channelusername and the intake half cannot use
// one, so a box with both boots clean and answers nothing. Refuse the
// config instead.
p := writeConfig(t, `{"telegram":{"bot_token":"t","chat_id":"@maven","intake":true}}`)
if _, err := Load(p); err == nil {
t.Fatal("Load succeeded for intake with an @-name chat id; want error")
}
}
func TestTelegramNamedChatOKWithoutIntake(t *testing.T) {
// Push-only is what the @-name is for, so nothing changes for a box that
// never turned intake on.
p := writeConfig(t, `{"telegram":{"bot_token":"t","chat_id":"@maven"}}`)
if _, err := Load(p); err != nil {
t.Fatalf("Load: %v", err)
}
}
func TestTelegramIntakeAcceptsNumericChat(t *testing.T) {
p := writeConfig(t, `{"telegram":{"bot_token":"t","chat_id":"-1001234567890","intake":true}}`)
if _, err := Load(p); err != nil {
t.Fatalf("Load: %v", err)
}
}
+11 -9
View File
@@ -94,20 +94,22 @@ func openTestStore(t *testing.T) *store.Store {
// attemptStatus reads one attempt row back. Returns ok=false when the row is
// gone, which would itself be a broken promise (a dropped attempt).
//
// It goes through ListDeliveryAttempts rather than raw SQL. This helper used to
// reach past the store into store.DB, which was the tell that the outbox was
// write-only; the reader landed in V-390 and this caller was not moved over.
func attemptStatus(t *testing.T, st *store.Store, id int64) (status string, completed bool, ok bool) {
t.Helper()
tx, err := st.DB(context.Background())
attempts, err := st.ListDeliveryAttempts(context.Background(), "", 200)
if err != nil {
t.Fatalf("read tx: %v", err)
t.Fatalf("ListDeliveryAttempts: %v", err)
}
defer func() { _ = tx.Rollback() }()
var completedTS *int64
err = tx.QueryRowContext(context.Background(),
`SELECT status, completed_ts FROM delivery_attempts WHERE id = ?`, id).Scan(&status, &completedTS)
if err != nil {
return "", false, false
for _, a := range attempts {
if a.ID == id {
return a.Status, a.HasComplete, true
}
}
return status, completedTS != nil, true
return "", false, false
}
// TestCrashBetweenBeginAndCompleteBecomesUnknown — simulate the crash window:
+24 -9
View File
@@ -49,6 +49,24 @@ type Poller struct {
offset int64
}
// ValidateIntakeChatID refuses a chat id the intake half cannot use. The push
// half accepts @channelusername as a destination. The intake half cannot: an
// inbound update names its chat by numeric id, so an @-name would match nothing
// and the poller would read the chat and answer none of it. Config validation
// calls this, so the box refuses to boot rather than running a dead reach —
// NewPoller returning an error is too late, because the daemon is already up.
func ValidateIntakeChatID(chatID string) error {
id := strings.TrimSpace(chatID)
if id == "" {
return errors.New("telegramsink: intake needs a chat id")
}
digits := strings.TrimPrefix(id, "-")
if digits == "" || strings.TrimLeft(digits, "0123456789") != "" {
return fmt.Errorf("telegramsink: intake needs the numeric chat id, not %s", chatID)
}
return nil
}
// NewPoller builds the intake half around an already-validated sink, so the
// token, the base URL and the relay are resolved in one place. turn is
// required; correct may be nil, and then the reply carries no buttons.
@@ -59,12 +77,8 @@ func NewPoller(s *Sink, turn Turn, correct Correct) (*Poller, error) {
if turn == nil {
return nil, errors.New("telegramsink: intake needs a turn handler")
}
// 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.
// 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)
if err := ValidateIntakeChatID(s.cfg.ChatID); err != nil {
return nil, err
}
// The sink's transport already carries the relay. Only the timeout differs,
// and it has to clear the long poll.
@@ -164,9 +178,10 @@ func (p *Poller) onMessage(ctx context.Context, m *message) {
}
}
// onCallback handles a tap on a correction button. Every path answers the
// callback: telegram spins a clock on the button until it is answered, and an
// unanswered tap reads as a gesture that was dropped.
// onCallback handles a tap on a correction button. Every path from the owner
// answers the callback: telegram spins a clock on the button until it is
// answered, and an unanswered tap reads as a gesture that was dropped. A tap
// from anyone else gets silence, the same as a message from a stranger.
func (p *Poller) onCallback(ctx context.Context, cb *callbackQuery) {
if !p.fromOwner(cb.Message.Chat.idString()) {
return
@@ -199,3 +199,18 @@ func TestNewPollerNeedsATurn(t *testing.T) {
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)
}
}
}
+192
View File
@@ -0,0 +1,192 @@
package ipc
import (
"context"
"errors"
"net"
"path/filepath"
"testing"
"time"
)
// A cancelled context has to abort a call that is already in flight. It did not
// until V-638: call checked ctx once before sending and then blocked in
// roundtrip with no connection deadline, so a daemon that read the frame and
// never answered parked the caller for as long as the socket stayed open.
//
// The server here is that daemon: it accepts, reads nothing, replies nothing.
func deafServer(t *testing.T) string {
t.Helper()
sock := filepath.Join(t.TempDir(), "deaf.sock")
ln, err := net.Listen("unix", sock)
if err != nil {
t.Fatalf("listen: %v", err)
}
t.Cleanup(func() { _ = ln.Close() })
go func() {
for {
conn, err := ln.Accept()
if err != nil {
return
}
// Hold it open and say nothing. Closed by the listener cleanup.
t.Cleanup(func() { _ = conn.Close() })
}
}()
return sock
}
func TestClientCancelAbortsAReadInFlight(t *testing.T) {
c, err := Dial(deafServer(t))
if err != nil {
t.Fatalf("dial: %v", err)
}
defer c.Close()
ctx, cancel := context.WithCancel(context.Background())
go func() {
time.Sleep(50 * time.Millisecond)
cancel()
}()
done := make(chan error, 1)
go func() {
_, err := c.Ping(ctx)
done <- err
}()
select {
case err := <-done:
// Ping is read-only, so the cancellation is reported as itself rather
// than as an ambiguous mutation.
if !errors.Is(err, context.Canceled) {
t.Errorf("got %v, want context.Canceled", err)
}
case <-time.After(5 * time.Second):
t.Fatal("a cancelled Ping did not return")
}
}
// A mutation cancelled while awaiting the reply may already have committed, so
// it is ErrAmbiguousOutcome and never a retry. That split is the invariant
// internal/ipc/maperr_test.go's neighbours rest on.
func TestClientCancelLeavesAMutationAmbiguous(t *testing.T) {
c, err := Dial(deafServer(t))
if err != nil {
t.Fatalf("dial: %v", err)
}
defer c.Close()
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
defer cancel()
done := make(chan error, 1)
go func() {
_, err := c.WriteFact(ctx, WriteFactReq{Key: "water", Value: "drank"})
done <- err
}()
select {
case err := <-done:
if !errors.Is(err, ErrAmbiguousOutcome) {
t.Errorf("got %v, want ErrAmbiguousOutcome", err)
}
case <-time.After(5 * time.Second):
t.Fatal("a cancelled WriteFact did not return")
}
}
// The deadline itself, with no cancellation: a call on a context with no
// deadline used to have no bound at all. This one has one and must respect it.
func TestClientDeadlineBoundsACall(t *testing.T) {
c, err := Dial(deafServer(t))
if err != nil {
t.Fatalf("dial: %v", err)
}
defer c.Close()
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
defer cancel()
start := time.Now()
if _, err := c.Ping(ctx); err == nil {
t.Fatal("a deaf server answered a Ping")
}
if elapsed := time.Since(start); elapsed > 3*time.Second {
t.Errorf("Ping took %v, want the context deadline to bound it", elapsed)
}
}
// blockingAPI parks Presence until its context is cancelled and records what
// cancelled it. Every other method is the unimplemented floor.
type blockingAPI struct {
UnimplementedCoreAPI
entered chan struct{}
err chan error
}
func (b *blockingAPI) Presence(ctx context.Context) (Presence, error) {
close(b.entered)
<-ctx.Done()
b.err <- ctx.Err()
return Presence{}, ctx.Err()
}
// serveConn dispatched under context.Background() until V-638, so Close could
// only abandon a dispatch in flight and never tell it to stop.
func TestServerCloseCancelsADispatchInFlight(t *testing.T) {
api := &blockingAPI{entered: make(chan struct{}), err: make(chan error, 1)}
srv, err := Listen(filepath.Join(t.TempDir(), "core.sock"), api)
if err != nil {
t.Fatalf("listen: %v", err)
}
served := make(chan struct{})
go func() { _ = srv.Serve(); close(served) }()
cli, err := Dial(srv.Path())
if err != nil {
t.Fatalf("dial: %v", err)
}
defer cli.Close()
go func() { _, _ = cli.Presence(context.Background()) }()
select {
case <-api.entered:
case <-time.After(5 * time.Second):
t.Fatal("the handler was never dispatched")
}
_ = srv.Close()
<-served
select {
case got := <-api.err:
if !errors.Is(got, context.Canceled) {
t.Errorf("handler saw %v, want context.Canceled", got)
}
case <-time.After(5 * time.Second):
t.Fatal("Close did not cancel the dispatch")
}
}
// The watchdog closes the conn, and it races the end of the call: a
// cancellation landing as the reply arrives can close a conn the call was
// already done with. That is survivable either way, because a write to a closed
// socket is errWriteLost and errWriteLost re-dials and retries, so this test
// passes with or without the drop in roundtrip's defer. What it pins is that
// the recovery is real and costs one round trip at most, never an error the
// caller sees.
func TestClientSurvivesACancelledCall(t *testing.T) {
_, _, cli, _ := newServerWithStore(t)
for i := 0; i < 20; i++ {
ctx, cancel := context.WithCancel(context.Background())
go cancel() // races the reply on purpose
_, _ = cli.Ping(ctx)
cancel()
if _, err := cli.Ping(context.Background()); err != nil {
t.Fatalf("call %d after a cancelled one: %v", i, err)
}
}
}
+94 -12
View File
@@ -26,9 +26,20 @@ type Client struct {
conn net.Conn
path string // the address as configured, kept for errors and logs
addr netaddr.Addr // parsed, so a dropped conn can be re-dialed (core restart)
mu sync.Mutex
mu sync.Mutex // one request at a time, so a frame and its reply pair up
// connMu guards the conn field alone, and is held only across an assignment
// or a read. It exists so Close and the cancellation watchdog can reach the
// connection without waiting for the call that is holding c.mu (V-638).
connMu sync.Mutex
}
// defaultCallTimeout bounds a call whose context carries no deadline. It is
// the same 120s internal/voice/client.go settles on: long enough for a model
// call on a cold resident model, short enough that a daemon which stopped
// answering does not park the caller forever.
const defaultCallTimeout = 120 * time.Second
// errWriteLost marks a conn drop while sending the request frame: the request
// never reached the server (or the server never saw a complete frame), so
// retrying is always safe regardless of method — nothing was applied to
@@ -103,11 +114,22 @@ func Dial(path string) (*Client, error) {
return &Client{conn: c, path: path, addr: addr}, nil
}
// Close closes the connection out from under a call in flight, on purpose: a
// shutdown must not wait out a parked read. It takes connMu and never c.mu, so
// it cannot block behind the call it is interrupting.
//
// The lock is taken and released by hand, around the two field accesses and
// nothing else. The socket close happens outside it, because a close on a tcp
// conn can block and connMu is on the path of every call.
func (c *Client) Close() error {
if c.conn == nil {
c.connMu.Lock()
conn := c.conn
c.conn = nil
c.connMu.Unlock()
if conn == nil {
return nil
}
return c.conn.Close()
return conn.Close()
}
// DialWait is Dial with patience: it retries with capped backoff until the
@@ -163,16 +185,26 @@ func (c *Client) call(ctx context.Context, m Method, params, result any) error {
}
var resp Response
err := c.roundtrip(m, raw, &resp)
err := c.roundtrip(ctx, m, raw, &resp)
switch {
case errors.Is(err, errWriteLost):
// The request never left; a duplicate send can't double-apply.
// Redial (roundtrip re-dials on a nil conn) and retry exactly once.
err = c.roundtrip(m, raw, &resp)
// Not when the caller has given up — a retry would only be a second
// frame nobody is waiting for.
if ctx.Err() == nil {
err = c.roundtrip(ctx, m, raw, &resp)
}
case errors.Is(err, errReadLost):
if readOnlyMethods[m] {
if ctx.Err() != nil {
// The caller cancelled the read it was waiting for. Nothing
// was applied, so this is the cancellation and not an
// ambiguity.
return ctx.Err()
}
// A duplicate read can't double-apply either — safe to replay.
err = c.roundtrip(m, raw, &resp)
err = c.roundtrip(ctx, m, raw, &resp)
} else {
// The mutation may have already committed server-side. Do not
// retry: report the ambiguity instead of guessing.
@@ -199,19 +231,55 @@ func (c *Client) call(ctx context.Context, m Method, params, result any) error {
// failure is wrapped in errReadLost (ambiguous — call() only retries it for
// read-only methods). Either way a failed conn is dropped so the next call
// re-dials clean. Caller holds c.mu.
func (c *Client) roundtrip(m Method, raw json.RawMessage, resp *Response) error {
if c.conn == nil {
conn, err := netaddr.Dial(c.addr)
//
// The connection carries a deadline derived from ctx, falling back to
// defaultCallTimeout, and a watchdog closes it if ctx is cancelled mid-call
// (V-638). Before that a daemon which stopped answering parked the caller for
// as long as the socket stayed open. The watchdog closes the conn rather than
// calling drop, because drop wants c.mu and the caller is holding it — the
// closed socket fails the read, and roundtrip drops it on the way out.
func (c *Client) roundtrip(ctx context.Context, m Method, raw json.RawMessage, resp *Response) error {
conn := c.currentConn()
if conn == nil {
dialed, err := netaddr.Dial(c.addr)
if err != nil {
return fmt.Errorf("%w: dial %s: %v", errWriteLost, c.addr, err)
}
c.conn = conn
c.setConn(dialed)
conn = dialed
}
if err := writeFrame(c.conn, Request{Method: m, Params: raw}); err != nil {
if dl, ok := ctx.Deadline(); ok {
_ = conn.SetDeadline(dl)
} else {
_ = conn.SetDeadline(time.Now().Add(defaultCallTimeout))
}
defer conn.SetDeadline(time.Time{})
// The watchdog and the end of the call race by construction: a cancellation
// landing just as the reply arrives can close a conn this call is already
// done with, and c.conn would still point at the closed socket. So a call
// whose context ended does not leave the conn behind for the next one,
// whichever of the two got there first.
done := make(chan struct{})
defer func() {
close(done)
if ctx.Err() != nil {
c.drop()
}
}()
go func() {
select {
case <-ctx.Done():
_ = conn.Close()
case <-done:
}
}()
if err := writeFrame(conn, Request{Method: m, Params: raw}); err != nil {
c.drop()
return fmt.Errorf("%w: %v", errWriteLost, err)
}
if err := readFrame(c.conn, resp); err != nil {
if err := readFrame(conn, resp); err != nil {
c.drop()
return fmt.Errorf("%w: %v", errReadLost, err)
}
@@ -220,12 +288,26 @@ func (c *Client) roundtrip(m Method, raw json.RawMessage, resp *Response) error
// drop closes and forgets the current conn so the next call re-dials.
func (c *Client) drop() {
c.connMu.Lock()
defer c.connMu.Unlock()
if c.conn != nil {
_ = c.conn.Close()
c.conn = nil
}
}
func (c *Client) currentConn() net.Conn {
c.connMu.Lock()
defer c.connMu.Unlock()
return c.conn
}
func (c *Client) setConn(conn net.Conn) {
c.connMu.Lock()
defer c.connMu.Unlock()
c.conn = conn
}
// hydrate rehydrates a wire RpcError into the matching package sentinel. The
// code↔sentinel table is the only place the wire "knows" about errors; keep it
// in sync with codeOf in wire.go.
+36 -8
View File
@@ -17,9 +17,11 @@ import (
// Server — the core side of the boundary. Listens on a unix domain socket,
// accepts module connections, frames requests to a CoreAPI and responses back.
// One Server per daemon process; concurrent connections are handled in their
// own goroutine but share the single CoreAPI (and therefore the single store
// writer — store is single-connection, SetMaxOpenConns(1), so serialization is
// already guaranteed at the db; the Server adds no locking of its own).
// own goroutine but share the single CoreAPI, and so the single store writer.
// The store opens at SetMaxOpenConns(1), so serialisation is already guaranteed
// at the database and the Server adds no locking of its own. That cap is an
// invariant this comment depends on, measured and kept on 07-08-2026 (V-642,
// docs/evals/2026-08-07-store-connection-cap.md).
type Server struct {
api atomic.Value // stores CoreAPI
path string
@@ -30,6 +32,14 @@ type Server struct {
done chan struct{}
accept sync.Mutex // guards wg.Add vs Close's wg.Wait sequence
// ctx — server-scoped, cancelled by Close, and the parent of every request
// context. serveConn dispatched under context.Background() until V-638, so
// a dispatch in flight during shutdown could not be told to stop and the
// closeGrace below could only abandon it. Cancelling gives a handler that
// respects its context the chance to return instead.
ctx context.Context
cancel context.CancelFunc
// conns — every accepted connection still being served. Close needs these
// because closing the listener does nothing to a connection already
// accepted: serveConn is parked in readFrame waiting for a peer that may
@@ -208,11 +218,14 @@ func Listen(path string, api CoreAPI) (*Server, error) {
if err != nil {
return nil, err
}
ctx, cancel := context.WithCancel(context.Background())
s := &Server{
path: path,
addr: addr,
ln: ln,
done: make(chan struct{}),
path: path,
addr: addr,
ln: ln,
done: make(chan struct{}),
ctx: ctx,
cancel: cancel,
}
s.api.Store(api)
return s, nil
@@ -250,7 +263,10 @@ func (s *Server) Serve() error {
func (s *Server) serveConn(c net.Conn) {
caller, callerOK := peerCaller(c)
ctx := context.Background()
// Derived from the server's, so Close cancels a dispatch in flight, and
// cancelled when this conn ends so nothing a handler spawned outlives it.
ctx, cancel := context.WithCancel(s.serverContext())
defer cancel()
if callerOK {
ctx = WithCaller(ctx, caller)
}
@@ -274,6 +290,15 @@ func (s *Server) serveConn(c net.Conn) {
}
}
// serverContext is s.ctx, or Background for a Server built as a zero value
// rather than by Listen (the wiring tests do that).
func (s *Server) serverContext() context.Context {
if s.ctx == nil {
return context.Background()
}
return s.ctx
}
func (s *Server) safeDispatch(ctx context.Context, req Request) (result json.RawMessage, err error) {
defer func() {
if r := recover(); r != nil {
@@ -720,6 +745,9 @@ func (s *Server) Close() error {
default:
close(s.done)
}
if s.cancel != nil {
s.cancel()
}
err := s.ln.Close()
// Closing the listener stops new connections; it does nothing to the ones
// already accepted. Close those too, or every serveConn parked in readFrame
+37 -6
View File
@@ -89,8 +89,8 @@ type Poller struct {
ranker Ranker
cfg Config
nextDue map[string]time.Time
seen map[string]map[string]bool // feed → item ID, for items with no date
polled map[string]bool // feed → polled at least once in THIS process
seen map[string]*seenIDs // feed → item IDs, for items with no date
polled map[string]bool // feed → polled at least once in THIS process
}
// NewPoller wires a poller. Returns nil when there is nothing to poll — a
@@ -122,7 +122,7 @@ func NewPoller(feeds []FeedConfig, fetch Fetcher, notes Notes, marks Marks, embe
feeds: valid, fetch: fetch, notes: notes, marks: marks,
embed: embed, ranker: ranker, cfg: cfg,
nextDue: map[string]time.Time{},
seen: map[string]map[string]bool{},
seen: map[string]*seenIDs{},
polled: map[string]bool{},
}
}
@@ -283,6 +283,38 @@ func (p *Poller) mark(ctx context.Context, feed string, now time.Time) (time.Tim
return at, true
}
// maxSeenPerFeed bounds the undated-item set. It has to stay comfortably above
// any one feed's front page, or an item still listed there would fall out of the
// set and be written a second time. A few hundred entries covers the largest
// page anyone publishes, and the set only has to span one poll window plus the
// resync guard, not all of history.
const maxSeenPerFeed = 512
// seenIDs is a bounded insertion-ordered set. The map answers the lookup, the
// slice remembers what to drop first, so an undated feed cannot grow the poller
// for as long as mavend runs.
type seenIDs struct {
ids map[string]bool
order []string
}
// add records id and reports whether it was new.
func (s *seenIDs) add(id string) bool {
if s.ids == nil {
s.ids = make(map[string]bool, maxSeenPerFeed)
}
if s.ids[id] {
return false
}
s.ids[id] = true
s.order = append(s.order, id)
if len(s.order) > maxSeenPerFeed {
delete(s.ids, s.order[0])
s.order = s.order[1:]
}
return true
}
// fresh — two dedup rules, because feeds are inconsistent about dates. A dated
// item must be newer than the mark; an undated one is kept once per process by
// ID.
@@ -308,12 +340,11 @@ func (p *Poller) fresh(f FeedConfig, it Item, mark, now time.Time, resync bool)
id = it.Title
}
if p.seen[f.Name] == nil {
p.seen[f.Name] = map[string]bool{}
p.seen[f.Name] = &seenIDs{}
}
if p.seen[f.Name][id] {
if !p.seen[f.Name].add(id) {
return false
}
p.seen[f.Name][id] = true
return !resync
}
+25
View File
@@ -3,6 +3,7 @@ package rss
import (
"context"
"errors"
"fmt"
"strings"
"testing"
"time"
@@ -208,3 +209,27 @@ func TestNoFeedsMeansNoPoller(t *testing.T) {
t.Fatal("a feed with no name or url is not a configuration")
}
}
// An undated feed used to grow p.seen for as long as mavend ran. The set is
// bounded now, and the bound must not cost the dedupe an item still on the
// front page — only ids far older than any page fall out.
func TestSeenIDsBounded(t *testing.T) {
var s seenIDs
for i := 0; i < maxSeenPerFeed*3; i++ {
if !s.add(fmt.Sprintf("item-%d", i)) {
t.Fatalf("item-%d read as already seen", i)
}
if len(s.ids) > maxSeenPerFeed || len(s.order) > maxSeenPerFeed {
t.Fatalf("after %d inserts: ids=%d order=%d, cap is %d",
i+1, len(s.ids), len(s.order), maxSeenPerFeed)
}
}
// The newest insert is still deduped; the oldest was evicted.
last := fmt.Sprintf("item-%d", maxSeenPerFeed*3-1)
if s.add(last) {
t.Fatalf("%s read as new, so the most recent id was dropped", last)
}
if !s.add("item-0") {
t.Fatal("item-0 survived, so nothing was evicted")
}
}
+154
View File
@@ -0,0 +1,154 @@
package store
import (
"context"
"database/sql"
"fmt"
"path/filepath"
"sort"
"sync"
"sync/atomic"
"testing"
"time"
)
// conncap_test.go measures whether a read queues behind a write at
// SetMaxOpenConns(1), which is what openAt sets (V-642). It is a measurement
// harness, not an assertion: the numbers it prints are the evidence, and the
// decision to move the cap or leave it belongs in docs/evals.
//
// Run it with -v, and note that it is skipped under -short because it spends
// seconds on purpose.
// openCapped opens a plaintext store at the given connection cap. In-package,
// so it can reach the handle openAt caps at 1.
func openCapped(t *testing.T, cap int) *Store {
t.Helper()
path := filepath.Join(t.TempDir(), "cap.db")
db, err := openAt(context.Background(), path)
if err != nil {
t.Fatalf("openAt: %v", err)
}
db.SetMaxOpenConns(cap)
s := &Store{db: db}
t.Cleanup(func() { _ = s.Close() })
return s
}
func percentile(d []time.Duration, p float64) time.Duration {
if len(d) == 0 {
return 0
}
i := int(float64(len(d)-1) * p)
return d[i]
}
// seedFacts writes n facts so a read has rows to decode.
func seedFacts(t *testing.T, s *Store, n int) {
t.Helper()
ctx := context.Background()
now := time.Now().UTC()
for i := 0; i < n; i++ {
key := fmt.Sprintf("seed_%d", i)
if _, err := s.SetValue(ctx, KindSelf, key, "tap:test",
map[string]int{"ml": i}, now.Add(time.Duration(i)*time.Millisecond)); err != nil {
t.Fatalf("seed %d: %v", i, err)
}
}
}
// measureReadsUnderWrites reports read latency percentiles while a writer
// writes at a fixed pace. The pace matters: an unpaced writer completes a
// different number of writes at each cap, because at a higher cap it competes
// with the readers for the write lock instead of taking turns on one
// connection. Two runs that did different work cannot be compared.
// It runs for a fixed wall-clock window rather than a fixed read count, so the
// paced writer does the same work at every cap. Tying the window to a read
// count made the faster configuration receive fewer writes.
func measureReadsUnderWrites(t *testing.T, s *Store, window, pace time.Duration) []time.Duration {
t.Helper()
ctx := context.Background()
var stop atomic.Bool
var writes atomic.Int64
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
now := time.Now().UTC()
for i := 0; !stop.Load(); i++ {
key := fmt.Sprintf("hot_%d", i%16)
if _, err := s.SetValue(ctx, KindSelf, key, "tap:test",
map[string]int{"n": i}, now.Add(time.Duration(i)*time.Millisecond)); err != nil {
t.Errorf("write: %v", err)
return
}
writes.Add(1)
time.Sleep(pace)
}
}()
var lat []time.Duration
deadline := time.Now().Add(window)
for time.Now().Before(deadline) {
start := time.Now()
if _, err := s.RecentFacts(ctx, 50); err != nil {
t.Fatalf("RecentFacts: %v", err)
}
lat = append(lat, time.Since(start))
}
stop.Store(true)
wg.Wait()
t.Logf("in %v: %d reads, %d writes", window, len(lat), writes.Load())
sort.Slice(lat, func(i, j int) bool { return lat[i] < lat[j] })
return lat
}
// TestConnCap_ReadLatencyUnderWrites is the V-642 measurement: read latency at
// cap 1 against cap 4, same workload, same schema, same driver.
func TestConnCap_ReadLatencyUnderWrites(t *testing.T) {
if testing.Short() {
t.Skip("measurement harness; runs for seconds")
}
for _, cap := range []int{1, 4} {
t.Run(fmt.Sprintf("cap=%d", cap), func(t *testing.T) {
s := openCapped(t, cap)
seedFacts(t, s, 500)
lat := measureReadsUnderWrites(t, s, 2*time.Second, 2*time.Millisecond)
t.Logf("cap=%d reads=%d p50=%v p95=%v max=%v",
cap, len(lat), percentile(lat, 0.50), percentile(lat, 0.95), lat[len(lat)-1])
})
}
}
// TestConnCap_ReadBlocksBehindOpenSnapshot is the sharper claim: at cap 1 an
// open read-only transaction holds the only connection, so an unrelated read
// cannot proceed until it commits. This is why the store exposes no way to
// begin one — `Store.DB` used to, and was deleted in V-642 with no caller. The
// test stays as the reason, so re-adding that seam fails a measurement rather
// than shipping a stall.
func TestConnCap_ReadBlocksBehindOpenSnapshot(t *testing.T) {
if testing.Short() {
t.Skip("measurement harness; waits on a timeout")
}
for _, cap := range []int{1, 4} {
t.Run(fmt.Sprintf("cap=%d", cap), func(t *testing.T) {
s := openCapped(t, cap)
seedFacts(t, s, 50)
tx, err := s.db.BeginTx(context.Background(), &sql.TxOptions{ReadOnly: true})
if err != nil {
t.Fatalf("BeginTx: %v", err)
}
defer func() { _ = tx.Rollback() }()
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
start := time.Now()
_, err = s.RecentFacts(ctx, 10)
t.Logf("cap=%d read alongside an open snapshot: waited %v, err=%v",
cap, time.Since(start).Round(time.Millisecond), err)
})
}
}
+135 -11
View File
@@ -68,6 +68,13 @@ func (m *MemoryStore) Insert(ctx context.Context, id string, vec []float32, meta
// Rows under memory.NonRecallPrefix are excluded in SQL. They are speaker
// voiceprints sharing this table, and note recall must not rank them; see that
// constant for why the previous arrangement only appeared to do this.
//
// Every row is still scored, because a full scan is what picks the winners.
// What the scan does NOT do is pay for a row it is about to discard: the score
// is read straight off the stored bytes without materializing a []float32, and
// the meta blob is copied and unmarshalled only for a row that has entered the
// topK. Losers cost one dot product and nothing else. Ranking is unchanged —
// same scores, same order, same ties.
func (m *MemoryStore) Search(ctx context.Context, vec []float32, topK int) ([]memory.Result, error) {
if topK <= 0 {
topK = 10
@@ -80,30 +87,127 @@ func (m *MemoryStore) Search(ctx context.Context, vec []float32, topK int) ([]me
}
defer rows.Close()
var out []memory.Result
// sql.RawBytes hands us the driver's own buffer, valid only until the next
// Next(). Nothing here outlives the row except what topK.offer copies on a
// survivor, so the three columns cost no allocation per row.
var id, blob, metaJSON sql.RawBytes
top := newTopK(topK)
for rows.Next() {
var id, metaJSON string
var blob []byte
if err := rows.Scan(&id, &blob, &metaJSON); err != nil {
return nil, fmt.Errorf("memory: row: %w", err)
}
meta := map[string]string{}
if err := json.Unmarshal([]byte(metaJSON), &meta); err != nil {
return nil, fmt.Errorf("memory: unmarshal meta for %q: %w", id, err)
}
out = append(out, memory.Result{ID: id, Score: dot(vec, decodeVec(blob)), Meta: meta})
top.offer(dotBlob(vec, blob), id, metaJSON)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("memory: rows: %w", err)
}
sort.Slice(out, func(i, j int) bool { return out[i].Score > out[j].Score })
if topK < len(out) {
out = out[:topK]
survivors := top.sorted()
out := make([]memory.Result, 0, len(survivors))
for _, c := range survivors {
meta := map[string]string{}
if err := json.Unmarshal(c.meta, &meta); err != nil {
return nil, fmt.Errorf("memory: unmarshal meta for %q: %w", c.id, err)
}
out = append(out, memory.Result{ID: c.id, Score: c.score, Meta: meta})
}
return out, nil
}
// candidate is one row that is currently in the topK: its score, its id, and
// its meta blob copied out of the driver's buffer. The copy is the price of
// surviving, and only survivors pay it.
type candidate struct {
score float64
id string
meta []byte
}
// topK keeps the k highest-scoring candidates seen so far as a min-heap, so the
// weakest survivor is always heap[0] and one comparison decides whether a new
// row is worth copying. k is 10 in practice, so the heap is tiny and the whole
// structure fits in cache.
//
// It is a plain slice with hand-written sift operations rather than
// container/heap, because that interface boxes every element into an `any` on
// Push and costs an allocation per surviving row.
type topK struct {
k int
heap []candidate
}
func newTopK(k int) *topK {
return &topK{k: k, heap: make([]candidate, 0, k)}
}
// offer admits a row if it beats the weakest survivor, or if the heap is not
// full yet. id and meta are the driver's buffers and are copied here, never
// retained.
//
// A row that only ties the weakest survivor does not displace it, so among
// equal scores the earliest k rows are kept. The full sort this replaced used
// sort.Slice, which is not stable, so it broke such a tie arbitrarily. That is
// the ONE observable difference between the two, and it is deliberate:
// deterministic beats arbitrary.
//
// It is not academic. Under the real embedder an exact tie means duplicate
// vectors and nothing in the recall eval moved (V-643). Under the hash
// embedder the eval's deterministic floor uses, ties are everywhere — it is
// bag-of-words, so every note sharing no word with the query scores exactly 0
// — and recall@3 on that run moved 74.1% to 81.5% purely because the zeros now
// come out in a fixed order. Neither number measures retrieval. recall@1 and
// false recall, which the eval actually asserts, are unchanged on both runs.
func (t *topK) offer(score float64, id, meta []byte) {
if t.k == 0 {
return
}
if len(t.heap) < t.k {
t.heap = append(t.heap, candidate{score: score, id: string(id), meta: append([]byte(nil), meta...)})
t.up(len(t.heap) - 1)
return
}
if score <= t.heap[0].score {
return
}
t.heap[0] = candidate{score: score, id: string(id), meta: append([]byte(nil), meta...)}
t.down(0)
}
func (t *topK) up(i int) {
for i > 0 {
parent := (i - 1) / 2
if t.heap[parent].score <= t.heap[i].score {
return
}
t.heap[parent], t.heap[i] = t.heap[i], t.heap[parent]
i = parent
}
}
func (t *topK) down(i int) {
for {
l, r, small := 2*i+1, 2*i+2, i
if l < len(t.heap) && t.heap[l].score < t.heap[small].score {
small = l
}
if r < len(t.heap) && t.heap[r].score < t.heap[small].score {
small = r
}
if small == i {
return
}
t.heap[small], t.heap[i] = t.heap[i], t.heap[small]
i = small
}
}
// sorted drains the heap into descending score order — what Search returns.
func (t *topK) sorted() []candidate {
out := t.heap
sort.Slice(out, func(i, j int) bool { return out[i].score > out[j].score })
return out
}
// ByPrefix returns every row whose id starts with prefix, vectors included.
//
// This is not a similarity query and deliberately does not score anything:
@@ -240,6 +344,26 @@ func decodeVec(b []byte) []float32 {
return v
}
// dotBlob is dot against a vector still in its stored encoding, so scoring a
// row the query is about to discard does not allocate the []float32 that
// decodeVec would build. Same arithmetic, same order of operations, so it
// returns bit-identical scores to dot(a, decodeVec(b)).
//
// A blob whose length isn't a multiple of 4 is truncated to the whole-element
// prefix, matching decodeVec, and a length mismatch is 0, matching dot.
func dotBlob(a []float32, b []byte) float64 {
n := len(b) / 4
if len(a) != n || n == 0 {
return 0
}
var sum float64
for i := 0; i < n; i++ {
f := math.Float32frombits(binary.LittleEndian.Uint32(b[4*i:]))
sum += float64(a[i]) * float64(f)
}
return sum
}
// dot is the cosine similarity for L2-normalized vectors (mismatched lengths ⇒
// 0, matching internal/memory's cosine).
func dot(a, b []float32) float64 {
+77
View File
@@ -0,0 +1,77 @@
package store
import (
"context"
"fmt"
"math"
"math/rand"
"path/filepath"
"testing"
)
// benchDim is the resident embedder's width (multilingual-e5-small, 384), so
// the per-row decode cost the benchmark measures is the real one.
const benchDim = 384
// seedMemVectors fills a fresh store with n L2-normalized rows carrying a meta
// blob the size recall actually stores — the note text plus its type — because
// the cost this benchmark exists to measure is unmarshalling that blob for
// every row when only topK survivors need it.
func seedMemVectors(tb testing.TB, n int) *MemoryStore {
tb.Helper()
path := filepath.Join(tb.TempDir(), "mem_bench.db")
st, err := Open(context.Background(), path)
if err != nil {
tb.Fatalf("Open: %v", err)
}
tb.Cleanup(func() { _ = st.Close() })
m := st.VectorMemory()
rng := rand.New(rand.NewSource(1))
ctx := context.Background()
for i := 0; i < n; i++ {
if err := m.Insert(ctx, fmt.Sprintf("note:%d", i), randUnitVec(rng, benchDim), map[string]string{
"type": "note",
"text": fmt.Sprintf("заметка номер %d о том, что надо не забыть сделать на неделе", i),
}); err != nil {
tb.Fatalf("Insert %d: %v", i, err)
}
}
return m
}
func randUnitVec(rng *rand.Rand, dim int) []float32 {
v := make([]float32, dim)
var norm float64
for i := range v {
f := rng.NormFloat64()
v[i] = float32(f)
norm += f * f
}
norm = math.Sqrt(norm)
for i := range v {
v[i] = float32(float64(v[i]) / norm)
}
return v
}
// BenchmarkMemoryStoreSearch measures one recall query against a store of n
// rows. Row counts bracket the documented scale: 1000 is a plausible today,
// 10000 is the "thousands, not millions" ceiling the type doc claims a full
// scan is fine at.
func BenchmarkMemoryStoreSearch(b *testing.B) {
for _, n := range []int{1000, 10000} {
b.Run(fmt.Sprintf("rows=%d", n), func(b *testing.B) {
m := seedMemVectors(b, n)
q := randUnitVec(rand.New(rand.NewSource(2)), benchDim)
ctx := context.Background()
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
if _, err := m.Search(ctx, q, 10); err != nil {
b.Fatal(err)
}
}
})
}
}
+107
View File
@@ -0,0 +1,107 @@
package store
import (
"context"
"fmt"
"math/rand"
"sort"
"testing"
"github.com/kami/maven/internal/memory"
)
// naiveSearch is the implementation Search replaced: score every row into a
// slice, sort the whole slice, truncate. It stays in the test file as the
// reference the bounded-heap version is judged against, because "recall must
// not change" is a claim about output, not about the code that produces it.
func naiveSearch(t *testing.T, m *MemoryStore, vec []float32, topK int) []memory.Result {
t.Helper()
rows, err := m.db.QueryContext(context.Background(),
`SELECT id, vec FROM memory_vectors WHERE id NOT LIKE ? ESCAPE '\'`,
escapeLike(memory.NonRecallPrefix)+"%")
if err != nil {
t.Fatalf("naive scan: %v", err)
}
defer rows.Close()
var out []memory.Result
for rows.Next() {
var id string
var blob []byte
if err := rows.Scan(&id, &blob); err != nil {
t.Fatalf("naive row: %v", err)
}
out = append(out, memory.Result{ID: id, Score: dot(vec, decodeVec(blob))})
}
if err := rows.Err(); err != nil {
t.Fatalf("naive rows: %v", err)
}
sort.Slice(out, func(i, j int) bool { return out[i].Score > out[j].Score })
if topK < len(out) {
out = out[:topK]
}
return out
}
// TestMemoryStoreSearchMatchesNaive is the constraint on V-643: the bounded
// heap must return exactly what a full scan and sort returned. Distinct random
// vectors, so no two scores tie and the ranking is total — a mismatch here is
// arithmetic or heap logic, not a tie-break difference.
func TestMemoryStoreSearchMatchesNaive(t *testing.T) {
ctx := context.Background()
m := newMemTestStore(t).VectorMemory()
rng := rand.New(rand.NewSource(7))
const rows, dim = 500, 64
for i := 0; i < rows; i++ {
if err := m.Insert(ctx, fmt.Sprintf("n%d", i), randUnitVec(rng, dim), map[string]string{
"text": fmt.Sprintf("note %d", i),
}); err != nil {
t.Fatalf("Insert %d: %v", i, err)
}
}
for _, topK := range []int{1, 3, 10, 50, rows, rows + 100} {
q := randUnitVec(rng, dim)
got, err := m.Search(ctx, q, topK)
if err != nil {
t.Fatalf("Search topK=%d: %v", topK, err)
}
want := naiveSearch(t, m, q, topK)
if len(got) != len(want) {
t.Fatalf("topK=%d: got %d results, naive returned %d", topK, len(got), len(want))
}
for i := range want {
if got[i].ID != want[i].ID {
t.Errorf("topK=%d rank %d: got %q, naive says %q", topK, i, got[i].ID, want[i].ID)
}
if got[i].Score != want[i].Score {
t.Errorf("topK=%d rank %d (%s): score %v, naive says %v",
topK, i, got[i].ID, got[i].Score, want[i].Score)
}
}
if len(got) > 0 && got[0].Meta["text"] == "" {
t.Errorf("topK=%d: survivor %s has no meta — it was never unmarshalled", topK, got[0].ID)
}
}
}
// TestDotBlobMatchesDot pins the claim in dotBlob's doc comment: reading the
// vector out of its stored bytes is bit-identical to decoding it first. Scores
// feed a gate with a 0.008 margin, so "close enough" is not the bar.
func TestDotBlobMatchesDot(t *testing.T) {
rng := rand.New(rand.NewSource(11))
for i := 0; i < 200; i++ {
a := randUnitVec(rng, 384)
b := randUnitVec(rng, 384)
if got, want := dotBlob(a, encodeVec(b)), dot(a, b); got != want {
t.Fatalf("dotBlob = %v, dot = %v", got, want)
}
}
// Length mismatch is 0 in both, and so is an empty vector.
if got := dotBlob([]float32{1, 0}, encodeVec([]float32{1, 0, 0})); got != 0 {
t.Errorf("mismatched lengths scored %v, want 0", got)
}
if got := dotBlob(nil, nil); got != 0 {
t.Errorf("empty scored %v, want 0", got)
}
}
+14 -8
View File
@@ -95,7 +95,20 @@ func openAt(ctx context.Context, path string) (*sql.DB, error) {
if err != nil {
return nil, fmt.Errorf("open %s: %w", path, err)
}
// single writer expected; the daemon is the only process touching the db.
// One connection, so every statement is serialised at the database and no
// caller above needs a lock of its own. internal/ipc's Server relies on
// exactly this, which is why the cap is an invariant rather than a tuning
// knob: raising it moves the serialisation guarantee somewhere it is not
// written down.
//
// Measured on 07-08-2026 (V-642, docs/evals/2026-08-07-store-connection-cap.md).
// WAL exists to let readers run beside one writer, and the cap gives that
// up, but reads do not queue: p50 594µs against 525µs at a cap of four,
// while write throughput more than halves. The one thing the cap cannot
// survive is a long-lived transaction, which holds the only connection and
// stalls every read for its lifetime. So the store begins none, and
// TestConnCap_ReadBlocksBehindOpenSnapshot is the standing measurement of
// what re-adding one would cost.
db.SetMaxOpenConns(1)
if _, err := db.ExecContext(ctx, schemaSQL); err != nil {
if closeErr := db.Close(); closeErr != nil {
@@ -133,13 +146,6 @@ func (s *Store) Close() error {
return s.enc.closeAndSeal(s.db)
}
// DB exposes the underlying handle for internal read-only snapshots.
// Used by the loop to take a consistent read under a single transaction.
// Modules never receive this handle — core mediates.
func (s *Store) DB(ctx context.Context) (*sql.Tx, error) {
return s.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true})
}
var (
// ErrNoFact — no non-voided row exists for this key.
ErrNoFact = errors.New("store: no fact for key")
+7 -2
View File
@@ -26,6 +26,8 @@
package voice
import (
"context"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/router"
)
@@ -38,8 +40,10 @@ import (
// decision's Intent + Slots + Clarify. The Intent largely names the reply
// shape (act/reminder/fact/note/query/clarify); the Slots carry the
// specifics that personalise it ("got it: water at 14:00").
// The context is the turn's, and it is the only bound an LLM-backed impl has
// besides the phraser timeout (V-638). A floor impl ignores it.
type Replier interface {
Reply(d router.Decision) string
Reply(ctx context.Context, d router.Decision) string
}
// StubReplier — the deterministic, no-model floor. Canned per intent;
@@ -54,7 +58,8 @@ func NewStubReplier() *StubReplier { return &StubReplier{} }
// Reply dispatches on Intent + Clarify. Each branch is short; the LLM impl
// will replace this with prompted text and the same dispatch shape.
func (s *StubReplier) Reply(d router.Decision) string {
// It makes no model call, so the context is unused.
func (s *StubReplier) Reply(_ context.Context, d router.Decision) string {
if d.Clarify {
return "не совсем поняла — можешь переформулировать?"
}
+22
View File
@@ -295,6 +295,27 @@ func (f *Fetcher) checkURL(u *url.URL) error {
return nil
}
// pruneHostsAbove is when pruneLocked bothers to walk the map. Below it the
// walk costs more than the entries do, and `crawl.on_demand` means the host set
// is whatever he names out loud, so it grows slowly.
const pruneHostsAbove = 64
// pruneLocked drops hosts whose last dial is further back than HostInterval.
// Such an entry cannot delay anything — waitTurn would let the next request
// through immediately — so keeping it only holds memory for the life of the
// process. Caller holds f.mu.
func (f *Fetcher) pruneLocked(now time.Time) {
if len(f.last) <= pruneHostsAbove {
return
}
cutoff := now.Add(-f.cfg.HostInterval)
for h, at := range f.last {
if at.Before(cutoff) {
delete(f.last, h)
}
}
}
// waitTurn blocks until this host's rate-limit interval has elapsed. It holds
// no lock while sleeping, so two hosts never wait on each other.
func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
@@ -304,6 +325,7 @@ func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
earliest := f.last[host].Add(f.cfg.HostInterval)
if !now.Before(earliest) {
f.last[host] = now
f.pruneLocked(now)
f.mu.Unlock()
return nil
}
+35
View File
@@ -4,6 +4,7 @@ import (
"bytes"
"context"
"errors"
"fmt"
"io"
"net"
"net/http"
@@ -318,3 +319,37 @@ func TestPostObeysDenylist(t *testing.T) {
t.Fatalf("error = %v, want ErrBlocked", err)
}
}
// f.last used to hold one entry per host ever dialed, for the life of the
// process. A host whose last dial is older than HostInterval cannot delay
// anything, so it is dropped once the map is worth walking.
func TestHostRateMapIsPruned(t *testing.T) {
f := New(Config{HostInterval: time.Minute, AllowPrivate: true})
stale := time.Now().Add(-time.Hour)
for i := 0; i < pruneHostsAbove*2; i++ {
f.last[fmt.Sprintf("h%d.example", i)] = stale
}
// One real turn is what triggers the sweep.
if err := f.waitTurn(context.Background(), "fresh.example"); err != nil {
t.Fatal(err)
}
if len(f.last) != 1 {
t.Fatalf("len(f.last) = %d after the sweep, want 1 (only the host just dialed)", len(f.last))
}
if _, ok := f.last["fresh.example"]; !ok {
t.Fatal("the host just dialed was pruned, so its own rate limit is lost")
}
// A host inside the interval is kept: pruning must not hand out a free turn.
f.last["recent.example"] = time.Now()
for i := 0; i < pruneHostsAbove*2; i++ {
f.last[fmt.Sprintf("g%d.example", i)] = stale
}
if err := f.waitTurn(context.Background(), "other.example"); err != nil {
t.Fatal(err)
}
if _, ok := f.last["recent.example"]; !ok {
t.Fatal("a host dialed inside HostInterval was pruned")
}
}