Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4be6852b94 | |||
| f7b76c572f | |||
| 05ddc5c92e | |||
| b55e68f98d | |||
| beaa24754c | |||
| aed8cac439 | |||
| af4eeceb6a |
@@ -101,9 +101,15 @@ protocol; the config in `deploy/mavend.json` (with `${VAR}` env expansion from g
|
||||
Count against compose, not against the table. Four of the nine daemons are absent, and each
|
||||
absence has a different reason.
|
||||
|
||||
`mavmaild` is commented out in compose, with the reason written beside it: it needs a mail
|
||||
account and this box has none. `mavcaldav` appears nowhere at all, and unlike the other
|
||||
three that is an oversight rather than a decision (V-644).
|
||||
`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`).
|
||||
|
||||
+35
-8
@@ -50,10 +50,10 @@ func run(args []string) error {
|
||||
socket := fs.String("socket", "", "core IPC socket path (required)")
|
||||
url := fs.String("url", "", "CalDAV calendar URL, e.g. http://localhost:5232/kami/personal (required)")
|
||||
user := fs.String("user", "", "CalDAV basic-auth username (required)")
|
||||
pass := fs.String("pass", "", "CalDAV basic-auth password (required)")
|
||||
passFile := fs.String("pass-file", "", "file holding the CalDAV basic-auth password (required — never passed as a flag value)")
|
||||
renderURL := fs.String("render-url", "", "CalDAV collection maven publishes her own reminders to; empty disables rendering")
|
||||
renderUser := fs.String("render-user", "", "basic-auth username for -render-url (defaults to -user)")
|
||||
renderPass := fs.String("render-pass", "", "basic-auth password for -render-url (defaults to -pass)")
|
||||
renderPassFile := fs.String("render-pass-file", "", "file holding the password for -render-url (defaults to -pass-file)")
|
||||
renderDur := fs.Duration("render-duration", calendar.DefaultReminderDuration, "how long a rendered reminder occupies")
|
||||
interval := fs.Duration("interval", 5*time.Minute, "poll cadence")
|
||||
timeout := fs.Duration("timeout", 10*time.Second, "per-request HTTP timeout")
|
||||
@@ -63,13 +63,22 @@ func run(args []string) error {
|
||||
if *socket == "" {
|
||||
return fmt.Errorf("-socket is required")
|
||||
}
|
||||
if *url == "" || *user == "" || *pass == "" {
|
||||
return fmt.Errorf("-url, -user, -pass are required")
|
||||
if *url == "" || *user == "" || *passFile == "" {
|
||||
return fmt.Errorf("-url, -user, -pass-file are required")
|
||||
}
|
||||
if err := checkRenderTarget([]string{*url}, *renderURL); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// The password is read from a file, never taken as a flag value: an argv
|
||||
// secret is visible in `ps` to every user on the box and lands in the compose
|
||||
// file and the shell history. Same rule mavmaild and mavpoll follow. Read
|
||||
// once at start, so a rotated password means a restart.
|
||||
pass, err := readSecret(*passFile)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||
defer stop()
|
||||
|
||||
@@ -85,17 +94,20 @@ func run(args []string) error {
|
||||
http: hc,
|
||||
url: strings.TrimRight(*url, "/"),
|
||||
user: *user,
|
||||
pass: *pass,
|
||||
pass: pass,
|
||||
}
|
||||
|
||||
var rend *renderer
|
||||
if *renderURL != "" {
|
||||
ru, rp := *renderUser, *renderPass
|
||||
ru, rp := *renderUser, pass
|
||||
if ru == "" {
|
||||
ru = *user
|
||||
}
|
||||
if rp == "" {
|
||||
rp = *pass
|
||||
if *renderPassFile != "" {
|
||||
rp, err = readSecret(*renderPassFile)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
rend = newRenderer(core, hc, *renderURL, ru, rp, *renderDur)
|
||||
log.Printf("mavcaldav: rendering reminders to %s", *renderURL)
|
||||
@@ -131,6 +143,21 @@ func run(args []string) error {
|
||||
// It takes the whole read set, not one URL. The guarantee in the package
|
||||
// comment is about every calendar maven reads, and a second read target added
|
||||
// later must not quietly fall outside the check.
|
||||
// readSecret reads one credential from a file and refuses an empty one. An
|
||||
// empty file is a deployment mistake, not a password, and CalDAV basic auth
|
||||
// would send it and get a 401 every poll.
|
||||
func readSecret(path string) (string, error) {
|
||||
raw, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("read password file: %w", err)
|
||||
}
|
||||
secret := strings.TrimSpace(string(raw))
|
||||
if secret == "" {
|
||||
return "", fmt.Errorf("password file %s is empty", path)
|
||||
}
|
||||
return secret, nil
|
||||
}
|
||||
|
||||
func checkRenderTarget(readURLs []string, renderURL string) error {
|
||||
if renderURL == "" {
|
||||
return nil
|
||||
|
||||
@@ -5,12 +5,38 @@ import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
|
||||
// The password comes from a file so it never reaches argv. An empty or missing
|
||||
// file must fail at start rather than authenticate as "" against his calendar.
|
||||
func TestReadSecret(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
good := filepath.Join(dir, "ok")
|
||||
if err := os.WriteFile(good, []byte(" hunter2\n"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got, err := readSecret(good); err != nil || got != "hunter2" {
|
||||
t.Fatalf("readSecret(good) = %q, %v; want \"hunter2\", nil", got, err)
|
||||
}
|
||||
|
||||
empty := filepath.Join(dir, "empty")
|
||||
if err := os.WriteFile(empty, []byte("\n \n"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := readSecret(empty); err == nil {
|
||||
t.Fatal("readSecret(empty) = nil error, want refusal")
|
||||
}
|
||||
if _, err := readSecret(filepath.Join(dir, "absent")); err == nil {
|
||||
t.Fatal("readSecret(absent) = nil error, want refusal")
|
||||
}
|
||||
}
|
||||
|
||||
type fakeCore struct {
|
||||
ipc.UnimplementedCoreAPI
|
||||
facts map[string]ipc.Fact // composite key "key|source" → Fact
|
||||
|
||||
@@ -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.
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
+37
-6
@@ -89,8 +89,8 @@ type Poller struct {
|
||||
ranker Ranker
|
||||
cfg Config
|
||||
nextDue map[string]time.Time
|
||||
seen map[string]map[string]bool // feed → item ID, for items with no date
|
||||
polled map[string]bool // feed → polled at least once in THIS process
|
||||
seen map[string]*seenIDs // feed → item IDs, for items with no date
|
||||
polled map[string]bool // feed → polled at least once in THIS process
|
||||
}
|
||||
|
||||
// NewPoller wires a poller. Returns nil when there is nothing to poll — a
|
||||
@@ -122,7 +122,7 @@ func NewPoller(feeds []FeedConfig, fetch Fetcher, notes Notes, marks Marks, embe
|
||||
feeds: valid, fetch: fetch, notes: notes, marks: marks,
|
||||
embed: embed, ranker: ranker, cfg: cfg,
|
||||
nextDue: map[string]time.Time{},
|
||||
seen: map[string]map[string]bool{},
|
||||
seen: map[string]*seenIDs{},
|
||||
polled: map[string]bool{},
|
||||
}
|
||||
}
|
||||
@@ -283,6 +283,38 @@ func (p *Poller) mark(ctx context.Context, feed string, now time.Time) (time.Tim
|
||||
return at, true
|
||||
}
|
||||
|
||||
// maxSeenPerFeed bounds the undated-item set. It has to stay comfortably above
|
||||
// any one feed's front page, or an item still listed there would fall out of the
|
||||
// set and be written a second time. A few hundred entries covers the largest
|
||||
// page anyone publishes, and the set only has to span one poll window plus the
|
||||
// resync guard, not all of history.
|
||||
const maxSeenPerFeed = 512
|
||||
|
||||
// seenIDs is a bounded insertion-ordered set. The map answers the lookup, the
|
||||
// slice remembers what to drop first, so an undated feed cannot grow the poller
|
||||
// for as long as mavend runs.
|
||||
type seenIDs struct {
|
||||
ids map[string]bool
|
||||
order []string
|
||||
}
|
||||
|
||||
// add records id and reports whether it was new.
|
||||
func (s *seenIDs) add(id string) bool {
|
||||
if s.ids == nil {
|
||||
s.ids = make(map[string]bool, maxSeenPerFeed)
|
||||
}
|
||||
if s.ids[id] {
|
||||
return false
|
||||
}
|
||||
s.ids[id] = true
|
||||
s.order = append(s.order, id)
|
||||
if len(s.order) > maxSeenPerFeed {
|
||||
delete(s.ids, s.order[0])
|
||||
s.order = s.order[1:]
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// fresh — two dedup rules, because feeds are inconsistent about dates. A dated
|
||||
// item must be newer than the mark; an undated one is kept once per process by
|
||||
// ID.
|
||||
@@ -308,12 +340,11 @@ func (p *Poller) fresh(f FeedConfig, it Item, mark, now time.Time, resync bool)
|
||||
id = it.Title
|
||||
}
|
||||
if p.seen[f.Name] == nil {
|
||||
p.seen[f.Name] = map[string]bool{}
|
||||
p.seen[f.Name] = &seenIDs{}
|
||||
}
|
||||
if p.seen[f.Name][id] {
|
||||
if !p.seen[f.Name].add(id) {
|
||||
return false
|
||||
}
|
||||
p.seen[f.Name][id] = true
|
||||
return !resync
|
||||
}
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@ package rss
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -208,3 +209,27 @@ func TestNoFeedsMeansNoPoller(t *testing.T) {
|
||||
t.Fatal("a feed with no name or url is not a configuration")
|
||||
}
|
||||
}
|
||||
|
||||
// An undated feed used to grow p.seen for as long as mavend ran. The set is
|
||||
// bounded now, and the bound must not cost the dedupe an item still on the
|
||||
// front page — only ids far older than any page fall out.
|
||||
func TestSeenIDsBounded(t *testing.T) {
|
||||
var s seenIDs
|
||||
for i := 0; i < maxSeenPerFeed*3; i++ {
|
||||
if !s.add(fmt.Sprintf("item-%d", i)) {
|
||||
t.Fatalf("item-%d read as already seen", i)
|
||||
}
|
||||
if len(s.ids) > maxSeenPerFeed || len(s.order) > maxSeenPerFeed {
|
||||
t.Fatalf("after %d inserts: ids=%d order=%d, cap is %d",
|
||||
i+1, len(s.ids), len(s.order), maxSeenPerFeed)
|
||||
}
|
||||
}
|
||||
// The newest insert is still deduped; the oldest was evicted.
|
||||
last := fmt.Sprintf("item-%d", maxSeenPerFeed*3-1)
|
||||
if s.add(last) {
|
||||
t.Fatalf("%s read as new, so the most recent id was dropped", last)
|
||||
}
|
||||
if !s.add("item-0") {
|
||||
t.Fatal("item-0 survived, so nothing was evicted")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,154 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// conncap_test.go measures whether a read queues behind a write at
|
||||
// SetMaxOpenConns(1), which is what openAt sets (V-642). It is a measurement
|
||||
// harness, not an assertion: the numbers it prints are the evidence, and the
|
||||
// decision to move the cap or leave it belongs in docs/evals.
|
||||
//
|
||||
// Run it with -v, and note that it is skipped under -short because it spends
|
||||
// seconds on purpose.
|
||||
|
||||
// openCapped opens a plaintext store at the given connection cap. In-package,
|
||||
// so it can reach the handle openAt caps at 1.
|
||||
func openCapped(t *testing.T, cap int) *Store {
|
||||
t.Helper()
|
||||
path := filepath.Join(t.TempDir(), "cap.db")
|
||||
db, err := openAt(context.Background(), path)
|
||||
if err != nil {
|
||||
t.Fatalf("openAt: %v", err)
|
||||
}
|
||||
db.SetMaxOpenConns(cap)
|
||||
s := &Store{db: db}
|
||||
t.Cleanup(func() { _ = s.Close() })
|
||||
return s
|
||||
}
|
||||
|
||||
func percentile(d []time.Duration, p float64) time.Duration {
|
||||
if len(d) == 0 {
|
||||
return 0
|
||||
}
|
||||
i := int(float64(len(d)-1) * p)
|
||||
return d[i]
|
||||
}
|
||||
|
||||
// seedFacts writes n facts so a read has rows to decode.
|
||||
func seedFacts(t *testing.T, s *Store, n int) {
|
||||
t.Helper()
|
||||
ctx := context.Background()
|
||||
now := time.Now().UTC()
|
||||
for i := 0; i < n; i++ {
|
||||
key := fmt.Sprintf("seed_%d", i)
|
||||
if _, err := s.SetValue(ctx, KindSelf, key, "tap:test",
|
||||
map[string]int{"ml": i}, now.Add(time.Duration(i)*time.Millisecond)); err != nil {
|
||||
t.Fatalf("seed %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// measureReadsUnderWrites reports read latency percentiles while a writer
|
||||
// writes at a fixed pace. The pace matters: an unpaced writer completes a
|
||||
// different number of writes at each cap, because at a higher cap it competes
|
||||
// with the readers for the write lock instead of taking turns on one
|
||||
// connection. Two runs that did different work cannot be compared.
|
||||
// It runs for a fixed wall-clock window rather than a fixed read count, so the
|
||||
// paced writer does the same work at every cap. Tying the window to a read
|
||||
// count made the faster configuration receive fewer writes.
|
||||
func measureReadsUnderWrites(t *testing.T, s *Store, window, pace time.Duration) []time.Duration {
|
||||
t.Helper()
|
||||
ctx := context.Background()
|
||||
|
||||
var stop atomic.Bool
|
||||
var writes atomic.Int64
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
now := time.Now().UTC()
|
||||
for i := 0; !stop.Load(); i++ {
|
||||
key := fmt.Sprintf("hot_%d", i%16)
|
||||
if _, err := s.SetValue(ctx, KindSelf, key, "tap:test",
|
||||
map[string]int{"n": i}, now.Add(time.Duration(i)*time.Millisecond)); err != nil {
|
||||
t.Errorf("write: %v", err)
|
||||
return
|
||||
}
|
||||
writes.Add(1)
|
||||
time.Sleep(pace)
|
||||
}
|
||||
}()
|
||||
|
||||
var lat []time.Duration
|
||||
deadline := time.Now().Add(window)
|
||||
for time.Now().Before(deadline) {
|
||||
start := time.Now()
|
||||
if _, err := s.RecentFacts(ctx, 50); err != nil {
|
||||
t.Fatalf("RecentFacts: %v", err)
|
||||
}
|
||||
lat = append(lat, time.Since(start))
|
||||
}
|
||||
stop.Store(true)
|
||||
wg.Wait()
|
||||
t.Logf("in %v: %d reads, %d writes", window, len(lat), writes.Load())
|
||||
|
||||
sort.Slice(lat, func(i, j int) bool { return lat[i] < lat[j] })
|
||||
return lat
|
||||
}
|
||||
|
||||
// TestConnCap_ReadLatencyUnderWrites is the V-642 measurement: read latency at
|
||||
// cap 1 against cap 4, same workload, same schema, same driver.
|
||||
func TestConnCap_ReadLatencyUnderWrites(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("measurement harness; runs for seconds")
|
||||
}
|
||||
for _, cap := range []int{1, 4} {
|
||||
t.Run(fmt.Sprintf("cap=%d", cap), func(t *testing.T) {
|
||||
s := openCapped(t, cap)
|
||||
seedFacts(t, s, 500)
|
||||
lat := measureReadsUnderWrites(t, s, 2*time.Second, 2*time.Millisecond)
|
||||
t.Logf("cap=%d reads=%d p50=%v p95=%v max=%v",
|
||||
cap, len(lat), percentile(lat, 0.50), percentile(lat, 0.95), lat[len(lat)-1])
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestConnCap_ReadBlocksBehindOpenSnapshot is the sharper claim: at cap 1 an
|
||||
// open read-only transaction holds the only connection, so an unrelated read
|
||||
// cannot proceed until it commits. This is why the store exposes no way to
|
||||
// begin one — `Store.DB` used to, and was deleted in V-642 with no caller. The
|
||||
// test stays as the reason, so re-adding that seam fails a measurement rather
|
||||
// than shipping a stall.
|
||||
func TestConnCap_ReadBlocksBehindOpenSnapshot(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("measurement harness; waits on a timeout")
|
||||
}
|
||||
for _, cap := range []int{1, 4} {
|
||||
t.Run(fmt.Sprintf("cap=%d", cap), func(t *testing.T) {
|
||||
s := openCapped(t, cap)
|
||||
seedFacts(t, s, 50)
|
||||
|
||||
tx, err := s.db.BeginTx(context.Background(), &sql.TxOptions{ReadOnly: true})
|
||||
if err != nil {
|
||||
t.Fatalf("BeginTx: %v", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||
defer cancel()
|
||||
start := time.Now()
|
||||
_, err = s.RecentFacts(ctx, 10)
|
||||
t.Logf("cap=%d read alongside an open snapshot: waited %v, err=%v",
|
||||
cap, time.Since(start).Round(time.Millisecond), err)
|
||||
})
|
||||
}
|
||||
}
|
||||
+14
-8
@@ -95,7 +95,20 @@ func openAt(ctx context.Context, path string) (*sql.DB, error) {
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("open %s: %w", path, err)
|
||||
}
|
||||
// single writer expected; the daemon is the only process touching the db.
|
||||
// One connection, so every statement is serialised at the database and no
|
||||
// caller above needs a lock of its own. internal/ipc's Server relies on
|
||||
// exactly this, which is why the cap is an invariant rather than a tuning
|
||||
// knob: raising it moves the serialisation guarantee somewhere it is not
|
||||
// written down.
|
||||
//
|
||||
// Measured on 07-08-2026 (V-642, docs/evals/2026-08-07-store-connection-cap.md).
|
||||
// WAL exists to let readers run beside one writer, and the cap gives that
|
||||
// up, but reads do not queue: p50 594µs against 525µs at a cap of four,
|
||||
// while write throughput more than halves. The one thing the cap cannot
|
||||
// survive is a long-lived transaction, which holds the only connection and
|
||||
// stalls every read for its lifetime. So the store begins none, and
|
||||
// TestConnCap_ReadBlocksBehindOpenSnapshot is the standing measurement of
|
||||
// what re-adding one would cost.
|
||||
db.SetMaxOpenConns(1)
|
||||
if _, err := db.ExecContext(ctx, schemaSQL); err != nil {
|
||||
if closeErr := db.Close(); closeErr != nil {
|
||||
@@ -133,13 +146,6 @@ func (s *Store) Close() error {
|
||||
return s.enc.closeAndSeal(s.db)
|
||||
}
|
||||
|
||||
// DB exposes the underlying handle for internal read-only snapshots.
|
||||
// Used by the loop to take a consistent read under a single transaction.
|
||||
// Modules never receive this handle — core mediates.
|
||||
func (s *Store) DB(ctx context.Context) (*sql.Tx, error) {
|
||||
return s.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true})
|
||||
}
|
||||
|
||||
var (
|
||||
// ErrNoFact — no non-voided row exists for this key.
|
||||
ErrNoFact = errors.New("store: no fact for key")
|
||||
|
||||
@@ -295,6 +295,27 @@ func (f *Fetcher) checkURL(u *url.URL) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// pruneHostsAbove is when pruneLocked bothers to walk the map. Below it the
|
||||
// walk costs more than the entries do, and `crawl.on_demand` means the host set
|
||||
// is whatever he names out loud, so it grows slowly.
|
||||
const pruneHostsAbove = 64
|
||||
|
||||
// pruneLocked drops hosts whose last dial is further back than HostInterval.
|
||||
// Such an entry cannot delay anything — waitTurn would let the next request
|
||||
// through immediately — so keeping it only holds memory for the life of the
|
||||
// process. Caller holds f.mu.
|
||||
func (f *Fetcher) pruneLocked(now time.Time) {
|
||||
if len(f.last) <= pruneHostsAbove {
|
||||
return
|
||||
}
|
||||
cutoff := now.Add(-f.cfg.HostInterval)
|
||||
for h, at := range f.last {
|
||||
if at.Before(cutoff) {
|
||||
delete(f.last, h)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// waitTurn blocks until this host's rate-limit interval has elapsed. It holds
|
||||
// no lock while sleeping, so two hosts never wait on each other.
|
||||
func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
|
||||
@@ -304,6 +325,7 @@ func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
|
||||
earliest := f.last[host].Add(f.cfg.HostInterval)
|
||||
if !now.Before(earliest) {
|
||||
f.last[host] = now
|
||||
f.pruneLocked(now)
|
||||
f.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
@@ -318,3 +319,37 @@ func TestPostObeysDenylist(t *testing.T) {
|
||||
t.Fatalf("error = %v, want ErrBlocked", err)
|
||||
}
|
||||
}
|
||||
|
||||
// f.last used to hold one entry per host ever dialed, for the life of the
|
||||
// process. A host whose last dial is older than HostInterval cannot delay
|
||||
// anything, so it is dropped once the map is worth walking.
|
||||
func TestHostRateMapIsPruned(t *testing.T) {
|
||||
f := New(Config{HostInterval: time.Minute, AllowPrivate: true})
|
||||
stale := time.Now().Add(-time.Hour)
|
||||
for i := 0; i < pruneHostsAbove*2; i++ {
|
||||
f.last[fmt.Sprintf("h%d.example", i)] = stale
|
||||
}
|
||||
|
||||
// One real turn is what triggers the sweep.
|
||||
if err := f.waitTurn(context.Background(), "fresh.example"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(f.last) != 1 {
|
||||
t.Fatalf("len(f.last) = %d after the sweep, want 1 (only the host just dialed)", len(f.last))
|
||||
}
|
||||
if _, ok := f.last["fresh.example"]; !ok {
|
||||
t.Fatal("the host just dialed was pruned, so its own rate limit is lost")
|
||||
}
|
||||
|
||||
// A host inside the interval is kept: pruning must not hand out a free turn.
|
||||
f.last["recent.example"] = time.Now()
|
||||
for i := 0; i < pruneHostsAbove*2; i++ {
|
||||
f.last[fmt.Sprintf("g%d.example", i)] = stale
|
||||
}
|
||||
if err := f.waitTurn(context.Background(), "other.example"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, ok := f.last["recent.example"]; !ok {
|
||||
t.Fatal("a host dialed inside HostInterval was pruned")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user