Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| be066a4b04 | |||
| ad074cea31 | |||
| 2c1b0eede0 |
@@ -8,6 +8,7 @@
|
||||
/mavcaldav
|
||||
/mavwaked
|
||||
/mavmaild
|
||||
/mavupdate
|
||||
|
||||
# Certs (private keys, don't commit)
|
||||
certs/
|
||||
|
||||
@@ -20,7 +20,7 @@ PIPER_ESPEAK := $(shell pwd)/deps/piper/espeak-ng-data
|
||||
|
||||
all: build
|
||||
|
||||
build: build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav build-mail
|
||||
build: build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav build-mail build-update
|
||||
|
||||
build-stt:
|
||||
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
|
||||
@@ -53,6 +53,12 @@ build-caldav:
|
||||
build-mail:
|
||||
$(GO) build $(GOFLAGS) -o mavmaild ./cmd/mavmaild/
|
||||
|
||||
# mavupdate is an operator CLI, not a daemon: nothing runs it but a human on the
|
||||
# box. It is built with the rest so a broken update path is caught by `make
|
||||
# build` rather than the first time it is needed.
|
||||
build-update:
|
||||
$(GO) build $(GOFLAGS) -o mavupdate ./cmd/mavupdate/
|
||||
|
||||
run-web: build-web
|
||||
./mavweb -addr :9200 -voice 127.0.0.1:9100
|
||||
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/crawl"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/memory"
|
||||
"github.com/kami/maven/internal/morning"
|
||||
@@ -78,6 +79,13 @@ var querySources = []querySource{
|
||||
{"embed", (*reactiveHandler).queryEmbed},
|
||||
{"memory", (*reactiveHandler).queryMemory},
|
||||
{"notes", (*reactiveHandler).queryNotes},
|
||||
// LAST before the model answers from memory, and that position is the whole
|
||||
// design (Vikunja #259): local sources first. The model, his own notes and
|
||||
// facts, and — once internal/kiwix is wired into this chain — the offline
|
||||
// ZIMs all get their turn before anything touches the network. This source
|
||||
// only claims a turn where he named a URL out loud, so it never competes
|
||||
// with a local answer.
|
||||
{"web", (*reactiveHandler).queryWeb},
|
||||
{"general-knowledge", (*reactiveHandler).queryGeneral},
|
||||
}
|
||||
|
||||
@@ -371,6 +379,57 @@ func (h *reactiveHandler) queryNotes(ctx context.Context, t *queryTurn) (string,
|
||||
return reply, true
|
||||
}
|
||||
|
||||
// webPageContextRunes — how much of a fetched page is handed to the phraser.
|
||||
// Less than the crawler keeps: the rest of the 4096-token window belongs to the
|
||||
// prompt, the persona block and the reply.
|
||||
const webPageContextRunes = 1500
|
||||
|
||||
// queryWeb — "посмотри https://example.org/x — что там?" (Vikunja #259).
|
||||
//
|
||||
// It claims a turn ONLY when he named a URL, which is what keeps a fallback from
|
||||
// becoming a habit: no URL, no fetch, and the model answers from what is local.
|
||||
// What leaves the box is the URL and nothing else — no note, no fact, no history
|
||||
// travels with it.
|
||||
func (h *reactiveHandler) queryWeb(ctx context.Context, t *queryTurn) (string, bool) {
|
||||
link, ok := router.FirstURL(t.dec.Utterance)
|
||||
if !ok {
|
||||
return "", false
|
||||
}
|
||||
if h.crawler == nil {
|
||||
// Claim rather than fall through: he asked about a specific page, and
|
||||
// letting the model answer from the URL's spelling alone is how a small
|
||||
// model invents a page's contents.
|
||||
return "я не читаю страницы — это не настроено.", true
|
||||
}
|
||||
ctxFetch, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||
defer cancel()
|
||||
page, err := h.crawler.Page(ctxFetch, link)
|
||||
if err != nil {
|
||||
if errors.Is(err, crawl.ErrRobots) {
|
||||
return "эта страница закрыта для чтения — robots.txt не разрешает.", true
|
||||
}
|
||||
log.Printf("voice: web: %v", err)
|
||||
return "не получилось прочитать страницу.", true
|
||||
}
|
||||
if page.Text == "" {
|
||||
return "страница открылась, но читать там нечего.", true
|
||||
}
|
||||
// The page is handed to the phraser the same way a note is: as context for
|
||||
// the question he actually asked. She answers the question, she does not
|
||||
// recite the page.
|
||||
snippet := page.Title + "\n" + crawl.TrimRunes(page.Text, webPageContextRunes)
|
||||
reply, perr := h.phraser.PhraseQuery(ctx, t.dec.Utterance, []string{snippet})
|
||||
if perr != nil {
|
||||
log.Printf("voice: web: phrase: %v", perr)
|
||||
}
|
||||
if reply == "" {
|
||||
// No phraser (or it failed): read back the top of the page rather than
|
||||
// pretend the fetch did not happen.
|
||||
return "вот что на странице: " + crawl.TrimRunes(page.Text, 300), true
|
||||
}
|
||||
return reply, true
|
||||
}
|
||||
|
||||
// queryGeneral — general knowledge from the phraser, the last source before
|
||||
// giving up. It always claims: either the model answers or Maven says she
|
||||
// doesn't know.
|
||||
|
||||
@@ -0,0 +1,184 @@
|
||||
// mavend/crawls.go — the driver for reading web pages (Vikunja #259,
|
||||
// docs/plans/14-web-crawler.md). The crawler is pure and lives in
|
||||
// internal/crawl; this is the impure half: the guarded fetcher, a ticker for the
|
||||
// scheduled watches, and the fact-backed dedup hashes.
|
||||
//
|
||||
// Two paths, one config block, both off unless configured:
|
||||
//
|
||||
// - ON DEMAND — he names a URL out loud and she reads it. That is the
|
||||
// `queryWeb` source in actions_query.go, LAST in the chain: after his
|
||||
// memory, after the notes, and (once Kiwix is wired into the chain) after
|
||||
// the local ZIMs. A local read costs nothing and leaks nothing; a fetch puts
|
||||
// a URL in someone's log, so it goes last.
|
||||
// - SCHEDULED — a watched page is re-read on its interval, and a page whose
|
||||
// text changed is written as a note. It does NOT announce itself. Same rule
|
||||
// as the feed poller: notes, never nudges.
|
||||
//
|
||||
// Only the URL goes out. Nothing here reads a note, a fact, the persona block or
|
||||
// the history, and internal/crawl has no access to the store at all.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/crawl"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/router"
|
||||
"github.com/kami/maven/internal/webfetch"
|
||||
)
|
||||
|
||||
// newCrawler builds the crawler from the `crawl` block, or returns nil when
|
||||
// there is none. Every caller checks for nil, and nil means no page is ever
|
||||
// fetched.
|
||||
func newCrawler(cfg *config.Config) *crawl.Crawler {
|
||||
if cfg.Crawl == nil {
|
||||
return nil
|
||||
}
|
||||
cc := cfg.Crawl
|
||||
|
||||
hosts := append([]string(nil), cc.AllowHosts...)
|
||||
// A watched page's own host is always reachable; otherwise an allowlist and
|
||||
// a watch list would have to be kept in sync by hand.
|
||||
for _, w := range cc.Watches {
|
||||
if u, err := url.Parse(w.URL); err == nil && u.Hostname() != "" {
|
||||
hosts = append(hosts, u.Hostname())
|
||||
}
|
||||
}
|
||||
// An allowlist plus on-demand is a contradiction worth logging rather than
|
||||
// silently resolving: he asked for arbitrary pages AND for a fixed list.
|
||||
// The allowlist wins, because it is the narrower instruction.
|
||||
if len(hosts) > 0 && cc.OnDemand && len(cc.AllowHosts) > 0 {
|
||||
log.Printf("crawl: allow_hosts is set, so on-demand reading is limited to those hosts")
|
||||
}
|
||||
ua := cc.UserAgent
|
||||
if ua == "" {
|
||||
ua = webfetch.DefaultUserAgent
|
||||
}
|
||||
fetcher := webfetch.New(webfetch.Config{
|
||||
AllowHosts: hosts,
|
||||
DenyHosts: cc.DenyHosts,
|
||||
Timeout: time.Duration(cc.Timeout),
|
||||
MaxBytes: cc.MaxBytes,
|
||||
UserAgent: ua,
|
||||
})
|
||||
// The user-agent handed to the crawler is the one the fetcher sends: obeying
|
||||
// robots rules written for a different name would be a lie.
|
||||
return crawl.New(&crawlFetcher{f: fetcher}, crawl.Config{
|
||||
UserAgent: ua,
|
||||
MaxRunes: cc.MaxRunes,
|
||||
})
|
||||
}
|
||||
|
||||
// onDemandCrawler returns a crawler for the answer path, or nil when on-demand
|
||||
// reading is off. The scheduled watches can be on while this is off: reading a
|
||||
// fixed list of pages on a timer and reading whatever URL is in an utterance are
|
||||
// different permissions, and the config keeps them separate.
|
||||
func onDemandCrawler(cfg *config.Config) *crawl.Crawler {
|
||||
if cfg.Crawl == nil || !cfg.Crawl.OnDemand {
|
||||
return nil
|
||||
}
|
||||
return newCrawler(cfg)
|
||||
}
|
||||
|
||||
// crawlWorker — ticker + watcher for the scheduled half.
|
||||
type crawlWorker struct {
|
||||
watcher *crawl.Watcher
|
||||
interval time.Duration
|
||||
}
|
||||
|
||||
// crawlTickInterval — how often the worker asks what is due. Per-watch cadence
|
||||
// is the watcher's business.
|
||||
const crawlTickInterval = 15 * time.Minute
|
||||
|
||||
// newCrawlWorker wires the scheduled crawls, or nil when nothing is watched.
|
||||
func newCrawlWorker(c *crawl.Crawler, api ipc.CoreAPI, emb router.Embedder, cfg *config.Config) *crawlWorker {
|
||||
if c == nil || cfg.Crawl == nil || len(cfg.Crawl.Watches) == 0 {
|
||||
return nil
|
||||
}
|
||||
watches := make([]crawl.WatchConfig, 0, len(cfg.Crawl.Watches))
|
||||
for _, w := range cfg.Crawl.Watches {
|
||||
watches = append(watches, crawl.WatchConfig{
|
||||
Name: w.Name,
|
||||
URL: w.URL,
|
||||
Interval: time.Duration(w.Interval),
|
||||
})
|
||||
}
|
||||
watcher := crawl.NewWatcher(c, watches, api, &factHashes{api: api},
|
||||
crawlEmbedder(emb), time.Duration(cfg.Crawl.Interval))
|
||||
if watcher == nil {
|
||||
log.Printf("crawl: configured but nothing watchable — scheduled crawls disabled")
|
||||
return nil
|
||||
}
|
||||
log.Printf("crawl: watching %d page(s), checking what is due every %s", len(watches), crawlTickInterval)
|
||||
return &crawlWorker{watcher: watcher, interval: crawlTickInterval}
|
||||
}
|
||||
|
||||
// run checks what is due until ctx is canceled. The first round runs
|
||||
// immediately; it writes notes only, so an early round startles nobody.
|
||||
func (w *crawlWorker) run(ctx context.Context) {
|
||||
w.watcher.CheckDue(ctx, time.Now())
|
||||
t := time.NewTicker(w.interval)
|
||||
defer t.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case now := <-t.C:
|
||||
w.watcher.CheckDue(ctx, now)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// crawlFetcher adapts webfetch to crawl.Fetcher, which is the seam that keeps
|
||||
// net/http out of the crawler package.
|
||||
type crawlFetcher struct{ f *webfetch.Fetcher }
|
||||
|
||||
func (a *crawlFetcher) Get(ctx context.Context, u string) (*crawl.Response, error) {
|
||||
resp, err := a.f.Get(ctx, u)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &crawl.Response{URL: resp.URL, ContentType: resp.ContentType, Body: resp.Body}, nil
|
||||
}
|
||||
|
||||
// factHashes stores each watch's last content hash as a config fact, so a
|
||||
// restart does not re-note an unchanged page. Same mechanism the feed reader
|
||||
// uses for its marks, and inspectable on /dash.
|
||||
type factHashes struct{ api ipc.CoreAPI }
|
||||
|
||||
func hashKey(name string) string { return "crawl:hash:" + name }
|
||||
|
||||
func (h *factHashes) LastHash(ctx context.Context, name string) (string, error) {
|
||||
f, err := h.api.LatestFact(ctx, hashKey(name))
|
||||
if err != nil {
|
||||
// No hash yet is not an error: the watcher treats "" as "never read".
|
||||
return "", nil
|
||||
}
|
||||
return f.Value, nil
|
||||
}
|
||||
|
||||
func (h *factHashes) SetHash(ctx context.Context, name, hash string) error {
|
||||
_, err := h.api.WriteFact(ctx, ipc.WriteFactReq{
|
||||
Ts: time.Now(),
|
||||
Kind: "config",
|
||||
Key: hashKey(name),
|
||||
Value: hash,
|
||||
Source: "poll:crawl",
|
||||
Confidence: 1.0,
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
// crawlEmbedder adapts router.Embedder for the watcher, embedding with
|
||||
// EmbedPassage (a page is text being searched FOR, and the e5 embedder is
|
||||
// asymmetric).
|
||||
func crawlEmbedder(emb router.Embedder) crawl.Embedder {
|
||||
if emb == nil {
|
||||
return nil
|
||||
}
|
||||
return passageEmbedder{emb}
|
||||
}
|
||||
@@ -0,0 +1,186 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/crawl"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/phraser"
|
||||
"github.com/kami/maven/internal/router"
|
||||
"github.com/kami/maven/internal/voice"
|
||||
)
|
||||
|
||||
// The default config reads nothing. This is the whole "off unless configured"
|
||||
// contract for the crawler, asserted at the wiring level rather than trusted.
|
||||
func TestCrawlOffByDefault(t *testing.T) {
|
||||
cfg := &config.Config{}
|
||||
if c := newCrawler(cfg); c != nil {
|
||||
t.Error("newCrawler with no crawl block returned a crawler")
|
||||
}
|
||||
if c := onDemandCrawler(cfg); c != nil {
|
||||
t.Error("onDemandCrawler with no crawl block returned a crawler")
|
||||
}
|
||||
if w := newCrawlWorker(nil, nil, nil, cfg); w != nil {
|
||||
t.Error("newCrawlWorker with no crawl block returned a worker")
|
||||
}
|
||||
// Watches configured but on_demand off ⇒ the answer path still reads
|
||||
// nothing: a timer over a fixed list is not permission for arbitrary URLs.
|
||||
withWatch := &config.Config{Crawl: &config.CrawlConfig{
|
||||
Watches: []config.CrawlWatchConfig{{Name: "p", URL: "https://example.org/p"}},
|
||||
}}
|
||||
if c := onDemandCrawler(withWatch); c != nil {
|
||||
t.Error("onDemandCrawler honoured a watch list as on-demand permission")
|
||||
}
|
||||
if c := newCrawler(withWatch); c == nil {
|
||||
t.Error("newCrawler returned nil for a configured watch")
|
||||
}
|
||||
}
|
||||
|
||||
// The wired fetcher must refuse a private address, because the crawler on this
|
||||
// box sits one hop from the whole homelab. Same guard the webfetch tests cover;
|
||||
// this asserts the daemon actually wires it.
|
||||
func TestCrawlerRefusesPrivateAddress(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "text/html")
|
||||
w.Write([]byte("<html><body>secret</body></html>"))
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
c := newCrawler(&config.Config{Crawl: &config.CrawlConfig{OnDemand: true}})
|
||||
if c == nil {
|
||||
t.Fatal("newCrawler returned nil for an on-demand config")
|
||||
}
|
||||
if _, err := c.Page(context.Background(), srv.URL); err == nil {
|
||||
t.Fatalf("reading %s succeeded; a loopback address must be refused", srv.URL)
|
||||
}
|
||||
}
|
||||
|
||||
func TestFactHashesRoundTrip(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
st := newTestStore(t)
|
||||
h := &factHashes{api: ipc.NewStoreAPI(st)}
|
||||
|
||||
got, err := h.LastHash(ctx, "page")
|
||||
if err != nil {
|
||||
t.Fatalf("LastHash on a fresh store: %v", err)
|
||||
}
|
||||
if got != "" {
|
||||
t.Errorf("LastHash = %q, want empty for a never-read page", got)
|
||||
}
|
||||
if err := h.SetHash(ctx, "page", "deadbeef"); err != nil {
|
||||
t.Fatalf("SetHash: %v", err)
|
||||
}
|
||||
got, err = h.LastHash(ctx, "page")
|
||||
if err != nil {
|
||||
t.Fatalf("LastHash: %v", err)
|
||||
}
|
||||
if got != "deadbeef" {
|
||||
t.Errorf("LastHash = %q, want deadbeef", got)
|
||||
}
|
||||
if key := hashKey("page"); key != "crawl:hash:page" {
|
||||
t.Errorf("hashKey = %q", key)
|
||||
}
|
||||
}
|
||||
|
||||
// stubCrawlFetcher serves one fixed page to every URL, so queryWeb can be
|
||||
// exercised without a network or an allowlist.
|
||||
type stubCrawlFetcher struct{ body, ctype string }
|
||||
|
||||
func (s *stubCrawlFetcher) Get(_ context.Context, u string) (*crawl.Response, error) {
|
||||
ct := s.ctype
|
||||
if ct == "" {
|
||||
ct = "text/html"
|
||||
}
|
||||
if strings.HasSuffix(u, "/robots.txt") {
|
||||
return &crawl.Response{URL: u, ContentType: "text/plain", Body: []byte("")}, nil
|
||||
}
|
||||
return &crawl.Response{URL: u, ContentType: ct, Body: []byte(s.body)}, nil
|
||||
}
|
||||
|
||||
func buildWebHandler(c *crawl.Crawler) *reactiveHandler {
|
||||
return &reactiveHandler{
|
||||
replier: voice.NewStubReplier(),
|
||||
phraser: phraser.NewStub(),
|
||||
crawler: c,
|
||||
}
|
||||
}
|
||||
|
||||
func askWeb(h *reactiveHandler, q string) (string, bool) {
|
||||
return h.queryWeb(context.Background(), &queryTurn{
|
||||
dec: router.Decision{Intent: router.IntentQuery, Utterance: q},
|
||||
})
|
||||
}
|
||||
|
||||
func TestQueryWebPassesWithoutAURL(t *testing.T) {
|
||||
h := buildWebHandler(crawl.New(&stubCrawlFetcher{body: "<html><body>x</body></html>"}, crawl.Config{}))
|
||||
if reply, ok := askWeb(h, "почему небо синее?"); ok {
|
||||
t.Errorf("the web source claimed a question with no URL: %q", reply)
|
||||
}
|
||||
}
|
||||
|
||||
// Not configured is said out loud rather than falling through, so a small model
|
||||
// never invents a page's contents from its URL.
|
||||
func TestQueryWebSaysWhenNotConfigured(t *testing.T) {
|
||||
h := buildWebHandler(nil)
|
||||
reply, ok := askWeb(h, "посмотри https://example.org/page")
|
||||
if !ok {
|
||||
t.Fatal("the web source did not claim a question with a URL")
|
||||
}
|
||||
if !strings.Contains(reply, "не настроено") {
|
||||
t.Errorf("reply = %q, want the not-configured answer", reply)
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueryWebReadsThePage(t *testing.T) {
|
||||
h := buildWebHandler(crawl.New(&stubCrawlFetcher{
|
||||
body: "<html><head><title>Заголовок</title></head><body><p>текст страницы</p></body></html>",
|
||||
}, crawl.Config{}))
|
||||
reply, ok := askWeb(h, "посмотри https://example.org/page — что там?")
|
||||
if !ok {
|
||||
t.Fatal("the web source did not claim a question with a URL")
|
||||
}
|
||||
if !strings.Contains(reply, "текст страницы") {
|
||||
t.Errorf("reply = %q, want the page text read back", reply)
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueryWebRefusesNonHTML(t *testing.T) {
|
||||
h := buildWebHandler(crawl.New(&stubCrawlFetcher{
|
||||
body: "\x00\x01binary", ctype: "application/octet-stream",
|
||||
}, crawl.Config{}))
|
||||
reply, ok := askWeb(h, "почитай https://example.org/blob.bin")
|
||||
if !ok {
|
||||
t.Fatal("the web source did not claim a question with a URL")
|
||||
}
|
||||
if !strings.Contains(reply, "не получилось") {
|
||||
t.Errorf("reply = %q, want the read-failed answer", reply)
|
||||
}
|
||||
}
|
||||
|
||||
// robots.txt is honoured on the answer path too, and she says so instead of
|
||||
// reporting a generic failure.
|
||||
func TestQueryWebObeysRobots(t *testing.T) {
|
||||
h := buildWebHandler(crawl.New(&robotsDenyFetcher{}, crawl.Config{}))
|
||||
reply, ok := askWeb(h, "посмотри https://example.org/private")
|
||||
if !ok {
|
||||
t.Fatal("the web source did not claim a question with a URL")
|
||||
}
|
||||
if !strings.Contains(reply, "robots.txt") {
|
||||
t.Errorf("reply = %q, want the robots answer", reply)
|
||||
}
|
||||
}
|
||||
|
||||
type robotsDenyFetcher struct{}
|
||||
|
||||
func (robotsDenyFetcher) Get(_ context.Context, u string) (*crawl.Response, error) {
|
||||
if strings.HasSuffix(u, "/robots.txt") {
|
||||
return &crawl.Response{URL: u, ContentType: "text/plain",
|
||||
Body: []byte("User-agent: *\nDisallow: /private\n")}, nil
|
||||
}
|
||||
return &crawl.Response{URL: u, ContentType: "text/html", Body: []byte("<html>nope</html>")}, nil
|
||||
}
|
||||
+1
-2
@@ -27,7 +27,6 @@ import (
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/email"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/llm"
|
||||
"github.com/kami/maven/internal/phraser"
|
||||
"github.com/kami/maven/internal/store"
|
||||
)
|
||||
@@ -66,7 +65,7 @@ func newMailIntake(st *store.Store, phr phraser.Phraser, cfg *config.Config) *ma
|
||||
if timeout <= 0 {
|
||||
timeout = config.DefaultEmailTimeout
|
||||
}
|
||||
ex := email.NewExtractor(llm.New(lp.BaseURL(), timeout), cfg.Email.MaxTasks, contextBlockFn(cfg, time.Now))
|
||||
ex := email.NewExtractor(llmClientFor(lp, timeout), cfg.Email.MaxTasks, contextBlockFn(cfg, time.Now))
|
||||
log.Printf("mail intake: enabled (max %d candidates per message, timeout %s)", cfg.Email.MaxTasks, timeout)
|
||||
return &mailIntake{st: st, ex: ex, timeout: timeout, now: time.Now}
|
||||
}
|
||||
|
||||
@@ -171,6 +171,7 @@ func run(args []string) error {
|
||||
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)
|
||||
)
|
||||
|
||||
if !locked {
|
||||
@@ -267,6 +268,7 @@ func run(args []string) error {
|
||||
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
|
||||
evalWorker = newMemoryEvalWorker(st, phr, cfg)
|
||||
feedWkr = newFeedWorker(ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
|
||||
crawlWkr = newCrawlWorker(newCrawler(cfg), ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
|
||||
|
||||
coreAPI = &daemonAPI{
|
||||
CoreAPI: ipc.NewStoreAPI(st),
|
||||
@@ -327,6 +329,7 @@ func run(args []string) error {
|
||||
// ipc.MethodIngestMail reports ErrUnknownMethod.
|
||||
if !locked {
|
||||
wireMailIntake(srv, st, phr, cfg)
|
||||
wireModelSwap(srv, phr, cfg)
|
||||
}
|
||||
|
||||
// WrapKeyFn — wraps the env key with a passkey credential public key and
|
||||
@@ -454,6 +457,7 @@ func run(args []string) error {
|
||||
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
|
||||
evalWorker = newMemoryEvalWorker(st, phr, cfg)
|
||||
feedWkr = newFeedWorker(ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
|
||||
crawlWkr = newCrawlWorker(newCrawler(cfg), ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
|
||||
|
||||
// Swap the CoreAPI from the locked placeholder to the real store adapter.
|
||||
newAPI := &daemonAPI{
|
||||
@@ -468,6 +472,7 @@ func run(args []string) error {
|
||||
srv.SetAPI(newAPI)
|
||||
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
|
||||
wireMailIntake(srv, st, phr, cfg)
|
||||
wireModelSwap(srv, phr, cfg)
|
||||
|
||||
// Start voice server.
|
||||
if voiceW != nil {
|
||||
@@ -506,6 +511,13 @@ func run(args []string) error {
|
||||
}()
|
||||
}
|
||||
|
||||
// Start the watched-page crawls (nil unless configured).
|
||||
if crawlWkr != nil {
|
||||
go func() {
|
||||
crawlWkr.run(ctx)
|
||||
}()
|
||||
}
|
||||
|
||||
dl.unlock()
|
||||
log.Printf("mavend: unlocked via passkey assertion")
|
||||
return nil
|
||||
@@ -558,6 +570,13 @@ func run(args []string) error {
|
||||
feedWkr.run(ctx)
|
||||
}()
|
||||
}
|
||||
if crawlWkr != nil {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
crawlWkr.run(ctx)
|
||||
}()
|
||||
}
|
||||
}
|
||||
|
||||
<-ctx.Done()
|
||||
|
||||
@@ -16,7 +16,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/llm"
|
||||
"github.com/kami/maven/internal/memeval"
|
||||
"github.com/kami/maven/internal/phraser"
|
||||
"github.com/kami/maven/internal/store"
|
||||
@@ -50,7 +49,7 @@ func newMemoryEvalWorker(st *store.Store, phr phraser.Phraser, cfg *config.Confi
|
||||
}
|
||||
// A generous per-request timeout: this is a long prompt to a Thinking model
|
||||
// and nobody is waiting on the answer.
|
||||
client := llm.New(lp.BaseURL(), 5*time.Minute)
|
||||
client := llmClientFor(lp, 5*time.Minute)
|
||||
ev := memeval.NewEvaluator(st, st, client, memeval.Config{
|
||||
MaxItems: cfg.MemoryEval.MaxItems,
|
||||
MinConfidence: cfg.MemoryEval.MinConfidence,
|
||||
|
||||
@@ -0,0 +1,113 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"path/filepath"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/llm"
|
||||
"github.com/kami/maven/internal/phraser"
|
||||
)
|
||||
|
||||
// Swapping the resident model while the daemon runs (Vikunja #250).
|
||||
//
|
||||
// Off unless configured: with no phraser.swap_models allowlist the two IPC
|
||||
// methods are never wired, so they answer ErrUnknownMethod. When it is wired the
|
||||
// swap method is AuthStepUp (internal/auth), which means an authed human surface
|
||||
// only — there is no act, no intent and no timer that reaches it. The daemon
|
||||
// never decides to change its own brain.
|
||||
//
|
||||
// The allowlist is exact-match against paths a human wrote in mavend.json. The
|
||||
// request carries a path and llama-server is started with it as `-m`, so
|
||||
// anything looser would turn "swap the model" into "load any file on my disk".
|
||||
func wireModelSwap(srv *ipc.Server, phr phraser.Phraser, cfg *config.Config) {
|
||||
if cfg.Phraser == nil || len(cfg.Phraser.SwapModels) == 0 {
|
||||
return
|
||||
}
|
||||
lp, ok := phr.(*phraser.LLMPhraser)
|
||||
if !ok {
|
||||
log.Printf("model swap: phraser.swap_models is set but there is no llama-server phraser — swap disabled")
|
||||
return
|
||||
}
|
||||
allowed := map[string]bool{}
|
||||
for _, m := range cfg.Phraser.SwapModels {
|
||||
allowed[filepath.Clean(m)] = true
|
||||
}
|
||||
// The configured model is always swappable back to, listed or not: the way
|
||||
// out of a bad swap must not depend on remembering to allowlist the model
|
||||
// you are already running.
|
||||
allowed[filepath.Clean(cfg.Phraser.ModelPath)] = true
|
||||
|
||||
srv.SwapModelFn = func(ctx context.Context, req ipc.SwapModelReq) (ipc.SwapModelResp, error) {
|
||||
path := filepath.Clean(req.ModelPath)
|
||||
if !allowed[path] {
|
||||
log.Printf("model swap: REFUSED %q — not in phraser.swap_models", req.ModelPath)
|
||||
return ipc.SwapModelResp{}, fmt.Errorf("%w: %q is not in phraser.swap_models", ipc.ErrForbidden, req.ModelPath)
|
||||
}
|
||||
res, err := lp.Swap(ctx, phraser.SwapSpec{
|
||||
ModelPath: path,
|
||||
NGpuLayers: req.NGpuLayers,
|
||||
NCtx: req.NCtx,
|
||||
})
|
||||
resp := ipc.SwapModelResp{
|
||||
Model: res.Model,
|
||||
ModelPath: res.ModelPath,
|
||||
BaseURL: res.BaseURL,
|
||||
RolledBack: res.RolledBack,
|
||||
TookMs: res.Took.Milliseconds(),
|
||||
}
|
||||
if err != nil {
|
||||
// A rolled-back swap is a failure that left a working daemon behind.
|
||||
// Both halves matter to the caller, so the response is filled in even
|
||||
// though the error is returned.
|
||||
log.Printf("model swap: %v", err)
|
||||
return resp, err
|
||||
}
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
srv.ModelStatusFn = func(ctx context.Context) (ipc.ModelStatusResp, error) {
|
||||
path, ngl, nctx := lp.LiveModel()
|
||||
base := lp.BaseURL()
|
||||
resp := ipc.ModelStatusResp{
|
||||
ModelPath: path,
|
||||
BaseURL: base,
|
||||
NGpuLayers: ngl,
|
||||
NCtx: nctx,
|
||||
Swappable: cfg.Phraser.SwapModels,
|
||||
}
|
||||
if base == "" {
|
||||
resp.Model = llm.UnknownModel
|
||||
return resp, nil
|
||||
}
|
||||
id, err := llm.ModelID(ctx, base)
|
||||
if err != nil {
|
||||
// Report the honest "I could not confirm it" rather than echoing the
|
||||
// configured filename as if the server had said it.
|
||||
resp.Model = llm.UnknownModel
|
||||
return resp, nil
|
||||
}
|
||||
resp.Model = id
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
log.Printf("model swap: enabled, %d allowlisted model(s) — step-up required", len(cfg.Phraser.SwapModels))
|
||||
}
|
||||
|
||||
// llmClientFor builds a completion client on the phraser's llama-server and
|
||||
// keeps it pointed at the right one across a model swap.
|
||||
//
|
||||
// Without the OnSwap registration every holder of a base URL — the LLM router,
|
||||
// the replier, the mail extractor, the memory evaluator — would keep talking to
|
||||
// the port of a server that no longer exists, and the daemon would degrade to
|
||||
// the classifier permanently after the first swap. The client is re-pointed, not
|
||||
// rebuilt, so nothing that holds it has to know a swap happened.
|
||||
func llmClientFor(lp *phraser.LLMPhraser, timeout time.Duration) *llm.Client {
|
||||
c := llm.New(lp.BaseURL(), timeout)
|
||||
lp.OnSwap(func(base string) { c.SetBaseURL(base) })
|
||||
return c
|
||||
}
|
||||
@@ -52,6 +52,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/audio"
|
||||
"github.com/kami/maven/internal/crawl"
|
||||
"github.com/kami/maven/internal/dialogue"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/memory"
|
||||
@@ -82,6 +83,10 @@ type reactiveHandler struct {
|
||||
replier voice.Replier
|
||||
now func() time.Time
|
||||
|
||||
// crawler reads a web page he names out loud (queryWeb). nil ⇒ on-demand
|
||||
// page reading is off, which is the default: no `crawl` block, no fetch.
|
||||
crawler *crawl.Crawler
|
||||
|
||||
// feedsOn — whether any RSS feed is configured (config.Feeds). It changes
|
||||
// only what she SAYS when asked and nothing is there: "ленты не настроены"
|
||||
// instead of "ничего нового", which are different truths.
|
||||
|
||||
+17
-12
@@ -148,7 +148,9 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
|
||||
// The replier uses the same llama-server as the phraser.
|
||||
var llmClient *llm.Client
|
||||
if lp, ok := phr.(*phraser.LLMPhraser); ok {
|
||||
llmClient = llm.New(lp.BaseURL(), 60*time.Second)
|
||||
// llmClientFor, not llm.New: this client must follow the phraser onto
|
||||
// the new llama-server when the resident model is swapped (Vikunja #250).
|
||||
llmClient = llmClientFor(lp, 60*time.Second)
|
||||
}
|
||||
// ----- router (the cascade; floor examples seed the classifier) -----
|
||||
// The act matcher's allowlist is exactly the enabled tool names — the
|
||||
@@ -201,17 +203,20 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
|
||||
|
||||
// ----- the handler (the reactive path; closes over stt / tts / router / coreAPI / memory) -----
|
||||
h := &reactiveHandler{
|
||||
stt: transcriber,
|
||||
tts: synthesizer,
|
||||
router: rtr,
|
||||
embedder: emb,
|
||||
api: coreAPI,
|
||||
tools: exec,
|
||||
matcher: matcher,
|
||||
replier: replier,
|
||||
phraser: phr,
|
||||
now: time.Now,
|
||||
feedsOn: cfg.Feeds != nil,
|
||||
stt: transcriber,
|
||||
tts: synthesizer,
|
||||
router: rtr,
|
||||
embedder: emb,
|
||||
api: coreAPI,
|
||||
tools: exec,
|
||||
matcher: matcher,
|
||||
replier: replier,
|
||||
phraser: phr,
|
||||
now: time.Now,
|
||||
feedsOn: cfg.Feeds != nil,
|
||||
// nil unless `crawl.on_demand` is on: reading a page he names is a
|
||||
// capability, and capabilities are off unless configured.
|
||||
crawler: onDemandCrawler(cfg),
|
||||
weatherProvider: weatherProvider,
|
||||
weatherLocation: weatherLocation,
|
||||
memStore: memStore,
|
||||
|
||||
@@ -0,0 +1,204 @@
|
||||
// Command mavupdate deploys a new build of Maven to the box she runs on, with
|
||||
// an automatic rollback when the new build does not come up (Vikunja #249).
|
||||
//
|
||||
// It is a CLI on purpose, and it is the ONLY trigger for the update path.
|
||||
//
|
||||
// The obvious design — an IPC method plus a button on the web UI behind the
|
||||
// step-up passkey gate, the way /tools works — was considered and refused. A
|
||||
// step-up gate protects against the wrong person clicking; it does not change
|
||||
// the fact that anything reachable over the network becomes, in the event of a
|
||||
// mavweb bug, a remote arbitrary-code path with a build system attached. An
|
||||
// update needs shell access on the host, which is a strictly higher bar than
|
||||
// the gate that guards the tool allowlist. That is deliberate and it is the
|
||||
// reason there is no MethodApplyUpdate anywhere in internal/ipc.
|
||||
//
|
||||
// Consequently: mavend does not import internal/update, nothing runs on a timer,
|
||||
// nothing checks a release server, and no act, intent, tool or LLM output can
|
||||
// reach any of this. She cannot update herself. She can be updated, by him.
|
||||
//
|
||||
// mavupdate -config deploy/mavend.json list # snapshots available to roll back to
|
||||
// mavupdate -config deploy/mavend.json verify # make build + make test, deploys nothing
|
||||
// mavupdate -config deploy/mavend.json apply -yes # the whole thing
|
||||
// mavupdate -config deploy/mavend.json rollback [id] # restore + restart (default: newest)
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/update"
|
||||
)
|
||||
|
||||
func main() {
|
||||
cfgPath := flag.String("config", "deploy/mavend.json", "path to mavend.json (the update block is read from it)")
|
||||
yes := flag.Bool("yes", false, "required by `apply` and `rollback`: yes, restart the daemon")
|
||||
flag.Usage = usage
|
||||
flag.Parse()
|
||||
|
||||
// The stdlib flag package stops parsing at the first non-flag argument, so a
|
||||
// `-yes` written after the subcommand (which is how anyone would type it, and
|
||||
// how the usage text shows it) lands in Args instead of the flag. Pick it out
|
||||
// by hand rather than silently treating "apply -yes" as an unconfirmed apply.
|
||||
var args []string
|
||||
for _, a := range flag.Args() {
|
||||
if a == "-yes" || a == "--yes" {
|
||||
*yes = true
|
||||
continue
|
||||
}
|
||||
args = append(args, a)
|
||||
}
|
||||
if len(args) == 0 {
|
||||
usage()
|
||||
os.Exit(2)
|
||||
}
|
||||
|
||||
cfg, err := config.Load(*cfgPath)
|
||||
if err != nil {
|
||||
die("config: %v", err)
|
||||
}
|
||||
if cfg.Update == nil {
|
||||
die("no `update` block in %s — the update capability is off unless configured.\nSee the package comment in internal/update for what it does and does not do.", *cfgPath)
|
||||
}
|
||||
|
||||
logf := func(format string, a ...any) {
|
||||
fmt.Fprintf(os.Stderr, "%s %s\n", time.Now().Format("15:04:05"), fmt.Sprintf(format, a...))
|
||||
}
|
||||
u, err := update.New(*cfg.Update, update.WithLogger(logf))
|
||||
if err != nil {
|
||||
die("%v", err)
|
||||
}
|
||||
|
||||
// Ctrl-C cancels the build or the health wait. It cannot cancel a rollback
|
||||
// midway into leaving the box in an unknown state, because the rollback runs
|
||||
// on its own context — see cmdApply.
|
||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||
defer stop()
|
||||
|
||||
switch args[0] {
|
||||
case "list":
|
||||
cmdList(u)
|
||||
case "verify":
|
||||
cmdVerify(ctx, u)
|
||||
case "apply":
|
||||
if !*yes {
|
||||
die("apply restarts mavend and can roll her back. Re-run with -yes if that is what you want.")
|
||||
}
|
||||
cmdApply(ctx, u)
|
||||
case "rollback":
|
||||
if !*yes {
|
||||
die("rollback restores the previous artifacts and restarts mavend. Re-run with -yes.")
|
||||
}
|
||||
id := ""
|
||||
if len(args) > 1 {
|
||||
id = args[1]
|
||||
}
|
||||
cmdRollback(ctx, u, id)
|
||||
default:
|
||||
usage()
|
||||
os.Exit(2)
|
||||
}
|
||||
}
|
||||
|
||||
func cmdList(u *update.Updater) {
|
||||
snaps, err := u.Snapshots()
|
||||
if err != nil {
|
||||
die("snapshots: %v", err)
|
||||
}
|
||||
if len(snaps) == 0 {
|
||||
fmt.Println("no snapshots yet — the first `apply` takes one before it builds anything")
|
||||
return
|
||||
}
|
||||
fmt.Printf("%-18s %-12s %s\n", "SNAPSHOT", "COMMIT", "FILES")
|
||||
for _, s := range snaps {
|
||||
commit := s.Commit
|
||||
if len(commit) > 12 {
|
||||
commit = commit[:12]
|
||||
}
|
||||
if commit == "" {
|
||||
commit = "-"
|
||||
}
|
||||
fmt.Printf("%-18s %-12s %d\n", s.ID, commit, len(s.Files))
|
||||
}
|
||||
fmt.Printf("\nrollback to the newest with: mavupdate rollback -yes\n")
|
||||
}
|
||||
|
||||
func cmdVerify(ctx context.Context, u *update.Updater) {
|
||||
steps, err := u.Verify(ctx)
|
||||
report(steps)
|
||||
if err != nil {
|
||||
die("%v", err)
|
||||
}
|
||||
fmt.Println("verified: the tree builds and passes its own tests. Nothing was deployed — run `apply -yes` for that.")
|
||||
}
|
||||
|
||||
func cmdApply(ctx context.Context, u *update.Updater) {
|
||||
res, err := u.Apply(ctx)
|
||||
report(res.Steps)
|
||||
summarize(res)
|
||||
switch {
|
||||
case err == nil:
|
||||
fmt.Println("\nupdate committed: she answers on the new build.")
|
||||
case errors.Is(err, update.ErrRollbackFailed):
|
||||
die("\n%v\n\nSHE IS PROBABLY DOWN. The previous artifacts are in the snapshot dir; copy them\nover the install dir and restart by hand.", err)
|
||||
case errors.Is(err, update.ErrRolledBack):
|
||||
die("\n%v\n\nShe is answering again on the previous build. Nothing was lost; fix the change and retry.", err)
|
||||
default:
|
||||
die("\n%v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func cmdRollback(ctx context.Context, u *update.Updater, id string) {
|
||||
res, err := u.Rollback(ctx, id)
|
||||
report(res.Steps)
|
||||
summarize(res)
|
||||
if err != nil && !errors.Is(err, update.ErrRolledBack) {
|
||||
die("\n%v", err)
|
||||
}
|
||||
fmt.Printf("\nrolled back to %s; she answers on it.\n", res.SnapshotID)
|
||||
}
|
||||
|
||||
func report(steps []update.Step) {
|
||||
for _, s := range steps {
|
||||
status := "ok"
|
||||
if s.Err != nil {
|
||||
status = "FAILED: " + s.Err.Error()
|
||||
}
|
||||
fmt.Printf(" %-8s %-8s %s\n", s.Name, s.Took.Round(time.Second), status)
|
||||
if s.Output != "" {
|
||||
fmt.Printf("---- %s output ----\n%s\n-------------------\n", s.Name, s.Output)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func summarize(res update.Result) {
|
||||
fmt.Printf("\nverified=%v snapshot=%s installed=%d restarted=%v healthy=%v rolled_back=%v rollback_healthy=%v took=%s\n",
|
||||
res.Verified, res.SnapshotID, len(res.Installed), res.Restarted, res.Healthy, res.RolledBack, res.RollbackHealthy, res.Took.Round(time.Second))
|
||||
}
|
||||
|
||||
func usage() {
|
||||
fmt.Fprint(os.Stderr, `mavupdate — deploy a new build of Maven, with rollback.
|
||||
|
||||
mavupdate [-config path] list
|
||||
mavupdate [-config path] verify
|
||||
mavupdate [-config path] apply -yes
|
||||
mavupdate [-config path] rollback [snapshot-id] -yes
|
||||
|
||||
apply is: health-check the running daemon, snapshot the deployed artifacts,
|
||||
make build, make test, install, restart, health-check — and restore the
|
||||
snapshot if any of that fails. It never fetches code and never runs by itself.
|
||||
|
||||
`)
|
||||
flag.PrintDefaults()
|
||||
}
|
||||
|
||||
func die(format string, a ...any) {
|
||||
fmt.Fprintf(os.Stderr, format+"\n", a...)
|
||||
os.Exit(1)
|
||||
}
|
||||
@@ -124,6 +124,7 @@ var sidebarSections = []struct {
|
||||
Label: "Settings",
|
||||
Pages: []struct{ Label, URL, Key string }{
|
||||
{Label: "Tools", URL: "/tools", Key: "tools"},
|
||||
{Label: "Model", URL: "/models", Key: "models"},
|
||||
{Label: "Passkey", URL: "/auth/passkey", Key: "passkey"},
|
||||
},
|
||||
},
|
||||
@@ -191,6 +192,8 @@ func pageIcon(key string) string {
|
||||
return `<svg class=icon width="14" height="14"><use href="/ethos-icons.svg#i-grid"/></svg>`
|
||||
case "tools":
|
||||
return `<svg class=icon width="14" height="14"><use href="/ethos-icons.svg#i-settings"/></svg>`
|
||||
case "models":
|
||||
return `<svg class=icon width="14" height="14"><use href="/ethos-icons.svg#i-wave"/></svg>`
|
||||
case "passkey":
|
||||
return `<svg class=icon width="14" height="14"><use href="/ethos-icons.svg#i-lock"/></svg>`
|
||||
default:
|
||||
@@ -225,6 +228,8 @@ func pageTitle(key string) string {
|
||||
return "Ecosystem"
|
||||
case "tools":
|
||||
return "Tools"
|
||||
case "models":
|
||||
return "Resident Model"
|
||||
case "passkey":
|
||||
return "Passkey"
|
||||
default:
|
||||
@@ -479,6 +484,12 @@ func main() {
|
||||
mux.HandleFunc("/routines", func(w http.ResponseWriter, r *http.Request) {
|
||||
handleRoutines(w, r, core, stepUpSession, *requireStepUp)
|
||||
})
|
||||
// /models — the resident-model surface (Vikunja #250). Same step-up gate as
|
||||
// /tools, and for a comparable reason: which model is loaded decides how every
|
||||
// utterance is routed and how every reply is worded. GET is read-only.
|
||||
mux.HandleFunc("/models", func(w http.ResponseWriter, r *http.Request) {
|
||||
handleModels(w, r, core, stepUpSession, *requireStepUp)
|
||||
})
|
||||
|
||||
// State-changing routes on this server, and their gate (Vikunja #317):
|
||||
//
|
||||
|
||||
@@ -0,0 +1,145 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"html/template"
|
||||
"log"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/webauthn"
|
||||
)
|
||||
|
||||
// The resident-model surface (Vikunja #250).
|
||||
//
|
||||
// GET shows which model llama-server actually has loaded and which files the
|
||||
// daemon is configured to allow. POST swaps to one of them, behind the same
|
||||
// step-up gate as POST /tools: the loaded model decides how every utterance is
|
||||
// routed and how every reply is worded, so it is an owner action.
|
||||
//
|
||||
// There is nothing on this page Maven can press. The swap is an IPC method rated
|
||||
// AuthStepUp in internal/auth, unreachable from an act, an intent or a timer.
|
||||
|
||||
// modelController — the two non-CoreAPI methods this page needs. *ipc.Client
|
||||
// satisfies it; a core without a swap allowlist answers ErrUnknownMethod, which
|
||||
// the page renders as "not configured" rather than an error.
|
||||
type modelController interface {
|
||||
ModelStatus(ctx context.Context) (ipc.ModelStatusResp, error)
|
||||
SwapModel(ctx context.Context, req ipc.SwapModelReq) (ipc.SwapModelResp, error)
|
||||
}
|
||||
|
||||
var modelsTmpl = template.Must(template.New("models").Funcs(shellFuncs()).Parse(shellTopHTML + modelsHTML + shellBottomHTML))
|
||||
|
||||
const modelsHTML = `{{template "shellTop" "models"}}
|
||||
<h1>Resident model</h1>
|
||||
<p class=hint>swapping requires step-up — <a href=/auth/passkey>assert a passkey</a> first. The old model is unloaded before the new one is loaded (one model fits the iGPU at a time), so turns during the load are refused and fall back to the classifier.</p>
|
||||
{{if .Msg}}<div class="msg msg-ok">{{.Msg}}</div>{{end}}
|
||||
{{if .Err}}<div class="msg msg-err">{{.Err}}</div>{{end}}
|
||||
{{if .Off}}
|
||||
<section class=card>
|
||||
<h2 class=card-title>swap not configured</h2>
|
||||
<p class=hint>this core has no <code>phraser.swap_models</code> allowlist, so there is nothing to swap to. Add the gguf paths you allow to <code>deploy/mavend.json</code> and restart once.</p>
|
||||
</section>
|
||||
{{else}}
|
||||
<section class=card>
|
||||
<h2 class=card-title>loaded now</h2>
|
||||
<div class=scroll><table>
|
||||
<tr><th>model</th><td><code>{{.Status.Model}}</code></td></tr>
|
||||
<tr><th>file</th><td><code>{{.Status.ModelPath}}</code></td></tr>
|
||||
<tr><th>server</th><td><code>{{.Status.BaseURL}}</code></td></tr>
|
||||
<tr><th>n_ctx</th><td>{{.Status.NCtx}}</td></tr>
|
||||
<tr><th>n_gpu_layers</th><td>{{.Status.NGpuLayers}}</td></tr>
|
||||
</table></div>
|
||||
<p class=hint>the model name is what llama-server reports for itself, not what the config says it should be.</p>
|
||||
</section>
|
||||
<section class=card>
|
||||
<h2 class=card-title>allowed models <span class=badge>{{len .Status.Swappable}}</span></h2>
|
||||
{{if .Status.Swappable}}<div class=scroll><table><tr><th>file</th><th></th></tr>
|
||||
{{range .Status.Swappable}}<tr><td><code>{{.}}</code></td>
|
||||
<td><form method=post action=/models class=inline-form>
|
||||
<input type=hidden name=model_path value="{{.}}">
|
||||
<button class=btn>load this one</button></form></td></tr>{{end}}
|
||||
</table></div>
|
||||
{{else}}<div class=empty><div>no models allowlisted</div></div>{{end}}
|
||||
</section>
|
||||
{{end}}
|
||||
{{template "shellBottom"}}`
|
||||
|
||||
type modelsPage struct {
|
||||
Msg string
|
||||
Err string
|
||||
Off bool
|
||||
Status ipc.ModelStatusResp
|
||||
}
|
||||
|
||||
// handleModels renders the model surface (GET) and applies a swap (POST).
|
||||
//
|
||||
// A failed swap is reported as a failure with the model that is still serving
|
||||
// named, because that is the state the operator needs: the daemon rolled back
|
||||
// and is answering turns, it just is not answering them with what he asked for.
|
||||
func handleModels(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI, session *webauthn.PasskeySession, requireStepUp bool) {
|
||||
if core == nil {
|
||||
http.Error(w, "models disabled (no -core)", http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
mc, ok := core.(modelController)
|
||||
if !ok {
|
||||
http.Error(w, "models unavailable: core connection does not support model swap", http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
ctx := r.Context()
|
||||
page := modelsPage{}
|
||||
|
||||
if r.Method == http.MethodPost {
|
||||
if !stepUpOK(session, requireStepUp) {
|
||||
http.Error(w, "step-up required: assert a passkey first", http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
path := strings.TrimSpace(r.FormValue("model_path"))
|
||||
if path == "" {
|
||||
http.Error(w, "model_path required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
req := ipc.SwapModelReq{ModelPath: path}
|
||||
if v, err := strconv.Atoi(r.FormValue("n_ctx")); err == nil {
|
||||
req.NCtx = v
|
||||
}
|
||||
res, err := mc.SwapModel(ctx, req)
|
||||
switch {
|
||||
case err == nil:
|
||||
page.Msg = "loaded " + res.Model + " (" + strconv.FormatInt(res.TookMs, 10) + "ms)"
|
||||
log.Printf("models: swapped to %s (%s) in %dms", res.ModelPath, res.Model, res.TookMs)
|
||||
case errors.Is(err, ipc.ErrForbidden):
|
||||
http.Error(w, "refused: that model is not in phraser.swap_models, or step-up was not asserted", http.StatusForbidden)
|
||||
return
|
||||
case errors.Is(err, ipc.ErrUnknownMethod):
|
||||
http.Error(w, "swap not configured on this core", http.StatusServiceUnavailable)
|
||||
return
|
||||
case res.RolledBack:
|
||||
page.Err = "swap failed, rolled back to " + res.Model + " — she is still answering, with the old model"
|
||||
log.Printf("models: swap to %s failed, rolled back: %v", path, err)
|
||||
default:
|
||||
page.Err = "swap failed: " + err.Error()
|
||||
log.Printf("models: swap to %s failed: %v", path, err)
|
||||
}
|
||||
}
|
||||
|
||||
st, err := mc.ModelStatus(ctx)
|
||||
if err != nil {
|
||||
if errors.Is(err, ipc.ErrUnknownMethod) {
|
||||
page.Off = true
|
||||
} else {
|
||||
log.Printf("models: status: %v", err)
|
||||
http.Error(w, "core read failed", http.StatusBadGateway)
|
||||
return
|
||||
}
|
||||
}
|
||||
page.Status = st
|
||||
w.Header().Set("Content-Type", "text/html; charset=utf-8")
|
||||
if err := modelsTmpl.Execute(w, page); err != nil {
|
||||
log.Printf("models render: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,150 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/webauthn"
|
||||
)
|
||||
|
||||
// fakeModelCore is a core that supports the two model methods. It records what
|
||||
// the page asked for, so the tests can assert the gate rather than the HTML.
|
||||
type fakeModelCore struct {
|
||||
ipc.UnimplementedCoreAPI
|
||||
|
||||
status ipc.ModelStatusResp
|
||||
statusErr error
|
||||
|
||||
swapResp ipc.SwapModelResp
|
||||
swapErr error
|
||||
swapped []ipc.SwapModelReq
|
||||
}
|
||||
|
||||
func (f *fakeModelCore) ModelStatus(ctx context.Context) (ipc.ModelStatusResp, error) {
|
||||
return f.status, f.statusErr
|
||||
}
|
||||
|
||||
func (f *fakeModelCore) SwapModel(ctx context.Context, req ipc.SwapModelReq) (ipc.SwapModelResp, error) {
|
||||
f.swapped = append(f.swapped, req)
|
||||
return f.swapResp, f.swapErr
|
||||
}
|
||||
|
||||
func modelsGET(t *testing.T, core ipc.CoreAPI) *httptest.ResponseRecorder {
|
||||
t.Helper()
|
||||
w := httptest.NewRecorder()
|
||||
handleModels(w, httptest.NewRequest(http.MethodGet, "/models", nil), core, nil, false)
|
||||
return w
|
||||
}
|
||||
|
||||
func modelsPOST(t *testing.T, core ipc.CoreAPI, session *webauthn.PasskeySession, requireStepUp bool, path string) *httptest.ResponseRecorder {
|
||||
t.Helper()
|
||||
r := httptest.NewRequest(http.MethodPost, "/models", strings.NewReader("model_path="+path))
|
||||
r.Header.Set("Content-Type", "application/x-www-form-urlencoded")
|
||||
w := httptest.NewRecorder()
|
||||
handleModels(w, r, core, session, requireStepUp)
|
||||
return w
|
||||
}
|
||||
|
||||
func TestModels_GETShowsTheLoadedModelAndTheAllowlist(t *testing.T) {
|
||||
core := &fakeModelCore{status: ipc.ModelStatusResp{
|
||||
Model: "Qwen3-1.7B-UD-Q4_K_XL",
|
||||
ModelPath: "/opt/maven/models/llm/qwen3.gguf",
|
||||
BaseURL: "http://127.0.0.1:18099",
|
||||
NCtx: 4096,
|
||||
Swappable: []string{"/opt/maven/models/llm/qwen3.gguf", "/opt/maven/models/llm/qwen3-cpt.gguf"},
|
||||
}}
|
||||
w := modelsGET(t, core)
|
||||
if w.Code != http.StatusOK {
|
||||
t.Fatalf("GET /models = %d; want 200", w.Code)
|
||||
}
|
||||
body := w.Body.String()
|
||||
for _, want := range []string{"Qwen3-1.7B-UD-Q4_K_XL", "qwen3-cpt.gguf", "4096"} {
|
||||
if !strings.Contains(body, want) {
|
||||
t.Errorf("page does not mention %q", want)
|
||||
}
|
||||
}
|
||||
if len(core.swapped) != 0 {
|
||||
t.Errorf("a GET swapped the model: %v", core.swapped)
|
||||
}
|
||||
}
|
||||
|
||||
func TestModels_POSTRequiresStepUpWhenFailingClosed(t *testing.T) {
|
||||
// No WebAuthn configured (nil session) + -require-stepup ⇒ deny, exactly
|
||||
// like POST /tools. Nothing reaches core.
|
||||
core := &fakeModelCore{}
|
||||
w := modelsPOST(t, core, nil, true, "/opt/maven/models/llm/qwen3.gguf")
|
||||
if w.Code != http.StatusForbidden {
|
||||
t.Fatalf("POST /models without assertable step-up = %d; want 403", w.Code)
|
||||
}
|
||||
if len(core.swapped) != 0 {
|
||||
t.Fatalf("a denied POST still called SwapModel: %v", core.swapped)
|
||||
}
|
||||
}
|
||||
|
||||
func TestModels_POSTSwapsAndReportsTheModelThatAnswered(t *testing.T) {
|
||||
core := &fakeModelCore{
|
||||
swapResp: ipc.SwapModelResp{Model: "qwen3-cpt", ModelPath: "/m/cpt.gguf", TookMs: 4200},
|
||||
status: ipc.ModelStatusResp{Model: "qwen3-cpt", ModelPath: "/m/cpt.gguf"},
|
||||
}
|
||||
w := modelsPOST(t, core, nil, false, "/m/cpt.gguf")
|
||||
if w.Code != http.StatusOK {
|
||||
t.Fatalf("POST /models = %d; want 200", w.Code)
|
||||
}
|
||||
if len(core.swapped) != 1 || core.swapped[0].ModelPath != "/m/cpt.gguf" {
|
||||
t.Fatalf("SwapModel calls = %v; want one for /m/cpt.gguf", core.swapped)
|
||||
}
|
||||
if !strings.Contains(w.Body.String(), "loaded qwen3-cpt") {
|
||||
t.Errorf("page does not report which model was loaded:\n%s", w.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestModels_RolledBackSwapSaysSheIsStillAnswering(t *testing.T) {
|
||||
core := &fakeModelCore{
|
||||
swapResp: ipc.SwapModelResp{Model: "qwen3", ModelPath: "/m/old.gguf", RolledBack: true},
|
||||
swapErr: errBrokenModel{},
|
||||
status: ipc.ModelStatusResp{Model: "qwen3", ModelPath: "/m/old.gguf"},
|
||||
}
|
||||
w := modelsPOST(t, core, nil, false, "/m/cpt.gguf")
|
||||
if w.Code != http.StatusOK {
|
||||
t.Fatalf("POST /models after a rollback = %d; want 200 with the failure rendered", w.Code)
|
||||
}
|
||||
body := w.Body.String()
|
||||
if !strings.Contains(body, "rolled back to qwen3") {
|
||||
t.Errorf("page does not say it rolled back:\n%s", body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestModels_RefusedPathIs403(t *testing.T) {
|
||||
core := &fakeModelCore{swapErr: ipc.ErrForbidden}
|
||||
w := modelsPOST(t, core, nil, false, "/etc/passwd")
|
||||
if w.Code != http.StatusForbidden {
|
||||
t.Fatalf("POST /models with a non-allowlisted path = %d; want 403", w.Code)
|
||||
}
|
||||
}
|
||||
|
||||
func TestModels_UnconfiguredCoreRendersOff(t *testing.T) {
|
||||
core := &fakeModelCore{statusErr: ipc.ErrUnknownMethod}
|
||||
w := modelsGET(t, core)
|
||||
if w.Code != http.StatusOK {
|
||||
t.Fatalf("GET /models against a core without the swap = %d; want 200", w.Code)
|
||||
}
|
||||
if !strings.Contains(w.Body.String(), "swap not configured") {
|
||||
t.Errorf("page does not say the capability is off:\n%s", w.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestModels_CoreWithoutTheMethodsIs503(t *testing.T) {
|
||||
// An in-process CoreAPI (no swap methods) must not 500 the page.
|
||||
w := modelsGET(t, ipc.UnimplementedCoreAPI{})
|
||||
if w.Code != http.StatusServiceUnavailable {
|
||||
t.Fatalf("GET /models on a core without the methods = %d; want 503", w.Code)
|
||||
}
|
||||
}
|
||||
|
||||
type errBrokenModel struct{}
|
||||
|
||||
func (errBrokenModel) Error() string { return "llm: server did not start" }
|
||||
@@ -67,6 +67,36 @@ What it does and does not do:
|
||||
- how far each feed was read is stored as a config fact `rss:latest:<name>`, so
|
||||
a restart does not re-note yesterday's headlines.
|
||||
|
||||
### Reading a page (`crawl`, also off by default)
|
||||
|
||||
There is no `crawl` block either, so no page is fetched. Two halves, separately
|
||||
switched:
|
||||
|
||||
```json
|
||||
"crawl": {
|
||||
"on_demand": true,
|
||||
"interval": "6h",
|
||||
"max_runes": 4000,
|
||||
"watches": [
|
||||
{ "name": "changelog", "url": "https://example.org/changelog", "interval": "12h" }
|
||||
]
|
||||
}
|
||||
```
|
||||
|
||||
- `on_demand` lets her read a page he names in the utterance: "посмотри
|
||||
https://example.org/x — что там?". The page becomes context for his question,
|
||||
and only the URL leaves the box. Without a URL nothing is fetched, so this is
|
||||
a fallback and not a habit;
|
||||
- `watches` re-reads a fixed list on its interval and writes a note when the
|
||||
text changed. Like the feeds, it announces nothing;
|
||||
- the answer path sits **last** in the query chain, behind his memory, his notes
|
||||
and (once wired) the local Kiwix ZIMs. A local read costs nothing;
|
||||
- `robots.txt` is fetched first and obeyed with no override; a `Disallow` is a
|
||||
refusal she says out loud. `Crawl-delay` is honoured;
|
||||
- same guarded fetcher as the feeds: allowlist/denylist, no private addresses,
|
||||
size cap, redirect cap, timeout, one request per host per second;
|
||||
- dedup state is the config fact `crawl:hash:<name>`.
|
||||
|
||||
## Not yet verified / host-dependent
|
||||
|
||||
This stack is correct-by-construction but has **not been build-tested here**
|
||||
@@ -84,3 +114,49 @@ build on the target host, most likely in one of these:
|
||||
work fine over the core socket.
|
||||
- **netdata** — `mavpoll` reaches it via `host.docker.internal`; adjust if
|
||||
netdata runs elsewhere.
|
||||
|
||||
## Updating her (`mavupdate`, Vikunja #249)
|
||||
|
||||
Off unless configured, and there is deliberately no button for it. There is no
|
||||
IPC method, no web route, no timer and no act that starts an update — the trigger
|
||||
is a human running `mavupdate` on the host, which needs shell access, a strictly
|
||||
higher bar than the step-up passkey gate that guards `/tools`. She cannot update
|
||||
herself; she can be updated. Nothing here ever fetches code: the new version is
|
||||
whatever you pulled into the working tree yourself.
|
||||
|
||||
Add an `update` block to `mavend.json` (mavend ignores it — only the CLI reads
|
||||
it), with paths as they exist **on the host**, not inside a container:
|
||||
|
||||
```json
|
||||
"update": {
|
||||
"source_dir": "/home/kami/apps/Maven",
|
||||
"install_dir": "/home/kami/apps/Maven",
|
||||
"snapshot_dir": "/var/lib/maven-snapshots",
|
||||
"binaries": ["mavend", "mavweb", "mavsttd", "mavttsd", "mavwaked",
|
||||
"mavenclient", "mavpoll", "mavcaldav", "mavmaild"],
|
||||
"config_files": ["deploy/mavend.json"],
|
||||
"restart_cmd": ["docker", "compose", "up", "-d", "--build"],
|
||||
"health_socket": "/var/lib/docker/volumes/maven_sockets/_data/mavend.sock",
|
||||
"health_timeout_sec": 120
|
||||
}
|
||||
```
|
||||
|
||||
`snapshot_dir` must be outside `install_dir` (a restore must not read from what
|
||||
the install writes) and `health_socket` is required: an update that cannot check
|
||||
its own result cannot roll itself back, so the config is refused without one.
|
||||
|
||||
Then:
|
||||
|
||||
```sh
|
||||
mavupdate -config deploy/mavend.json verify # make build + make test, deploys nothing
|
||||
mavupdate -config deploy/mavend.json apply -yes # snapshot, verify, install, restart, health-check
|
||||
mavupdate -config deploy/mavend.json list # what you can roll back to
|
||||
mavupdate -config deploy/mavend.json rollback -yes # restore the previous artifacts and restart
|
||||
```
|
||||
|
||||
`apply` refuses to start if she is not already answering — otherwise a failed
|
||||
update and a box that was already broken are indistinguishable afterwards. On any
|
||||
failure after the install it restores the snapshot, restarts, and checks again;
|
||||
if that also fails it says so loudly and names the directory to copy back by hand.
|
||||
The database is never snapshotted or rolled back (see the package comment in
|
||||
`internal/update`); schema compatibility is `store.Migrate`'s job.
|
||||
|
||||
@@ -28,3 +28,41 @@
|
||||
7. Add IPC methods `MethodTriggerCrawl(name)`, `MethodListCrawls`, `MethodGetCrawlResult(name)`
|
||||
8. Add `crawls` block to `config.Config` and `deploy/mavend.json`
|
||||
9. Test with a static HTML page — verify extraction matches expected values, verify scheduling fires correctly
|
||||
|
||||
## Shipped 2026-08-01 (#259)
|
||||
|
||||
Built as `internal/crawl` (pure: robots, extraction, watcher) plus
|
||||
`cmd/mavend/crawls.go` (fetcher, ticker, dedup facts), on top of the guarded
|
||||
`internal/webfetch` door added with the feed reader (#258). Off unless
|
||||
configured, in two separately-switched halves: `crawl.on_demand` for a URL he
|
||||
names, `crawl.watches` for a scheduled re-read.
|
||||
|
||||
**Limits are code, not documentation** (`internal/webfetch`, tested one test per
|
||||
limit): host allowlist/denylist, no private addresses (loopback, RFC1918 —
|
||||
hence the LAN and the `10.42.0.0/24` wg range —, link-local incl. cloud
|
||||
metadata, CGNAT, v6 ULA) enforced in the dialer's `Control` hook so DNS
|
||||
rebinding and every redirect hop are covered, response size cap, redirect cap,
|
||||
timeout, one request per host per second. `robots.txt` is fetched first, cached
|
||||
per host, and a `Disallow` is refused with no override.
|
||||
|
||||
Deliberate deviations from the plan above:
|
||||
|
||||
- **No CSS selectors and no LLM structured extraction** (steps 2). The output is
|
||||
plaintext handed to the phraser as context for the question he asked. A 1.7B
|
||||
extracting a JSON price table from 4000 runes is a worse bet than reading, and
|
||||
`goquery` is not vendored.
|
||||
- **No `crawl` act verb and no new IPC methods** (steps 5, 7). Reading a page is
|
||||
a query source (`queryWeb` in `actions_query.go`, last in the chain, behind
|
||||
Kiwix once that is wired), not an action he commands. Nothing needs a new wire
|
||||
method to work.
|
||||
- **Notes, not facts.** A page's text is not a fact about him. Only the dedup
|
||||
hash is a fact (`crawl:hash:<name>`, kind `config`, source `poll:crawl`).
|
||||
- **Nothing is dispatched.** A changed page writes a note; it does not nudge.
|
||||
Not a nag.
|
||||
- **No `/tools` crawl history page.** The notes and the hash facts are already
|
||||
visible on `/dash`.
|
||||
|
||||
**No new dependency.** The vendored tree has no `x/net/html`, no `goquery` and
|
||||
no `temoto/robotstxt`, so robots parsing and HTML-to-text are stdlib
|
||||
(`regexp`, `html`) — RE2 has no backreferences, hence the `pairsRE` builder in
|
||||
`extract.go`.
|
||||
|
||||
@@ -393,3 +393,26 @@ func mustWriteFactParams(source string) []byte {
|
||||
}
|
||||
return b
|
||||
}
|
||||
|
||||
// TestRequirement_SwapModel — loading a different resident model is an owner
|
||||
// action at the same rung as mutating the tool allowlist: it decides how every
|
||||
// utterance is routed and how every reply is worded. The read side is not.
|
||||
func TestRequirement_SwapModel(t *testing.T) {
|
||||
if got := Requirement(ipc.MethodSwapModel); got != AuthStepUp {
|
||||
t.Errorf("SwapModel authority = %v; want AuthStepUp", got)
|
||||
}
|
||||
if got := Requirement(ipc.MethodModelStatus); got != AuthRead {
|
||||
t.Errorf("ModelStatus authority = %v; want AuthRead", got)
|
||||
}
|
||||
// A surface that cannot carry a passkey gesture cannot swap the model, no
|
||||
// matter what it is enrolled as — this is the "never through voice" property.
|
||||
voice := Scope{Surface: SurfaceVoice, Module: "voice", SourceScope: []string{"*"}}
|
||||
if err := Can(ipc.MethodSwapModel, voice, nil); !errors.Is(err, ErrForbidden) {
|
||||
t.Errorf("voice swapping the model = %v; want ErrForbidden", err)
|
||||
}
|
||||
// And with no step-up session asserted, the gate refuses even a capable surface.
|
||||
noSession := &Gate{Enrollment: NewFloorEnrollment()}
|
||||
if err := noSession.Check(context.Background(), ipc.MethodSwapModel, nil); !errors.Is(err, ipc.ErrForbidden) {
|
||||
t.Errorf("SwapModel with no asserted step-up = %v; want ErrForbidden", err)
|
||||
}
|
||||
}
|
||||
|
||||
+11
-1
@@ -53,6 +53,13 @@ func Requirement(m ipc.Method) Authority {
|
||||
// asserted — never a module or the voice/chat path. maven can propose
|
||||
// (MethodProposeTool, no step-up: she has no passkey) but never en/disable.
|
||||
return AuthStepUp
|
||||
case ipc.MethodSwapModel:
|
||||
// Swapping the resident model changes what routes every utterance and
|
||||
// what words every reply. It is the owner's call, from a surface that can
|
||||
// carry a passkey gesture — the same rung as mutating the tool allowlist,
|
||||
// and for the same reason: nothing Maven says or does may reach it.
|
||||
// MethodModelStatus is only the read side, so it stays at AuthRead.
|
||||
return AuthStepUp
|
||||
case ipc.MethodWriteFact:
|
||||
return AuthWrite
|
||||
case ipc.MethodAssertStepUp:
|
||||
@@ -79,7 +86,10 @@ func Requirement(m ipc.Method) Authority {
|
||||
// produce: candidate tasks and nothing else. It cannot write a fact, set a
|
||||
// reminder, or touch the tool allowlist, so a compromised mail reader can
|
||||
// at worst put junk on a review page he clears in one click.
|
||||
ipc.MethodIngestMail:
|
||||
ipc.MethodIngestMail,
|
||||
// The read side of the model swap: which model is resident, which ones are
|
||||
// allowlisted. It loads nothing and changes nothing.
|
||||
ipc.MethodModelStatus:
|
||||
return AuthRead
|
||||
}
|
||||
// Unknown method ⇒ AuthRead, but ipc.dispatch returns ErrUnknownMethod
|
||||
|
||||
@@ -23,6 +23,7 @@ import (
|
||||
"github.com/kami/maven/internal/delivery/ntfysink"
|
||||
"github.com/kami/maven/internal/delivery/telegramsink"
|
||||
"github.com/kami/maven/internal/morning"
|
||||
"github.com/kami/maven/internal/update"
|
||||
"github.com/robfig/cron/v3"
|
||||
)
|
||||
|
||||
@@ -108,6 +109,16 @@ type Config struct {
|
||||
// calls its /v1/chat/completions endpoint to phrase nudges and reminders.
|
||||
Phraser *PhraserConfig `json:"phraser,omitempty"`
|
||||
|
||||
// Update — how THIS box deploys a new build of Maven (Vikunja #249). nil ⇒
|
||||
// the update capability does not exist, which is the state to leave it in
|
||||
// unless the operator has read internal/update's package comment.
|
||||
//
|
||||
// mavend never reads this block: the daemon does not import internal/update
|
||||
// and cannot update itself. It lives here because cmd/mavupdate — a CLI the
|
||||
// owner runs on the host, the only trigger there is — reads the same config
|
||||
// file to find the socket it health-checks.
|
||||
Update *update.Config `json:"update,omitempty"`
|
||||
|
||||
// Voice — the client↔core surface + the stt/tts modules the daemon
|
||||
// wires. nil ⇒ the daemon doesn't wire voice: the TCP listener stays
|
||||
// down, the dispatcher's Voice slot stays nil (the routing table's
|
||||
@@ -161,6 +172,10 @@ type Config struct {
|
||||
// the weather and telegram. See FeedsConfig.
|
||||
Feeds *FeedsConfig `json:"feeds,omitempty"`
|
||||
|
||||
// Crawl — reading a web page (Vikunja #259). nil / absent ⇒ Maven never
|
||||
// fetches a page: not on request, not on a schedule. See CrawlConfig.
|
||||
Crawl *CrawlConfig `json:"crawl,omitempty"`
|
||||
|
||||
// Praxis — the ecosystem attention-state service. When configured, maven
|
||||
// calls the Praxis HTTP tools API for attention listing and item lifecycle.
|
||||
// Maven never touches Praxis's database directly (ecosystem invariant: no
|
||||
@@ -454,6 +469,62 @@ type FeedSourceConfig struct {
|
||||
Exclude []string `json:"exclude,omitempty"` // drop items containing any of these
|
||||
}
|
||||
|
||||
// CrawlConfig — the web crawler (Vikunja #259, docs/plans/14-web-crawler.md).
|
||||
//
|
||||
// Absent ⇒ off, and off means no page is ever fetched. Present with neither
|
||||
// `on_demand` nor a `watches` entry is also off: there would be nothing to do.
|
||||
//
|
||||
// The crawler is the LAST place an answer is looked for, behind the model, his
|
||||
// own memory and the local Kiwix ZIMs. That ordering lives in the query-source
|
||||
// chain (cmd/mavend/actions_query.go), not here, but it is the reason this block
|
||||
// is small: it is a fallback, not a search engine.
|
||||
//
|
||||
// Only the URL leaves the box. His notes, facts, persona block and history are
|
||||
// never part of a request — the crawler package cannot even read the store.
|
||||
type CrawlConfig struct {
|
||||
// OnDemand — may he ask her to read a page he names out loud
|
||||
// ("посмотри https://… — что там пишут?"). false ⇒ the on-demand answer
|
||||
// source stays off and only the watches below run.
|
||||
OnDemand bool `json:"on_demand,omitempty"`
|
||||
|
||||
// Watches — pages re-read on a schedule. A page whose text changed is
|
||||
// written as a note (source "crawl:<name>"); nothing is announced.
|
||||
Watches []CrawlWatchConfig `json:"watches,omitempty"`
|
||||
|
||||
// Interval — default watch cadence. 0 ⇒ crawl.DefaultWatchInterval (6h).
|
||||
Interval Duration `json:"interval,omitempty"`
|
||||
|
||||
// AllowHosts — when set, the ONLY hosts the crawler may reach (subdomains
|
||||
// included). Watched pages' own hosts are added automatically. Setting this
|
||||
// is how "she may read the arch wiki and nothing else" is expressed.
|
||||
AllowHosts []string `json:"allow_hosts,omitempty"`
|
||||
|
||||
// DenyHosts — never reachable, checked first. Private addresses do not need
|
||||
// to be listed: they are refused unconditionally (see internal/webfetch).
|
||||
DenyHosts []string `json:"deny_hosts,omitempty"`
|
||||
|
||||
// UserAgent — sent on every request AND matched against robots.txt groups.
|
||||
// Empty ⇒ webfetch.DefaultUserAgent.
|
||||
UserAgent string `json:"user_agent,omitempty"`
|
||||
|
||||
// Timeout — per-request budget. 0 ⇒ webfetch.DefaultTimeout.
|
||||
Timeout Duration `json:"timeout,omitempty"`
|
||||
|
||||
// MaxBytes — response size cap. 0 ⇒ webfetch.DefaultMaxBytes (2 MiB).
|
||||
MaxBytes int64 `json:"max_bytes,omitempty"`
|
||||
|
||||
// MaxRunes — how much extracted text is kept. 0 ⇒ crawl.DefaultMaxRunes
|
||||
// (4000), which is what fits a 4096-token context alongside a prompt.
|
||||
MaxRunes int `json:"max_runes,omitempty"`
|
||||
}
|
||||
|
||||
// CrawlWatchConfig — one page kept an eye on.
|
||||
type CrawlWatchConfig struct {
|
||||
Name string `json:"name"` // note source is "crawl:<name>"
|
||||
URL string `json:"url"`
|
||||
Interval Duration `json:"interval,omitempty"` // 0 ⇒ CrawlConfig.Interval
|
||||
}
|
||||
|
||||
// MemoryEvalConfig — the background memory-evaluation loop (Vikunja #248).
|
||||
// Absent ⇒ off, like every other capability that costs something the owner did
|
||||
// not ask for. Each evaluation is a full LLM round-trip on the one resident
|
||||
@@ -519,6 +590,21 @@ type PhraserConfig struct {
|
||||
// persona and invented units). Chat, query and reminder phrasing always go
|
||||
// through the model regardless. See phraser.Config.LLMNudges.
|
||||
LLMNudges bool `json:"llm_nudges,omitempty"`
|
||||
|
||||
// SwapModels — the gguf files the running daemon is allowed to swap to
|
||||
// without a restart (Vikunja #250). Empty (the default) means the swap
|
||||
// capability does not exist: ipc.MethodSwapModel answers ErrUnknownMethod,
|
||||
// exactly like an unconfigured weather or telegram block.
|
||||
//
|
||||
// It is an allowlist and not a directory on purpose. The request carries a
|
||||
// path, and llama-server is started with it as `-m`; anything short of an
|
||||
// exact match against a list a human wrote in this file would make "swap the
|
||||
// model" mean "load a file of your choosing off my disk". ModelPath is
|
||||
// always swappable back to whether or not it is listed.
|
||||
//
|
||||
// Paths must be absolute — the daemon's working directory is not the
|
||||
// operator's, and a relative path here would resolve somewhere surprising.
|
||||
SwapModels []string `json:"swap_models,omitempty"`
|
||||
}
|
||||
|
||||
// EmbedderConfig — paths for the ONNX multilingual embedder. The daemon
|
||||
@@ -697,6 +783,12 @@ func (c *Config) applyDefaults() {
|
||||
c.Feeds = nil
|
||||
}
|
||||
|
||||
// Same rule for the crawler: a block that neither answers on demand nor
|
||||
// watches anything has nothing to do, so it is normalised to "off".
|
||||
if c.Crawl != nil && !c.Crawl.OnDemand && len(c.Crawl.Watches) == 0 {
|
||||
c.Crawl = nil
|
||||
}
|
||||
|
||||
if c.Voice != nil {
|
||||
if c.Voice.RouterThreshold <= 0 {
|
||||
c.Voice.RouterThreshold = DefaultRouterThreshold
|
||||
@@ -754,6 +846,22 @@ func (c *Config) validate() error {
|
||||
if c.Phraser.ModelPath == "" {
|
||||
return errors.New("phraser.model_path is required")
|
||||
}
|
||||
// A relative entry in the swap allowlist would resolve against the
|
||||
// daemon's working directory, so the path a human reads in this file
|
||||
// would not be the path llama-server is handed. Fail at startup.
|
||||
for _, m := range c.Phraser.SwapModels {
|
||||
if !filepath.IsAbs(m) {
|
||||
return fmt.Errorf("phraser.swap_models: %q must be an absolute path", m)
|
||||
}
|
||||
}
|
||||
}
|
||||
// The update block is validated here even though mavend never acts on it: a
|
||||
// half-written update config that is only noticed by cmd/mavupdate is noticed
|
||||
// at the worst possible moment, halfway through deploying a new build.
|
||||
if c.Update != nil {
|
||||
if err := c.Update.Validate(); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if c.Voice != nil && c.Voice.Enabled {
|
||||
if c.Voice.Bind == "" {
|
||||
|
||||
@@ -293,3 +293,76 @@ func TestPatternProposalNotifyDefaultsOff(t *testing.T) {
|
||||
t.Errorf("cooldown = %v, want 6h", c.PatternProposals.Cooldown)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSwapModelsAbsentMeansOff — the swap capability does not exist unless the
|
||||
// operator lists the models he allows (Vikunja #250).
|
||||
func TestSwapModelsAbsentMeansOff(t *testing.T) {
|
||||
c, err := Load(writeConfig(t, `{"phraser": {"model_path": "/m/qwen.gguf"}}`))
|
||||
if err != nil {
|
||||
t.Fatalf("Load: %v", err)
|
||||
}
|
||||
if len(c.Phraser.SwapModels) != 0 {
|
||||
t.Errorf("swap_models = %v; want empty when unconfigured", c.Phraser.SwapModels)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSwapModelsParsedAndMustBeAbsolute(t *testing.T) {
|
||||
c, err := Load(writeConfig(t, `{"phraser": {
|
||||
"model_path": "/m/qwen.gguf",
|
||||
"swap_models": ["/m/qwen.gguf", "/m/qwen-cpt.gguf"]
|
||||
}}`))
|
||||
if err != nil {
|
||||
t.Fatalf("Load: %v", err)
|
||||
}
|
||||
if len(c.Phraser.SwapModels) != 2 {
|
||||
t.Fatalf("swap_models = %v; want 2 entries", c.Phraser.SwapModels)
|
||||
}
|
||||
// A relative entry would resolve against the daemon's cwd, not the operator's.
|
||||
if _, err := Load(writeConfig(t, `{"phraser": {
|
||||
"model_path": "/m/qwen.gguf",
|
||||
"swap_models": ["models/llm/qwen.gguf"]
|
||||
}}`)); err == nil {
|
||||
t.Error("Load accepted a relative swap_models entry; want a startup failure")
|
||||
}
|
||||
}
|
||||
|
||||
// TestUpdateBlockAbsentMeansOff — mavend never updates itself; the block only
|
||||
// exists so cmd/mavupdate can find the deployment it is asked to update
|
||||
// (Vikunja #249). Absent is the normal state.
|
||||
func TestUpdateBlockAbsentMeansOff(t *testing.T) {
|
||||
c, err := Load(writeConfig(t, `{}`))
|
||||
if err != nil {
|
||||
t.Fatalf("Load: %v", err)
|
||||
}
|
||||
if c.Update != nil {
|
||||
t.Errorf("update = %+v; want nil when unconfigured", c.Update)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateBlockValidatedAtStartup(t *testing.T) {
|
||||
good := `{"update": {
|
||||
"source_dir": "/srv/maven",
|
||||
"install_dir": "/srv/maven",
|
||||
"snapshot_dir": "/var/lib/maven/snapshots",
|
||||
"binaries": ["mavend", "mavweb"],
|
||||
"restart_cmd": ["docker", "compose", "up", "-d", "--build", "mavend"],
|
||||
"health_socket": "/run/maven/mavend.sock"
|
||||
}}`
|
||||
c, err := Load(writeConfig(t, good))
|
||||
if err != nil {
|
||||
t.Fatalf("Load: %v", err)
|
||||
}
|
||||
if c.Update == nil || len(c.Update.Binaries) != 2 {
|
||||
t.Fatalf("update block = %+v; want it parsed", c.Update)
|
||||
}
|
||||
// A block with no health check cannot detect its own failure, so it cannot
|
||||
// roll back — refused at load, not halfway through a deploy.
|
||||
noHealth := `{"update": {
|
||||
"source_dir": "/srv/maven", "install_dir": "/srv/maven",
|
||||
"snapshot_dir": "/var/lib/maven/snapshots",
|
||||
"binaries": ["mavend"], "restart_cmd": ["true"]
|
||||
}}`
|
||||
if _, err := Load(writeConfig(t, noHealth)); err == nil {
|
||||
t.Error("Load accepted an update block with no health_socket")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,174 @@
|
||||
// Package crawl reads a web page: fetch, robots check, HTML to text.
|
||||
//
|
||||
// It is the LAST place Maven looks for an answer, and that ordering is the whole
|
||||
// design. "Never phones home" is deprecated, but what replaced it puts local
|
||||
// sources first: the resident model, then his own memory, then the Kiwix ZIMs on
|
||||
// the box (internal/kiwix), and only then the network. A local read costs
|
||||
// nothing and leaks nothing; a fetch costs a round-trip and puts a URL in
|
||||
// someone's access log. So this package exists to be the fallback, not the
|
||||
// front door — see the querySources chain in cmd/mavend/actions_query.go for
|
||||
// where it actually sits.
|
||||
//
|
||||
// What never leaves the box: his notes, his facts, the persona block, the
|
||||
// conversation history. Only the URL is requested and, for the on-demand path,
|
||||
// only because he said it out loud. Nothing here reads the store.
|
||||
//
|
||||
// The limits are not in this package — they are in internal/webfetch, which is
|
||||
// the only way anything here touches a socket: http(s) only, host allow/deny,
|
||||
// private-address refusal, size cap, redirect cap, per-host rate limit. What
|
||||
// this package adds is politeness (robots.txt) and dedup.
|
||||
package crawl
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Errors callers distinguish.
|
||||
var (
|
||||
ErrRobots = errors.New("crawl: robots.txt disallows this path")
|
||||
ErrNotHTML = errors.New("crawl: response is not html or text")
|
||||
)
|
||||
|
||||
// Fetcher is the guarded HTTP door (internal/webfetch adapted by the daemon). An
|
||||
// interface so this package constructs no http.Client of its own and can be
|
||||
// tested without a network.
|
||||
type Fetcher interface {
|
||||
Get(ctx context.Context, url string) (*Response, error)
|
||||
}
|
||||
|
||||
// Response is the minimum a crawl needs from a fetch.
|
||||
type Response struct {
|
||||
URL string
|
||||
ContentType string
|
||||
Body []byte
|
||||
}
|
||||
|
||||
// Config — crawler knobs.
|
||||
type Config struct {
|
||||
// UserAgent is the name matched against robots.txt groups. It must be the
|
||||
// same string the fetcher sends, or Maven would be claiming one identity
|
||||
// and obeying the rules for another.
|
||||
UserAgent string
|
||||
// MaxRunes caps extracted text. 0 ⇒ DefaultMaxRunes.
|
||||
MaxRunes int
|
||||
// RobotsTTL — how long a parsed robots.txt is trusted. 0 ⇒ 1h.
|
||||
RobotsTTL time.Duration
|
||||
// Now is injectable for tests. nil ⇒ time.Now.
|
||||
Now func() time.Time
|
||||
}
|
||||
|
||||
// Crawler fetches and extracts pages. Safe for concurrent use.
|
||||
type Crawler struct {
|
||||
fetch Fetcher
|
||||
cfg Config
|
||||
robots *robotsCache
|
||||
}
|
||||
|
||||
// New builds a crawler. Returns nil when there is no fetcher, which is how the
|
||||
// daemon expresses "crawling is off unless configured".
|
||||
func New(fetch Fetcher, cfg Config) *Crawler {
|
||||
if fetch == nil {
|
||||
return nil
|
||||
}
|
||||
if cfg.UserAgent == "" {
|
||||
cfg.UserAgent = "Maven"
|
||||
}
|
||||
if cfg.MaxRunes <= 0 {
|
||||
cfg.MaxRunes = DefaultMaxRunes
|
||||
}
|
||||
if cfg.RobotsTTL <= 0 {
|
||||
cfg.RobotsTTL = time.Hour
|
||||
}
|
||||
if cfg.Now == nil {
|
||||
cfg.Now = time.Now
|
||||
}
|
||||
return &Crawler{fetch: fetch, cfg: cfg, robots: newRobotsCache(cfg.RobotsTTL)}
|
||||
}
|
||||
|
||||
// Page fetches rawURL and returns its text. It checks robots.txt first and
|
||||
// refuses a disallowed path with ErrRobots — there is no override.
|
||||
func (c *Crawler) Page(ctx context.Context, rawURL string) (Page, error) {
|
||||
u, err := url.Parse(strings.TrimSpace(rawURL))
|
||||
if err != nil {
|
||||
return Page{}, fmt.Errorf("crawl: bad url %q: %w", rawURL, err)
|
||||
}
|
||||
ok, err := c.allowed(ctx, u)
|
||||
if err != nil {
|
||||
return Page{}, err
|
||||
}
|
||||
if !ok {
|
||||
return Page{}, fmt.Errorf("%w: %s", ErrRobots, u.Path)
|
||||
}
|
||||
resp, err := c.fetch.Get(ctx, u.String())
|
||||
if err != nil {
|
||||
return Page{}, err
|
||||
}
|
||||
// A PDF or an image is bytes Maven cannot read; saying so beats storing
|
||||
// binary garbage as a "note".
|
||||
ct := strings.ToLower(resp.ContentType)
|
||||
if ct != "" && !strings.Contains(ct, "html") && !strings.Contains(ct, "text/") &&
|
||||
!strings.Contains(ct, "xml") && !strings.Contains(ct, "json") {
|
||||
return Page{}, fmt.Errorf("%w: %s", ErrNotHTML, resp.ContentType)
|
||||
}
|
||||
return Extract(resp.URL, resp.Body, c.cfg.MaxRunes), nil
|
||||
}
|
||||
|
||||
// allowed consults robots.txt for u's host, reading it at most once per TTL.
|
||||
//
|
||||
// A robots.txt that cannot be fetched (404, a timeout, a blocked host) means
|
||||
// allow, per the standard. The one thing that is NOT fail-open is an explicit
|
||||
// Disallow.
|
||||
func (c *Crawler) allowed(ctx context.Context, u *url.URL) (bool, error) {
|
||||
host := u.Host
|
||||
now := c.cfg.Now()
|
||||
rules, ok := c.robots.get(host, now)
|
||||
if !ok {
|
||||
robotsURL := u.Scheme + "://" + host + "/robots.txt"
|
||||
resp, err := c.fetch.Get(ctx, robotsURL)
|
||||
switch {
|
||||
case err != nil:
|
||||
// Note what is NOT swallowed: a refusal from the guarded fetcher.
|
||||
// If webfetch says this host is denied or private, the page fetch
|
||||
// would fail the same way, and reporting the real reason beats
|
||||
// reporting a robots verdict we never got.
|
||||
if isFatalFetchError(err) {
|
||||
return false, err
|
||||
}
|
||||
rules = Rules{}
|
||||
default:
|
||||
rules = ParseRobots(string(resp.Body), c.cfg.UserAgent)
|
||||
}
|
||||
c.robots.put(host, rules, now)
|
||||
}
|
||||
path := u.EscapedPath()
|
||||
if u.RawQuery != "" {
|
||||
path += "?" + u.RawQuery
|
||||
}
|
||||
return rules.Allowed(path), nil
|
||||
}
|
||||
|
||||
// isFatalFetchError — a fetch failure that means "this host is off limits"
|
||||
// rather than "there is no robots.txt here". The sentinel set is webfetch's, but
|
||||
// this package must not import it (the interface exists precisely so it does
|
||||
// not), so the check is on the message. Ugly and honest: the alternative is a
|
||||
// dependency inversion for two strings.
|
||||
func isFatalFetchError(err error) bool {
|
||||
s := err.Error()
|
||||
return strings.Contains(s, "not allowed") || strings.Contains(s, "private address") ||
|
||||
strings.Contains(s, "only http and https")
|
||||
}
|
||||
|
||||
// Hash is the dedup key for a crawl result: the sha256 of the extracted text,
|
||||
// hex, first 16 chars. Text and not raw HTML, because a page whose only change
|
||||
// is a rotating ad slot or a CSRF token has not changed.
|
||||
func Hash(text string) string {
|
||||
sum := sha256.Sum256([]byte(strings.TrimSpace(text)))
|
||||
return hex.EncodeToString(sum[:])[:16]
|
||||
}
|
||||
@@ -0,0 +1,155 @@
|
||||
package crawl
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// fakeFetcher serves canned pages by URL and counts requests, so a test can
|
||||
// assert that robots.txt was read once and that a refusal never reached the page.
|
||||
type fakeFetcher struct {
|
||||
pages map[string]Response
|
||||
err error
|
||||
calls []string
|
||||
}
|
||||
|
||||
func (f *fakeFetcher) Get(_ context.Context, u string) (*Response, error) {
|
||||
f.calls = append(f.calls, u)
|
||||
if f.err != nil {
|
||||
return nil, f.err
|
||||
}
|
||||
r, ok := f.pages[u]
|
||||
if !ok {
|
||||
return nil, errors.New("http 404")
|
||||
}
|
||||
if r.URL == "" {
|
||||
r.URL = u
|
||||
}
|
||||
if r.ContentType == "" {
|
||||
r.ContentType = "text/html; charset=utf-8"
|
||||
}
|
||||
return &r, nil
|
||||
}
|
||||
|
||||
const htmlPage = `<html><head><title>Почему небо синее</title>
|
||||
<style>body{color:red}</style><script>track()</script></head>
|
||||
<body><nav>меню</nav><h1>Небо</h1>
|
||||
<p>Свет рассеивается на молекулах воздуха.</p>
|
||||
<p>Короткие волны рассеиваются сильнее.</p>
|
||||
<footer>© 2026</footer></body></html>`
|
||||
|
||||
func newTestCrawler(f *fakeFetcher) *Crawler {
|
||||
return New(f, Config{UserAgent: "Maven/1.0", Now: func() time.Time { return time.Unix(0, 0) }})
|
||||
}
|
||||
|
||||
func TestPageExtractsText(t *testing.T) {
|
||||
f := &fakeFetcher{pages: map[string]Response{
|
||||
"https://example.org/sky": {Body: []byte(htmlPage)},
|
||||
}}
|
||||
page, err := newTestCrawler(f).Page(context.Background(), "https://example.org/sky")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if page.Title != "Почему небо синее" {
|
||||
t.Errorf("title = %q", page.Title)
|
||||
}
|
||||
if !strings.Contains(page.Text, "Свет рассеивается") {
|
||||
t.Errorf("body text missing: %q", page.Text)
|
||||
}
|
||||
for _, junk := range []string{"track()", "color:red", "меню", "© 2026"} {
|
||||
if strings.Contains(page.Text, junk) {
|
||||
t.Errorf("%q survived extraction: %q", junk, page.Text)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRobotsIsCheckedAndObeyed(t *testing.T) {
|
||||
f := &fakeFetcher{pages: map[string]Response{
|
||||
"https://example.org/robots.txt": {Body: []byte("User-agent: *\nDisallow: /secret\n"), ContentType: "text/plain"},
|
||||
"https://example.org/secret/x": {Body: []byte(htmlPage)},
|
||||
"https://example.org/open": {Body: []byte(htmlPage)},
|
||||
}}
|
||||
c := newTestCrawler(f)
|
||||
if _, err := c.Page(context.Background(), "https://example.org/secret/x"); !errors.Is(err, ErrRobots) {
|
||||
t.Fatalf("error = %v, want ErrRobots", err)
|
||||
}
|
||||
for _, u := range f.calls {
|
||||
if strings.Contains(u, "/secret") {
|
||||
t.Fatal("the disallowed page was fetched anyway")
|
||||
}
|
||||
}
|
||||
if _, err := c.Page(context.Background(), "https://example.org/open"); err != nil {
|
||||
t.Fatalf("allowed page: %v", err)
|
||||
}
|
||||
// robots.txt was read once for the host, not once per page.
|
||||
robotsReads := 0
|
||||
for _, u := range f.calls {
|
||||
if strings.HasSuffix(u, "/robots.txt") {
|
||||
robotsReads++
|
||||
}
|
||||
}
|
||||
if robotsReads != 1 {
|
||||
t.Fatalf("robots.txt read %d times, want 1", robotsReads)
|
||||
}
|
||||
}
|
||||
|
||||
// No robots.txt means allow — that is the standard, and the alternative makes
|
||||
// most of the web unreadable.
|
||||
func TestMissingRobotsAllows(t *testing.T) {
|
||||
f := &fakeFetcher{pages: map[string]Response{
|
||||
"https://example.org/page": {Body: []byte(htmlPage)},
|
||||
}}
|
||||
if _, err := newTestCrawler(f).Page(context.Background(), "https://example.org/page"); err != nil {
|
||||
t.Fatalf("err = %v, want the page", err)
|
||||
}
|
||||
}
|
||||
|
||||
// A refusal from the guarded fetcher must surface as itself, not be laundered
|
||||
// into "no robots.txt, go ahead".
|
||||
func TestFetcherRefusalIsNotSwallowed(t *testing.T) {
|
||||
f := &fakeFetcher{err: errors.New("webfetch: refusing to connect to a private address: 127.0.0.1")}
|
||||
_, err := newTestCrawler(f).Page(context.Background(), "http://127.0.0.1:9100/mcp")
|
||||
if err == nil || !strings.Contains(err.Error(), "private address") {
|
||||
t.Fatalf("error = %v, want the fetcher's refusal", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNonTextIsRefused(t *testing.T) {
|
||||
f := &fakeFetcher{pages: map[string]Response{
|
||||
"https://example.org/f.pdf": {Body: []byte("%PDF-1.7"), ContentType: "application/pdf"},
|
||||
}}
|
||||
if _, err := newTestCrawler(f).Page(context.Background(), "https://example.org/f.pdf"); !errors.Is(err, ErrNotHTML) {
|
||||
t.Fatalf("error = %v, want ErrNotHTML", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMaxRunesCapsText(t *testing.T) {
|
||||
long := "<html><body><p>" + strings.Repeat("привет ", 2000) + "</p></body></html>"
|
||||
f := &fakeFetcher{pages: map[string]Response{"https://example.org/l": {Body: []byte(long)}}}
|
||||
c := New(f, Config{MaxRunes: 50})
|
||||
page, err := c.Page(context.Background(), "https://example.org/l")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if n := len([]rune(page.Text)); n > 51 {
|
||||
t.Fatalf("text = %d runes, want the 50-rune cap", n)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewWithoutFetcherIsNil(t *testing.T) {
|
||||
if New(nil, Config{}) != nil {
|
||||
t.Fatal("a crawler with no fetcher must be nil — crawling is off unless configured")
|
||||
}
|
||||
}
|
||||
|
||||
func TestHashIgnoresNothingButText(t *testing.T) {
|
||||
if Hash("a") == Hash("b") {
|
||||
t.Fatal("different text hashed the same")
|
||||
}
|
||||
if Hash(" same \n") != Hash("same") {
|
||||
t.Fatal("surrounding whitespace changed the hash")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,106 @@
|
||||
package crawl
|
||||
|
||||
import (
|
||||
"html"
|
||||
"regexp"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// HTML → text, with a regexp and no tokenizer.
|
||||
//
|
||||
// golang.org/x/net/html is not vendored and the network is not assumed, so this
|
||||
// is stdlib. That is less of a compromise than it sounds: the unit of context
|
||||
// here is a few hundred words for a 4096-token model to read, exactly like the
|
||||
// Kiwix snippet, so what matters is dropping script/style/nav noise and keeping
|
||||
// paragraph boundaries. A DOM would buy correctness on malformed markup that is
|
||||
// then thrown away by truncation anyway.
|
||||
//
|
||||
// What this deliberately does NOT do: run JavaScript, follow links, or extract
|
||||
// structured fields with CSS selectors or an LLM prompt. The plan's step 2 asked
|
||||
// for the last of those; see docs/plans/14-web-crawler.md for why it was left
|
||||
// out for now.
|
||||
|
||||
var (
|
||||
// RE2 has no backreferences, so each tag pair is spelled out rather than
|
||||
// captured and matched against itself.
|
||||
dropRE = regexp.MustCompile(pairsRE("script", "style", "noscript", "svg", "head", "nav", "footer", "form"))
|
||||
titleRE = regexp.MustCompile(`(?is)<title\b[^>]*>(.*?)</title>`)
|
||||
h1RE = regexp.MustCompile(`(?is)<h1\b[^>]*>(.*?)</h1>`)
|
||||
// Block-level tags become newlines so paragraphs survive as paragraphs.
|
||||
blockRE = regexp.MustCompile(`(?is)</?(p|div|br|li|tr|h[1-6]|section|article|blockquote|pre)\b[^>]*>`)
|
||||
tagRE = regexp.MustCompile(`(?s)<[^>]*>`)
|
||||
commentRE = regexp.MustCompile(`(?s)<!--.*?-->`)
|
||||
spaceRE = regexp.MustCompile(`[ \t\f\v]+`)
|
||||
blankRE = regexp.MustCompile(`\n{2,}`)
|
||||
)
|
||||
|
||||
// pairsRE builds `(?is)<tag …>…</tag>|…` for the given tags.
|
||||
func pairsRE(tags ...string) string {
|
||||
parts := make([]string, 0, len(tags))
|
||||
for _, t := range tags {
|
||||
parts = append(parts, `<`+t+`\b[^>]*>.*?</`+t+`>`)
|
||||
}
|
||||
return `(?is)` + strings.Join(parts, "|")
|
||||
}
|
||||
|
||||
// Page is an extracted page.
|
||||
type Page struct {
|
||||
URL string
|
||||
Title string
|
||||
Text string // plain text, paragraphs separated by single newlines
|
||||
}
|
||||
|
||||
// Extract turns a fetched HTML document into a Page. maxRunes caps the text (0 ⇒
|
||||
// DefaultMaxRunes); the cap is on runes, not bytes, because a Russian page cut
|
||||
// at a byte boundary ends in half a letter.
|
||||
func Extract(url string, body []byte, maxRunes int) Page {
|
||||
if maxRunes <= 0 {
|
||||
maxRunes = DefaultMaxRunes
|
||||
}
|
||||
s := string(body)
|
||||
s = commentRE.ReplaceAllString(s, " ")
|
||||
|
||||
title := firstGroup(titleRE, s)
|
||||
if title == "" {
|
||||
title = firstGroup(h1RE, s)
|
||||
}
|
||||
|
||||
s = dropRE.ReplaceAllString(s, "\n")
|
||||
s = blockRE.ReplaceAllString(s, "\n")
|
||||
s = tagRE.ReplaceAllString(s, " ")
|
||||
s = html.UnescapeString(s)
|
||||
s = spaceRE.ReplaceAllString(s, " ")
|
||||
|
||||
var lines []string
|
||||
for _, l := range strings.Split(s, "\n") {
|
||||
if l = strings.TrimSpace(l); l != "" {
|
||||
lines = append(lines, l)
|
||||
}
|
||||
}
|
||||
text := blankRE.ReplaceAllString(strings.Join(lines, "\n"), "\n")
|
||||
|
||||
return Page{URL: url, Title: title, Text: TrimRunes(text, maxRunes)}
|
||||
}
|
||||
|
||||
// DefaultMaxRunes — how much of a page is kept. ~4000 runes is a long answer's
|
||||
// worth of context and still leaves room in a 4096-token window for the prompt
|
||||
// and the reply.
|
||||
const DefaultMaxRunes = 4000
|
||||
|
||||
func firstGroup(re *regexp.Regexp, s string) string {
|
||||
m := re.FindStringSubmatch(s)
|
||||
if len(m) < 2 {
|
||||
return ""
|
||||
}
|
||||
t := tagRE.ReplaceAllString(m[1], " ")
|
||||
return strings.TrimSpace(strings.Join(strings.Fields(html.UnescapeString(t)), " "))
|
||||
}
|
||||
|
||||
// TrimRunes cuts s to at most max runes, on a rune boundary.
|
||||
func TrimRunes(s string, max int) string {
|
||||
r := []rune(s)
|
||||
if len(r) <= max {
|
||||
return s
|
||||
}
|
||||
return strings.TrimSpace(string(r[:max])) + "…"
|
||||
}
|
||||
@@ -0,0 +1,211 @@
|
||||
package crawl
|
||||
|
||||
import (
|
||||
"regexp"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// robots.txt, parsed the small way: no wildcards beyond the two the standard
|
||||
// actually defines (`*` inside a path and `$` at the end), no sitemaps, no
|
||||
// crawl-delay-per-agent gymnastics. A personal assistant reading a handful of
|
||||
// pages does not need a spec-complete implementation; it needs to not be rude,
|
||||
// and to be auditable in one sitting.
|
||||
//
|
||||
// Two rules worth stating because they are choices, not accidents:
|
||||
//
|
||||
// - a missing or unreadable robots.txt means ALLOW. That is what the standard
|
||||
// says (404 ⇒ unrestricted), and the alternative would make a site that
|
||||
// simply has no robots.txt unreadable;
|
||||
// - an explicit Disallow means REFUSE, and Maven does not offer an override.
|
||||
// There is no "but he asked me to" flag: the page is not read.
|
||||
|
||||
// Rules is a parsed robots.txt for one user-agent.
|
||||
type Rules struct {
|
||||
allow []string
|
||||
disallow []string
|
||||
// Delay is Crawl-delay in seconds when the group named one, 0 otherwise.
|
||||
// The fetcher's own per-host rate limit is the floor; this can only make
|
||||
// Maven slower, never faster.
|
||||
Delay time.Duration
|
||||
}
|
||||
|
||||
// ParseRobots reads robots.txt and returns the rules that apply to agent.
|
||||
//
|
||||
// Group selection follows the standard: the most specific matching group wins,
|
||||
// which here means an exact user-agent match beats `*`. Lines that are neither
|
||||
// are ignored rather than guessed at.
|
||||
func ParseRobots(body string, agent string) Rules {
|
||||
agent = strings.ToLower(agent)
|
||||
|
||||
type group struct {
|
||||
agents []string
|
||||
allow []string
|
||||
disallow []string
|
||||
delay time.Duration
|
||||
}
|
||||
var groups []group
|
||||
var cur *group
|
||||
// startNew tracks whether the next User-agent line opens a new group or
|
||||
// joins the current one: consecutive User-agent lines share their rules.
|
||||
startNew := true
|
||||
|
||||
for _, raw := range strings.Split(body, "\n") {
|
||||
line := raw
|
||||
if i := strings.IndexByte(line, '#'); i >= 0 {
|
||||
line = line[:i]
|
||||
}
|
||||
line = strings.TrimSpace(line)
|
||||
if line == "" {
|
||||
continue
|
||||
}
|
||||
key, val, ok := strings.Cut(line, ":")
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
key = strings.ToLower(strings.TrimSpace(key))
|
||||
val = strings.TrimSpace(val)
|
||||
|
||||
switch key {
|
||||
case "user-agent":
|
||||
if startNew || cur == nil {
|
||||
groups = append(groups, group{})
|
||||
cur = &groups[len(groups)-1]
|
||||
startNew = false
|
||||
}
|
||||
cur.agents = append(cur.agents, strings.ToLower(val))
|
||||
case "disallow":
|
||||
if cur == nil {
|
||||
continue
|
||||
}
|
||||
startNew = true
|
||||
// "Disallow:" with an empty value allows everything, and is not a
|
||||
// path rule at all.
|
||||
if val != "" {
|
||||
cur.disallow = append(cur.disallow, val)
|
||||
}
|
||||
case "allow":
|
||||
if cur == nil {
|
||||
continue
|
||||
}
|
||||
startNew = true
|
||||
if val != "" {
|
||||
cur.allow = append(cur.allow, val)
|
||||
}
|
||||
case "crawl-delay":
|
||||
if cur == nil {
|
||||
continue
|
||||
}
|
||||
startNew = true
|
||||
if d, err := time.ParseDuration(val + "s"); err == nil && d > 0 {
|
||||
cur.delay = d
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
var star, exact *group
|
||||
for i := range groups {
|
||||
for _, a := range groups[i].agents {
|
||||
if a == "*" && star == nil {
|
||||
star = &groups[i]
|
||||
}
|
||||
// A robots.txt names "maven", we send "Maven/1.0 (…)": match on
|
||||
// prefix, which is how every crawler reads this field.
|
||||
if a != "*" && a != "" && strings.HasPrefix(agent, a) {
|
||||
exact = &groups[i]
|
||||
}
|
||||
}
|
||||
}
|
||||
g := exact
|
||||
if g == nil {
|
||||
g = star
|
||||
}
|
||||
if g == nil {
|
||||
return Rules{}
|
||||
}
|
||||
return Rules{allow: g.allow, disallow: g.disallow, Delay: g.delay}
|
||||
}
|
||||
|
||||
// Allowed reports whether path may be fetched. Longest matching rule wins, and
|
||||
// Allow beats Disallow at equal length — the standard's tie-break, and the one
|
||||
// that makes "Disallow: /" plus "Allow: /public" mean what it looks like.
|
||||
func (r Rules) Allowed(path string) bool {
|
||||
if path == "" {
|
||||
path = "/"
|
||||
}
|
||||
best, allowed := -1, true
|
||||
for _, p := range r.disallow {
|
||||
if n, ok := matchPath(p, path); ok && n > best {
|
||||
best, allowed = n, false
|
||||
}
|
||||
}
|
||||
for _, p := range r.allow {
|
||||
if n, ok := matchPath(p, path); ok && n >= best {
|
||||
best, allowed = n, true
|
||||
}
|
||||
}
|
||||
return allowed
|
||||
}
|
||||
|
||||
// matchPath applies a robots path pattern and returns the pattern's length as
|
||||
// the specificity score. `*` matches any run of characters, `$` anchors the end.
|
||||
// A pattern is a PREFIX match otherwise, which is what "Disallow: /admin" means.
|
||||
func matchPath(pattern, path string) (int, bool) {
|
||||
score := len(pattern)
|
||||
re, err := robotsRegexp(pattern)
|
||||
if err != nil {
|
||||
return 0, false
|
||||
}
|
||||
return score, re.MatchString(path)
|
||||
}
|
||||
|
||||
// robotsRegexp turns a robots path pattern into an anchored-at-the-start
|
||||
// regexp. Everything but `*` and a trailing `$` is a literal, so the pattern is
|
||||
// quoted first and the two metacharacters are put back afterwards.
|
||||
func robotsRegexp(pattern string) (*regexp.Regexp, error) {
|
||||
end := ""
|
||||
if strings.HasSuffix(pattern, "$") {
|
||||
pattern = strings.TrimSuffix(pattern, "$")
|
||||
end = "$"
|
||||
}
|
||||
parts := strings.Split(pattern, "*")
|
||||
for i, p := range parts {
|
||||
parts[i] = regexp.QuoteMeta(p)
|
||||
}
|
||||
return regexp.Compile("^" + strings.Join(parts, ".*") + end)
|
||||
}
|
||||
|
||||
// robotsCache holds parsed rules per host so a crawl of ten pages on one site
|
||||
// reads robots.txt once. TTL because a site may change its mind, and a daemon
|
||||
// that runs for weeks would otherwise never notice.
|
||||
type robotsCache struct {
|
||||
ttl time.Duration
|
||||
mu sync.Mutex
|
||||
m map[string]robotsEntry
|
||||
}
|
||||
|
||||
type robotsEntry struct {
|
||||
rules Rules
|
||||
at time.Time
|
||||
}
|
||||
|
||||
func newRobotsCache(ttl time.Duration) *robotsCache {
|
||||
return &robotsCache{ttl: ttl, m: map[string]robotsEntry{}}
|
||||
}
|
||||
|
||||
func (c *robotsCache) get(host string, now time.Time) (Rules, bool) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
e, ok := c.m[host]
|
||||
if !ok || now.Sub(e.at) > c.ttl {
|
||||
return Rules{}, false
|
||||
}
|
||||
return e.rules, true
|
||||
}
|
||||
|
||||
func (c *robotsCache) put(host string, r Rules, now time.Time) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
c.m[host] = robotsEntry{rules: r, at: now}
|
||||
}
|
||||
@@ -0,0 +1,84 @@
|
||||
package crawl
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
const robotsBody = `# a comment
|
||||
User-agent: *
|
||||
Disallow: /private
|
||||
Disallow: /tmp/
|
||||
Crawl-delay: 5
|
||||
|
||||
User-agent: Maven
|
||||
Disallow: /
|
||||
Allow: /public
|
||||
`
|
||||
|
||||
func TestParseRobotsPicksTheMostSpecificGroup(t *testing.T) {
|
||||
// The Maven group applies to us even though we send a longer UA string.
|
||||
r := ParseRobots(robotsBody, "Maven/1.0 (self-hosted personal assistant)")
|
||||
if r.Allowed("/anything") {
|
||||
t.Error("Disallow: / in our own group was ignored")
|
||||
}
|
||||
if !r.Allowed("/public/page") {
|
||||
t.Error("Allow: /public must beat the shorter Disallow: /")
|
||||
}
|
||||
|
||||
// A different agent falls into the * group.
|
||||
star := ParseRobots(robotsBody, "SomeoneElse/2")
|
||||
if !star.Allowed("/anything") {
|
||||
t.Error("the * group disallows nothing but /private and /tmp/")
|
||||
}
|
||||
if star.Allowed("/private/x") || star.Allowed("/tmp/") {
|
||||
t.Error("the * group's disallows were not applied")
|
||||
}
|
||||
if star.Delay != 5*time.Second {
|
||||
t.Errorf("crawl-delay = %v, want 5s", star.Delay)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseRobotsEmptyMeansAllowAll(t *testing.T) {
|
||||
for _, body := range []string{"", "# nothing here\n", "User-agent: *\nDisallow:\n"} {
|
||||
if !ParseRobots(body, "Maven").Allowed("/whatever") {
|
||||
t.Errorf("body %q must allow everything", body)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRobotsWildcards(t *testing.T) {
|
||||
r := ParseRobots("User-agent: *\nDisallow: /*.pdf$\nDisallow: /a/*/secret\n", "Maven")
|
||||
if r.Allowed("/docs/manual.pdf") {
|
||||
t.Error("*.pdf$ did not match")
|
||||
}
|
||||
if !r.Allowed("/docs/manual.pdf.html") {
|
||||
t.Error("$ must anchor at the end")
|
||||
}
|
||||
if r.Allowed("/a/b/secret") {
|
||||
t.Error("/a/*/secret did not match")
|
||||
}
|
||||
if !r.Allowed("/a/b/public") {
|
||||
t.Error("unrelated path was refused")
|
||||
}
|
||||
}
|
||||
|
||||
// Consecutive User-agent lines share one group, which is common in the wild.
|
||||
func TestRobotsSharedGroup(t *testing.T) {
|
||||
r := ParseRobots("User-agent: Googlebot\nUser-agent: Maven\nDisallow: /x\n", "Maven/1.0")
|
||||
if r.Allowed("/x/y") {
|
||||
t.Fatal("a shared group's rules were not applied to the second agent")
|
||||
}
|
||||
}
|
||||
|
||||
func TestRobotsCacheTTL(t *testing.T) {
|
||||
c := newRobotsCache(time.Minute)
|
||||
now := time.Now()
|
||||
c.put("example.com", ParseRobots("User-agent: *\nDisallow: /\n", "Maven"), now)
|
||||
if _, ok := c.get("example.com", now.Add(30*time.Second)); !ok {
|
||||
t.Error("a fresh entry must be served from cache")
|
||||
}
|
||||
if _, ok := c.get("example.com", now.Add(2*time.Minute)); ok {
|
||||
t.Error("an expired entry must be re-read")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,177 @@
|
||||
package crawl
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Scheduled crawls: a page is re-read on an interval, and when its TEXT changed
|
||||
// the new text is written as a note. Nothing is dispatched — same rule as the
|
||||
// feed poller (Vikunja #258). A page that announced its own change would be a
|
||||
// nag, and "the docs page changed" is not worth interrupting anyone for.
|
||||
//
|
||||
// Dedup is by content hash, so a page that re-renders identically writes nothing
|
||||
// and a rotating ad slot does not count as news.
|
||||
|
||||
// WatchConfig — one page to keep an eye on.
|
||||
type WatchConfig struct {
|
||||
Name string // note source is "crawl:<Name>"
|
||||
URL string // http(s), guarded by the fetcher
|
||||
Interval time.Duration // 0 ⇒ Watcher's default
|
||||
}
|
||||
|
||||
// Notes is core's note-writing half (same shape as ipc.CoreAPI's method).
|
||||
type Notes interface {
|
||||
WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error)
|
||||
}
|
||||
|
||||
// Hashes remembers the last text hash per watch, durably, so a restart does not
|
||||
// re-note an unchanged page. The daemon backs this with config facts
|
||||
// ("crawl:hash:<name>").
|
||||
type Hashes interface {
|
||||
LastHash(ctx context.Context, name string) (string, error)
|
||||
SetHash(ctx context.Context, name, hash string) error
|
||||
}
|
||||
|
||||
// Embedder embeds a note on its way into the store. nil ⇒ no vector.
|
||||
type Embedder interface {
|
||||
Embed(ctx context.Context, text string) ([]float32, error)
|
||||
}
|
||||
|
||||
// DefaultWatchInterval — pages change slowly, and every check is a request in
|
||||
// someone's log.
|
||||
const DefaultWatchInterval = 6 * time.Hour
|
||||
|
||||
// Watcher re-reads watched pages on their interval.
|
||||
type Watcher struct {
|
||||
c *Crawler
|
||||
watches []WatchConfig
|
||||
notes Notes
|
||||
hashes Hashes
|
||||
embed Embedder
|
||||
interval time.Duration
|
||||
nextDue map[string]time.Time
|
||||
}
|
||||
|
||||
// NewWatcher wires the scheduled half, or returns nil when there is nothing to
|
||||
// watch. Callers check for nil: no watches, no goroutine, no request.
|
||||
func NewWatcher(c *Crawler, watches []WatchConfig, notes Notes, hashes Hashes, embed Embedder, defaultInterval time.Duration) *Watcher {
|
||||
if c == nil || notes == nil {
|
||||
return nil
|
||||
}
|
||||
var valid []WatchConfig
|
||||
for _, w := range watches {
|
||||
if strings.TrimSpace(w.Name) == "" || strings.TrimSpace(w.URL) == "" {
|
||||
log.Printf("crawl: skipping a watch with no name or no url")
|
||||
continue
|
||||
}
|
||||
valid = append(valid, w)
|
||||
}
|
||||
if len(valid) == 0 {
|
||||
return nil
|
||||
}
|
||||
if defaultInterval <= 0 {
|
||||
defaultInterval = DefaultWatchInterval
|
||||
}
|
||||
return &Watcher{
|
||||
c: c, watches: valid, notes: notes, hashes: hashes, embed: embed,
|
||||
interval: defaultInterval, nextDue: map[string]time.Time{},
|
||||
}
|
||||
}
|
||||
|
||||
// Watches returns the configured watches.
|
||||
func (w *Watcher) Watches() []WatchConfig { return w.watches }
|
||||
|
||||
// CheckDue re-reads every watch whose interval elapsed and returns how many
|
||||
// notes were written. Errors are logged per watch, never returned: one dead page
|
||||
// must not stop the others.
|
||||
func (w *Watcher) CheckDue(ctx context.Context, now time.Time) int {
|
||||
written := 0
|
||||
for _, watch := range w.watches {
|
||||
if due, ok := w.nextDue[watch.Name]; ok && now.Before(due) {
|
||||
continue
|
||||
}
|
||||
interval := watch.Interval
|
||||
if interval <= 0 {
|
||||
interval = w.interval
|
||||
}
|
||||
w.nextDue[watch.Name] = now.Add(interval)
|
||||
changed, err := w.Check(ctx, watch, now)
|
||||
if err != nil {
|
||||
log.Printf("crawl: watch %s: %v", watch.Name, err)
|
||||
continue
|
||||
}
|
||||
if changed {
|
||||
log.Printf("crawl: watch %s: page changed, noted", watch.Name)
|
||||
written++
|
||||
}
|
||||
}
|
||||
return written
|
||||
}
|
||||
|
||||
// Check re-reads one watch now and reports whether it wrote a note.
|
||||
func (w *Watcher) Check(ctx context.Context, watch WatchConfig, now time.Time) (bool, error) {
|
||||
page, err := w.c.Page(ctx, watch.URL)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
// Title included: a page whose headline changed has changed.
|
||||
h := Hash(page.Title + "\n" + page.Text)
|
||||
if w.hashes != nil {
|
||||
prev, err := w.hashes.LastHash(ctx, watch.Name)
|
||||
if err != nil {
|
||||
log.Printf("crawl: watch %s: read hash: %v", watch.Name, err)
|
||||
}
|
||||
if prev == h {
|
||||
return false, nil
|
||||
}
|
||||
}
|
||||
text := NoteText(watch, page)
|
||||
var vec []float32
|
||||
if w.embed != nil {
|
||||
v, err := w.embed.Embed(ctx, text)
|
||||
if err != nil {
|
||||
log.Printf("crawl: watch %s: embed: %v", watch.Name, err)
|
||||
} else {
|
||||
vec = v
|
||||
}
|
||||
}
|
||||
if _, err := w.notes.WriteNote(ctx, now, text, vec, SourceFor(watch.Name)); err != nil {
|
||||
return false, fmt.Errorf("write note: %w", err)
|
||||
}
|
||||
if w.hashes != nil {
|
||||
if err := w.hashes.SetHash(ctx, watch.Name, h); err != nil {
|
||||
log.Printf("crawl: watch %s: save hash: %v", watch.Name, err)
|
||||
}
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// SourceFor is the note source for a watch, and SourcePrefix is what the answer
|
||||
// path matches to recognise one.
|
||||
func SourceFor(name string) string { return SourcePrefix + name }
|
||||
|
||||
// SourcePrefix — provenance for anything read off the network on a schedule.
|
||||
const SourcePrefix = "crawl:"
|
||||
|
||||
// noteRunes — how much of a watched page goes into a note. Shorter than what the
|
||||
// on-demand path reads: a note is a record of a change, not an archive.
|
||||
const noteRunes = 800
|
||||
|
||||
// NoteText renders a watched page as a note body.
|
||||
func NoteText(watch WatchConfig, page Page) string {
|
||||
var b strings.Builder
|
||||
if page.Title != "" {
|
||||
b.WriteString(page.Title)
|
||||
} else {
|
||||
b.WriteString(watch.Name)
|
||||
}
|
||||
b.WriteString("\n")
|
||||
b.WriteString(TrimRunes(page.Text, noteRunes))
|
||||
b.WriteString("\n")
|
||||
b.WriteString(watch.URL)
|
||||
return b.String()
|
||||
}
|
||||
@@ -0,0 +1,117 @@
|
||||
package crawl
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
type note struct {
|
||||
text string
|
||||
source string
|
||||
}
|
||||
|
||||
type fakeNotes struct{ notes []note }
|
||||
|
||||
func (n *fakeNotes) WriteNote(_ context.Context, _ time.Time, text string, _ []float32, source string) (int64, error) {
|
||||
n.notes = append(n.notes, note{text, source})
|
||||
return int64(len(n.notes)), nil
|
||||
}
|
||||
|
||||
type fakeHashes struct{ m map[string]string }
|
||||
|
||||
func newHashes() *fakeHashes { return &fakeHashes{m: map[string]string{}} }
|
||||
func (f *fakeHashes) LastHash(_ context.Context, name string) (string, error) {
|
||||
return f.m[name], nil
|
||||
}
|
||||
func (f *fakeHashes) SetHash(_ context.Context, name, h string) error { f.m[name] = h; return nil }
|
||||
|
||||
var t0 = time.Date(2026, 8, 1, 9, 0, 0, 0, time.UTC)
|
||||
|
||||
func TestWatchNotesAChangedPage(t *testing.T) {
|
||||
f := &fakeFetcher{pages: map[string]Response{
|
||||
"https://example.org/docs": {Body: []byte(htmlPage)},
|
||||
}}
|
||||
notes := &fakeNotes{}
|
||||
hashes := newHashes()
|
||||
w := NewWatcher(newTestCrawler(f), []WatchConfig{{Name: "docs", URL: "https://example.org/docs"}},
|
||||
notes, hashes, nil, time.Hour)
|
||||
if w == nil {
|
||||
t.Fatal("NewWatcher returned nil for a configured watch")
|
||||
}
|
||||
if n := w.CheckDue(context.Background(), t0); n != 1 {
|
||||
t.Fatalf("first check wrote %d notes, want 1", n)
|
||||
}
|
||||
if notes.notes[0].source != "crawl:docs" {
|
||||
t.Errorf("source = %q, want crawl:docs", notes.notes[0].source)
|
||||
}
|
||||
if !strings.Contains(notes.notes[0].text, "https://example.org/docs") {
|
||||
t.Errorf("note does not carry the url: %q", notes.notes[0].text)
|
||||
}
|
||||
|
||||
// Unchanged page, interval elapsed: nothing written.
|
||||
if n := w.CheckDue(context.Background(), t0.Add(2*time.Hour)); n != 0 {
|
||||
t.Fatalf("an unchanged page wrote %d notes", n)
|
||||
}
|
||||
|
||||
// Changed page: one note.
|
||||
f.pages["https://example.org/docs"] = Response{Body: []byte(strings.Replace(htmlPage, "синее", "серое", 1))}
|
||||
if n := w.CheckDue(context.Background(), t0.Add(4*time.Hour)); n != 1 {
|
||||
t.Fatalf("a changed page wrote %d notes, want 1", n)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWatchIntervalIsRespected(t *testing.T) {
|
||||
f := &fakeFetcher{pages: map[string]Response{"https://example.org/d": {Body: []byte(htmlPage)}}}
|
||||
w := NewWatcher(newTestCrawler(f), []WatchConfig{{Name: "d", URL: "https://example.org/d", Interval: time.Hour}},
|
||||
&fakeNotes{}, newHashes(), nil, 0)
|
||||
w.CheckDue(context.Background(), t0)
|
||||
before := len(f.calls)
|
||||
w.CheckDue(context.Background(), t0.Add(time.Minute))
|
||||
if len(f.calls) != before {
|
||||
t.Fatal("the page was re-read inside its interval")
|
||||
}
|
||||
}
|
||||
|
||||
// The hash is durable so a restart does not re-note an unchanged page.
|
||||
func TestWatchHashSurvivesRestart(t *testing.T) {
|
||||
f := &fakeFetcher{pages: map[string]Response{"https://example.org/d": {Body: []byte(htmlPage)}}}
|
||||
hashes := newHashes()
|
||||
watches := []WatchConfig{{Name: "d", URL: "https://example.org/d"}}
|
||||
NewWatcher(newTestCrawler(f), watches, &fakeNotes{}, hashes, nil, time.Hour).CheckDue(context.Background(), t0)
|
||||
|
||||
notes2 := &fakeNotes{}
|
||||
NewWatcher(newTestCrawler(f), watches, notes2, hashes, nil, time.Hour).CheckDue(context.Background(), t0.Add(time.Hour))
|
||||
if len(notes2.notes) != 0 {
|
||||
t.Fatalf("a fresh watcher re-noted an unchanged page: %q", notes2.notes[0].text)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWatchDeadPageDoesNotStopTheOthers(t *testing.T) {
|
||||
f := &fakeFetcher{pages: map[string]Response{"https://example.org/live": {Body: []byte(htmlPage)}}}
|
||||
notes := &fakeNotes{}
|
||||
w := NewWatcher(newTestCrawler(f), []WatchConfig{
|
||||
{Name: "dead", URL: "https://example.org/gone"},
|
||||
{Name: "live", URL: "https://example.org/live"},
|
||||
}, notes, newHashes(), nil, time.Hour)
|
||||
if n := w.CheckDue(context.Background(), t0); n != 1 {
|
||||
t.Fatalf("wrote %d notes, want 1 (the live page)", n)
|
||||
}
|
||||
if notes.notes[0].source != "crawl:live" {
|
||||
t.Fatalf("source = %q", notes.notes[0].source)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNoWatchesMeansNoWatcher(t *testing.T) {
|
||||
c := newTestCrawler(&fakeFetcher{})
|
||||
if NewWatcher(c, nil, &fakeNotes{}, nil, nil, 0) != nil {
|
||||
t.Fatal("no watches must mean no watcher")
|
||||
}
|
||||
if NewWatcher(nil, []WatchConfig{{Name: "a", URL: "u"}}, &fakeNotes{}, nil, nil, 0) != nil {
|
||||
t.Fatal("no crawler must mean no watcher")
|
||||
}
|
||||
if NewWatcher(c, []WatchConfig{{Name: "", URL: ""}}, &fakeNotes{}, nil, nil, 0) != nil {
|
||||
t.Fatal("a watch with no name or url is not a configuration")
|
||||
}
|
||||
}
|
||||
@@ -176,6 +176,50 @@ type IngestMailResp struct {
|
||||
Skipped bool `json:"skipped,omitempty"`
|
||||
}
|
||||
|
||||
// SwapModelReq — load another resident model without restarting the daemon
|
||||
// (Vikunja #250). ModelPath must be one of the paths in phraser.swap_models;
|
||||
// anything else is ErrForbidden, and an unconfigured allowlist makes the whole
|
||||
// method ErrUnknownMethod.
|
||||
//
|
||||
// NGpuLayers and NCtx are zero for "keep what is loaded now", which is the
|
||||
// normal case — the same laptop iGPU, a different gguf.
|
||||
//
|
||||
// This is an owner action. It is AuthStepUp in the authority table, it is not on
|
||||
// CoreAPI, and no act, intent or timer can reach it: swapping the model is not
|
||||
// something Maven does to herself.
|
||||
type SwapModelReq struct {
|
||||
ModelPath string `json:"model_path"`
|
||||
NGpuLayers int `json:"n_gpu_layers,omitempty"`
|
||||
NCtx int `json:"n_ctx,omitempty"`
|
||||
}
|
||||
|
||||
// SwapModelResp — what the daemon ended up serving. Model is the identity the
|
||||
// new llama-server reported for itself, not an echo of the request: if the file
|
||||
// was not the model the operator thought it was, this is where it shows.
|
||||
//
|
||||
// RolledBack is true when the requested model failed to load or would not answer
|
||||
// and the previous one was put back. In that case the call also returns an error
|
||||
// — the swap did not happen — and Model names the model still serving.
|
||||
type SwapModelResp struct {
|
||||
Model string `json:"model"`
|
||||
ModelPath string `json:"model_path"`
|
||||
BaseURL string `json:"base_url"`
|
||||
RolledBack bool `json:"rolled_back,omitempty"`
|
||||
TookMs int64 `json:"took_ms"`
|
||||
}
|
||||
|
||||
// ModelStatusResp — which model is resident and which ones may be swapped in.
|
||||
// Read-only; the authed page renders it. Swappable is the configured allowlist,
|
||||
// so an empty list means the capability is off.
|
||||
type ModelStatusResp struct {
|
||||
Model string `json:"model"`
|
||||
ModelPath string `json:"model_path"`
|
||||
BaseURL string `json:"base_url"`
|
||||
NGpuLayers int `json:"n_gpu_layers"`
|
||||
NCtx int `json:"n_ctx"`
|
||||
Swappable []string `json:"swappable,omitempty"`
|
||||
}
|
||||
|
||||
type listTasksReq struct {
|
||||
Status string `json:"status"` // "" all | "live" | candidate|open|done|dropped
|
||||
}
|
||||
|
||||
@@ -459,6 +459,28 @@ func (c *Client) IngestMail(ctx context.Context, req IngestMailReq) (IngestMailR
|
||||
return r, nil
|
||||
}
|
||||
|
||||
// SwapModel asks core to load another resident model (Vikunja #250).
|
||||
// ErrUnknownMethod means core has no phraser.swap_models allowlist configured;
|
||||
// ErrForbidden means the path is not on it, or step-up was not asserted. A
|
||||
// non-nil error with RolledBack set means nothing changed — the old model is
|
||||
// still serving.
|
||||
func (c *Client) SwapModel(ctx context.Context, req SwapModelReq) (SwapModelResp, error) {
|
||||
var r SwapModelResp
|
||||
if err := c.call(ctx, MethodSwapModel, req, &r); err != nil {
|
||||
return SwapModelResp{}, err
|
||||
}
|
||||
return r, nil
|
||||
}
|
||||
|
||||
// ModelStatus reports the resident model and the swap allowlist. Read-only.
|
||||
func (c *Client) ModelStatus(ctx context.Context) (ModelStatusResp, error) {
|
||||
var r ModelStatusResp
|
||||
if err := c.call(ctx, MethodModelStatus, nil, &r); err != nil {
|
||||
return ModelStatusResp{}, err
|
||||
}
|
||||
return r, nil
|
||||
}
|
||||
|
||||
func (c *Client) DismissProposedRoutine(ctx context.Context, id int64) error {
|
||||
return c.call(ctx, MethodDismissProposedRoutine, dismissProposedRoutineReq{ID: id}, nil)
|
||||
}
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"os"
|
||||
@@ -598,3 +599,44 @@ func TestIngestMail_Hook(t *testing.T) {
|
||||
t.Errorf("req across the wire = %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSwapModel_OffUnlessConfigured — no allowlist in the config means the
|
||||
// daemon never sets the hook, and the method does not exist. That is what "off
|
||||
// unless configured" looks like at the wire for the model swap (Vikunja #250).
|
||||
func TestSwapModel_OffUnlessConfigured(t *testing.T) {
|
||||
_, _, cli, _ := newServerWithStore(t)
|
||||
if _, err := cli.SwapModel(context.Background(), SwapModelReq{ModelPath: "/m/x.gguf"}); !errors.Is(err, ErrUnknownMethod) {
|
||||
t.Fatalf("SwapModel error = %v, want ErrUnknownMethod", err)
|
||||
}
|
||||
if _, err := cli.ModelStatus(context.Background()); !errors.Is(err, ErrUnknownMethod) {
|
||||
t.Fatalf("ModelStatus error = %v, want ErrUnknownMethod", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSwapModel_Hook — the request crosses the boundary intact and the reported
|
||||
// identity comes back. A refusal from the daemon's allowlist arrives as
|
||||
// ErrForbidden, which is what a caller keys its error message off.
|
||||
func TestSwapModel_Hook(t *testing.T) {
|
||||
_, srv, cli, _ := newServerWithStore(t)
|
||||
var got SwapModelReq
|
||||
srv.SwapModelFn = func(_ context.Context, req SwapModelReq) (SwapModelResp, error) {
|
||||
got = req
|
||||
if req.ModelPath != "/m/allowed.gguf" {
|
||||
return SwapModelResp{}, fmt.Errorf("%w: not allowlisted", ErrForbidden)
|
||||
}
|
||||
return SwapModelResp{Model: "allowed", ModelPath: req.ModelPath, BaseURL: "http://127.0.0.1:9", TookMs: 12}, nil
|
||||
}
|
||||
resp, err := cli.SwapModel(context.Background(), SwapModelReq{ModelPath: "/m/allowed.gguf", NCtx: 4096})
|
||||
if err != nil {
|
||||
t.Fatalf("SwapModel: %v", err)
|
||||
}
|
||||
if resp.Model != "allowed" || resp.TookMs != 12 {
|
||||
t.Errorf("resp = %+v", resp)
|
||||
}
|
||||
if got.NCtx != 4096 {
|
||||
t.Errorf("req across the wire = %+v", got)
|
||||
}
|
||||
if _, err := cli.SwapModel(context.Background(), SwapModelReq{ModelPath: "/etc/shadow"}); !errors.Is(err, ErrForbidden) {
|
||||
t.Fatalf("swap to a non-allowlisted path = %v; want ErrForbidden", err)
|
||||
}
|
||||
}
|
||||
|
||||
+46
-2
@@ -432,6 +432,19 @@ type Server struct {
|
||||
// every CoreAPI implementation has to carry.
|
||||
IngestMailFn IngestMailFunc
|
||||
|
||||
// SwapModelFn / ModelStatusFn — the on-the-fly resident model swap (Vikunja
|
||||
// #250) and its read side. Set by the daemon only when phraser.swap_models
|
||||
// lists at least one model AND the phraser owns a llama-server; nil ⇒ both
|
||||
// methods answer ErrUnknownMethod, which is what "off unless configured"
|
||||
// looks like at the wire.
|
||||
//
|
||||
// They bypass CoreAPI for the same reason IngestMailFn does: this is not a
|
||||
// store operation, it needs the daemon's llama-server, and no other CoreAPI
|
||||
// implementation should have to carry it. MethodSwapModel is AuthStepUp in
|
||||
// internal/auth — owner-triggered, never an act and never a timer.
|
||||
SwapModelFn SwapModelFunc
|
||||
ModelStatusFn ModelStatusFunc
|
||||
|
||||
// UnlockFn — unwraps the store encryption key from the wrapped blob using
|
||||
// the passkey credential public key, opens the encrypted store, and wires
|
||||
// the rest of the daemon (voice, loop, delivery). Set by the daemon when
|
||||
@@ -450,6 +463,12 @@ type WrapKeyFunc func(ctx context.Context, publicKey []byte) error
|
||||
// public key and completes daemon initialization.
|
||||
type UnlockFunc func(ctx context.Context, publicKey []byte) error
|
||||
|
||||
// SwapModelFunc — loads another resident model in place of the live one.
|
||||
type SwapModelFunc func(ctx context.Context, req SwapModelReq) (SwapModelResp, error)
|
||||
|
||||
// ModelStatusFunc — reports the resident model and the swap allowlist.
|
||||
type ModelStatusFunc func(ctx context.Context) (ModelStatusResp, error)
|
||||
|
||||
// IngestMailFunc — core-side mail extraction. Returns what was captured.
|
||||
type IngestMailFunc func(ctx context.Context, req IngestMailReq) (IngestMailResp, error)
|
||||
|
||||
@@ -620,8 +639,9 @@ func withoutParams[R any](fn func(ctx context.Context, api CoreAPI) (R, error))
|
||||
// existed) as an argument — so SetAPI's runtime swap (the unlock transition)
|
||||
// is still honored on the very next request with no extra plumbing here.
|
||||
//
|
||||
// MethodAssertStepUp, MethodStoreEncryptionKey, MethodUnlock and
|
||||
// MethodIngestMail are NOT in this table: they bypass CoreAPI entirely
|
||||
// MethodAssertStepUp, MethodStoreEncryptionKey, MethodUnlock,
|
||||
// MethodIngestMail, MethodSwapModel and MethodModelStatus are NOT in this
|
||||
// table: they bypass CoreAPI entirely
|
||||
// (s.StepUp / s.WrapKeyFn / s.UnlockFn / s.IngestMailFn), so dispatch
|
||||
// special-cases them before consulting the table.
|
||||
var methodTable = map[Method]handlerFunc{
|
||||
@@ -874,6 +894,30 @@ func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, er
|
||||
return marshalResult(resp), nil
|
||||
}
|
||||
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
|
||||
|
||||
case MethodSwapModel:
|
||||
if s.SwapModelFn != nil {
|
||||
var p SwapModelReq
|
||||
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
resp, err := s.SwapModelFn(ctx, p)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(resp), nil
|
||||
}
|
||||
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
|
||||
|
||||
case MethodModelStatus:
|
||||
if s.ModelStatusFn != nil {
|
||||
resp, err := s.ModelStatusFn(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return marshalResult(resp), nil
|
||||
}
|
||||
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
|
||||
}
|
||||
|
||||
h, ok := methodTable[req.Method]
|
||||
|
||||
@@ -51,6 +51,8 @@ const (
|
||||
MethodListTasks Method = "list_tasks"
|
||||
MethodSetTaskStatus Method = "set_task_status"
|
||||
MethodIngestMail Method = "ingest_mail"
|
||||
MethodSwapModel Method = "swap_model"
|
||||
MethodModelStatus Method = "model_status"
|
||||
)
|
||||
|
||||
// Request — one frame from module to core. Params is the JSON-encoded argument
|
||||
|
||||
+27
-1
@@ -10,10 +10,17 @@ import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Client struct {
|
||||
// mu guards base only. The base URL changes when the daemon swaps the
|
||||
// resident model (Vikunja #250): llama-server is relaunched on a fresh
|
||||
// port, and every holder of this client — the LLM router, the replier, the
|
||||
// mail extractor — must follow without being rebuilt. One mutexed field is
|
||||
// the whole mechanism; a swap re-points the client, it does not replace it.
|
||||
mu sync.RWMutex
|
||||
base string
|
||||
http *http.Client
|
||||
}
|
||||
@@ -22,6 +29,25 @@ func New(baseURL string, timeout time.Duration) *Client {
|
||||
return &Client{base: baseURL, http: &http.Client{Timeout: timeout}}
|
||||
}
|
||||
|
||||
// SetBaseURL re-points the client at another llama-server. Safe to call while
|
||||
// requests are in flight: a request that already read the old base finishes
|
||||
// against the old base (or fails, and every caller of Complete has a fallback),
|
||||
// and the next one uses the new base. It is deliberately NOT a queue-and-retry —
|
||||
// the phraser quiesces around a swap, so the window is small and a lost turn
|
||||
// degrades to the classifier rather than hanging.
|
||||
func (c *Client) SetBaseURL(base string) {
|
||||
c.mu.Lock()
|
||||
c.base = base
|
||||
c.mu.Unlock()
|
||||
}
|
||||
|
||||
// BaseURL is the server this client currently talks to.
|
||||
func (c *Client) BaseURL() string {
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
return c.base
|
||||
}
|
||||
|
||||
type Req struct {
|
||||
System string
|
||||
User string
|
||||
@@ -63,7 +89,7 @@ func (c *Client) Complete(ctx context.Context, r Req) (string, error) {
|
||||
RepeatPenalty: r.RepeatPenalty,
|
||||
Stop: r.Stop,
|
||||
})
|
||||
req, err := http.NewRequestWithContext(ctx, "POST", c.base+"/v1/chat/completions", bytes.NewReader(b))
|
||||
req, err := http.NewRequestWithContext(ctx, "POST", c.BaseURL()+"/v1/chat/completions", bytes.NewReader(b))
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
@@ -57,3 +57,37 @@ func TestComplete(t *testing.T) {
|
||||
t.Errorf("got %q, want %q", got, "ok")
|
||||
}
|
||||
}
|
||||
|
||||
// TestSetBaseURL — a model swap re-points every holder of the client rather than
|
||||
// rebuilding the router, the replier and the extractors (Vikunja #250).
|
||||
func TestSetBaseURL(t *testing.T) {
|
||||
var hit string
|
||||
srvA := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
hit = "A"
|
||||
w.Write([]byte(`{"choices":[{"message":{"content":"a"}}]}`))
|
||||
}))
|
||||
defer srvA.Close()
|
||||
srvB := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
hit = "B"
|
||||
w.Write([]byte(`{"choices":[{"message":{"content":"b"}}]}`))
|
||||
}))
|
||||
defer srvB.Close()
|
||||
|
||||
c := New(srvA.URL, 5*time.Second)
|
||||
if _, err := c.Complete(context.Background(), Req{User: "x"}); err != nil {
|
||||
t.Fatalf("Complete against A: %v", err)
|
||||
}
|
||||
if hit != "A" {
|
||||
t.Fatalf("first request went to %q; want A", hit)
|
||||
}
|
||||
c.SetBaseURL(srvB.URL)
|
||||
if got := c.BaseURL(); got != srvB.URL {
|
||||
t.Errorf("BaseURL = %q; want %q", got, srvB.URL)
|
||||
}
|
||||
if _, err := c.Complete(context.Background(), Req{User: "x"}); err != nil {
|
||||
t.Fatalf("Complete against B: %v", err)
|
||||
}
|
||||
if hit != "B" {
|
||||
t.Errorf("request after the swap went to %q; want B", hit)
|
||||
}
|
||||
}
|
||||
|
||||
+150
-34
@@ -27,14 +27,46 @@ var listenRE = regexp.MustCompile(`listening on (https?://\S+)`)
|
||||
type LLMPhraser struct {
|
||||
cfg Config
|
||||
client *http.Client
|
||||
port string
|
||||
cmd *exec.Cmd
|
||||
cancel context.CancelFunc
|
||||
wg sync.WaitGroup
|
||||
|
||||
// tmpl — the hand-written Russian nudges. Default path for nudges; see
|
||||
// Config.LLMNudges. nil only if the template file failed to load.
|
||||
tmpl *NudgeTemplates
|
||||
|
||||
// spawnCtx — the parent of every llama-server this phraser starts, i.e. the
|
||||
// daemon's own context. Deliberately NOT the per-request context of the call
|
||||
// that asked for a model swap: that one is cancelled the moment the request
|
||||
// returns, which would kill the model it had just loaded.
|
||||
spawnCtx context.Context
|
||||
cancel context.CancelFunc
|
||||
|
||||
// launch / probe — the two side effects of a swap, injectable so the swap
|
||||
// logic is testable without a real llama-server and a real model file.
|
||||
// launch is nil when this phraser does not own its server (NewLLMPhraserAt),
|
||||
// which is also what makes Swap refuse there.
|
||||
launch func(ctx context.Context, cfg Config) (backend, error)
|
||||
probe func(ctx context.Context, base string) (string, error)
|
||||
|
||||
// swapMu — single-flight around Swap. Held for the whole swap, including the
|
||||
// model load, so two concurrent swap requests can never both be loading.
|
||||
swapMu sync.Mutex
|
||||
|
||||
// mu guards everything below: the live backend, the swap gate and the
|
||||
// in-flight request count. See acquire/quiesce in swap.go.
|
||||
mu sync.Mutex
|
||||
be backend
|
||||
live liveModel
|
||||
swapping bool
|
||||
inflight int
|
||||
observers []func(baseURL string)
|
||||
}
|
||||
|
||||
// liveModel — what is actually loaded right now. Distinct from Config, which
|
||||
// stays immutable after construction: a swap changes these three fields and
|
||||
// nothing else, so no reader of cfg (prompts, grammar, timeouts) races a swap.
|
||||
type liveModel struct {
|
||||
ModelPath string
|
||||
NGpuLayers int
|
||||
NCtx int
|
||||
}
|
||||
|
||||
type Config struct {
|
||||
@@ -85,15 +117,21 @@ func DefaultConfig(modelPath string) Config {
|
||||
func NewLLMPhraser(ctx context.Context, cfg Config) (*LLMPhraser, error) {
|
||||
ctx, cancel := context.WithCancel(ctx)
|
||||
p := &LLMPhraser{
|
||||
cfg: cfg,
|
||||
client: &http.Client{Timeout: cfg.Timeout},
|
||||
cancel: cancel,
|
||||
tmpl: loadNudgeTemplates(),
|
||||
cfg: cfg,
|
||||
client: &http.Client{Timeout: cfg.Timeout},
|
||||
tmpl: loadNudgeTemplates(),
|
||||
spawnCtx: ctx,
|
||||
cancel: cancel,
|
||||
launch: spawnLlamaServer,
|
||||
probe: defaultProbe,
|
||||
live: liveModel{ModelPath: cfg.ModelPath, NGpuLayers: cfg.NGpuLayers, NCtx: cfg.NCtx},
|
||||
}
|
||||
if err := p.start(ctx); err != nil {
|
||||
be, err := p.launch(ctx, cfg)
|
||||
if err != nil {
|
||||
cancel()
|
||||
return nil, err
|
||||
}
|
||||
p.be = be
|
||||
return p, nil
|
||||
}
|
||||
|
||||
@@ -106,11 +144,17 @@ func NewLLMPhraser(ctx context.Context, cfg Config) (*LLMPhraser, error) {
|
||||
// still uses NewLLMPhraser and still owns its own child process.
|
||||
func NewLLMPhraserAt(baseURL string, cfg Config) *LLMPhraser {
|
||||
return &LLMPhraser{
|
||||
cfg: cfg,
|
||||
client: &http.Client{Timeout: cfg.Timeout},
|
||||
port: strings.TrimSuffix(baseURL, "/"),
|
||||
cancel: func() {},
|
||||
tmpl: loadNudgeTemplates(),
|
||||
cfg: cfg,
|
||||
client: &http.Client{Timeout: cfg.Timeout},
|
||||
tmpl: loadNudgeTemplates(),
|
||||
spawnCtx: context.Background(),
|
||||
cancel: func() {},
|
||||
probe: defaultProbe,
|
||||
// launch stays nil: we did not start this server, so we must not stop it.
|
||||
// Swap therefore refuses here (ErrSwapNotOwned) instead of killing a
|
||||
// server another process depends on.
|
||||
be: borrowedBackend(strings.TrimSuffix(baseURL, "/")),
|
||||
live: liveModel{ModelPath: cfg.ModelPath, NGpuLayers: cfg.NGpuLayers, NCtx: cfg.NCtx},
|
||||
}
|
||||
}
|
||||
|
||||
@@ -126,16 +170,66 @@ func loadNudgeTemplates() *NudgeTemplates {
|
||||
return nt
|
||||
}
|
||||
|
||||
func (p *LLMPhraser) start(ctx context.Context) error {
|
||||
// backend — one llama-server this phraser talks to. Two implementations: a
|
||||
// llamaProc we spawned and must reap, and a borrowedBackend someone else owns.
|
||||
type backend interface {
|
||||
BaseURL() string
|
||||
Close() error
|
||||
}
|
||||
|
||||
// borrowedBackend — a server started and owned by someone else (the phrasing
|
||||
// scorer's shared llama-server). Closing it is a no-op by construction.
|
||||
type borrowedBackend string
|
||||
|
||||
func (b borrowedBackend) BaseURL() string { return string(b) }
|
||||
func (b borrowedBackend) Close() error { return nil }
|
||||
|
||||
// llamaProc — a llama-server child process plus the goroutine reading its
|
||||
// stderr. Close kills and reaps it; see the Pdeathsig note in spawnLlamaServer.
|
||||
type llamaProc struct {
|
||||
base string
|
||||
cmd *exec.Cmd
|
||||
cancel context.CancelFunc
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
func (l *llamaProc) BaseURL() string { return l.base }
|
||||
|
||||
func (l *llamaProc) Close() error {
|
||||
l.cancel()
|
||||
if l.cmd != nil && l.cmd.Process != nil {
|
||||
_ = l.cmd.Process.Kill()
|
||||
_ = l.cmd.Wait() // reap the process — without Wait, the child becomes a zombie
|
||||
}
|
||||
l.wg.Wait()
|
||||
return nil
|
||||
}
|
||||
|
||||
// spawnLlamaServer starts one llama-server for cfg and waits until it says which
|
||||
// address it is listening on. ctx owns the process lifetime, so it must be the
|
||||
// daemon's context, not a request's.
|
||||
func spawnLlamaServer(ctx context.Context, cfg Config) (backend, error) {
|
||||
ctx, cancel := context.WithCancel(ctx)
|
||||
p, err := startLlamaProc(ctx, cfg)
|
||||
if err != nil {
|
||||
cancel()
|
||||
return nil, err
|
||||
}
|
||||
p.cancel = cancel
|
||||
return p, nil
|
||||
}
|
||||
|
||||
func startLlamaProc(ctx context.Context, cfg Config) (*llamaProc, error) {
|
||||
p := &llamaProc{}
|
||||
args := []string{
|
||||
"-m", p.cfg.ModelPath,
|
||||
"-m", cfg.ModelPath,
|
||||
"--host", "127.0.0.1",
|
||||
"--port", extractPort(p.cfg.Listen),
|
||||
"-c", fmt.Sprintf("%d", p.cfg.NCtx),
|
||||
"-ngl", fmt.Sprintf("%d", p.cfg.NGpuLayers),
|
||||
"--port", extractPort(cfg.Listen),
|
||||
"-c", fmt.Sprintf("%d", cfg.NCtx),
|
||||
"-ngl", fmt.Sprintf("%d", cfg.NGpuLayers),
|
||||
"--no-webui",
|
||||
}
|
||||
cmd := exec.CommandContext(ctx, p.cfg.BinPath, args...)
|
||||
cmd := exec.CommandContext(ctx, cfg.BinPath, args...)
|
||||
// Pdeathsig: the kernel SIGKILLs llama-server the moment mavend dies — by
|
||||
// ANY means, including SIGKILL/OOM/panic where our Close() never runs. Without
|
||||
// it a hard-killed mavend orphans its llama-server (reparented to init, keeps
|
||||
@@ -148,12 +242,12 @@ func (p *LLMPhraser) start(ctx context.Context) error {
|
||||
|
||||
stderr, err := cmd.StderrPipe()
|
||||
if err != nil {
|
||||
return fmt.Errorf("llm: stderr pipe: %w", err)
|
||||
return nil, fmt.Errorf("llm: stderr pipe: %w", err)
|
||||
}
|
||||
|
||||
if err := cmd.Start(); err != nil {
|
||||
stderr.Close()
|
||||
return fmt.Errorf("llm: start: %w", err)
|
||||
return nil, fmt.Errorf("llm: start: %w", err)
|
||||
}
|
||||
|
||||
portCh := make(chan string, 1)
|
||||
@@ -186,32 +280,44 @@ func (p *LLMPhraser) start(ctx context.Context) error {
|
||||
|
||||
select {
|
||||
case addr := <-portCh:
|
||||
p.port = addr
|
||||
return nil
|
||||
p.base = addr
|
||||
return p, nil
|
||||
case err := <-errCh:
|
||||
_ = cmd.Process.Kill()
|
||||
_ = cmd.Wait()
|
||||
return fmt.Errorf("llm: server output: %w", err)
|
||||
return nil, fmt.Errorf("llm: server output: %w", err)
|
||||
case <-ctx.Done():
|
||||
_ = cmd.Process.Kill()
|
||||
_ = cmd.Wait()
|
||||
return ctx.Err()
|
||||
return nil, ctx.Err()
|
||||
case <-time.After(60 * time.Second):
|
||||
_ = cmd.Process.Kill()
|
||||
_ = cmd.Wait()
|
||||
return fmt.Errorf("llm: server did not start within 60s")
|
||||
return nil, fmt.Errorf("llm: server did not start within 60s")
|
||||
}
|
||||
}
|
||||
|
||||
func (p *LLMPhraser) BaseURL() string { return p.port }
|
||||
// BaseURL is the llama-server this phraser talks to right now. It changes when
|
||||
// the model is swapped, so callers that cache it must register an observer
|
||||
// (OnSwap) rather than keeping the string forever.
|
||||
func (p *LLMPhraser) BaseURL() string {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
if p.be == nil {
|
||||
return ""
|
||||
}
|
||||
return p.be.BaseURL()
|
||||
}
|
||||
|
||||
func (p *LLMPhraser) Close() error {
|
||||
p.cancel()
|
||||
if p.cmd != nil && p.cmd.Process != nil {
|
||||
_ = p.cmd.Process.Kill()
|
||||
_ = p.cmd.Wait() // reap the process — without Wait, the child becomes a zombie
|
||||
p.mu.Lock()
|
||||
be := p.be
|
||||
p.be = nil
|
||||
p.mu.Unlock()
|
||||
if be != nil {
|
||||
return be.Close()
|
||||
}
|
||||
p.wg.Wait()
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -350,6 +456,11 @@ func chatSystemPrompt(block func() string) string {
|
||||
// the LLM completion endpoint. Like chatWithSystem but for an arbitrary message
|
||||
// slice — the caller owns the system prompt placement.
|
||||
func (p *LLMPhraser) chatWithMessages(ctx context.Context, msgs []chatMsg, maxTokens int) (string, error) {
|
||||
base, release, err := p.acquire()
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer release()
|
||||
req := chatReq{
|
||||
Messages: msgs,
|
||||
Temperature: 0.7,
|
||||
@@ -360,7 +471,7 @@ func (p *LLMPhraser) chatWithMessages(ctx context.Context, msgs []chatMsg, maxTo
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("llm: marshal: %w", err)
|
||||
}
|
||||
httpReq, err := http.NewRequestWithContext(ctx, "POST", p.port+"/v1/chat/completions", bytes.NewReader(body))
|
||||
httpReq, err := http.NewRequestWithContext(ctx, "POST", base+"/v1/chat/completions", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("llm: request: %w", err)
|
||||
}
|
||||
@@ -493,6 +604,11 @@ func (p *LLMPhraser) chat(ctx context.Context, userPrompt string) (string, error
|
||||
}
|
||||
|
||||
func (p *LLMPhraser) chatWithSystem(ctx context.Context, system, user string, maxTokens int) (string, error) {
|
||||
base, release, err := p.acquire()
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer release()
|
||||
req := chatReq{
|
||||
Messages: []chatMsg{
|
||||
{Role: "system", Content: system},
|
||||
@@ -507,7 +623,7 @@ func (p *LLMPhraser) chatWithSystem(ctx context.Context, system, user string, ma
|
||||
return "", fmt.Errorf("llm: marshal: %w", err)
|
||||
}
|
||||
|
||||
httpReq, err := http.NewRequestWithContext(ctx, "POST", p.port+"/v1/chat/completions", bytes.NewReader(body))
|
||||
httpReq, err := http.NewRequestWithContext(ctx, "POST", base+"/v1/chat/completions", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("llm: request: %w", err)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,330 @@
|
||||
package phraser
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/llm"
|
||||
)
|
||||
|
||||
// Swapping the resident model without restarting the daemon (Vikunja #250).
|
||||
//
|
||||
// Three properties this file exists to hold, in order of importance:
|
||||
//
|
||||
// 1. NEVER two models resident at once. The deploy target is a laptop iGPU
|
||||
// with the whole 1.7B offloaded to it (`n_gpu_layers: 99`); loading a second
|
||||
// model beside the first is how you OOM the box, and a blue/green swap that
|
||||
// "keeps the old one warm until the new one answers" does exactly that. So
|
||||
// the old server is killed FIRST and the new one loaded after. The cost of
|
||||
// that ordering is a window with no model at all, which is why:
|
||||
//
|
||||
// 2. A swap is atomic from a turn's point of view. An in-flight turn finishes
|
||||
// on the old model — Swap waits for the last one to return before killing
|
||||
// anything. A turn that arrives during the swap is REFUSED immediately with
|
||||
// ErrSwapping rather than blocked: every phrasing path already has a
|
||||
// fallback (templates, "вот что я нашла", the classifier for routing), so a
|
||||
// fast refusal degrades one turn instead of hanging it for the length of a
|
||||
// model load. No turn ever gets half of one model and half of another.
|
||||
//
|
||||
// 3. A failed load rolls back to the model that was working. The new server is
|
||||
// probed (it must say which model it loaded) before it is published; if the
|
||||
// launch or the probe fails, the previous config is relaunched and the
|
||||
// phraser goes back to serving. Only if the rollback ALSO fails is the
|
||||
// phraser left without a backend, and then it says so loudly and every turn
|
||||
// degrades rather than breaks.
|
||||
//
|
||||
// Not here, deliberately: nothing calls Swap on a timer, and no act or intent can
|
||||
// reach it. It is an IPC method behind the step-up gate, i.e. owner-triggered.
|
||||
|
||||
var (
|
||||
// ErrSwapping — a turn arrived while the model was being swapped. Callers
|
||||
// treat it like any other LLM error and use their fallback.
|
||||
ErrSwapping = errors.New("phraser: model swap in progress")
|
||||
|
||||
// ErrSwapNotOwned — this phraser did not start its llama-server, so it must
|
||||
// not stop one (NewLLMPhraserAt: the eval harness shares a server).
|
||||
ErrSwapNotOwned = errors.New("phraser: llama-server is not ours to swap")
|
||||
|
||||
// ErrNoBackend — no model is loaded at all. Only reachable after a failed
|
||||
// swap whose rollback also failed.
|
||||
ErrNoBackend = errors.New("phraser: no llama-server loaded")
|
||||
|
||||
// ErrSwapBusy — a turn was still running when the drain deadline expired, so
|
||||
// the swap was abandoned. Nothing was killed; ask again.
|
||||
ErrSwapBusy = errors.New("phraser: turns still in flight, swap abandoned")
|
||||
)
|
||||
|
||||
// SwapSpec — what to load. Zero NGpuLayers/NCtx keep whatever is live, so the
|
||||
// common case ("same settings, different gguf") is one field.
|
||||
type SwapSpec struct {
|
||||
ModelPath string
|
||||
NGpuLayers int
|
||||
NCtx int
|
||||
}
|
||||
|
||||
// SwapResult — what happened. Model is the identity the NEW server reported, so
|
||||
// it is evidence rather than an echo of the request: if the file at ModelPath is
|
||||
// not what the operator thought it was, this is where that shows up.
|
||||
type SwapResult struct {
|
||||
Model string
|
||||
BaseURL string
|
||||
ModelPath string
|
||||
RolledBack bool
|
||||
Took time.Duration
|
||||
}
|
||||
|
||||
// drainTimeout — how long Swap waits for in-flight turns before giving up. A
|
||||
// turn is at most Config.Timeout (30s in deploy) plus the model's own latency;
|
||||
// 90s covers a slow Thinking generation without wedging the caller forever.
|
||||
const drainTimeout = 90 * time.Second
|
||||
|
||||
// probeTimeout — how long the new server gets to answer "which model do you
|
||||
// have". The load itself is bounded by spawnLlamaServer's own 60s wait.
|
||||
const probeTimeout = 30 * time.Second
|
||||
|
||||
// defaultProbe asks the server which model it has loaded. This is the health
|
||||
// check: a server that answers /v1/models has finished loading weights and is
|
||||
// serving, and its answer is the identity we report back.
|
||||
func defaultProbe(ctx context.Context, base string) (string, error) {
|
||||
return llm.ModelID(ctx, base)
|
||||
}
|
||||
|
||||
// OnSwap registers a callback fired with the new base URL every time the live
|
||||
// backend changes, including after a rollback. Holders of an *llm.Client (the
|
||||
// LLM router, the replier, the mail extractor) register SetBaseURL here so a
|
||||
// swap re-points them without rebuilding the router or the handler.
|
||||
//
|
||||
// Callbacks run with no lock held, in registration order.
|
||||
func (p *LLMPhraser) OnSwap(fn func(baseURL string)) {
|
||||
if fn == nil {
|
||||
return
|
||||
}
|
||||
p.mu.Lock()
|
||||
p.observers = append(p.observers, fn)
|
||||
p.mu.Unlock()
|
||||
}
|
||||
|
||||
// LiveModel is the model file currently loaded (and its load settings). Empty
|
||||
// ModelPath means no model is loaded.
|
||||
func (p *LLMPhraser) LiveModel() (path string, nGpuLayers, nCtx int) {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
return p.live.ModelPath, p.live.NGpuLayers, p.live.NCtx
|
||||
}
|
||||
|
||||
// acquire reserves a slot for one request and returns the base URL to use.
|
||||
// Every request path must call it and must call the returned release exactly
|
||||
// once — that count is what Swap drains.
|
||||
func (p *LLMPhraser) acquire() (string, func(), error) {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
if p.swapping {
|
||||
return "", nil, ErrSwapping
|
||||
}
|
||||
if p.be == nil {
|
||||
return "", nil, ErrNoBackend
|
||||
}
|
||||
p.inflight++
|
||||
base := p.be.BaseURL()
|
||||
var once bool
|
||||
return base, func() {
|
||||
p.mu.Lock()
|
||||
if !once {
|
||||
once = true
|
||||
p.inflight--
|
||||
}
|
||||
p.mu.Unlock()
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Swap loads another model in place of the live one. See the file comment for
|
||||
// the properties it guarantees. Returns the new model's reported identity, or
|
||||
// an error plus RolledBack=true when the old model was put back.
|
||||
//
|
||||
// ctx bounds the drain and the probe. It does NOT own the new server's lifetime
|
||||
// — that is the daemon's context, captured at construction — so a swap survives
|
||||
// the request that asked for it.
|
||||
func (p *LLMPhraser) Swap(ctx context.Context, spec SwapSpec) (SwapResult, error) {
|
||||
if spec.ModelPath == "" {
|
||||
return SwapResult{}, fmt.Errorf("phraser: swap needs a model path")
|
||||
}
|
||||
p.swapMu.Lock()
|
||||
defer p.swapMu.Unlock()
|
||||
|
||||
if p.launch == nil {
|
||||
return SwapResult{}, ErrSwapNotOwned
|
||||
}
|
||||
|
||||
started := time.Now()
|
||||
oldLive := p.liveSnapshot()
|
||||
newLive := liveModel{
|
||||
ModelPath: spec.ModelPath,
|
||||
NGpuLayers: pickInt(spec.NGpuLayers, oldLive.NGpuLayers),
|
||||
NCtx: pickInt(spec.NCtx, oldLive.NCtx),
|
||||
}
|
||||
if newLive == oldLive && p.BaseURL() != "" {
|
||||
// Already serving exactly this. Report the live identity rather than
|
||||
// pointlessly unloading and reloading the same weights.
|
||||
base := p.BaseURL()
|
||||
id, err := p.probeWith(ctx, base)
|
||||
if err != nil {
|
||||
return SwapResult{}, err
|
||||
}
|
||||
return SwapResult{Model: id, BaseURL: base, ModelPath: oldLive.ModelPath, Took: time.Since(started)}, nil
|
||||
}
|
||||
|
||||
if err := p.quiesce(ctx); err != nil {
|
||||
return SwapResult{}, err
|
||||
}
|
||||
defer p.resume()
|
||||
|
||||
// Property 1: the old model leaves the GPU before the new one arrives.
|
||||
p.mu.Lock()
|
||||
old := p.be
|
||||
p.be = nil
|
||||
p.mu.Unlock()
|
||||
if old != nil {
|
||||
_ = old.Close()
|
||||
}
|
||||
|
||||
be, err := p.loadAndProbe(ctx, newLive)
|
||||
if err != nil {
|
||||
log.Printf("phraser: swap to %s FAILED (%v) — rolling back to %s", newLive.ModelPath, err, oldLive.ModelPath)
|
||||
rb, rbErr := p.loadAndProbe(ctx, oldLive)
|
||||
if rbErr != nil {
|
||||
log.Printf("phraser: ROLLBACK to %s ALSO FAILED (%v) — no model is loaded, every phrasing path is on its fallback and routing is on the classifier until the daemon is restarted", oldLive.ModelPath, rbErr)
|
||||
return SwapResult{RolledBack: true, Took: time.Since(started)},
|
||||
fmt.Errorf("phraser: swap failed (%w) and rollback failed too: %v", err, rbErr)
|
||||
}
|
||||
p.publish(rb, oldLive)
|
||||
return SwapResult{
|
||||
Model: rb.id, BaseURL: rb.be.BaseURL(), ModelPath: oldLive.ModelPath,
|
||||
RolledBack: true, Took: time.Since(started),
|
||||
},
|
||||
fmt.Errorf("phraser: swap to %s failed, rolled back to %s: %w", newLive.ModelPath, oldLive.ModelPath, err)
|
||||
}
|
||||
p.publish(be, newLive)
|
||||
log.Printf("phraser: model swapped to %s (%s) at %s in %s", newLive.ModelPath, be.id, be.be.BaseURL(), time.Since(started).Round(time.Millisecond))
|
||||
return SwapResult{
|
||||
Model: be.id, BaseURL: be.be.BaseURL(), ModelPath: newLive.ModelPath,
|
||||
Took: time.Since(started),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// loaded — a started server plus the identity it reported.
|
||||
type loaded struct {
|
||||
be backend
|
||||
id string
|
||||
}
|
||||
|
||||
// loadAndProbe starts a server for lm and verifies it answers. A server that
|
||||
// starts but will not say what it loaded is treated as a failed load and is
|
||||
// killed here — publishing it would hand every turn to a backend we could not
|
||||
// confirm.
|
||||
func (p *LLMPhraser) loadAndProbe(ctx context.Context, lm liveModel) (loaded, error) {
|
||||
cfg := p.cfg
|
||||
cfg.ModelPath = lm.ModelPath
|
||||
cfg.NGpuLayers = lm.NGpuLayers
|
||||
cfg.NCtx = lm.NCtx
|
||||
// p.spawnCtx, not ctx: the process must outlive the request asking for it.
|
||||
be, err := p.launch(p.spawnCtx, cfg)
|
||||
if err != nil {
|
||||
return loaded{}, err
|
||||
}
|
||||
id, err := p.probeWith(ctx, be.BaseURL())
|
||||
if err != nil {
|
||||
_ = be.Close()
|
||||
return loaded{}, fmt.Errorf("phraser: %s started but would not answer: %w", lm.ModelPath, err)
|
||||
}
|
||||
return loaded{be: be, id: id}, nil
|
||||
}
|
||||
|
||||
func (p *LLMPhraser) probeWith(ctx context.Context, base string) (string, error) {
|
||||
probe := p.probe
|
||||
if probe == nil {
|
||||
probe = defaultProbe
|
||||
}
|
||||
pctx, cancel := context.WithTimeout(ctx, probeTimeout)
|
||||
defer cancel()
|
||||
return probe(pctx, base)
|
||||
}
|
||||
|
||||
// quiesce closes the door on new turns and waits for the ones already running.
|
||||
// Polling rather than a sync.Cond: the wait happens once per swap, a 25ms poll
|
||||
// is invisible next to a model load, and a poll cannot deadlock on a release
|
||||
// path that panicked.
|
||||
func (p *LLMPhraser) quiesce(ctx context.Context) error {
|
||||
p.mu.Lock()
|
||||
if p.swapping {
|
||||
p.mu.Unlock()
|
||||
return ErrSwapping
|
||||
}
|
||||
p.swapping = true
|
||||
inflight := p.inflight
|
||||
p.mu.Unlock()
|
||||
if inflight == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
deadline := time.Now().Add(drainTimeout)
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
p.resume()
|
||||
return ctx.Err()
|
||||
case <-time.After(25 * time.Millisecond):
|
||||
}
|
||||
p.mu.Lock()
|
||||
inflight = p.inflight
|
||||
p.mu.Unlock()
|
||||
if inflight == 0 {
|
||||
return nil
|
||||
}
|
||||
if time.Now().After(deadline) {
|
||||
// Nothing has been killed yet, so abandoning is free: reopen the door
|
||||
// and let the operator try again rather than cutting a live turn off
|
||||
// mid-generation.
|
||||
p.resume()
|
||||
return fmt.Errorf("%w (%d still running after %s)", ErrSwapBusy, inflight, drainTimeout)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (p *LLMPhraser) resume() {
|
||||
p.mu.Lock()
|
||||
p.swapping = false
|
||||
p.mu.Unlock()
|
||||
}
|
||||
|
||||
// publish makes l the live backend and tells everyone holding a base URL.
|
||||
func (p *LLMPhraser) publish(l loaded, lm liveModel) {
|
||||
p.mu.Lock()
|
||||
p.be = l.be
|
||||
p.live = lm
|
||||
obs := make([]func(string), len(p.observers))
|
||||
copy(obs, p.observers)
|
||||
p.mu.Unlock()
|
||||
base := l.be.BaseURL()
|
||||
for _, fn := range obs {
|
||||
fn(base)
|
||||
}
|
||||
}
|
||||
|
||||
func (p *LLMPhraser) liveSnapshot() liveModel {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
return p.live
|
||||
}
|
||||
|
||||
// pickInt returns v when the caller set it, and fallback otherwise. 0 is the
|
||||
// "unset" value: -1 already means "offload every layer" and deploy uses 99, so
|
||||
// nothing legitimate asks for exactly zero GPU layers through this path.
|
||||
func pickInt(v, fallback int) int {
|
||||
if v == 0 {
|
||||
return fallback
|
||||
}
|
||||
return v
|
||||
}
|
||||
@@ -0,0 +1,311 @@
|
||||
package phraser
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// fakeModel — a stand-in llama-server. It answers /v1/models with its own name
|
||||
// and /v1/chat/completions with a phrasing-contract reply that names itself, so
|
||||
// a test can tell WHICH model answered a turn — the property the swap is about.
|
||||
type fakeModel struct {
|
||||
srv *httptest.Server
|
||||
name string
|
||||
closed atomic.Bool
|
||||
}
|
||||
|
||||
func newFakeModel(t *testing.T, name string) *fakeModel {
|
||||
t.Helper()
|
||||
f := &fakeModel{name: name}
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Write([]byte(`{"data":[{"id":"/models/` + name + `.gguf"}]}`))
|
||||
})
|
||||
mux.HandleFunc("/v1/chat/completions", func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Write([]byte(`{"choices":[{"message":{"content":"{\"response\":\"` + name + `\",\"mood\":\"neutral\"}"}}]}`))
|
||||
})
|
||||
f.srv = httptest.NewServer(mux)
|
||||
t.Cleanup(f.srv.Close)
|
||||
return f
|
||||
}
|
||||
|
||||
func (f *fakeModel) BaseURL() string { return f.srv.URL }
|
||||
func (f *fakeModel) Close() error { f.closed.Store(true); return nil }
|
||||
|
||||
// fakeFleet is the injected launcher: it hands out a prepared fakeModel per
|
||||
// model path, and refuses paths the test did not prepare (that is what a bad
|
||||
// gguf looks like from here). It also asserts the invariant that matters on a
|
||||
// laptop iGPU: never two servers alive at the same time.
|
||||
type fakeFleet struct {
|
||||
mu sync.Mutex
|
||||
models map[string]string // model path → fake name
|
||||
live int
|
||||
maxLive int
|
||||
launch int
|
||||
}
|
||||
|
||||
func (fl *fakeFleet) launcher(t *testing.T) func(context.Context, Config) (backend, error) {
|
||||
return func(ctx context.Context, cfg Config) (backend, error) {
|
||||
fl.mu.Lock()
|
||||
name, ok := fl.models[cfg.ModelPath]
|
||||
fl.launch++
|
||||
if !ok {
|
||||
fl.mu.Unlock()
|
||||
return nil, errors.New("no such model file: " + cfg.ModelPath)
|
||||
}
|
||||
fl.live++
|
||||
if fl.live > fl.maxLive {
|
||||
fl.maxLive = fl.live
|
||||
}
|
||||
fl.mu.Unlock()
|
||||
f := newFakeModel(t, name)
|
||||
return &fleetBackend{fleet: fl, model: f}, nil
|
||||
}
|
||||
}
|
||||
|
||||
type fleetBackend struct {
|
||||
fleet *fakeFleet
|
||||
model *fakeModel
|
||||
once sync.Once
|
||||
}
|
||||
|
||||
func (b *fleetBackend) BaseURL() string { return b.model.BaseURL() }
|
||||
func (b *fleetBackend) Close() error {
|
||||
b.once.Do(func() {
|
||||
b.fleet.mu.Lock()
|
||||
b.fleet.live--
|
||||
b.fleet.mu.Unlock()
|
||||
})
|
||||
return b.model.Close()
|
||||
}
|
||||
|
||||
// newSwapPhraser builds an LLMPhraser with an injected launcher, so the swap
|
||||
// path is exercised without a gguf or a GPU.
|
||||
func newSwapPhraser(t *testing.T, fl *fakeFleet, modelPath string) *LLMPhraser {
|
||||
t.Helper()
|
||||
cfg := DefaultConfig(modelPath)
|
||||
cfg.Timeout = 5 * time.Second
|
||||
p := &LLMPhraser{
|
||||
cfg: cfg,
|
||||
client: &http.Client{Timeout: cfg.Timeout},
|
||||
spawnCtx: context.Background(),
|
||||
cancel: func() {},
|
||||
launch: fl.launcher(t),
|
||||
probe: defaultProbe,
|
||||
live: liveModel{ModelPath: modelPath, NGpuLayers: cfg.NGpuLayers, NCtx: cfg.NCtx},
|
||||
}
|
||||
be, err := p.launch(p.spawnCtx, cfg)
|
||||
if err != nil {
|
||||
t.Fatalf("initial launch: %v", err)
|
||||
}
|
||||
p.be = be
|
||||
t.Cleanup(func() { p.Close() })
|
||||
return p
|
||||
}
|
||||
|
||||
func TestSwap_LoadsNewModelAndRepointsHolders(t *testing.T) {
|
||||
fl := &fakeFleet{models: map[string]string{"/m/old.gguf": "old", "/m/new.gguf": "new"}}
|
||||
p := newSwapPhraser(t, fl, "/m/old.gguf")
|
||||
|
||||
// A holder of the base URL (the LLM router's client, in the daemon).
|
||||
var seen []string
|
||||
p.OnSwap(func(base string) { seen = append(seen, base) })
|
||||
|
||||
before, err := p.PhraseChat(context.Background(), "привет", nil)
|
||||
if err != nil || before != "old" {
|
||||
t.Fatalf("before swap: %q, %v; want the old model to answer", before, err)
|
||||
}
|
||||
|
||||
res, err := p.Swap(context.Background(), SwapSpec{ModelPath: "/m/new.gguf"})
|
||||
if err != nil {
|
||||
t.Fatalf("Swap: %v", err)
|
||||
}
|
||||
if res.Model != "new" {
|
||||
t.Errorf("res.Model = %q; want the identity the NEW server reported (%q)", res.Model, "new")
|
||||
}
|
||||
if res.RolledBack {
|
||||
t.Errorf("res.RolledBack = true on a successful swap")
|
||||
}
|
||||
after, err := p.PhraseChat(context.Background(), "привет", nil)
|
||||
if err != nil || after != "new" {
|
||||
t.Fatalf("after swap: %q, %v; want the new model to answer", after, err)
|
||||
}
|
||||
if path, _, _ := p.LiveModel(); path != "/m/new.gguf" {
|
||||
t.Errorf("LiveModel = %q; want /m/new.gguf", path)
|
||||
}
|
||||
if len(seen) != 1 || seen[0] != p.BaseURL() {
|
||||
t.Errorf("observers saw %v; want exactly one call with the new base %q", seen, p.BaseURL())
|
||||
}
|
||||
if fl.maxLive > 1 {
|
||||
t.Errorf("%d servers were alive at once; the iGPU only fits one model", fl.maxLive)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSwap_FailedLoadRollsBackToTheWorkingModel(t *testing.T) {
|
||||
fl := &fakeFleet{models: map[string]string{"/m/old.gguf": "old"}}
|
||||
p := newSwapPhraser(t, fl, "/m/old.gguf")
|
||||
|
||||
res, err := p.Swap(context.Background(), SwapSpec{ModelPath: "/m/broken.gguf"})
|
||||
if err == nil {
|
||||
t.Fatal("Swap to a model that will not load returned nil error")
|
||||
}
|
||||
if !res.RolledBack {
|
||||
t.Errorf("res.RolledBack = false; a failed swap must say it rolled back")
|
||||
}
|
||||
if res.Model != "old" {
|
||||
t.Errorf("res.Model = %q; want the old model back", res.Model)
|
||||
}
|
||||
// The point of the rollback: turns keep working.
|
||||
got, err := p.PhraseChat(context.Background(), "привет", nil)
|
||||
if err != nil || got != "old" {
|
||||
t.Fatalf("after rollback: %q, %v; want the old model serving again", got, err)
|
||||
}
|
||||
if path, _, _ := p.LiveModel(); path != "/m/old.gguf" {
|
||||
t.Errorf("LiveModel = %q; want the old model", path)
|
||||
}
|
||||
if fl.maxLive > 1 {
|
||||
t.Errorf("%d servers alive at once during a rollback", fl.maxLive)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSwap_ProbeFailureIsTreatedAsAFailedLoad(t *testing.T) {
|
||||
// A server that starts but will not say what it loaded must never be
|
||||
// published — we would be serving turns from a backend we cannot confirm.
|
||||
fl := &fakeFleet{models: map[string]string{"/m/old.gguf": "old", "/m/mute.gguf": "mute"}}
|
||||
p := newSwapPhraser(t, fl, "/m/old.gguf")
|
||||
// Fail the probe once — for the newly launched server — and let the
|
||||
// rollback's probe through.
|
||||
calls := 0
|
||||
p.probe = func(ctx context.Context, base string) (string, error) {
|
||||
calls++
|
||||
if calls == 1 {
|
||||
return "", errors.New("no answer from the new server")
|
||||
}
|
||||
return defaultProbe(ctx, base)
|
||||
}
|
||||
|
||||
_, err := p.Swap(context.Background(), SwapSpec{ModelPath: "/m/mute.gguf"})
|
||||
if err == nil {
|
||||
t.Fatal("Swap published a server that failed its probe")
|
||||
}
|
||||
if path, _, _ := p.LiveModel(); path != "/m/old.gguf" {
|
||||
t.Errorf("LiveModel = %q; want the old model after a failed probe", path)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSwap_RollbackFailureLeavesNoBackendAndDegrades(t *testing.T) {
|
||||
fl := &fakeFleet{models: map[string]string{"/m/old.gguf": "old"}}
|
||||
p := newSwapPhraser(t, fl, "/m/old.gguf")
|
||||
// Make the rollback fail too: the old file "disappears" mid-swap.
|
||||
fl.mu.Lock()
|
||||
delete(fl.models, "/m/old.gguf")
|
||||
fl.mu.Unlock()
|
||||
|
||||
_, err := p.Swap(context.Background(), SwapSpec{ModelPath: "/m/broken.gguf"})
|
||||
if err == nil {
|
||||
t.Fatal("Swap returned nil when both the load and the rollback failed")
|
||||
}
|
||||
// Nothing is loaded, and the request path says so rather than panicking.
|
||||
if _, _, aerr := p.acquire(); !errors.Is(aerr, ErrNoBackend) {
|
||||
t.Errorf("acquire error = %v; want ErrNoBackend", aerr)
|
||||
}
|
||||
// Phrasing degrades to its fallback instead of failing the turn.
|
||||
got, err := p.PhraseChat(context.Background(), "привет", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("PhraseChat after a total failure returned an error: %v", err)
|
||||
}
|
||||
if got == "" {
|
||||
t.Error("PhraseChat returned empty; the fallback must still say something")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSwap_WaitsForInFlightTurnAndRefusesNewOnes(t *testing.T) {
|
||||
fl := &fakeFleet{models: map[string]string{"/m/old.gguf": "old", "/m/new.gguf": "new"}}
|
||||
p := newSwapPhraser(t, fl, "/m/old.gguf")
|
||||
|
||||
// Hold one turn open by taking a slot directly — the same slot every
|
||||
// request path takes.
|
||||
base, release, err := p.acquire()
|
||||
if err != nil {
|
||||
t.Fatalf("acquire: %v", err)
|
||||
}
|
||||
if base == "" {
|
||||
t.Fatal("acquire returned an empty base URL")
|
||||
}
|
||||
|
||||
swapped := make(chan error, 1)
|
||||
go func() { _, e := p.Swap(context.Background(), SwapSpec{ModelPath: "/m/new.gguf"}); swapped <- e }()
|
||||
|
||||
// While the swap waits to drain, a NEW turn is refused immediately rather
|
||||
// than blocked for the length of a model load.
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for {
|
||||
_, rel, aerr := p.acquire()
|
||||
if rel != nil {
|
||||
rel()
|
||||
}
|
||||
if errors.Is(aerr, ErrSwapping) {
|
||||
break
|
||||
}
|
||||
if time.Now().After(deadline) {
|
||||
t.Fatalf("new turns were never refused during a swap (last error: %v)", aerr)
|
||||
}
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
|
||||
// The swap cannot have completed while our turn was still in flight.
|
||||
select {
|
||||
case e := <-swapped:
|
||||
t.Fatalf("Swap finished before the in-flight turn released: %v", e)
|
||||
case <-time.After(50 * time.Millisecond):
|
||||
}
|
||||
|
||||
release()
|
||||
if e := <-swapped; e != nil {
|
||||
t.Fatalf("Swap after drain: %v", e)
|
||||
}
|
||||
got, err := p.PhraseChat(context.Background(), "привет", nil)
|
||||
if err != nil || got != "new" {
|
||||
t.Fatalf("after swap: %q, %v; want the new model", got, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSwap_RefusedWhenWeDoNotOwnTheServer(t *testing.T) {
|
||||
// NewLLMPhraserAt points at a shared server the eval harness owns. Swapping
|
||||
// there would kill a server another process depends on.
|
||||
p := NewLLMPhraserAt("http://127.0.0.1:1/", DefaultConfig("/m/old.gguf"))
|
||||
if _, err := p.Swap(context.Background(), SwapSpec{ModelPath: "/m/new.gguf"}); !errors.Is(err, ErrSwapNotOwned) {
|
||||
t.Fatalf("Swap on a borrowed server = %v; want ErrSwapNotOwned", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSwap_SameModelIsANoOp(t *testing.T) {
|
||||
fl := &fakeFleet{models: map[string]string{"/m/old.gguf": "old"}}
|
||||
p := newSwapPhraser(t, fl, "/m/old.gguf")
|
||||
launchesBefore := fl.launch
|
||||
|
||||
res, err := p.Swap(context.Background(), SwapSpec{ModelPath: "/m/old.gguf"})
|
||||
if err != nil {
|
||||
t.Fatalf("Swap to the live model: %v", err)
|
||||
}
|
||||
if res.Model != "old" {
|
||||
t.Errorf("res.Model = %q; want old", res.Model)
|
||||
}
|
||||
if fl.launch != launchesBefore {
|
||||
t.Errorf("%d extra launches; swapping to the live model must not reload weights", fl.launch-launchesBefore)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSwap_EmptyModelPathRefused(t *testing.T) {
|
||||
fl := &fakeFleet{models: map[string]string{"/m/old.gguf": "old"}}
|
||||
p := newSwapPhraser(t, fl, "/m/old.gguf")
|
||||
if _, err := p.Swap(context.Background(), SwapSpec{}); err == nil {
|
||||
t.Fatal("Swap with no model path returned nil error")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,39 @@
|
||||
package router
|
||||
|
||||
import (
|
||||
"regexp"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// Finding a URL in an utterance (Vikunja #259).
|
||||
//
|
||||
// This is deliberately strict: a scheme is required. "посмотри на example.org"
|
||||
// is not treated as a fetch request, because a bare dotted word is also how
|
||||
// people say file names, versions and Russian abbreviations, and the cost of a
|
||||
// false positive here is an outbound request nobody asked for.
|
||||
//
|
||||
// Note where this runs: an utterance from STT. Whisper will mangle a spoken URL,
|
||||
// which is fine — the URL that survives is one he pasted into the web chat, and
|
||||
// a mangled one simply fails to match.
|
||||
var urlRE = regexp.MustCompile(`(?i)\bhttps?://[^\s<>"']+`)
|
||||
|
||||
// FirstURL returns the first http(s) URL in text.
|
||||
//
|
||||
// Trailing punctuation is trimmed: he ends sentences, and "…/page." is not a
|
||||
// path component. A closing bracket is only trimmed when it has no opener,
|
||||
// because a wikipedia URL legitimately ends in one.
|
||||
func FirstURL(text string) (string, bool) {
|
||||
m := urlRE.FindString(text)
|
||||
if m == "" {
|
||||
return "", false
|
||||
}
|
||||
m = strings.TrimRight(m, ".,;:!?…")
|
||||
if strings.HasSuffix(m, ")") && strings.Count(m, "(") == 0 {
|
||||
m = strings.TrimSuffix(m, ")")
|
||||
}
|
||||
// A scheme with nothing after it is not a URL.
|
||||
if rest := strings.SplitN(m, "//", 2); len(rest) < 2 || rest[1] == "" {
|
||||
return "", false
|
||||
}
|
||||
return m, true
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
package router
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestFirstURL(t *testing.T) {
|
||||
cases := []struct {
|
||||
text string
|
||||
want string
|
||||
}{
|
||||
{"посмотри https://example.org/page — что там?", "https://example.org/page"},
|
||||
{"почитай http://example.org/a/b?x=1 и скажи", "http://example.org/a/b?x=1"},
|
||||
{"вот ссылка: https://example.org/page.", "https://example.org/page"},
|
||||
{"https://ru.wikipedia.org/wiki/Небо_(значения)", "https://ru.wikipedia.org/wiki/Небо_(значения)"},
|
||||
// No scheme ⇒ no fetch. A bare dotted word is not an instruction to
|
||||
// reach out to the network.
|
||||
{"посмотри на example.org", ""},
|
||||
{"открой файл config.json", ""},
|
||||
{"что нового?", ""},
|
||||
{"https://", ""},
|
||||
{"", ""},
|
||||
}
|
||||
for _, c := range cases {
|
||||
got, ok := FirstURL(c.text)
|
||||
if c.want == "" {
|
||||
if ok {
|
||||
t.Errorf("FirstURL(%q) = %q, want no match", c.text, got)
|
||||
}
|
||||
continue
|
||||
}
|
||||
if !ok || got != c.want {
|
||||
t.Errorf("FirstURL(%q) = %q, %v; want %q", c.text, got, ok, c.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,208 @@
|
||||
package update
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Result — the full account of one Apply. Every field is filled in on the
|
||||
// failure paths too, because "what state is my box in" is the only question that
|
||||
// matters after a failed update.
|
||||
type Result struct {
|
||||
Verified bool
|
||||
SnapshotID string // the rollback target; named even when the rollback failed
|
||||
Installed []string
|
||||
Restarted bool
|
||||
Healthy bool
|
||||
RolledBack bool
|
||||
// RollbackHealthy — whether she answered again after the restore. False with
|
||||
// RolledBack true is the manual-recovery case.
|
||||
RollbackHealthy bool
|
||||
Steps []Step
|
||||
Took time.Duration
|
||||
}
|
||||
|
||||
// Apply is the whole update, in the only order that is safe.
|
||||
//
|
||||
// It is called by a human running cmd/mavupdate on the box. Nothing else calls
|
||||
// it: no timer, no IPC method, no web route, no act. See the package comment.
|
||||
func (u *Updater) Apply(ctx context.Context) (Result, error) {
|
||||
start := u.now()
|
||||
res := Result{}
|
||||
defer func() { res.Took = u.now().Sub(start) }()
|
||||
|
||||
// 0. She has to be answering before we start. Otherwise a failed update and
|
||||
// a box that was already broken look identical afterwards, and the rollback
|
||||
// has no baseline to prove itself against.
|
||||
u.log("preflight: checking the running daemon")
|
||||
if err := u.health(ctx, u.cfg.HealthSocket); err != nil {
|
||||
return res, fmt.Errorf("%w: %v", ErrUnhealthyBefore, err)
|
||||
}
|
||||
|
||||
// 1. Snapshot what is deployed now, BEFORE the build.
|
||||
//
|
||||
// The order matters and it is not the obvious one. `make build` writes its
|
||||
// binaries into the working tree, and on the docker deployment the working
|
||||
// tree IS the install dir — so snapshotting after the build would snapshot
|
||||
// the new artifacts and leave nothing to roll back to. The snapshot is the
|
||||
// only thing standing between a bad build and a box that needs a screwdriver,
|
||||
// so it is taken first, while the deployed bytes are still the old ones.
|
||||
names := append(append([]string{}, u.cfg.Binaries...), u.cfg.ConfigFiles...)
|
||||
snap, err := u.store.Save(u.cfg.InstallDir, names, u.gitHead(ctx), "pre-update")
|
||||
if err != nil {
|
||||
return res, err
|
||||
}
|
||||
res.SnapshotID = snap.ID
|
||||
u.log("snapshot: %s (%d files) in %s", snap.ID, len(snap.Files), snap.Dir())
|
||||
|
||||
// 2. Build and test before anything is deployed. A broken tree costs time
|
||||
// and nothing else — but `make build` has already overwritten the binaries in
|
||||
// the tree, so restore them: otherwise a later restart by hand would deploy
|
||||
// code that failed its own tests. Nothing has been restarted, so this is a
|
||||
// file restore with no restart and no health check.
|
||||
steps, err := u.Verify(ctx)
|
||||
res.Steps = append(res.Steps, steps...)
|
||||
if err != nil {
|
||||
if rerr := snap.Restore(u.cfg.InstallDir); rerr != nil {
|
||||
u.log("verify failed and the artifacts could not be put back: %v — the previous ones are in %s", rerr, snap.Dir())
|
||||
} else {
|
||||
res.RolledBack = true
|
||||
u.log("verify failed; the previously deployed artifacts are back in place, she was never restarted")
|
||||
}
|
||||
return res, err
|
||||
}
|
||||
res.Verified = true
|
||||
|
||||
// 3. Install. Per-file temp+rename, so an interruption leaves whole files.
|
||||
// Config is snapshotted but never overwritten — an update does not get to
|
||||
// replace the operator's config.
|
||||
installed, err := u.install()
|
||||
res.Installed = installed
|
||||
if err != nil {
|
||||
// Files may be half-swapped across the set, so restore before returning
|
||||
// even though nothing has been restarted yet.
|
||||
u.log("install failed: %v — restoring", err)
|
||||
return u.rollback(ctx, snap, res, err)
|
||||
}
|
||||
u.log("install: %d artifact(s) into %s", len(installed), u.cfg.InstallDir)
|
||||
|
||||
// 4. Restart, then 5. prove she answers.
|
||||
if err := u.restart(ctx, &res); err != nil {
|
||||
return u.rollback(ctx, snap, res, err)
|
||||
}
|
||||
u.log("restart: ok, waiting for her to answer (up to %s)", u.cfg.healthTimeout())
|
||||
if err := u.waitHealthy(ctx, u.cfg.healthTimeout()); err != nil {
|
||||
return u.rollback(ctx, snap, res, err)
|
||||
}
|
||||
res.Healthy = true
|
||||
u.log("health: she answers on %s — update committed", u.cfg.HealthSocket)
|
||||
|
||||
if err := u.store.Prune(u.cfg.KeepSnapshots); err != nil {
|
||||
u.log("prune: %v (harmless)", err)
|
||||
}
|
||||
return res, nil
|
||||
}
|
||||
|
||||
// Rollback restores a snapshot by id (empty = the newest) and restarts. Exposed
|
||||
// separately so the operator can undo an update that verified, restarted and
|
||||
// answered a Presence call but is wrong in a way no health check can see.
|
||||
func (u *Updater) Rollback(ctx context.Context, id string) (Result, error) {
|
||||
var snap Snapshot
|
||||
var err error
|
||||
if id == "" {
|
||||
snaps, lerr := u.store.List()
|
||||
if lerr != nil {
|
||||
return Result{}, lerr
|
||||
}
|
||||
if len(snaps) == 0 {
|
||||
return Result{}, errors.New("update: no snapshots to roll back to")
|
||||
}
|
||||
snap = snaps[0]
|
||||
} else if snap, err = u.store.Load(id); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
res := Result{SnapshotID: snap.ID}
|
||||
return u.rollback(ctx, snap, res, errors.New("operator asked for a rollback"))
|
||||
}
|
||||
|
||||
// rollback restores the snapshot and restarts, then reports whether that worked.
|
||||
// It depends on nothing that the update changed: file copies out of the snapshot
|
||||
// dir and the same restart command. No build, no migration, no cooperation from
|
||||
// the code being replaced.
|
||||
func (u *Updater) rollback(ctx context.Context, snap Snapshot, res Result, cause error) (Result, error) {
|
||||
// A rollback interrupted halfway is the one outcome worse than the failure
|
||||
// that triggered it, so it does not inherit the caller's cancellation: a
|
||||
// Ctrl-C during the health wait must not abandon the restore mid-restart.
|
||||
ctx = context.WithoutCancel(ctx)
|
||||
res.RolledBack = true
|
||||
u.log("rollback: restoring snapshot %s over %s", snap.ID, u.cfg.InstallDir)
|
||||
if err := snap.Restore(u.cfg.InstallDir); err != nil {
|
||||
u.log("rollback: RESTORE FAILED: %v", err)
|
||||
return res, fmt.Errorf("%w: %v (after %v); the previous artifacts are in %s — copy them back by hand", ErrRollbackFailed, err, cause, snap.Dir())
|
||||
}
|
||||
// A restore with no restart leaves the failed process running, so a failed
|
||||
// restart here is still the manual-recovery case.
|
||||
if err := u.restart(ctx, &res); err != nil {
|
||||
u.log("rollback: RESTART FAILED: %v", err)
|
||||
return res, fmt.Errorf("%w: restored %s but the restart failed: %v (after %v)", ErrRollbackFailed, snap.ID, err, cause)
|
||||
}
|
||||
if err := u.waitHealthy(ctx, u.cfg.healthTimeout()); err != nil {
|
||||
u.log("rollback: she still does not answer: %v", err)
|
||||
return res, fmt.Errorf("%w: restored %s and restarted but she does not answer: %v (after %v)", ErrRollbackFailed, snap.ID, err, cause)
|
||||
}
|
||||
res.RollbackHealthy = true
|
||||
u.log("rollback: she answers again on the previous build (%s)", snap.ID)
|
||||
return res, fmt.Errorf("%w to %s: %v", ErrRolledBack, snap.ID, cause)
|
||||
}
|
||||
|
||||
// install copies the freshly built binaries from SourceDir into InstallDir.
|
||||
//
|
||||
// When the two are the same directory — the docker deployment builds the image
|
||||
// from the working tree — this is a no-op by design rather than by accident: the
|
||||
// artifacts are already where they belong and the restart command rebuilds the
|
||||
// image from them.
|
||||
func (u *Updater) install() ([]string, error) {
|
||||
if filepath.Clean(u.cfg.SourceDir) == filepath.Clean(u.cfg.InstallDir) {
|
||||
return u.cfg.Binaries, nil
|
||||
}
|
||||
var done []string
|
||||
for _, name := range u.cfg.Binaries {
|
||||
src := filepath.Join(u.cfg.SourceDir, name)
|
||||
fi, err := os.Stat(src)
|
||||
if err != nil {
|
||||
return done, fmt.Errorf("update: install %s: %w (did `make build` produce it?)", name, err)
|
||||
}
|
||||
if _, err := copyFile(src, filepath.Join(u.cfg.InstallDir, name), fi.Mode().Perm()); err != nil {
|
||||
return done, fmt.Errorf("update: install %s: %w", name, err)
|
||||
}
|
||||
done = append(done, name)
|
||||
}
|
||||
return done, nil
|
||||
}
|
||||
|
||||
func (u *Updater) restart(ctx context.Context, res *Result) error {
|
||||
u.log("restart: %v", u.cfg.RestartCmd)
|
||||
out, err := u.run(ctx, u.cfg.SourceDir, u.cfg.RestartCmd)
|
||||
if err != nil {
|
||||
res.Steps = append(res.Steps, Step{Name: "restart", Argv: u.cfg.RestartCmd, Err: err, Output: tail(out, 4000)})
|
||||
return fmt.Errorf("update: restart %v: %w", u.cfg.RestartCmd, err)
|
||||
}
|
||||
res.Restarted = true
|
||||
res.Steps = append(res.Steps, Step{Name: "restart", Argv: u.cfg.RestartCmd})
|
||||
return nil
|
||||
}
|
||||
|
||||
// gitHead records which commit produced a snapshot, for the operator's benefit.
|
||||
// Best-effort: a tree without git is not a reason to refuse to snapshot.
|
||||
func (u *Updater) gitHead(ctx context.Context) string {
|
||||
out, err := u.run(ctx, u.cfg.SourceDir, []string{"git", "rev-parse", "HEAD"})
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
return strings.TrimSpace(out)
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
package update
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
|
||||
// The health check is the whole basis for rolling back, so it has to mean
|
||||
// something. "The process is running" does not: mavend can be up with a dead
|
||||
// store, a socket it never bound, or a config it failed to parse. What is
|
||||
// checked instead is that she answers a real read over the real IPC socket —
|
||||
// which exercises the socket, the dispatch table and the store in one call.
|
||||
//
|
||||
// Presence is the method used because it is read-only (safe to retry), needs no
|
||||
// arguments, and touches the store. It cannot write anything, so a health check
|
||||
// never leaves a trace in her memory.
|
||||
|
||||
// DialHealth connects to the mavend socket and performs one read.
|
||||
func DialHealth(ctx context.Context, socket string) error {
|
||||
c, err := ipc.Dial(socket)
|
||||
if err != nil {
|
||||
return fmt.Errorf("update: health dial: %w", err)
|
||||
}
|
||||
defer c.Close()
|
||||
if _, err := c.Presence(ctx); err != nil {
|
||||
return fmt.Errorf("update: health read: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// waitHealthy retries the health check until it passes or the timeout elapses.
|
||||
// A restart is not instantaneous — she loads a 1.7B on boot — so the first few
|
||||
// failures are expected and are not a reason to roll back.
|
||||
func (u *Updater) waitHealthy(ctx context.Context, timeout time.Duration) error {
|
||||
deadline := u.now().Add(timeout)
|
||||
delay := 500 * time.Millisecond
|
||||
var last error
|
||||
for {
|
||||
attemptCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
|
||||
err := u.health(attemptCtx, u.cfg.HealthSocket)
|
||||
cancel()
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
last = err
|
||||
if u.now().After(deadline) {
|
||||
return fmt.Errorf("update: not healthy after %s: %w", timeout, last)
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-time.After(delay):
|
||||
}
|
||||
if delay < 5*time.Second {
|
||||
delay *= 2
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,267 @@
|
||||
package update
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"time"
|
||||
)
|
||||
|
||||
// A snapshot is a byte-for-byte copy of the deployed artifacts plus a manifest
|
||||
// of their sha256 sums, taken before an install.
|
||||
//
|
||||
// It is copies, not hardlinks and not a git stash, for one reason: the restore
|
||||
// path must work when everything else is broken. A hardlink into the install dir
|
||||
// would be clobbered by the very install it exists to undo, and a git-based
|
||||
// undo needs a toolchain, a clean tree, and a rebuild — three things a failed
|
||||
// update is likely to have taken away. Copying two dozen megabytes of Go
|
||||
// binaries costs a second and needs nothing but the filesystem.
|
||||
//
|
||||
// The sums are what make a restore verifiable rather than hopeful: Restore
|
||||
// re-hashes every file it writes, so "the old bytes are back" is checked, not
|
||||
// assumed.
|
||||
|
||||
// FileRec — one file in a snapshot.
|
||||
type FileRec struct {
|
||||
Name string `json:"name"` // relative name inside the install dir
|
||||
SHA256 string `json:"sha256"` // of the snapshotted bytes
|
||||
Mode os.FileMode `json:"mode"`
|
||||
Size int64 `json:"size"`
|
||||
}
|
||||
|
||||
// Snapshot — the manifest. Written last, so a directory without a readable
|
||||
// manifest.json is an aborted snapshot and is never offered as a rollback target.
|
||||
type Snapshot struct {
|
||||
ID string `json:"id"` // sortable timestamp, also the directory name
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
Commit string `json:"commit,omitempty"` // git HEAD of the tree that produced it, when known
|
||||
Note string `json:"note,omitempty"`
|
||||
Files []FileRec `json:"files"`
|
||||
|
||||
dir string // absolute path, filled in by List/Load
|
||||
}
|
||||
|
||||
// Dir — where this snapshot's file copies live.
|
||||
func (s Snapshot) Dir() string { return s.dir }
|
||||
|
||||
const manifestName = "manifest.json"
|
||||
|
||||
// Store is a directory of snapshots.
|
||||
type Store struct {
|
||||
Dir string
|
||||
now func() time.Time
|
||||
}
|
||||
|
||||
func (st *Store) clock() time.Time {
|
||||
if st.now != nil {
|
||||
return st.now()
|
||||
}
|
||||
return time.Now()
|
||||
}
|
||||
|
||||
// Save copies names (relative to srcDir) into a new snapshot and writes the
|
||||
// manifest. A name that does not exist is skipped rather than fatal: the first
|
||||
// ever run happens on a box where some artifact may legitimately be missing, and
|
||||
// refusing to snapshot then would mean refusing to update.
|
||||
func (st *Store) Save(srcDir string, names []string, commit, note string) (Snapshot, error) {
|
||||
ts := st.clock().UTC()
|
||||
snap := Snapshot{
|
||||
ID: ts.Format("20060102-150405"),
|
||||
CreatedAt: ts,
|
||||
Commit: commit,
|
||||
Note: note,
|
||||
}
|
||||
snap.dir = filepath.Join(st.Dir, snap.ID)
|
||||
if err := os.MkdirAll(snap.dir, 0o700); err != nil {
|
||||
return Snapshot{}, fmt.Errorf("update: snapshot dir: %w", err)
|
||||
}
|
||||
for _, name := range names {
|
||||
src := filepath.Join(srcDir, name)
|
||||
fi, err := os.Stat(src)
|
||||
if err != nil {
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
continue
|
||||
}
|
||||
return Snapshot{}, fmt.Errorf("update: snapshot %s: %w", name, err)
|
||||
}
|
||||
if fi.IsDir() {
|
||||
return Snapshot{}, fmt.Errorf("update: snapshot %s: is a directory (only files are deployable artifacts)", name)
|
||||
}
|
||||
dst := filepath.Join(snap.dir, name)
|
||||
if err := os.MkdirAll(filepath.Dir(dst), 0o700); err != nil {
|
||||
return Snapshot{}, err
|
||||
}
|
||||
sum, err := copyFile(src, dst, fi.Mode().Perm())
|
||||
if err != nil {
|
||||
return Snapshot{}, fmt.Errorf("update: snapshot %s: %w", name, err)
|
||||
}
|
||||
snap.Files = append(snap.Files, FileRec{Name: name, SHA256: sum, Mode: fi.Mode().Perm(), Size: fi.Size()})
|
||||
}
|
||||
if len(snap.Files) == 0 {
|
||||
os.RemoveAll(snap.dir)
|
||||
return Snapshot{}, fmt.Errorf("update: snapshot of %s is empty — none of the listed artifacts exist", srcDir)
|
||||
}
|
||||
// Manifest last: its presence is what makes the snapshot usable.
|
||||
blob, err := json.MarshalIndent(snap, "", " ")
|
||||
if err != nil {
|
||||
return Snapshot{}, err
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(snap.dir, manifestName), blob, 0o600); err != nil {
|
||||
return Snapshot{}, fmt.Errorf("update: snapshot manifest: %w", err)
|
||||
}
|
||||
return snap, nil
|
||||
}
|
||||
|
||||
// List returns the complete snapshots, newest first.
|
||||
func (st *Store) List() ([]Snapshot, error) {
|
||||
ents, err := os.ReadDir(st.Dir)
|
||||
if err != nil {
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
return nil, nil
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
var out []Snapshot
|
||||
for _, e := range ents {
|
||||
if !e.IsDir() {
|
||||
continue
|
||||
}
|
||||
s, err := st.Load(e.Name())
|
||||
if err != nil {
|
||||
continue // aborted or hand-mangled: not a rollback target
|
||||
}
|
||||
out = append(out, s)
|
||||
}
|
||||
sort.Slice(out, func(i, j int) bool { return out[i].ID > out[j].ID })
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// Load reads one snapshot's manifest.
|
||||
func (st *Store) Load(id string) (Snapshot, error) {
|
||||
dir := filepath.Join(st.Dir, id)
|
||||
blob, err := os.ReadFile(filepath.Join(dir, manifestName))
|
||||
if err != nil {
|
||||
return Snapshot{}, err
|
||||
}
|
||||
var s Snapshot
|
||||
if err := json.Unmarshal(blob, &s); err != nil {
|
||||
return Snapshot{}, fmt.Errorf("update: manifest %s: %w", id, err)
|
||||
}
|
||||
s.dir = dir
|
||||
return s, nil
|
||||
}
|
||||
|
||||
// Restore copies a snapshot's files back over dstDir and verifies every write
|
||||
// against the manifest sum. Only the named files are touched; anything else in
|
||||
// dstDir is left alone.
|
||||
//
|
||||
// This is the function the whole package exists to be able to run. It uses the
|
||||
// filesystem and nothing else — no toolchain, no build, no cooperation from the
|
||||
// code being replaced.
|
||||
func (s Snapshot) Restore(dstDir string) error {
|
||||
if s.dir == "" {
|
||||
return errors.New("update: snapshot has no directory (load it through the store)")
|
||||
}
|
||||
for _, f := range s.Files {
|
||||
src := filepath.Join(s.dir, f.Name)
|
||||
sum, err := hashFile(src)
|
||||
if err != nil {
|
||||
return fmt.Errorf("update: restore %s: %w", f.Name, err)
|
||||
}
|
||||
if sum != f.SHA256 {
|
||||
return fmt.Errorf("update: restore %s: snapshot is corrupt (sha256 %s, manifest says %s)", f.Name, sum, f.SHA256)
|
||||
}
|
||||
dst := filepath.Join(dstDir, f.Name)
|
||||
if err := os.MkdirAll(filepath.Dir(dst), 0o755); err != nil {
|
||||
return err
|
||||
}
|
||||
got, err := copyFile(src, dst, f.Mode)
|
||||
if err != nil {
|
||||
return fmt.Errorf("update: restore %s: %w", f.Name, err)
|
||||
}
|
||||
if got != f.SHA256 {
|
||||
return fmt.Errorf("update: restore %s: wrote the wrong bytes (sha256 %s)", f.Name, got)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Prune keeps the newest keep snapshots and removes the rest. The newest is
|
||||
// never pruned regardless of keep — it is the rollback target.
|
||||
func (st *Store) Prune(keep int) error {
|
||||
if keep < 1 {
|
||||
keep = 1
|
||||
}
|
||||
snaps, err := st.List()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, s := range snaps[min(keep, len(snaps)):] {
|
||||
if err := os.RemoveAll(s.dir); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// copyFile writes src to dst atomically (temp + rename, so a reader never sees a
|
||||
// half file and an interrupted copy leaves the old one intact) and returns the
|
||||
// sha256 of what was written.
|
||||
func copyFile(src, dst string, mode os.FileMode) (string, error) {
|
||||
in, err := os.Open(src)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer in.Close()
|
||||
if mode == 0 {
|
||||
mode = 0o644
|
||||
}
|
||||
tmp, err := os.CreateTemp(filepath.Dir(dst), ".update-*")
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
tmpName := tmp.Name()
|
||||
defer os.Remove(tmpName) // no-op once the rename succeeds
|
||||
h := sha256.New()
|
||||
if _, err := io.Copy(io.MultiWriter(tmp, h), in); err != nil {
|
||||
tmp.Close()
|
||||
return "", err
|
||||
}
|
||||
// fsync before the rename: a binary that is renamed into place but whose
|
||||
// bytes are still in the page cache is exactly the file a power cut turns
|
||||
// into an unbootable daemon.
|
||||
if err := tmp.Sync(); err != nil {
|
||||
tmp.Close()
|
||||
return "", err
|
||||
}
|
||||
if err := tmp.Chmod(mode); err != nil {
|
||||
tmp.Close()
|
||||
return "", err
|
||||
}
|
||||
if err := tmp.Close(); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := os.Rename(tmpName, dst); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return hex.EncodeToString(h.Sum(nil)), nil
|
||||
}
|
||||
|
||||
func hashFile(p string) (string, error) {
|
||||
f, err := os.Open(p)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer f.Close()
|
||||
h := sha256.New()
|
||||
if _, err := io.Copy(h, f); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return hex.EncodeToString(h.Sum(nil)), nil
|
||||
}
|
||||
@@ -0,0 +1,270 @@
|
||||
// Package update applies a new build of Maven to the box she runs on, with a
|
||||
// verified-before-committed install and an automatic rollback (Vikunja #249).
|
||||
//
|
||||
// # What this package refuses to be
|
||||
//
|
||||
// This is the highest-risk capability in the backlog — code that changes the
|
||||
// running system — so the refusals are as much of the design as the features,
|
||||
// and they are enforced here rather than described in a doc:
|
||||
//
|
||||
// - It is never automatic and never on a timer. There is no checker, no
|
||||
// channel, no "check for updates" call and nothing that fires from the tick
|
||||
// loop. Apply runs exactly when a human runs cmd/mavupdate on the box.
|
||||
// - The daemon cannot update itself. mavend does not import this package and
|
||||
// there is no IPC method and no web route that reaches it, so no act, no
|
||||
// intent, no tool and no LLM output can start an update. The trigger needs
|
||||
// shell access to the host, which is a strictly higher bar than the step-up
|
||||
// passkey gate that guards /tools — an update is not a thing to expose to
|
||||
// anything reachable over the network.
|
||||
// - It does not fetch code. Nothing here talks to a release server, a
|
||||
// registry, or GitHub. The new version is whatever is in the working tree
|
||||
// the operator points it at, which he pulled himself. Downloading and
|
||||
// running code on the strength of a checksum in the same download is not a
|
||||
// property we can verify on one box.
|
||||
// - It does not supervise its own death. The plan asked for an in-process
|
||||
// crash-loop detector; a process cannot reliably notice that it keeps
|
||||
// dying, and one that thinks it can is worse than nothing. Restart-on-crash
|
||||
// belongs to whatever starts mavend (compose `restart: unless-stopped`,
|
||||
// systemd `Restart=`). What this package guarantees instead is narrower and
|
||||
// real: within one Apply, the new build is proven to answer before the old
|
||||
// one is considered replaced, and if it does not answer the old bytes go
|
||||
// back and are proven to answer again.
|
||||
//
|
||||
// # The order of operations, and why
|
||||
//
|
||||
// Apply is: health-check the CURRENT daemon → build → test → snapshot → install
|
||||
// → restart → health-check → rollback on any failure.
|
||||
//
|
||||
// The first health check is not ceremony. If she is already not answering, a
|
||||
// failed update and a broken box are indistinguishable afterwards, and the
|
||||
// rollback has nothing to prove itself against — so Apply refuses to start.
|
||||
//
|
||||
// Build and test run BEFORE anything is written to the install dir, so a broken
|
||||
// tree costs nothing but time. Install is per-file write-temp-then-rename, so a
|
||||
// crash mid-install leaves whole files, not half ones.
|
||||
//
|
||||
// The rollback path deliberately depends on nothing that just changed: it copies
|
||||
// byte-for-byte from a snapshot taken before the install and re-runs the same
|
||||
// restart command. It does not ask the new binary to do anything, does not run
|
||||
// a migration, and does not need the update to have gotten far enough to leave
|
||||
// a working anything behind.
|
||||
//
|
||||
// # What is out of scope on purpose
|
||||
//
|
||||
// The database is not snapshotted or rolled back. It is encrypted, live, and
|
||||
// often larger than the disk headroom; a store rolled back under a schema that
|
||||
// already migrated forward loses writes silently, which is worse than a failed
|
||||
// update. Schema compatibility is store.Migrate's job. A snapshot here is the
|
||||
// deployable artifacts only: binaries and config.
|
||||
package update
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
var (
|
||||
// ErrNotConfigured — no update block in the config. The capability does not
|
||||
// exist unless the operator described his own deployment.
|
||||
ErrNotConfigured = errors.New("update: not configured")
|
||||
|
||||
// ErrUnhealthyBefore — the daemon was already not answering when Apply
|
||||
// started. Refused: see the package comment.
|
||||
ErrUnhealthyBefore = errors.New("update: the running daemon is not healthy — refusing to update on top of a broken box")
|
||||
|
||||
// ErrVerifyFailed — build or test failed. Nothing was installed.
|
||||
ErrVerifyFailed = errors.New("update: verification failed")
|
||||
|
||||
// ErrRolledBack — the new build was installed and did not come up healthy,
|
||||
// so the previous snapshot was restored. Wraps the underlying failure.
|
||||
ErrRolledBack = errors.New("update: rolled back")
|
||||
|
||||
// ErrRollbackFailed — the worst case: the new build failed AND the restore
|
||||
// did not bring her back. The operator has to fix the box by hand; the
|
||||
// snapshot directory is named in the result so he knows what to copy.
|
||||
ErrRollbackFailed = errors.New("update: ROLLBACK FAILED — manual recovery required")
|
||||
)
|
||||
|
||||
// Config — the operator's description of his own deployment. Every path is
|
||||
// absolute and validated; nothing is guessed, because guessing wrong here means
|
||||
// overwriting the wrong file.
|
||||
type Config struct {
|
||||
// SourceDir — the git working tree to build. The operator pulls it himself;
|
||||
// this package never fetches.
|
||||
SourceDir string `json:"source_dir"`
|
||||
|
||||
// InstallDir — where the built binaries are copied to. On the docker
|
||||
// deployment this is the tree the image is built from, so it is usually the
|
||||
// same as SourceDir and Install is a no-op copy; on a bare-metal deployment
|
||||
// it is /opt/maven/bin.
|
||||
InstallDir string `json:"install_dir"`
|
||||
|
||||
// SnapshotDir — where the pre-install copies live. Must not be inside
|
||||
// InstallDir: a restore reading from a directory the install is writing to
|
||||
// is not a restore.
|
||||
SnapshotDir string `json:"snapshot_dir"`
|
||||
|
||||
// Binaries — the artifact names to snapshot and install, relative to
|
||||
// SourceDir (built) and InstallDir (deployed). Listed explicitly rather than
|
||||
// globbed so a stray file in the tree never gets deployed.
|
||||
Binaries []string `json:"binaries"`
|
||||
|
||||
// ConfigFiles — extra files to snapshot alongside the binaries, relative to
|
||||
// InstallDir. Snapshotted, never overwritten by an install: the operator's
|
||||
// config is not something an update gets to replace.
|
||||
ConfigFiles []string `json:"config_files,omitempty"`
|
||||
|
||||
// RestartCmd — how this deployment restarts mavend, e.g.
|
||||
// ["docker","compose","up","-d","--build","mavend"] or
|
||||
// ["systemctl","restart","mavend"]. Run in SourceDir. Required: there is no
|
||||
// portable default and picking one would mean restarting the wrong thing.
|
||||
RestartCmd []string `json:"restart_cmd"`
|
||||
|
||||
// HealthSocket — mavend's IPC socket, used to prove she answers after a
|
||||
// restart. Required: without a health check there is no signal to roll back
|
||||
// on, and an update that cannot detect its own failure is not what this
|
||||
// package is for.
|
||||
HealthSocket string `json:"health_socket"`
|
||||
|
||||
// HealthTimeoutSec — how long to wait for the restarted daemon to answer.
|
||||
// Default 90s; she loads a 1.7B on boot, so this is not a couple of seconds.
|
||||
HealthTimeoutSec int `json:"health_timeout_sec,omitempty"`
|
||||
|
||||
// VerifyTimeoutMin — cap on `make build` + `make test`. Default 20m.
|
||||
VerifyTimeoutMin int `json:"verify_timeout_min,omitempty"`
|
||||
|
||||
// KeepSnapshots — how many snapshots to retain. Default 5, minimum 1: the
|
||||
// most recent one is the rollback target and is never pruned.
|
||||
KeepSnapshots int `json:"keep_snapshots,omitempty"`
|
||||
}
|
||||
|
||||
// Validate — fail at startup, not halfway through an install.
|
||||
func (c Config) Validate() error {
|
||||
if c.SourceDir == "" || c.InstallDir == "" || c.SnapshotDir == "" {
|
||||
return errors.New("update: source_dir, install_dir and snapshot_dir are all required")
|
||||
}
|
||||
for _, p := range []string{c.SourceDir, c.InstallDir, c.SnapshotDir} {
|
||||
if !filepath.IsAbs(p) {
|
||||
return fmt.Errorf("update: %q must be an absolute path", p)
|
||||
}
|
||||
}
|
||||
if within(c.SnapshotDir, c.InstallDir) {
|
||||
return fmt.Errorf("update: snapshot_dir %q is inside install_dir %q — a restore must not read from what the install writes", c.SnapshotDir, c.InstallDir)
|
||||
}
|
||||
if len(c.Binaries) == 0 {
|
||||
return errors.New("update: binaries is empty — nothing to install")
|
||||
}
|
||||
for _, b := range append(append([]string{}, c.Binaries...), c.ConfigFiles...) {
|
||||
if filepath.IsAbs(b) || strings.Contains(b, "..") {
|
||||
return fmt.Errorf("update: %q must be a plain relative name", b)
|
||||
}
|
||||
}
|
||||
if len(c.RestartCmd) == 0 {
|
||||
return errors.New("update: restart_cmd is required — there is no safe default for restarting someone else's deployment")
|
||||
}
|
||||
if c.HealthSocket == "" {
|
||||
return errors.New("update: health_socket is required — an update that cannot check its own result cannot roll back on failure")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c Config) withDefaults() Config {
|
||||
if c.HealthTimeoutSec <= 0 {
|
||||
c.HealthTimeoutSec = 90
|
||||
}
|
||||
if c.VerifyTimeoutMin <= 0 {
|
||||
c.VerifyTimeoutMin = 20
|
||||
}
|
||||
if c.KeepSnapshots < 1 {
|
||||
c.KeepSnapshots = 5
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
func (c Config) healthTimeout() time.Duration {
|
||||
return time.Duration(c.HealthTimeoutSec) * time.Second
|
||||
}
|
||||
|
||||
func (c Config) verifyTimeout() time.Duration {
|
||||
return time.Duration(c.VerifyTimeoutMin) * time.Minute
|
||||
}
|
||||
|
||||
// within reports whether p is dir or lives under it.
|
||||
func within(p, dir string) bool {
|
||||
p, dir = filepath.Clean(p), filepath.Clean(dir)
|
||||
if p == dir {
|
||||
return true
|
||||
}
|
||||
rel, err := filepath.Rel(dir, p)
|
||||
return err == nil && rel != ".." && !strings.HasPrefix(rel, ".."+string(filepath.Separator))
|
||||
}
|
||||
|
||||
// Runner runs one command and returns its combined output. Injected so the
|
||||
// tests can drive build/test/restart failures without a toolchain, a container
|
||||
// or a real daemon to break.
|
||||
type Runner func(ctx context.Context, dir string, argv []string) (string, error)
|
||||
|
||||
// ExecRunner is the real one.
|
||||
func ExecRunner(ctx context.Context, dir string, argv []string) (string, error) {
|
||||
cmd := exec.CommandContext(ctx, argv[0], argv[1:]...)
|
||||
cmd.Dir = dir
|
||||
out, err := cmd.CombinedOutput()
|
||||
return string(out), err
|
||||
}
|
||||
|
||||
// HealthCheck proves the daemon at socket answers. Injected for the same reason
|
||||
// as Runner.
|
||||
type HealthCheck func(ctx context.Context, socket string) error
|
||||
|
||||
// Logger receives one line per step. The CLI prints these as they happen: an
|
||||
// update that goes quiet for four minutes during `make test` reads as a hang.
|
||||
type Logger func(format string, args ...any)
|
||||
|
||||
// Updater is the whole capability. Construct with New and call Apply or
|
||||
// Rollback; there is no background goroutine and nothing starts on its own.
|
||||
type Updater struct {
|
||||
cfg Config
|
||||
store *Store
|
||||
run Runner
|
||||
health HealthCheck
|
||||
log Logger
|
||||
now func() time.Time
|
||||
}
|
||||
|
||||
// New builds an Updater. Every seam has a real default; the tests replace them.
|
||||
func New(cfg Config, opts ...Option) (*Updater, error) {
|
||||
if err := cfg.Validate(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
u := &Updater{
|
||||
cfg: cfg.withDefaults(),
|
||||
store: &Store{Dir: cfg.SnapshotDir},
|
||||
run: ExecRunner,
|
||||
health: DialHealth,
|
||||
log: func(string, ...any) {},
|
||||
now: time.Now,
|
||||
}
|
||||
for _, o := range opts {
|
||||
o(u)
|
||||
}
|
||||
u.store.now = u.now
|
||||
return u, nil
|
||||
}
|
||||
|
||||
// Option — a constructor seam.
|
||||
type Option func(*Updater)
|
||||
|
||||
func WithRunner(r Runner) Option { return func(u *Updater) { u.run = r } }
|
||||
func WithHealth(h HealthCheck) Option { return func(u *Updater) { u.health = h } }
|
||||
func WithLogger(l Logger) Option { return func(u *Updater) { u.log = l } }
|
||||
func WithClock(f func() time.Time) Option {
|
||||
return func(u *Updater) { u.now = f }
|
||||
}
|
||||
|
||||
// Snapshots lists what is available to roll back to, newest first.
|
||||
func (u *Updater) Snapshots() ([]Snapshot, error) { return u.store.List() }
|
||||
@@ -0,0 +1,437 @@
|
||||
package update
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// The tests drive the whole orchestration against a fake box: a directory tree
|
||||
// standing in for the install dir, an injected Runner standing in for
|
||||
// make/git/docker, and an injected HealthCheck standing in for mavend. That is
|
||||
// what makes the failure paths — the ones that matter — testable at all: you
|
||||
// cannot ask a real deployment to fail its health check on demand, and the
|
||||
// rollback path is exactly the path nobody exercises by hand.
|
||||
|
||||
type fakeBox struct {
|
||||
t *testing.T
|
||||
root string
|
||||
|
||||
// what the fake `make build` writes into the source tree
|
||||
newBytes string
|
||||
// scripted failures
|
||||
buildErr error
|
||||
testErr error
|
||||
restartErr error
|
||||
|
||||
// health: fails until the Nth call, then follows healthy
|
||||
healthErrs int // remaining failures to serve
|
||||
healthy bool
|
||||
healthChecks int
|
||||
// deployedAtRestart records the installed bytes each time restart runs, so a
|
||||
// test can prove the rollback put the old bytes back BEFORE restarting.
|
||||
deployedAtRestart []string
|
||||
|
||||
ran []string
|
||||
}
|
||||
|
||||
func newFakeBox(t *testing.T) *fakeBox {
|
||||
t.Helper()
|
||||
root := t.TempDir()
|
||||
for _, d := range []string{"src", "install", "snapshots"} {
|
||||
if err := os.MkdirAll(filepath.Join(root, d), 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
// The currently deployed build, and a config file next to it.
|
||||
write(t, filepath.Join(root, "install", "mavend"), "OLD-BUILD")
|
||||
write(t, filepath.Join(root, "install", "mavend.json"), `{"tick_interval":"60s"}`)
|
||||
// The source tree already contains a stale binary; `make build` overwrites it.
|
||||
write(t, filepath.Join(root, "src", "mavend"), "STALE")
|
||||
return &fakeBox{t: t, root: root, newBytes: "NEW-BUILD", healthy: true}
|
||||
}
|
||||
|
||||
func (b *fakeBox) cfg() Config {
|
||||
return Config{
|
||||
SourceDir: filepath.Join(b.root, "src"),
|
||||
InstallDir: filepath.Join(b.root, "install"),
|
||||
SnapshotDir: filepath.Join(b.root, "snapshots"),
|
||||
Binaries: []string{"mavend"},
|
||||
ConfigFiles: []string{"mavend.json"},
|
||||
RestartCmd: []string{"restart-the-thing"},
|
||||
HealthSocket: filepath.Join(b.root, "mavend.sock"),
|
||||
HealthTimeoutSec: 1,
|
||||
KeepSnapshots: 3,
|
||||
}
|
||||
}
|
||||
|
||||
func (b *fakeBox) run(ctx context.Context, dir string, argv []string) (string, error) {
|
||||
b.ran = append(b.ran, strings.Join(argv, " "))
|
||||
switch strings.Join(argv, " ") {
|
||||
case "make build":
|
||||
if b.buildErr != nil {
|
||||
return "ld: undefined reference to everything", b.buildErr
|
||||
}
|
||||
// A real build writes its artifacts into the working tree — the behaviour
|
||||
// the snapshot-before-build ordering exists to survive.
|
||||
write(b.t, filepath.Join(b.root, "src", "mavend"), b.newBytes)
|
||||
return "built", nil
|
||||
case "make test":
|
||||
if b.testErr != nil {
|
||||
return "--- FAIL: TestSomething", b.testErr
|
||||
}
|
||||
return "ok", nil
|
||||
case "git rev-parse HEAD":
|
||||
return "cafebabecafebabecafebabecafebabecafebabe\n", nil
|
||||
case "restart-the-thing":
|
||||
b.deployedAtRestart = append(b.deployedAtRestart, read(b.t, filepath.Join(b.root, "install", "mavend")))
|
||||
if b.restartErr != nil {
|
||||
return "no such container", b.restartErr
|
||||
}
|
||||
return "restarted", nil
|
||||
}
|
||||
return "", errors.New("unexpected command: " + strings.Join(argv, " "))
|
||||
}
|
||||
|
||||
func (b *fakeBox) health(ctx context.Context, socket string) error {
|
||||
b.healthChecks++
|
||||
if b.healthErrs > 0 {
|
||||
b.healthErrs--
|
||||
return errors.New("connection refused")
|
||||
}
|
||||
if !b.healthy {
|
||||
return errors.New("she does not answer")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b *fakeBox) updater(t *testing.T, extra ...Option) *Updater {
|
||||
t.Helper()
|
||||
opts := append([]Option{WithRunner(b.run), WithHealth(b.health)}, extra...)
|
||||
u, err := New(b.cfg(), opts...)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return u
|
||||
}
|
||||
|
||||
func (b *fakeBox) deployed() string { return read(b.t, filepath.Join(b.root, "install", "mavend")) }
|
||||
|
||||
func write(t *testing.T, path, content string) {
|
||||
t.Helper()
|
||||
if err := os.WriteFile(path, []byte(content), 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func read(t *testing.T, path string) string {
|
||||
t.Helper()
|
||||
b, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return string(b)
|
||||
}
|
||||
|
||||
func TestApply_HappyPath(t *testing.T) {
|
||||
b := newFakeBox(t)
|
||||
res, err := b.updater(t).Apply(context.Background())
|
||||
if err != nil {
|
||||
t.Fatalf("Apply: %v", err)
|
||||
}
|
||||
if !res.Verified || !res.Restarted || !res.Healthy || res.RolledBack {
|
||||
t.Fatalf("result = %+v; want verified+restarted+healthy and no rollback", res)
|
||||
}
|
||||
if got := b.deployed(); got != "NEW-BUILD" {
|
||||
t.Errorf("deployed binary = %q; want the new build", got)
|
||||
}
|
||||
// The order is the property: health, snapshot, build, test, install, restart.
|
||||
want := []string{"git rev-parse HEAD", "make build", "make test", "restart-the-thing"}
|
||||
if strings.Join(b.ran, "|") != strings.Join(want, "|") {
|
||||
t.Errorf("commands ran = %v; want %v", b.ran, want)
|
||||
}
|
||||
if res.SnapshotID == "" {
|
||||
t.Error("no snapshot was taken")
|
||||
}
|
||||
}
|
||||
|
||||
func TestApply_RefusesWhenSheIsAlreadyDown(t *testing.T) {
|
||||
// A box that is already broken has no baseline for the rollback to prove
|
||||
// itself against, so the update never starts.
|
||||
b := newFakeBox(t)
|
||||
b.healthy = false
|
||||
res, err := b.updater(t).Apply(context.Background())
|
||||
if !errors.Is(err, ErrUnhealthyBefore) {
|
||||
t.Fatalf("Apply on an unhealthy box = %v; want ErrUnhealthyBefore", err)
|
||||
}
|
||||
if len(b.ran) != 0 {
|
||||
t.Errorf("a refused update still ran %v", b.ran)
|
||||
}
|
||||
if res.SnapshotID != "" {
|
||||
t.Error("a refused update still took a snapshot")
|
||||
}
|
||||
}
|
||||
|
||||
func TestApply_TestFailureDeploysNothingAndPutsTheTreeBack(t *testing.T) {
|
||||
b := newFakeBox(t)
|
||||
b.testErr = errors.New("exit status 1")
|
||||
res, err := b.updater(t).Apply(context.Background())
|
||||
if !errors.Is(err, ErrVerifyFailed) {
|
||||
t.Fatalf("Apply with failing tests = %v; want ErrVerifyFailed", err)
|
||||
}
|
||||
if res.Verified {
|
||||
t.Error("result claims verified after a failing test suite")
|
||||
}
|
||||
for _, c := range b.ran {
|
||||
if c == "restart-the-thing" {
|
||||
t.Fatal("a failed verification restarted the daemon")
|
||||
}
|
||||
}
|
||||
if got := b.deployed(); got != "OLD-BUILD" {
|
||||
t.Errorf("deployed binary = %q; want the old build untouched", got)
|
||||
}
|
||||
// The failing output is kept so the operator can see why.
|
||||
var found bool
|
||||
for _, s := range res.Steps {
|
||||
if s.Name == "test" && strings.Contains(s.Output, "FAIL") {
|
||||
found = true
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
t.Error("the failing test output was not retained")
|
||||
}
|
||||
}
|
||||
|
||||
func TestApply_BuildFailureIsCaughtBeforeTheTests(t *testing.T) {
|
||||
b := newFakeBox(t)
|
||||
b.buildErr = errors.New("exit status 2")
|
||||
if _, err := b.updater(t).Apply(context.Background()); !errors.Is(err, ErrVerifyFailed) {
|
||||
t.Fatalf("Apply with a failing build = %v; want ErrVerifyFailed", err)
|
||||
}
|
||||
for _, c := range b.ran {
|
||||
if c == "make test" {
|
||||
t.Error("ran the test suite after the build failed")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestApply_UnhealthyAfterRestartRollsBackToTheOldBytes(t *testing.T) {
|
||||
// The case the package exists for: everything verifies, the new build
|
||||
// installs, and then she does not come up.
|
||||
b := newFakeBox(t)
|
||||
b.healthErrs = 1 // the preflight check passes, then she stops answering
|
||||
b.healthy = false
|
||||
u := b.updater(t)
|
||||
// Once the rollback restores the old build, she answers again.
|
||||
restored := false
|
||||
u.health = func(ctx context.Context, socket string) error {
|
||||
b.healthChecks++
|
||||
if b.deployed() == "OLD-BUILD" && restored {
|
||||
return nil
|
||||
}
|
||||
if b.healthChecks == 1 {
|
||||
return nil // preflight: the old build is up
|
||||
}
|
||||
if b.deployed() == "OLD-BUILD" {
|
||||
restored = true
|
||||
return nil
|
||||
}
|
||||
return errors.New("she does not answer on the new build")
|
||||
}
|
||||
res, err := u.Apply(context.Background())
|
||||
if !errors.Is(err, ErrRolledBack) {
|
||||
t.Fatalf("Apply with a dead new build = %v; want ErrRolledBack", err)
|
||||
}
|
||||
if !res.RolledBack || !res.RollbackHealthy || res.Healthy {
|
||||
t.Fatalf("result = %+v; want rolled back and healthy again on the old build", res)
|
||||
}
|
||||
if got := b.deployed(); got != "OLD-BUILD" {
|
||||
t.Errorf("deployed binary after the rollback = %q; want OLD-BUILD", got)
|
||||
}
|
||||
// And the restore happened BEFORE the second restart, not after it.
|
||||
if len(b.deployedAtRestart) != 2 {
|
||||
t.Fatalf("restarts = %v; want two (the update and the rollback)", b.deployedAtRestart)
|
||||
}
|
||||
if b.deployedAtRestart[0] != "NEW-BUILD" || b.deployedAtRestart[1] != "OLD-BUILD" {
|
||||
t.Errorf("bytes in place at each restart = %v; want [NEW-BUILD OLD-BUILD]", b.deployedAtRestart)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApply_RestartFailureRollsBack(t *testing.T) {
|
||||
b := newFakeBox(t)
|
||||
b.restartErr = errors.New("exit status 1")
|
||||
res, err := b.updater(t).Apply(context.Background())
|
||||
// The rollback's own restart fails too, so this is the manual-recovery case —
|
||||
// and it says so instead of reporting a tidy rollback.
|
||||
if !errors.Is(err, ErrRollbackFailed) {
|
||||
t.Fatalf("Apply with a broken restart command = %v; want ErrRollbackFailed", err)
|
||||
}
|
||||
if !res.RolledBack || res.RollbackHealthy {
|
||||
t.Fatalf("result = %+v; want rolled back but not healthy", res)
|
||||
}
|
||||
if got := b.deployed(); got != "OLD-BUILD" {
|
||||
t.Errorf("deployed binary = %q; want the old bytes restored even so", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApply_RollbackNeedsNoBuildAndNoNewCode(t *testing.T) {
|
||||
// The rollback must not depend on the toolchain, the source tree, or the
|
||||
// code it is replacing. Prove it: delete the source tree's binary and make
|
||||
// every command except the restart fail, then roll back.
|
||||
b := newFakeBox(t)
|
||||
if _, err := b.updater(t).Apply(context.Background()); err != nil {
|
||||
t.Fatalf("setup Apply: %v", err)
|
||||
}
|
||||
if b.deployed() != "NEW-BUILD" {
|
||||
t.Fatal("setup did not deploy")
|
||||
}
|
||||
os.RemoveAll(filepath.Join(b.root, "src"))
|
||||
if err := os.MkdirAll(filepath.Join(b.root, "src"), 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
b.buildErr = errors.New("no toolchain here")
|
||||
b.testErr = errors.New("no toolchain here")
|
||||
b.ran = nil
|
||||
|
||||
res, err := b.updater(t).Rollback(context.Background(), "")
|
||||
if err != nil && !errors.Is(err, ErrRolledBack) {
|
||||
t.Fatalf("Rollback: %v", err)
|
||||
}
|
||||
if !res.RollbackHealthy {
|
||||
t.Fatalf("result = %+v; want a healthy rollback", res)
|
||||
}
|
||||
if got := b.deployed(); got != "OLD-BUILD" {
|
||||
t.Errorf("deployed binary = %q; want OLD-BUILD", got)
|
||||
}
|
||||
for _, c := range b.ran {
|
||||
if strings.HasPrefix(c, "make") {
|
||||
t.Errorf("the rollback ran %q — it must not need a build", c)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestApply_ConfigIsSnapshottedButNeverOverwritten(t *testing.T) {
|
||||
b := newFakeBox(t)
|
||||
// A config in the source tree must not be deployed over the operator's.
|
||||
write(t, filepath.Join(b.root, "src", "mavend.json"), `{"tick_interval":"1s"}`)
|
||||
if _, err := b.updater(t).Apply(context.Background()); err != nil {
|
||||
t.Fatalf("Apply: %v", err)
|
||||
}
|
||||
if got := read(t, filepath.Join(b.root, "install", "mavend.json")); !strings.Contains(got, "60s") {
|
||||
t.Errorf("installed config = %q; an update must not replace his config", got)
|
||||
}
|
||||
snaps, err := b.updater(t).Snapshots()
|
||||
if err != nil || len(snaps) == 0 {
|
||||
t.Fatalf("Snapshots: %v %v", snaps, err)
|
||||
}
|
||||
var names []string
|
||||
for _, f := range snaps[0].Files {
|
||||
names = append(names, f.Name)
|
||||
}
|
||||
if len(names) != 2 {
|
||||
t.Errorf("snapshot files = %v; want the binary and the config", names)
|
||||
}
|
||||
if snaps[0].Commit == "" {
|
||||
t.Error("the snapshot did not record which commit produced it")
|
||||
}
|
||||
}
|
||||
|
||||
func TestRollback_CorruptSnapshotIsRefusedNotRestored(t *testing.T) {
|
||||
b := newFakeBox(t)
|
||||
if _, err := b.updater(t).Apply(context.Background()); err != nil {
|
||||
t.Fatalf("setup Apply: %v", err)
|
||||
}
|
||||
snaps, _ := b.updater(t).Snapshots()
|
||||
// Something ate the snapshot. Restoring it would deploy garbage.
|
||||
write(t, filepath.Join(snaps[0].Dir(), "mavend"), "CORRUPT")
|
||||
_, err := b.updater(t).Rollback(context.Background(), snaps[0].ID)
|
||||
if !errors.Is(err, ErrRollbackFailed) || !strings.Contains(err.Error(), "corrupt") {
|
||||
t.Fatalf("Rollback of a corrupt snapshot = %v; want a refusal naming the corruption", err)
|
||||
}
|
||||
if got := b.deployed(); got != "NEW-BUILD" {
|
||||
t.Errorf("deployed binary = %q; a refused restore must change nothing", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRollback_NoSnapshots(t *testing.T) {
|
||||
b := newFakeBox(t)
|
||||
if _, err := b.updater(t).Rollback(context.Background(), ""); err == nil {
|
||||
t.Error("Rollback with no snapshots succeeded; want an error")
|
||||
}
|
||||
}
|
||||
|
||||
func TestPrune_KeepsTheNewestAsTheRollbackTarget(t *testing.T) {
|
||||
b := newFakeBox(t)
|
||||
st := &Store{Dir: filepath.Join(b.root, "snapshots")}
|
||||
base := time.Date(2026, 8, 1, 3, 0, 0, 0, time.UTC)
|
||||
for i := 0; i < 4; i++ {
|
||||
i := i
|
||||
st.now = func() time.Time { return base.Add(time.Duration(i) * time.Minute) }
|
||||
if _, err := st.Save(filepath.Join(b.root, "install"), []string{"mavend"}, "", ""); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if err := st.Prune(0); err != nil { // 0 is clamped to 1, never to zero
|
||||
t.Fatal(err)
|
||||
}
|
||||
snaps, err := st.List()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(snaps) != 1 {
|
||||
t.Fatalf("kept %d snapshots; want 1", len(snaps))
|
||||
}
|
||||
if snaps[0].ID != "20260801-030300" {
|
||||
t.Errorf("kept %s; want the newest", snaps[0].ID)
|
||||
}
|
||||
}
|
||||
|
||||
func TestList_IgnoresSnapshotsWithNoManifest(t *testing.T) {
|
||||
// An interrupted snapshot has files but no manifest. It must never be offered
|
||||
// as a rollback target — restoring a half-copied binary is the worst outcome
|
||||
// in the package.
|
||||
b := newFakeBox(t)
|
||||
dir := filepath.Join(b.root, "snapshots", "20260801-000000")
|
||||
if err := os.MkdirAll(dir, 0o700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
write(t, filepath.Join(dir, "mavend"), "HALF")
|
||||
snaps, err := (&Store{Dir: filepath.Join(b.root, "snapshots")}).List()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(snaps) != 0 {
|
||||
t.Errorf("List returned %d snapshots; want none (no manifest)", len(snaps))
|
||||
}
|
||||
}
|
||||
|
||||
func TestConfigValidate(t *testing.T) {
|
||||
ok := (&fakeBox{root: t.TempDir()}).cfg()
|
||||
if err := ok.Validate(); err != nil {
|
||||
t.Fatalf("valid config rejected: %v", err)
|
||||
}
|
||||
bad := map[string]func(c Config) Config{
|
||||
"relative source": func(c Config) Config { c.SourceDir = "src"; return c },
|
||||
"no restart command": func(c Config) Config { c.RestartCmd = nil; return c },
|
||||
"no health socket": func(c Config) Config { c.HealthSocket = ""; return c },
|
||||
"no binaries": func(c Config) Config { c.Binaries = nil; return c },
|
||||
"escaping artifact name": func(c Config) Config { c.Binaries = []string{"../../etc/passwd"}; return c },
|
||||
"absolute artifact name": func(c Config) Config { c.Binaries = []string{"/usr/bin/mavend"}; return c },
|
||||
"snapshots inside install": func(c Config) Config { c.SnapshotDir = filepath.Join(c.InstallDir, "snaps"); return c },
|
||||
}
|
||||
for name, mutate := range bad {
|
||||
if err := mutate(ok).Validate(); err == nil {
|
||||
t.Errorf("%s was accepted; want a startup failure", name)
|
||||
}
|
||||
}
|
||||
// And New refuses an invalid config outright rather than half-configuring.
|
||||
if _, err := New(mutate(ok, "no health socket", bad)); err == nil {
|
||||
t.Error("New accepted a config with no health socket")
|
||||
}
|
||||
}
|
||||
|
||||
func mutate(c Config, key string, m map[string]func(Config) Config) Config { return m[key](c) }
|
||||
@@ -0,0 +1,65 @@
|
||||
package update
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Verification is "does this tree build and does it pass its own tests", run
|
||||
// before a single byte is written to the install dir.
|
||||
//
|
||||
// It is `make build` and `make test`, not `go build`: the CGO daemons need the
|
||||
// vendored toolchain and the whisper/piper include and library paths wired
|
||||
// through the Makefile, and a bare `go build` on them fails in a way that has
|
||||
// nothing to do with the change being deployed. `make test` is the -race suite
|
||||
// with the CGO env set, and it is the only evidence available on a single box
|
||||
// that the new code does what the old code did.
|
||||
//
|
||||
// This is not a substitute for a second environment. A test suite that passes
|
||||
// says the code is self-consistent; it does not say the new build will start
|
||||
// against this machine's actual models, sockets and encrypted store. That is
|
||||
// what the post-restart health check is for, and it is why the install is
|
||||
// reversible rather than merely careful.
|
||||
|
||||
// Step — one verification or orchestration step and how it went. Kept so the CLI
|
||||
// can print a truthful account of what was done, including on the failure path.
|
||||
type Step struct {
|
||||
Name string
|
||||
Argv []string
|
||||
Took time.Duration
|
||||
Err error
|
||||
Output string // combined output, only retained for failures
|
||||
}
|
||||
|
||||
// Verify runs the build and the test suite in SourceDir.
|
||||
func (u *Updater) Verify(ctx context.Context) ([]Step, error) {
|
||||
ctx, cancel := context.WithTimeout(ctx, u.cfg.verifyTimeout())
|
||||
defer cancel()
|
||||
var steps []Step
|
||||
for _, argv := range [][]string{{"make", "build"}, {"make", "test"}} {
|
||||
u.log("verify: %v (this takes a while)", argv)
|
||||
start := u.now()
|
||||
out, err := u.run(ctx, u.cfg.SourceDir, argv)
|
||||
st := Step{Name: argv[len(argv)-1], Argv: argv, Took: u.now().Sub(start), Err: err}
|
||||
if err != nil {
|
||||
st.Output = tail(out, 4000)
|
||||
}
|
||||
steps = append(steps, st)
|
||||
if err != nil {
|
||||
u.log("verify: %v FAILED after %s", argv, st.Took.Round(time.Second))
|
||||
return steps, fmt.Errorf("%w: %v: %v", ErrVerifyFailed, argv, err)
|
||||
}
|
||||
u.log("verify: %v ok in %s", argv, st.Took.Round(time.Second))
|
||||
}
|
||||
return steps, nil
|
||||
}
|
||||
|
||||
// tail keeps the last n bytes — a failing `make test` prints far more than is
|
||||
// useful, and the failure is always at the end.
|
||||
func tail(s string, n int) string {
|
||||
if len(s) <= n {
|
||||
return s
|
||||
}
|
||||
return "…" + s[len(s)-n:]
|
||||
}
|
||||
Reference in New Issue
Block a user