Compare commits

..

2 Commits

Author SHA1 Message Date
kami cb3641e7bb Read RSS and Atom feeds, and speak about them only when asked (#258)
internal/rss parses RSS 2.0 and Atom, and polls each configured feed on its own
interval; internal/webfetch is the one door either of them uses to touch the
network. The poller writes items as notes with source "rss:<feed>" and nothing
else: the answer path reads them back when he asks "что нового в лентах?", and
nothing is announced on arrival. A feed that dispatched would be a nag, which is
why the plan's breaking-news rule was left out rather than built.

webfetch is where the limits live, as code rather than a paragraph: http(s)
only, an allowlist (the configured feeds' hosts) and a denylist, a 2 MiB body
cap, a 3-redirect cap, one request per host per second, and a refusal to connect
to any private address — checked in the dialer's Control hook so it holds for
every resolved address and every redirect hop, not just for a literal IP.

Off unless configured: no "feeds" block, no poller, no outbound request. How far
a feed was read is a config fact (rss:latest:<name>), so a restart does not
re-note yesterday's headlines.
2026-08-01 03:27:45 +04:00
kami ee7bec11e3 Add mavmaild, the read-only IMAP poller that feeds mail intake (#246)
The extraction seam landed on the previous branch but nothing fed it. This
adds the daemon that does: every interval it opens one mailbox read-only
(EXAMINE + BODY.PEEK, so reading leaves no \Seen behind), fetches the UIDs
it has not handed over yet, and posts each message to core over
ingest_mail. Core runs the model and writes task candidates; this daemon
writes nothing and cannot create a reminder.

It is a separate daemon because of the credential. mavpoll set the
precedent with the zenmoney token (#125): the module talking to the third
party holds the secret, reads it from a file so it never lands in argv, in
docker-compose.yml or in shell history, and core never sees it. There is
deliberately no -password flag, and a test asserts that.

Off unless configured at both ends: without -password-file the daemon
refuses to start, and if core has no email block the first ingest returns
ErrUnknownMethod, which disables the reader instead of hammering a socket
that will keep refusing. A seen-UID state file (0600, atomic write) keeps a
restart from re-extracting the whole lookback window; correctness does not
depend on it, since capture dedupes on normalised text. Logs are counts and
UIDs — no subject, sender or body.

Verified with an in-process IMAP server and a fake core: bulk mail is
filtered before core is asked, seen UIDs are not re-fetched, a failed
ingest is retried next poll, ErrUnknownMethod stops at the first message,
and state survives a restart. The live half is untested by design — no IMAP
credential exists on this box; setup is written up as QA steps.

Vikunja #246
2026-08-01 03:13:18 +04:00
26 changed files with 2705 additions and 16 deletions
+3
View File
@@ -7,6 +7,7 @@
/mavpoll
/mavcaldav
/mavwaked
/mavmaild
# Certs (private keys, don't commit)
certs/
@@ -36,6 +37,8 @@ deploy/db_key.env
deploy/telegram.env
# zenmoney API token, read by mavpoll (never in argv, never committed)
deploy/zenmoney.token
# IMAP password, read by mavmaild (never in argv, never committed)
deploy/imap.password
# Temp files
/tmp/
+2 -1
View File
@@ -34,7 +34,7 @@ CGO daemons (`mavend`, `mavsttd`, `mavttsd`, `mavenclient`) need the vendored to
and libs wired through the Makefile — **do not** call `go build` on them bare, use `make`:
```sh
make build # all 8 binaries
make build # all 9 binaries
make build-web # single daemon (pure-Go ones: web/waked/poll/caldav build without CGO)
make test # go test -race across ./internal/... ./cmd/... with CGO env set
```
@@ -62,6 +62,7 @@ Pure-Go packages (`router`, `memory`, `mavweb`, …) run under a plain `go test
| `mavenclient` | Voice loop client (mic → stt → core → tts). |
| `mavpoll` | Telegram long-poll reach. |
| `mavcaldav` | CalDAV calendar sync. |
| `mavmaild` | Mail reader (IMAP, read-only). Holds the IMAP password; core never sees it. |
Daemons are wired socket-to-socket, not linked. `internal/ipc` is the client/server wire
protocol; the config in `deploy/mavend.json` (with `${VAR}` env expansion from gitignored
+2 -1
View File
@@ -51,7 +51,8 @@ RUN go build -o /out/mavend ./cmd/mavend && \
go build -o /out/mavttsd ./cmd/mavttsd && \
go build -o /out/mavweb ./cmd/mavweb && \
go build -o /out/mavpoll ./cmd/mavpoll && \
go build -o /out/mavcaldav ./cmd/mavcaldav
go build -o /out/mavcaldav ./cmd/mavcaldav && \
go build -o /out/mavmaild ./cmd/mavmaild
# llama.cpp Vulkan build — the phraser/router LFM engine (llama-server). Built
# from source (not a prebuilt vendored blob) so the binary's glibc/GLIBCXX match
+5 -2
View File
@@ -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: build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav build-mail
build-stt:
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
@@ -50,6 +50,9 @@ build-poll:
build-caldav:
$(GO) build $(GOFLAGS) -o mavcaldav ./cmd/mavcaldav/
build-mail:
$(GO) build $(GOFLAGS) -o mavmaild ./cmd/mavmaild/
run-web: build-web
./mavweb -addr :9200 -voice 127.0.0.1:9100
@@ -185,4 +188,4 @@ download-embedder:
@echo ' sudo cp onnxruntime-linux-x64-1.15.1/lib/libonnxruntime.so* /usr/local/lib/'
clean:
rm -f mavend mavenclient mavsttd mavttsd mavweb mavpoll mavcaldav mavwaked
rm -f mavend mavenclient mavsttd mavttsd mavweb mavpoll mavcaldav mavwaked mavmaild
+63
View File
@@ -12,6 +12,7 @@ import (
"github.com/kami/maven/internal/memory"
"github.com/kami/maven/internal/morning"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/rss"
"github.com/kami/maven/internal/weather"
)
@@ -67,6 +68,11 @@ var querySources = []querySource{
// answer it from whatever he once said about spending. Its matcher needs a
// money noun plus an actual ask, so "я потратил весь день" is untouched.
{"money", (*reactiveHandler).queryMoney},
// Before the recall sources and before general knowledge: "что нового?" is
// a question about the feeds she reads, and general knowledge would answer
// it by inventing news. Its matcher needs a feed noun plus an ask, so
// "у меня новая лента в инстаграме" is untouched.
{"feeds", (*reactiveHandler).queryFeeds},
{"calendar", (*reactiveHandler).queryCalendar},
{"weather", (*reactiveHandler).queryWeather},
{"embed", (*reactiveHandler).queryEmbed},
@@ -180,6 +186,63 @@ func (h *reactiveHandler) queryHabits(ctx context.Context, t *queryTurn) (string
return profile.FormatOverallRU(), true
}
// feedNoteWindow — how many recent notes are scanned for feed items, and
// feedReadOut — how many headlines she actually reads back. She summarises the
// top of the pile, she does not recite a river.
const (
feedNoteWindow = 200
feedReadOut = 3
)
// queryFeeds — "что нового в лентах?", "что нового по технологиям?"
// (Vikunja #258).
//
// This is the ONLY way a feed item reaches him. The poller writes notes and
// never speaks; asking is the trigger. If that ever changes, the thing that
// changed is "Maven is not a nag", not a detail of this file.
func (h *reactiveHandler) queryFeeds(ctx context.Context, t *queryTurn) (string, bool) {
q, ok := router.ParseFeedQuery(t.dec.Utterance)
if !ok {
return "", false
}
if !h.feedsOn {
// Claim the turn rather than fall through: "не читаю ленты" is true, and
// letting general knowledge answer "что нового?" would be an invented
// news bulletin.
return "я пока не читаю ленты — они не настроены.", true
}
notes, err := h.api.RecentNotes(ctx, feedNoteWindow)
if err != nil {
log.Printf("voice: feeds: recent notes: %v", err)
return "не получилось посмотреть ленты.", true
}
var picked []string
for _, n := range notes {
if !strings.HasPrefix(n.Source, rss.SourcePrefix) {
continue
}
if !router.CategoryMatches(n.Text, q.Category) {
continue
}
// The note carries title, summary and link; she reads the title.
title := n.Text
if i := strings.IndexByte(title, '\n'); i > 0 {
title = title[:i]
}
picked = append(picked, strings.TrimSpace(title))
if len(picked) == feedReadOut {
break
}
}
if len(picked) == 0 {
if q.Category != "" {
return "по этой теме в лентах пока ничего.", true
}
return "в лентах пока ничего нового.", true
}
return "вот что нового: " + strings.Join(picked, "; "), true
}
// queryCalendar — "что у меня сегодня?", "планы на завтра?"
// h.now(), not time.Now(): the handler's clock is the injected one, so this
// source can be tested at a fixed time like the rest.
+182
View File
@@ -0,0 +1,182 @@
// mavend/feeds.go — the driver for RSS/Atom reading (Vikunja #258,
// docs/plans/13-rss-news-feeds.md). The reader itself is pure and lives in
// internal/rss; this is the impure half: a ticker, the guarded fetcher, and the
// two adapters that let a pure package talk to the store.
//
// Why in-core rather than its own daemon like mavmaild and mavpoll: those two
// hold a CREDENTIAL (an IMAP password, a zenmoney token), and the reason they
// are separate processes is that core must never see it. A feed URL is public,
// there is no secret to isolate, and a whole extra binary and compose service
// would buy nothing. The other half of the mavpoll precedent — off unless
// configured — is kept: no `feeds` block, no poller, no outbound request.
//
// It is its own goroutine, not a step on the tick: the tick has a delivery
// deadline behind it, and a feed read is a network round-trip that nobody is
// waiting on.
//
// Nothing here dispatches. A feed that announced itself would be a nag, so the
// only output is notes with source "rss:<feed>", which the answer path reads
// when he asks ("что нового в лентах?" — see queryFeeds in actions_query.go).
package main
import (
"context"
"log"
"net/url"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/rss"
"github.com/kami/maven/internal/webfetch"
)
// feedWorker — ticker + poller.
type feedWorker struct {
poller *rss.Poller
interval time.Duration
}
// feedTickInterval — how often the worker asks the poller what is due. Per-feed
// cadence is the poller's business; this is just the granularity.
const feedTickInterval = 5 * time.Minute
// newFeedWorker wires feed reading, or returns nil when it must not run:
// no `feeds` block (the normal case), or nothing valid in it. Every caller
// checks for nil.
func newFeedWorker(api ipc.CoreAPI, emb router.Embedder, cfg *config.Config) *feedWorker {
if cfg.Feeds == nil {
return nil
}
fc := cfg.Feeds
feeds := make([]rss.FeedConfig, 0, len(fc.Sources))
hosts := append([]string(nil), fc.AllowHosts...)
for _, s := range fc.Sources {
feeds = append(feeds, rss.FeedConfig{
Name: s.Name,
URL: s.URL,
Category: s.Category,
Interval: time.Duration(s.Interval),
Include: s.Include,
Exclude: s.Exclude,
})
// Each configured feed's own host is allowed. The allowlist is then
// exactly "the feeds he asked for", so a redirect off to somewhere else
// is refused by the fetcher rather than followed.
if u, err := url.Parse(s.URL); err == nil && u.Hostname() != "" {
hosts = append(hosts, u.Hostname())
}
}
fetcher := webfetch.New(webfetch.Config{
AllowHosts: hosts,
Timeout: time.Duration(fc.Timeout),
MaxBytes: fc.MaxBytes,
})
poller := rss.NewPoller(feeds, &feedFetcher{f: fetcher}, api, &factMarks{api: api},
embedderFor(emb), nil, rss.Config{
DefaultInterval: time.Duration(fc.PollInterval),
MaxItems: fc.MaxItems,
MaxAge: time.Duration(fc.MaxAge),
})
if poller == nil {
log.Printf("feeds: configured but nothing pollable — feed reading disabled")
return nil
}
log.Printf("feeds: reading %d feed(s), checking what is due every %s", len(feeds), feedTickInterval)
return &feedWorker{poller: poller, interval: feedTickInterval}
}
// run polls what is due until ctx is canceled. The first round runs immediately
// so a restart does not blind her for the first interval; it writes notes only,
// so an early round cannot startle anyone.
func (w *feedWorker) run(ctx context.Context) {
w.poller.PollDue(ctx, time.Now())
t := time.NewTicker(w.interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case now := <-t.C:
w.poller.PollDue(ctx, now)
}
}
}
// embedderOf — the voice wiring's embedder, or nil when voice is not wired.
// Feed notes are embedded with the SAME model the rest of the store uses, or not
// at all; a second embedder would write vectors nothing can search.
func embedderOf(w *voiceWiring) router.Embedder {
if w == nil {
return nil
}
return w.embedder
}
// feedFetcher adapts webfetch to rss.Fetcher — the pure package names the two
// fields it needs and stays free of net/http.
type feedFetcher struct{ f *webfetch.Fetcher }
func (a *feedFetcher) Get(ctx context.Context, u string) (*rss.Body, error) {
resp, err := a.f.Get(ctx, u)
if err != nil {
return nil, err
}
return &rss.Body{Bytes: resp.Body}, nil
}
// factMarks stores "how far this feed was read" as a config fact, the same
// mechanism the plan named and the same one the pattern tick uses for its own
// bookkeeping. Durable, inspectable on /dash, and cheap.
type factMarks struct{ api ipc.CoreAPI }
func markKey(feed string) string { return "rss:latest:" + feed }
func (m *factMarks) LastMark(ctx context.Context, feed string) (time.Time, error) {
f, err := m.api.LatestFact(ctx, markKey(feed))
if err != nil {
// No mark yet is not an error worth propagating: the poller treats a
// zero time as a cold start.
return time.Time{}, nil
}
t, err := time.Parse(time.RFC3339, f.Value)
if err != nil {
return time.Time{}, nil
}
return t, nil
}
func (m *factMarks) SetMark(ctx context.Context, feed string, at time.Time) error {
_, err := m.api.WriteFact(ctx, ipc.WriteFactReq{
Ts: time.Now(),
Kind: "config",
Key: markKey(feed),
Value: at.UTC().Format(time.RFC3339),
Source: "poll:rss",
Confidence: 1.0,
})
return err
}
// embedderFor adapts router.Embedder to rss.Embedder, and returns nil when
// there is none — a note without a vector is still a note the recent-notes path
// can read.
//
// EmbedPassage, not Embed: a feed item is text being searched FOR, and the e5
// embedder is asymmetric. Getting this backwards makes the item unfindable by
// the question that should have matched it.
func embedderFor(emb router.Embedder) rss.Embedder {
if emb == nil {
return nil
}
return passageEmbedder{emb}
}
type passageEmbedder struct{ e router.Embedder }
func (p passageEmbedder) Embed(ctx context.Context, text string) ([]float32, error) {
return router.EmbedPassage(ctx, p.e, text)
}
+177
View File
@@ -0,0 +1,177 @@
package main
import (
"context"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/rss"
"github.com/kami/maven/internal/voice"
)
// buildFeedHandler — a handler with the given feed notes already stored. No
// embedder: the feed source answers from recent notes by source, which is what
// makes it work for notes written before an embedder existed.
func buildFeedHandler(t *testing.T, feedsOn bool, notes ...ipc.Note) *reactiveHandler {
t.Helper()
ctx := context.Background()
st := newTestStore(t)
now := time.Now()
for i, n := range notes {
ts := now.Add(time.Duration(i) * time.Minute)
if _, err := st.WriteNote(ctx, ts, n.Text, nil, n.Source); err != nil {
t.Fatalf("WriteNote: %v", err)
}
}
return &reactiveHandler{
api: ipc.NewStoreAPI(st),
replier: voice.NewStubReplier(),
phraser: phraser.NewStub(),
now: func() time.Time { return now },
feedsOn: feedsOn,
embedder: nil,
}
}
func askFeeds(t *testing.T, h *reactiveHandler, q string) (string, bool) {
t.Helper()
return h.queryFeeds(context.Background(), &queryTurn{
dec: router.Decision{Intent: router.IntentQuery, Utterance: q},
})
}
func TestQueryFeedsReadsFeedNotes(t *testing.T) {
h := buildFeedHandler(t, true,
ipc.Note{Text: "Новая уязвимость в ядре [технологии]\nпатч вышел\nhttps://example.org/a", Source: "rss:habr"},
ipc.Note{Text: "что-то он сам сказал", Source: "tap:voice"},
)
reply, ok := askFeeds(t, h, "что нового в лентах?")
if !ok {
t.Fatal("the feed source did not claim the question")
}
if !strings.Contains(reply, "уязвимость") {
t.Errorf("reply = %q, want the headline", reply)
}
if strings.Contains(reply, "он сам сказал") {
t.Errorf("a note he dictated leaked into the feed answer: %q", reply)
}
// She reads the headline, not the summary and not the URL.
if strings.Contains(reply, "https://") || strings.Contains(reply, "патч вышел") {
t.Errorf("reply = %q, want the title line only", reply)
}
}
func TestQueryFeedsByCategory(t *testing.T) {
h := buildFeedHandler(t, true,
ipc.Note{Text: "Релиз ядра [технологии]", Source: "rss:habr"},
ipc.Note{Text: "Выборы отложены [политика]", Source: "rss:news"},
)
reply, ok := askFeeds(t, h, "что нового по технологиям?")
if !ok {
t.Fatal("not claimed")
}
if !strings.Contains(reply, "ядра") || strings.Contains(reply, "Выборы") {
t.Fatalf("reply = %q, want only the технологии item", reply)
}
reply, _ = askFeeds(t, h, "что нового по спорту?")
if !strings.Contains(reply, "ничего") {
t.Fatalf("reply = %q, want an honest empty answer for an unread category", reply)
}
}
// "не настроены" and "ничего нового" are different truths, and neither may be
// answered by the model inventing a bulletin.
func TestQueryFeedsOffAndEmptyDiffer(t *testing.T) {
off := buildFeedHandler(t, false)
reply, ok := askFeeds(t, off, "что нового?")
if !ok || !strings.Contains(reply, "не настроены") {
t.Fatalf("feeds off: reply = %q, ok = %v", reply, ok)
}
on := buildFeedHandler(t, true)
reply, ok = askFeeds(t, on, "что нового?")
if !ok || !strings.Contains(reply, "ничего нового") {
t.Fatalf("feeds on but empty: reply = %q, ok = %v", reply, ok)
}
}
func TestQueryFeedsPassesOnANonFeedQuestion(t *testing.T) {
h := buildFeedHandler(t, true)
if reply, ok := askFeeds(t, h, "напомни полить цветы"); ok {
t.Fatalf("claimed an unrelated question with %q", reply)
}
}
// The mark is what stops a restart from re-noting yesterday's headlines, so the
// fact round-trip is worth a test of its own.
func TestFactMarksRoundTrip(t *testing.T) {
st := newTestStore(t)
m := &factMarks{api: ipc.NewStoreAPI(st)}
ctx := context.Background()
at, err := m.LastMark(ctx, "habr")
if err != nil || !at.IsZero() {
t.Fatalf("no mark yet: got %v, %v — want zero time and no error", at, err)
}
want := time.Date(2026, 7, 28, 10, 0, 0, 0, time.UTC)
if err := m.SetMark(ctx, "habr", want); err != nil {
t.Fatal(err)
}
got, err := m.LastMark(ctx, "habr")
if err != nil {
t.Fatal(err)
}
if !got.Equal(want) {
t.Fatalf("mark = %v, want %v", got, want)
}
}
// Off unless configured, checked at the wiring seam: no `feeds` block ⇒ no
// worker ⇒ no outbound request is possible.
func TestNewFeedWorkerOffByDefault(t *testing.T) {
st := newTestStore(t)
api := ipc.NewStoreAPI(st)
if w := newFeedWorker(api, nil, &config.Config{}); w != nil {
t.Fatal("a config with no feeds block wired a feed worker")
}
// An empty sources list is normalised to "off" by config.Load; the worker
// refuses it too, so a hand-built Config cannot switch it on by accident.
if w := newFeedWorker(api, nil, &config.Config{Feeds: &config.FeedsConfig{}}); w != nil {
t.Fatal("an empty sources list wired a feed worker")
}
cfg := &config.Config{Feeds: &config.FeedsConfig{Sources: []config.FeedSourceConfig{
{Name: "habr", URL: "https://example.org/rss"},
}}}
w := newFeedWorker(api, nil, cfg)
if w == nil {
t.Fatal("a configured feed did not wire a worker")
}
if got := w.poller.Feeds(); len(got) != 1 || got[0].Name != "habr" {
t.Fatalf("feeds = %+v", got)
}
}
// The fetcher the worker builds must be allowlisted to the configured feeds and
// nothing else — the crawler's SSRF guards are only worth as much as the
// allowlist handed to them.
func TestFeedWorkerFetcherIsAllowlisted(t *testing.T) {
cfg := &config.Config{Feeds: &config.FeedsConfig{Sources: []config.FeedSourceConfig{
{Name: "habr", URL: "https://feeds.example.org/rss"},
}}}
w := newFeedWorker(ipc.NewStoreAPI(newTestStore(t)), nil, cfg)
if w == nil {
t.Fatal("no worker")
}
// PollFeed goes through the guarded fetcher; a feed URL pointing at the box
// itself must fail rather than be read.
_, err := w.poller.PollFeed(context.Background(), rss.FeedConfig{
Name: "evil", URL: "http://127.0.0.1:9100/mcp",
}, time.Now())
if err == nil {
t.Fatal("the poller fetched a private address")
}
}
+17
View File
@@ -170,6 +170,7 @@ func run(args []string) error {
eco *ecosystemWiring
factWorker *factEnrichmentWorker
evalWorker *memoryEvalWorker // nil ⇒ memory evaluation off (the default)
feedWkr *feedWorker // nil ⇒ no feed is read (the default)
)
if !locked {
@@ -265,6 +266,7 @@ func run(args []string) error {
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines), cfg.PatternProposals)
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
evalWorker = newMemoryEvalWorker(st, phr, cfg)
feedWkr = newFeedWorker(ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
coreAPI = &daemonAPI{
CoreAPI: ipc.NewStoreAPI(st),
@@ -451,6 +453,7 @@ func run(args []string) error {
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines), cfg.PatternProposals)
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
evalWorker = newMemoryEvalWorker(st, phr, cfg)
feedWkr = newFeedWorker(ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
// Swap the CoreAPI from the locked placeholder to the real store adapter.
newAPI := &daemonAPI{
@@ -496,6 +499,13 @@ func run(args []string) error {
}()
}
// Start feed reading (nil unless configured).
if feedWkr != nil {
go func() {
feedWkr.run(ctx)
}()
}
dl.unlock()
log.Printf("mavend: unlocked via passkey assertion")
return nil
@@ -541,6 +551,13 @@ func run(args []string) error {
evalWorker.run(ctx)
}()
}
if feedWkr != nil {
wg.Add(1)
go func() {
defer wg.Done()
feedWkr.run(ctx)
}()
}
}
<-ctx.Done()
+5
View File
@@ -82,6 +82,11 @@ type reactiveHandler struct {
replier voice.Replier
now func() time.Time
// 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.
feedsOn bool
weatherProvider weather.Provider
weatherLocation string // default location for weather queries
+1
View File
@@ -211,6 +211,7 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
replier: replier,
phraser: phr,
now: time.Now,
feedsOn: cfg.Feeds != nil,
weatherProvider: weatherProvider,
weatherLocation: weatherLocation,
memStore: memStore,
+332
View File
@@ -0,0 +1,332 @@
// mavmaild — the mail reader module (Vikunja #246,
// docs/plans/01-email-reader.md).
//
// Every so often it opens one IMAP mailbox read-only, fetches the messages it
// has not read yet, and hands each one to core over ipc.MethodIngestMail. Core
// runs the extraction on the resident model and writes what comes back as task
// CANDIDATES he reviews on /tasks. Nothing here writes to the store, nothing
// here can create a reminder, and nothing here speaks.
//
// Why a separate daemon rather than a loop inside mavend, when extraction has
// to happen in mavend anyway: the credential. mavpoll set the precedent with the
// zenmoney token (#125) — the module that talks to a third party holds the
// secret, reads it from a FILE so it never appears in `ps`, in
// docker-compose.yml or in shell history, and core never sees it. Core learns
// that mail exists only as message text on one IPC method; it cannot connect to
// the mailbox even if it wanted to, and a compromised core yields no mail
// password.
//
// Off unless configured: without -password-file there is nothing to run, and
// the daemon says so and exits. If core has no `email` block the very first
// ingest comes back ErrUnknownMethod and this daemon stops polling instead of
// hammering a socket that will keep refusing.
//
// Mail is personal, so the log is counts and UIDs: how many messages were
// fetched, how many were bulk, how many candidates came back. No subject, no
// sender, no body, ever — reviewing a candidate is what /tasks is for.
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"log"
"os"
"os/signal"
"path/filepath"
"sort"
"strings"
"syscall"
"time"
"github.com/kami/maven/internal/email"
"github.com/kami/maven/internal/ipc"
)
func main() {
if err := run(os.Args[1:]); err != nil {
fmt.Fprintln(os.Stderr, "mavmaild:", err)
os.Exit(1)
}
}
func run(args []string) error {
fs := flag.NewFlagSet("mavmaild", flag.ContinueOnError)
socket := fs.String("socket", "", "core IPC socket path (required)")
server := fs.String("imap", "", "IMAP server, host or host:993 (required)")
user := fs.String("user", "", "IMAP username (required)")
passFile := fs.String("password-file", "", "file holding the IMAP password (required — never passed as a flag value)")
mailbox := fs.String("mailbox", "INBOX", "mailbox to read, read-only")
interval := fs.Duration("interval", 15*time.Minute, "how often to read the mailbox")
lookback := fs.Duration("lookback", 72*time.Hour, "how far back to search on each poll")
max := fs.Int("max", 25, "most messages to fetch in one poll")
timeout := fs.Duration("timeout", 30*time.Second, "IMAP network timeout")
statePath := fs.String("state", "", "file remembering which UIDs were read (default: none — every poll re-reads the window)")
if err := fs.Parse(args); err != nil {
return err
}
if *socket == "" {
return fmt.Errorf("-socket is required")
}
if *server == "" || *user == "" || *passFile == "" {
return fmt.Errorf("mail reading is off unless configured: set -imap, -user and -password-file")
}
// The password is read from a file, never taken as a flag value: an argv
// secret is visible in `ps` to every user on the box and lands in the compose
// file and the shell history. Read once at start — a rotated password means a
// restart, which is cheaper than re-reading his credential every quarter hour.
raw, err := os.ReadFile(*passFile)
if err != nil {
return fmt.Errorf("read password file: %w", err)
}
password := strings.TrimSpace(string(raw))
if password == "" {
return fmt.Errorf("password file %s is empty", *passFile)
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
core, err := ipc.DialWait(*socket, 60*time.Second)
if err != nil {
return err
}
defer core.Close()
r := &reader{
core: core,
addr: *server,
user: *user,
mailbox: *mailbox,
lookback: *lookback,
max: *max,
timeout: *timeout,
state: newSeenState(*statePath),
}
if err := r.state.load(); err != nil {
// A missing or corrupt state file must not stop mail from being read: the
// worst case is re-reading the window, and capture dedupes on text.
log.Printf("mavmaild: state: %v (starting from an empty seen-set)", err)
}
// The password is never logged, not even its length.
log.Printf("mavmaild: reading %s on %s every %s (lookback %s, max %d/poll)",
*mailbox, *server, *interval, *lookback, *max)
r.pollOnce(ctx, password) // don't idle a full interval on start
t := time.NewTicker(*interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
log.Printf("mavmaild: bye")
return nil
case <-t.C:
if r.disabled {
// Core told us mail ingestion is not configured. Nothing will change
// without a core restart, and a restart restarts us too.
log.Printf("mavmaild: core does not accept mail — idling")
return nil
}
r.pollOnce(ctx, password)
}
}
}
// mailIngester — the slice of core this daemon uses. One method: hand over a
// message. It cannot write a fact, create a reminder or read the store, and the
// interface says so.
type mailIngester interface {
IngestMail(ctx context.Context, req ipc.IngestMailReq) (ipc.IngestMailResp, error)
}
type reader struct {
core mailIngester
addr string
user string
mailbox string
lookback time.Duration
max int
timeout time.Duration
state *seenState
// dial — connection seam for the tests; nil ⇒ implicit TLS.
dial func(addr string, timeout time.Duration) (*email.Conn, error)
// disabled — core answered ErrUnknownMethod, i.e. it has no email block.
disabled bool
}
// pollOnce — one read of the mailbox, then one ingest per message.
//
// A fetch error aborts this poll and nothing else; the next tick tries again.
// An ingest error for one message does not skip the rest — one mail the model
// choked on should not hide the four behind it.
func (r *reader) pollOnce(ctx context.Context, password string) {
msgs, err := r.fetch(password)
if err != nil {
// The error may name a UID; it never names a subject or a sender.
log.Printf("mavmaild: fetch: %v", err)
if len(msgs) == 0 {
return
}
}
var junk, candidates, created int
for _, m := range msgs {
if ctx.Err() != nil {
return
}
if m.Junk {
junk++
// Marked seen without a model call: the header filter already decided,
// and re-classifying it every quarter hour would be pure waste.
r.state.mark(m.UID)
continue
}
resp, err := r.core.IngestMail(ctx, ipc.IngestMailReq{
Mailbox: r.mailbox,
UID: m.UID,
From: m.From,
Subject: m.Subject,
Date: m.Date,
Body: m.Body,
})
if errors.Is(err, ipc.ErrUnknownMethod) {
log.Printf("mavmaild: core has no email block configured — mail ingestion is off; stopping")
r.disabled = true
return
}
if err != nil {
// Not marked seen: an ingest that failed should be retried next poll.
log.Printf("mavmaild: ingest uid %d: %v", m.UID, err)
continue
}
r.state.mark(m.UID)
candidates += len(resp.TaskIDs)
created += resp.Created
}
if err := r.state.save(); err != nil {
log.Printf("mavmaild: state: %v", err)
}
log.Printf("mavmaild: %s: %d read, %d bulk, %d candidate(s), %d new", r.mailbox, len(msgs), junk, candidates, created)
}
// fetch reads the mailbox. Messages already in the seen-set are not fetched at
// all, so a steady mailbox costs one SEARCH per poll and nothing else.
func (r *reader) fetch(password string) ([]email.Message, error) {
f := email.FetchSince{
Addr: r.addr,
User: r.user,
Mailbox: r.mailbox,
Timeout: r.timeout,
Since: time.Now().Add(-r.lookback),
Max: r.max,
Skip: r.state.seen,
}
return f.RunWith(password, r.dial)
}
// ---- seen state ------------------------------------------------------------
// seenState — the UIDs already handed to core, persisted so a restart does not
// re-read (and re-extract, at multi-second LLM cost) the whole lookback window.
//
// Correctness does not depend on it: ipc.CaptureTask dedupes on normalised text
// among live tasks, so a re-read produces no duplicate rows. This exists to save
// the model's time, which is why a broken state file is a log line rather than a
// failure.
//
// UIDs are per-mailbox and monotonic, so the set is kept as a high-water mark
// plus the stragglers above it. If the server ever changes UIDVALIDITY, UIDs
// reset and the window is simply re-read once — dedupe absorbs it.
type seenState struct {
path string
high uint32
set map[uint32]bool
dirty bool
}
func newSeenState(path string) *seenState {
return &seenState{path: path, set: map[uint32]bool{}}
}
type seenFile struct {
High uint32 `json:"high"`
UIDs []uint32 `json:"uids,omitempty"`
}
func (s *seenState) seen(uid uint32) bool {
return uid <= s.high || s.set[uid]
}
func (s *seenState) mark(uid uint32) {
if s.seen(uid) {
return
}
s.set[uid] = true
s.dirty = true
// Advance the high-water mark through any contiguous run, so the explicit set
// stays small on a mailbox read in order.
for {
next := s.high + 1
if !s.set[next] {
break
}
delete(s.set, next)
s.high = next
}
}
func (s *seenState) load() error {
if s.path == "" {
return nil
}
b, err := os.ReadFile(s.path)
if errors.Is(err, os.ErrNotExist) {
return nil // first run
}
if err != nil {
return err
}
var f seenFile
if err := json.Unmarshal(b, &f); err != nil {
return fmt.Errorf("parse %s: %w", s.path, err)
}
s.high = f.High
for _, u := range f.UIDs {
s.set[u] = true
}
return nil
}
// save writes the state atomically (temp file + rename), 0600: it is a list of
// message ids from his mailbox, which is metadata about his mail.
func (s *seenState) save() error {
if s.path == "" || !s.dirty {
return nil
}
uids := make([]uint32, 0, len(s.set))
for u := range s.set {
uids = append(uids, u)
}
sort.Slice(uids, func(i, j int) bool { return uids[i] < uids[j] })
b, err := json.Marshal(seenFile{High: s.high, UIDs: uids})
if err != nil {
return err
}
tmp := s.path + ".tmp"
if err := os.MkdirAll(filepath.Dir(s.path), 0o700); err != nil {
return err
}
if err := os.WriteFile(tmp, b, 0o600); err != nil {
return err
}
if err := os.Rename(tmp, s.path); err != nil {
return err
}
s.dirty = false
return nil
}
+254
View File
@@ -0,0 +1,254 @@
package main
import (
"bufio"
"context"
"fmt"
"net"
"os"
"path/filepath"
"strconv"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/email"
"github.com/kami/maven/internal/ipc"
)
// ---- a scripted IMAP server, same shape internal/email's tests use ---------
type fakeIMAP struct {
msgs map[uint32]string
uids []uint32
cmds []string
}
func (f *fakeIMAP) serve(c net.Conn) {
defer c.Close()
fmt.Fprint(c, "* OK fake ready\r\n")
r := bufio.NewReader(c)
for {
line, err := r.ReadString('\n')
if err != nil {
return
}
parts := strings.SplitN(strings.TrimRight(line, "\r\n"), " ", 2)
if len(parts) != 2 {
return
}
tag, cmd := parts[0], parts[1]
f.cmds = append(f.cmds, cmd)
upper := strings.ToUpper(cmd)
switch {
case strings.HasPrefix(upper, "LOGIN"), strings.HasPrefix(upper, "EXAMINE"):
fmt.Fprintf(c, "%s OK\r\n", tag)
case strings.HasPrefix(upper, "UID SEARCH"):
var ids []string
for _, u := range f.uids {
ids = append(ids, strconv.FormatUint(uint64(u), 10))
}
fmt.Fprintf(c, "* SEARCH %s\r\n%s OK\r\n", strings.Join(ids, " "), tag)
case strings.HasPrefix(upper, "UID FETCH"):
uid, _ := strconv.ParseUint(strings.Fields(cmd)[2], 10, 32)
if raw, ok := f.msgs[uint32(uid)]; ok {
fmt.Fprintf(c, "* 1 FETCH (UID %d BODY[] {%d}\r\n%s)\r\n", uid, len(raw), raw)
}
fmt.Fprintf(c, "%s OK\r\n", tag)
case strings.HasPrefix(upper, "LOGOUT"):
fmt.Fprintf(c, "* BYE\r\n%s OK\r\n", tag)
return
default:
fmt.Fprintf(c, "%s BAD\r\n", tag)
}
}
}
func (f *fakeIMAP) dial(_ string, timeout time.Duration) (*email.Conn, error) {
cli, srv := net.Pipe()
go f.serve(srv)
return email.NewConn(cli, timeout)
}
// ---- a fake core -----------------------------------------------------------
type fakeCore struct {
got []ipc.IngestMailReq
resp ipc.IngestMailResp
err error
}
func (c *fakeCore) IngestMail(_ context.Context, req ipc.IngestMailReq) (ipc.IngestMailResp, error) {
c.got = append(c.got, req)
if c.err != nil {
return ipc.IngestMailResp{}, c.err
}
return c.resp, nil
}
func mail(subject, body string, extraHeaders ...string) string {
h := "Subject: " + subject + "\r\nFrom: a@b.c\r\nContent-Type: text/plain; charset=utf-8\r\n"
for _, e := range extraHeaders {
h += e + "\r\n"
}
return h + "\r\n" + body + "\r\n"
}
func newTestReader(t *testing.T, f *fakeIMAP, core *fakeCore, statePath string) *reader {
t.Helper()
return &reader{
core: core, addr: "mail.example:993", user: "kami", mailbox: "INBOX",
lookback: 72 * time.Hour, max: 25, timeout: 5 * time.Second,
state: newSeenState(statePath),
dial: f.dial,
}
}
func TestPollHandsMessagesToCore(t *testing.T) {
f := &fakeIMAP{
uids: []uint32{1, 2},
msgs: map[uint32]string{
1: mail("Счёт", "Оплатить до 5 августа."),
2: mail("Скидки", "Sale!", "List-Unsubscribe: <mailto:u@x>"),
},
}
core := &fakeCore{resp: ipc.IngestMailResp{TaskIDs: []int64{1}, Created: 1}}
r := newTestReader(t, f, core, "")
r.pollOnce(context.Background(), "secret")
// The newsletter is filtered before core is asked: only the real mail crosses.
if len(core.got) != 1 {
t.Fatalf("core saw %d messages, want 1 (the bulk one must not cross): %+v", len(core.got), core.got)
}
got := core.got[0]
if got.UID != 1 || got.Mailbox != "INBOX" || got.Subject != "Счёт" {
t.Errorf("ingest req = %+v", got)
}
if !strings.Contains(got.Body, "Оплатить") {
t.Errorf("body = %q", got.Body)
}
}
// A second poll must not re-send what core already saw — extraction is a
// multi-second LLM call per message.
func TestPollSkipsSeenUIDs(t *testing.T) {
f := &fakeIMAP{uids: []uint32{5}, msgs: map[uint32]string{5: mail("Счёт", "текст")}}
core := &fakeCore{}
r := newTestReader(t, f, core, "")
r.pollOnce(context.Background(), "secret")
r.pollOnce(context.Background(), "secret")
if len(core.got) != 1 {
t.Errorf("core saw %d messages over two polls, want 1", len(core.got))
}
}
// An ingest that failed is NOT marked seen: the next poll retries it.
func TestPollRetriesFailedIngest(t *testing.T) {
f := &fakeIMAP{uids: []uint32{5}, msgs: map[uint32]string{5: mail("Счёт", "текст")}}
core := &fakeCore{err: fmt.Errorf("llama-server is warming up")}
r := newTestReader(t, f, core, "")
r.pollOnce(context.Background(), "secret")
core.err = nil
r.pollOnce(context.Background(), "secret")
if len(core.got) != 2 {
t.Errorf("core saw %d attempts, want 2 (a failed ingest is retried)", len(core.got))
}
}
// Core without an email block ⇒ stop, don't hammer the socket.
func TestPollStopsWhenCoreRefusesMail(t *testing.T) {
f := &fakeIMAP{uids: []uint32{1, 2}, msgs: map[uint32]string{1: mail("a", "b"), 2: mail("c", "d")}}
core := &fakeCore{err: fmt.Errorf("call: %w", ipc.ErrUnknownMethod)}
r := newTestReader(t, f, core, "")
r.pollOnce(context.Background(), "secret")
if !r.disabled {
t.Error("ErrUnknownMethod must disable the reader")
}
if len(core.got) != 1 {
t.Errorf("core saw %d messages, want 1 — stop at the first refusal", len(core.got))
}
}
func TestSeenStatePersists(t *testing.T) {
path := filepath.Join(t.TempDir(), "state", "seen.json")
f := &fakeIMAP{uids: []uint32{9}, msgs: map[uint32]string{9: mail("Счёт", "текст")}}
core := &fakeCore{}
r := newTestReader(t, f, core, path)
r.pollOnce(context.Background(), "secret")
fi, err := os.Stat(path)
if err != nil {
t.Fatalf("state file: %v", err)
}
// A list of message ids from his mailbox is metadata about his mail.
if perm := fi.Mode().Perm(); perm != 0o600 {
t.Errorf("state file mode = %v, want 0600", perm)
}
// A fresh reader with the same state file must not re-read the message.
core2 := &fakeCore{}
r2 := newTestReader(t, f, core2, path)
if err := r2.state.load(); err != nil {
t.Fatalf("load: %v", err)
}
r2.pollOnce(context.Background(), "secret")
if len(core2.got) != 0 {
t.Errorf("after a restart core saw %d messages, want 0", len(core2.got))
}
}
func TestSeenStateHighWaterMark(t *testing.T) {
s := newSeenState("")
s.mark(1)
s.mark(3)
s.mark(2)
if s.high != 3 {
t.Errorf("high = %d, want 3 (contiguous run collapses)", s.high)
}
if len(s.set) != 0 {
t.Errorf("explicit set = %v, want empty", s.set)
}
if !s.seen(2) || s.seen(4) {
t.Errorf("seen(2)=%v seen(4)=%v", s.seen(2), s.seen(4))
}
}
func TestSeenStateCorruptFileIsNotFatal(t *testing.T) {
path := filepath.Join(t.TempDir(), "seen.json")
if err := os.WriteFile(path, []byte("{not json"), 0o600); err != nil {
t.Fatal(err)
}
s := newSeenState(path)
if err := s.load(); err == nil {
t.Error("a corrupt state file should report an error the caller logs")
}
if s.seen(1) {
t.Error("a corrupt state file must leave an empty seen-set, not a poisoned one")
}
}
// Off unless configured, and the credential is never a flag value.
func TestRunRequiresConfig(t *testing.T) {
if err := run([]string{}); err == nil {
t.Error("no -socket must be an error")
}
if err := run([]string{"-socket", "/tmp/nope.sock"}); err == nil {
t.Error("no mailbox configuration must be an error, not a default mailbox")
}
// There is no -password flag at all: only -password-file.
if err := run([]string{"-socket", "/x", "-imap", "h", "-user", "u", "-password", "p"}); err == nil ||
!strings.Contains(err.Error(), "flag provided but not defined") {
t.Errorf("a -password flag must not exist; err = %v", err)
}
}
func TestRunRejectsEmptyPasswordFile(t *testing.T) {
path := filepath.Join(t.TempDir(), "pass")
if err := os.WriteFile(path, []byte(" \n"), 0o600); err != nil {
t.Fatal(err)
}
err := run([]string{"-socket", "/x/y.sock", "-imap", "h", "-user", "u", "-password-file", path})
if err == nil || !strings.Contains(err.Error(), "empty") {
t.Errorf("an empty password file must be refused before dialling; err = %v", err)
}
}
+31
View File
@@ -36,6 +36,37 @@ present (see `.dockerignore`).
| `/var/lib/maven` (volume) | encrypted db at rest |
| `/dev/shm` (tmpfs) | decrypted db working copy (RAM only) |
## Reading the outside world (off by default)
`mavend.json` ships without a `feeds` block, which means no RSS/Atom feed is
fetched and no outbound request is made. Switching it on is adding the block:
```json
"feeds": {
"poll_interval": "30m",
"max_items": 5,
"max_age": "24h",
"sources": [
{ "name": "habr", "url": "https://habr.com/ru/rss/best/daily/",
"category": "технологии", "exclude": ["реклама"] }
]
}
```
What it does and does not do:
- items are written as notes with source `rss:<name>`, visible on `/dash`;
- **nothing is announced.** She reads them back when asked — "что нового в
лентах?", "что нового по технологиям?" — and never on arrival. There is no
severity or channel knob here on purpose;
- the fetcher is allowlisted to the hosts of the configured feeds, plus any
`allow_hosts`. It refuses non-http(s) schemes and every private address
(loopback, the LAN, the `10.42.0.0/24` wg range, cloud metadata). It caps the
response at 2 MiB and redirects at 3, and makes at most one request per host
per second. See `internal/webfetch`;
- 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.
## Not yet verified / host-dependent
This stack is correct-by-construction but has **not been build-tested here**
+25
View File
@@ -113,6 +113,31 @@ services:
- sockets:/run/maven
# - ./deploy/zenmoney.token:/run/secrets/zenmoney.token:ro
# The mail reader (Vikunja #246) is OFF and commented out: it needs an IMAP
# account, and there is none on this box. mavmaild reads the password from a
# FILE so it never appears in `ps`, in this file, or in shell history — the
# same rule mavpoll follows for the zenmoney token. Core never sees the
# password: the reader hands core message text on one IPC method, and core
# writes what the model extracts as task CANDIDATES he reviews on /tasks.
# Nothing here can create a reminder, so a misread mail cannot fire.
#
# To enable: write the password to deploy/imap.password (0600, gitignored),
# add an "email": {} block to deploy/mavend.json, and uncomment this service.
# mavmaild:
# <<: *image
# command: ["mavmaild", "-socket", "/run/maven/mavend.sock",
# "-imap", "imap.example.org:993",
# "-user", "kami@example.org",
# "-password-file", "/run/secrets/imap.password",
# "-mailbox", "INBOX",
# "-interval", "15m",
# "-state", "/var/lib/maven/mail-seen.json"]
# depends_on: [mavend]
# volumes:
# - sockets:/run/maven
# - dbdata:/var/lib/maven
# - ./deploy/imap.password:/run/secrets/imap.password:ro
volumes:
dbdata:
sockets:
+25
View File
@@ -27,3 +27,28 @@
7. Add voice query handler — `"что нового?"` queries `RecentNotes` filtered by source prefix `rss:` and phrases via `phraser.PhraseQuery`
8. Add `feeds` block to `config.Config` and `deploy/mavend.json`
9. Test with a live RSS feed (e.g., `https://news.ycombinator.com/rss`) — verify items appear in notes table
---
## Shipped 2026-08-01 (#258)
`internal/webfetch` (the guarded HTTP door: scheme, allow/deny hosts, private-address
refusal in the dialer, size cap, redirect cap, per-host rate limit), `internal/rss`
(RSS 2.0 + Atom parser, poller with durable marks and a keyword filter),
`cmd/mavend/feeds.go` (ticker, fetcher adapter, `rss:latest:<feed>` fact marks),
config block `feeds`, and the `feeds` query source with `router.ParseFeedQuery`.
Deviations from the plan above, both deliberate:
- **Step 5 (breaking-news nudges) was not built.** A feed that dispatches is a nag,
and the one thing Maven is not is a nag. Items are read when asked and nowhere else.
If breaking news is ever wanted, it belongs behind the existing delivery policy
(severity, quiet hours, digest), not in the poller.
- **Step 3 (embedder relevance) is a seam, not an implementation.** `rss.Ranker`
exists and is wired nil. Scoring items against an "interest profile" needs a
profile, and there is none yet; a threshold with nothing to compare against is a
random filter with a confident name. The filter that runs is the per-feed
include/exclude keyword list, which he can read and predict.
No new dependency: stdlib `encoding/xml`, no gofeed. Stock deploy config has no
`feeds` block, so the capability is off.
+58
View File
@@ -156,6 +156,11 @@ type Config struct {
// live in the reader (cmd/mavmaild), never here.
Email *EmailConfig `json:"email,omitempty"`
// Feeds — RSS/Atom feed reading (Vikunja #258). nil / absent ⇒ no feed is
// ever fetched: reading the outside world is off unless configured, like
// the weather and telegram. See FeedsConfig.
Feeds *FeedsConfig `json:"feeds,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
@@ -402,6 +407,53 @@ func (p *PatternProposalConfig) AnnounceProposals() bool {
return p != nil && p.Notify
}
// FeedsConfig — the RSS/Atom reader (Vikunja #258, docs/plans/13-rss-news-feeds.md).
//
// Absent ⇒ off, and off means no outbound request at all. Present with an empty
// `sources` list is also off — a poller with nothing to poll is not wired.
//
// What a feed may NOT do here: speak. Items are written as notes with source
// "rss:<name>" and read back when he asks; nothing is dispatched, nudged or
// announced on arrival. That is the "not a nag" constraint, and it is why there
// is no severity or channel field in this block to reach for.
type FeedsConfig struct {
// Sources — the feeds to read. Empty ⇒ the reader stays down.
Sources []FeedSourceConfig `json:"sources,omitempty"`
// PollInterval — default per-feed cadence. 0 ⇒ rss.DefaultPollInterval (30m).
PollInterval Duration `json:"poll_interval,omitempty"`
// MaxItems — most items kept from one feed in one poll. 0 ⇒
// rss.DefaultMaxItems (5). This is the "не завали мне /dash" knob.
MaxItems int `json:"max_items,omitempty"`
// MaxAge — on a first poll (no saved mark), how far back to take items.
// 0 ⇒ rss.DefaultMaxAge (24h), so switching a feed on imports today, not
// the archive.
MaxAge Duration `json:"max_age,omitempty"`
// AllowHosts — when set, the reader may only connect to these hosts (and
// their subdomains). The feed URLs' own hosts are added automatically, so
// this is only needed to be stricter than that.
AllowHosts []string `json:"allow_hosts,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"`
}
// FeedSourceConfig — one feed.
type FeedSourceConfig struct {
Name string `json:"name"` // note source is "rss:<name>"
URL string `json:"url"` // http(s) only
Category string `json:"category,omitempty"` // "технологии" — what "что нового по X?" matches
Interval Duration `json:"interval,omitempty"` // 0 ⇒ FeedsConfig.PollInterval
Include []string `json:"include,omitempty"` // keep only items containing one of these
Exclude []string `json:"exclude,omitempty"` // drop items containing any of these
}
// 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
@@ -639,6 +691,12 @@ func (c *Config) applyDefaults() {
c.Email.Timeout = Duration(DefaultEmailTimeout)
}
// A feeds block with no sources is the same as no block: nothing to poll,
// nothing wired. Normalising it to nil keeps that "off" in one place.
if c.Feeds != nil && len(c.Feeds.Sources) == 0 {
c.Feeds = nil
}
if c.Voice != nil {
if c.Voice.RouterThreshold <= 0 {
c.Voice.RouterThreshold = DefaultRouterThreshold
+8 -6
View File
@@ -26,21 +26,23 @@ type FetchSince struct {
Since time.Time
Max int
Skip func(uid uint32) bool
// dial is the connection seam. nil means Dial (implicit TLS); the tests set
// it to an in-process fake. Unexported so no configuration path can point
// the reader at a non-TLS transport.
dial func(addr string, timeout time.Duration) (*Conn, error)
}
// Run performs one read. password is passed here, not stored in the struct, so
// the configuration of a mailbox and the secret for it are never the same value
// sitting in the same place.
func (f FetchSince) Run(password string) ([]Message, error) {
return f.RunWith(password, nil)
}
// RunWith is Run with an explicit connection function, which is how the reader
// daemon and the tests substitute an in-process server. nil ⇒ Dial, i.e.
// implicit TLS with certificate verification; there is no configuration path
// that reaches this, so no deployment can end up talking cleartext IMAP.
func (f FetchSince) RunWith(password string, dial func(addr string, timeout time.Duration) (*Conn, error)) ([]Message, error) {
if f.Addr == "" || f.User == "" || f.Mailbox == "" {
return nil, fmt.Errorf("email: mailbox not configured (addr/user/mailbox)")
}
dial := f.dial
if dial == nil {
dial = Dial
}
+5 -6
View File
@@ -21,13 +21,12 @@ func TestFetchSinceRun(t *testing.T) {
Since: time.Date(2026, 7, 30, 0, 0, 0, 0, time.UTC),
Max: 2,
Skip: func(uid uint32) bool { return uid == 3 },
dial: func(addr string, timeout time.Duration) (*Conn, error) {
cli, srv := net.Pipe()
go f.serve(t, srv)
return NewConn(cli, timeout)
},
}
msgs, err := fs.Run("secret")
msgs, err := fs.RunWith("secret", func(addr string, timeout time.Duration) (*Conn, error) {
cli, srv := net.Pipe()
go f.serve(t, srv)
return NewConn(cli, timeout)
})
if err != nil {
t.Fatalf("run: %v", err)
}
+109
View File
@@ -0,0 +1,109 @@
package router
import "strings"
// Feed queries — "что нового в лентах?", "что нового по технологиям?"
// (Vikunja #258).
//
// Deterministic matching, like the calendar, plan and habit matchers above it:
// the LLM router says this is a query, and this decides whether it is a question
// about the feeds. A model deciding that would occasionally answer "что нового?"
// out of world knowledge, which is the one thing a feed reader exists to avoid.
// FeedQuery — a parsed "what's new" question. Category is the topic he named
// ("технологии"), empty when he asked about the feeds in general.
type FeedQuery struct {
Category string
}
// feedNouns — the words that make a question about the feeds themselves.
var feedNouns = []string{
"лента", "ленте", "ленты", "лентах", "лентам",
"новости", "новостей", "новостях", "новостям",
"новое", "нового", "новенького",
"feed", "feeds", "news", "headlines",
}
// newnessMarkers — the "что нового" half. "нового" alone is in feedNouns
// because it carries the question on its own ("что нового?"); a bare "лента"
// needs the ask, which is what askMarkers below is for.
var askMarkers = []string{
"что", "какие", "какое", "расскажи", "почитай", "прочитай", "покажи",
"what", "any", "tell", "show", "read",
}
// ParseFeedQuery reports whether an utterance asks what is new in the feeds, and
// which topic if it names one after "по"/"о"/"про"/"about".
//
// Both a feed noun and an ask are required. "у меня новая лента в инстаграме" is
// a statement and must not be read as a request to recite headlines.
func ParseFeedQuery(text string) (FeedQuery, bool) {
toks := planTokens(text)
noun, ask := false, false
for _, t := range toks {
for _, n := range feedNouns {
if t == n {
noun = true
break
}
}
for _, a := range askMarkers {
if t == a {
ask = true
break
}
}
}
if !noun || !ask {
return FeedQuery{}, false
}
return FeedQuery{Category: feedCategory(toks)}, true
}
// categoryPreps — the prepositions a topic follows. Russian marks the topic with
// a preposition ("по технологиям", "про политику"), so the word after one is the
// category; there is no stemming here, and the match against the configured
// category is a prefix comparison for exactly that reason.
var categoryPreps = map[string]bool{"по": true, "о": true, "об": true, "про": true, "about": true, "on": true}
func feedCategory(toks []string) string {
for i, t := range toks {
if categoryPreps[t] && i+1 < len(toks) {
next := toks[i+1]
// "по новостям" names no topic, it repeats the noun.
for _, n := range feedNouns {
if next == n {
return ""
}
}
return next
}
}
return ""
}
// CategoryMatches reports whether a note's text plausibly belongs to the
// category he named. Russian inflects the topic ("технологиям" vs the configured
// "технологии"), and there is no stemmer in this repo, so the comparison is on a
// common prefix — long enough that "полит" and "погод" stay apart, short enough
// to survive a case ending.
func CategoryMatches(text, category string) bool {
if category == "" {
return true
}
stem := categoryStem(category)
if stem == "" {
return false
}
return strings.Contains(strings.ToLower(text), stem)
}
// categoryStem cuts a word down to the part inflection leaves alone. 5 runes is
// the compromise: shorter words are used whole.
func categoryStem(word string) string {
r := []rune(strings.ToLower(strings.TrimSpace(word)))
if len(r) > 5 {
r = r[:5]
}
return string(r)
}
+48
View File
@@ -0,0 +1,48 @@
package router
import "testing"
func TestParseFeedQuery(t *testing.T) {
cases := []struct {
text string
ok bool
category string
}{
{"что нового в лентах?", true, ""},
{"что нового?", true, ""},
{"какие новости?", true, ""},
{"что нового по технологиям?", true, "технологиям"},
{"расскажи новости про политику", true, "политику"},
{"что нового по новостям", true, ""},
{"what's new in the feeds?", true, ""},
{"any news about kubernetes", true, "kubernetes"},
// Statements, not requests.
{"у меня новая лента в инстаграме", false, ""},
{"новости меня утомили", false, ""},
{"напомни полить цветы", false, ""},
{"", false, ""},
}
for _, c := range cases {
q, ok := ParseFeedQuery(c.text)
if ok != c.ok {
t.Errorf("ParseFeedQuery(%q) ok = %v, want %v", c.text, ok, c.ok)
continue
}
if ok && q.Category != c.category {
t.Errorf("ParseFeedQuery(%q) category = %q, want %q", c.text, q.Category, c.category)
}
}
}
func TestCategoryMatches(t *testing.T) {
// The inflected form he says must match the form the config spells.
if !CategoryMatches("Новый релиз [технологии]", "технологиям") {
t.Error("inflected category did not match")
}
if CategoryMatches("Новый релиз [технологии]", "политику") {
t.Error("unrelated category matched")
}
if !CategoryMatches("anything", "") {
t.Error("an empty category must match everything")
}
}
+179
View File
@@ -0,0 +1,179 @@
// Package rss reads RSS 2.0 and Atom feeds, and does nothing else with them.
//
// Parsing and polling are split from delivery on purpose: a feed is a source
// Maven can be ASKED about, not a thing that speaks. Nothing in this package
// dispatches, nudges or notifies — the poller writes notes, and the answer path
// reads them when he asks "что нового в лентах?". "Not a nag" is the oldest
// constraint in the spec, and a news feed is the single most tempting way to
// break it.
//
// Stdlib only (encoding/xml). Feeds are XML from strangers, so the parser takes
// what it recognises and ignores the rest rather than failing a whole feed over
// one malformed item.
package rss
import (
"encoding/xml"
"fmt"
"html"
"io"
"regexp"
"strings"
"time"
)
// Item is one feed entry, normalised across RSS and Atom.
type Item struct {
Title string
Link string
Summary string // plain text, tags stripped, entities decoded
Published time.Time // zero when the feed did not say
ID string // guid / atom id, falling back to the link
}
// Feed is a parsed document.
type Feed struct {
Title string
Items []Item
}
// feedDoc covers both dialects in one struct. RSS puts items under
// channel>item, Atom puts entries at the top level, and the field names barely
// overlap — so both sets are declared and whichever the document filled in wins.
type feedDoc struct {
ChannelTitle string `xml:"channel>title"`
AtomTitle string `xml:"title"`
Items []struct {
Title string `xml:"title"`
Link string `xml:"link"`
Description string `xml:"description"`
Encoded string `xml:"encoded"` // content:encoded
GUID string `xml:"guid"`
PubDate string `xml:"pubDate"`
Date string `xml:"date"` // dc:date
} `xml:"channel>item"`
Entries []struct {
Title string `xml:"title"`
Links []struct {
Href string `xml:"href,attr"`
Rel string `xml:"rel,attr"`
} `xml:"link"`
Summary string `xml:"summary"`
Content string `xml:"content"`
ID string `xml:"id"`
Updated string `xml:"updated"`
Published string `xml:"published"`
} `xml:"entry"`
}
// Parse reads a feed document.
func Parse(r io.Reader) (Feed, error) {
var doc feedDoc
dec := xml.NewDecoder(r)
// Feeds in the wild declare windows-1251 and worse. We only ever read
// UTF-8; a charset we cannot decode is a feed we do not read, which is
// better than mojibake in his notes.
dec.Strict = false
if err := dec.Decode(&doc); err != nil {
return Feed{}, fmt.Errorf("rss: bad xml: %w", err)
}
f := Feed{Title: strings.TrimSpace(doc.ChannelTitle)}
if f.Title == "" {
f.Title = strings.TrimSpace(doc.AtomTitle)
}
for _, it := range doc.Items {
item := Item{
Title: PlainText(it.Title),
Link: strings.TrimSpace(it.Link),
Summary: PlainText(firstNonEmpty(it.Description, it.Encoded)),
Published: parseTime(firstNonEmpty(it.PubDate, it.Date)),
ID: strings.TrimSpace(firstNonEmpty(it.GUID, it.Link)),
}
if item.Title != "" || item.Link != "" {
f.Items = append(f.Items, item)
}
}
for _, e := range doc.Entries {
link := ""
for _, l := range e.Links {
if l.Rel == "" || l.Rel == "alternate" {
link = strings.TrimSpace(l.Href)
break
}
}
if link == "" && len(e.Links) > 0 {
link = strings.TrimSpace(e.Links[0].Href)
}
item := Item{
Title: PlainText(e.Title),
Link: link,
Summary: PlainText(firstNonEmpty(e.Summary, e.Content)),
Published: parseTime(firstNonEmpty(e.Published, e.Updated)),
ID: strings.TrimSpace(firstNonEmpty(e.ID, link)),
}
if item.Title != "" || item.Link != "" {
f.Items = append(f.Items, item)
}
}
return f, nil
}
func firstNonEmpty(vals ...string) string {
for _, v := range vals {
if strings.TrimSpace(v) != "" {
return v
}
}
return ""
}
// timeLayouts — RFC1123/822 for RSS, RFC3339 for Atom, plus the near-misses
// real feeds ship (no seconds, numeric zone where a name is expected).
var timeLayouts = []string{
time.RFC1123Z,
time.RFC1123,
time.RFC822Z,
time.RFC822,
time.RFC3339,
"2006-01-02T15:04:05Z0700",
"2006-01-02 15:04:05",
"2006-01-02",
"Mon, 02 Jan 2006 15:04:05 -0700",
"Mon, 2 Jan 2006 15:04:05 -0700",
"Mon, 2 Jan 2006 15:04:05 MST",
}
// parseTime returns the zero time on anything it cannot read. An undated item
// is still an item; the poller dedupes by ID, so a missing date costs nothing.
func parseTime(s string) time.Time {
s = strings.TrimSpace(s)
if s == "" {
return time.Time{}
}
for _, l := range timeLayouts {
if t, err := time.Parse(l, s); err == nil {
return t.UTC()
}
}
return time.Time{}
}
var (
// RE2 has no backreferences, so the two tags are spelled out rather than
// captured and matched against themselves.
scriptRE = regexp.MustCompile(`(?is)<script\b[^>]*>.*?</script>|<style\b[^>]*>.*?</style>`)
tagRE = regexp.MustCompile(`(?s)<[^>]*>`)
)
// PlainText strips markup and decodes entities — feed summaries are HTML, and
// what reaches a note (and possibly the TTS) must be text. Exported because the
// crawler's extractor needs exactly this on a bigger input.
func PlainText(s string) string {
s = scriptRE.ReplaceAllString(s, " ")
s = tagRE.ReplaceAllString(s, " ")
s = html.UnescapeString(s)
return strings.TrimSpace(strings.Join(strings.Fields(s), " "))
}
+112
View File
@@ -0,0 +1,112 @@
package rss
import (
"strings"
"testing"
"time"
)
const rss2 = `<?xml version="1.0"?>
<rss version="2.0">
<channel>
<title>Хабр</title>
<item>
<title>Новая уязвимость в ядре</title>
<link>https://example.org/a</link>
<description>&lt;p&gt;Патч уже &lt;b&gt;вышел&lt;/b&gt;.&lt;/p&gt;</description>
<guid>tag:example.org,a</guid>
<pubDate>Mon, 28 Jul 2026 10:00:00 +0000</pubDate>
</item>
<item>
<title>Без даты</title>
<link>https://example.org/b</link>
</item>
</channel>
</rss>`
const atom = `<?xml version="1.0" encoding="utf-8"?>
<feed xmlns="http://www.w3.org/2005/Atom">
<title>Example Atom</title>
<entry>
<title>Release 2.0</title>
<link rel="alternate" href="https://example.com/rel"/>
<link rel="edit" href="https://example.com/edit"/>
<id>urn:uuid:1</id>
<updated>2026-07-30T12:30:00Z</updated>
<summary>Ships &amp; works</summary>
</entry>
</feed>`
func TestParseRSS2(t *testing.T) {
f, err := Parse(strings.NewReader(rss2))
if err != nil {
t.Fatal(err)
}
if f.Title != "Хабр" {
t.Fatalf("title = %q", f.Title)
}
if len(f.Items) != 2 {
t.Fatalf("items = %d, want 2", len(f.Items))
}
it := f.Items[0]
if it.Title != "Новая уязвимость в ядре" {
t.Errorf("title = %q", it.Title)
}
if it.Summary != "Патч уже вышел ." && it.Summary != "Патч уже вышел." {
t.Errorf("summary = %q — tags must be stripped and entities decoded", it.Summary)
}
if it.ID != "tag:example.org,a" {
t.Errorf("id = %q", it.ID)
}
if want := time.Date(2026, 7, 28, 10, 0, 0, 0, time.UTC); !it.Published.Equal(want) {
t.Errorf("published = %v, want %v", it.Published, want)
}
if !f.Items[1].Published.IsZero() {
t.Errorf("undated item got a date: %v", f.Items[1].Published)
}
if f.Items[1].ID != "https://example.org/b" {
t.Errorf("id falls back to the link, got %q", f.Items[1].ID)
}
}
func TestParseAtom(t *testing.T) {
f, err := Parse(strings.NewReader(atom))
if err != nil {
t.Fatal(err)
}
if f.Title != "Example Atom" || len(f.Items) != 1 {
t.Fatalf("feed = %+v", f)
}
it := f.Items[0]
if it.Link != "https://example.com/rel" {
t.Errorf("link = %q, want the alternate link", it.Link)
}
if it.Summary != "Ships & works" {
t.Errorf("summary = %q", it.Summary)
}
if want := time.Date(2026, 7, 30, 12, 30, 0, 0, time.UTC); !it.Published.Equal(want) {
t.Errorf("published = %v, want %v", it.Published, want)
}
}
func TestParseGarbage(t *testing.T) {
if _, err := Parse(strings.NewReader("<html><body>not a feed")); err == nil {
t.Fatal("want an error on a non-feed document")
}
// A feed with an item that has neither title nor link contributes nothing
// rather than an empty note.
f, err := Parse(strings.NewReader(`<rss><channel><item><description>x</description></item></channel></rss>`))
if err != nil {
t.Fatal(err)
}
if len(f.Items) != 0 {
t.Fatalf("items = %d, want 0", len(f.Items))
}
}
func TestPlainTextDropsScript(t *testing.T) {
got := PlainText(`<p>hi</p><script>alert("x")</script><style>b{}</style> there`)
if got != "hi there" {
t.Fatalf("got %q", got)
}
}
+324
View File
@@ -0,0 +1,324 @@
package rss
import (
"context"
"fmt"
"log"
"strings"
"time"
)
// FeedConfig — one feed to read. A feed with no Name or no URL is ignored.
type FeedConfig struct {
Name string // short id; the note source is "rss:<Name>"
URL string // http(s) only, enforced by the fetcher
Category string // free text ("технологии"), used to answer "что по X?"
Interval time.Duration // 0 ⇒ the poller's default
Include []string // when non-empty, keep only items matching one of these
Exclude []string // drop items matching any of these, even if included
}
// Fetcher is the guarded HTTP door (internal/webfetch). An interface so the
// poller is testable without a network and so it CANNOT fetch by any other
// means: no http.Client is constructed in this package.
type Fetcher interface {
Get(ctx context.Context, url string) (*Body, error)
}
// Body is the minimum the poller needs from a response.
type Body struct{ Bytes []byte }
// Notes is core's note-writing half. Same shape as ipc.CoreAPI's method, so the
// daemon passes its API straight in.
type Notes interface {
WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error)
}
// Marks remembers how far a feed was read. Durable, because the alternative is
// re-writing yesterday's headlines as fresh notes after every restart. The
// daemon backs this with config facts (key "rss:latest:<feed>").
type Marks interface {
LastMark(ctx context.Context, feed string) (time.Time, error)
SetMark(ctx context.Context, feed string, at time.Time) error
}
// Embedder embeds a note on its way into the store so recall can find it. nil ⇒
// notes are written without a vector (still readable by the recent-notes path).
type Embedder interface {
Embed(ctx context.Context, text string) ([]float32, error)
}
// Ranker is the relevance seam. The plan called for scoring each item against
// an interest profile built from his notes; that profile does not exist yet, and
// a threshold over an embedder with no profile to compare to is a random filter
// with a confident name. So the seam is here, nil in the daemon, and the filter
// that actually runs is the per-feed keyword one — a rule he can read and
// predict. When there IS a profile, implement this and pass it.
//
// Note what a Ranker must NOT be: anything that sends his notes outward. The
// scoring happens locally against a local embedder; the feed item is the input,
// his memory is never the payload.
type Ranker interface {
Relevant(ctx context.Context, text string) (bool, error)
}
// Config — poller-wide settings.
type Config struct {
DefaultInterval time.Duration // 0 ⇒ DefaultPollInterval
MaxItems int // most notes written per feed per poll; 0 ⇒ DefaultMaxItems
MaxAge time.Duration // ignore items older than this on a cold start; 0 ⇒ DefaultMaxAge
}
// Defaults chosen to be quiet: a feed read every half hour, at most a handful of
// items kept, and a cold start that does not import a month of history.
const (
DefaultPollInterval = 30 * time.Minute
DefaultMaxItems = 5
DefaultMaxAge = 24 * time.Hour
)
// Poller reads feeds on a schedule and writes what survives filtering as notes.
type Poller struct {
feeds []FeedConfig
fetch Fetcher
notes Notes
marks Marks
embed Embedder
ranker Ranker
cfg Config
nextDue map[string]time.Time
seen map[string]map[string]bool // feed → item ID, for items with no date
}
// NewPoller wires a poller. Returns nil when there is nothing to poll — a
// capability is off unless configured, and callers check for nil.
func NewPoller(feeds []FeedConfig, fetch Fetcher, notes Notes, marks Marks, embed Embedder, ranker Ranker, cfg Config) *Poller {
var valid []FeedConfig
for _, f := range feeds {
if strings.TrimSpace(f.Name) == "" || strings.TrimSpace(f.URL) == "" {
log.Printf("rss: skipping a feed with no name or no url")
continue
}
valid = append(valid, f)
}
if len(valid) == 0 || fetch == nil || notes == nil {
return nil
}
if cfg.DefaultInterval <= 0 {
cfg.DefaultInterval = DefaultPollInterval
}
if cfg.MaxItems <= 0 {
cfg.MaxItems = DefaultMaxItems
}
if cfg.MaxAge <= 0 {
cfg.MaxAge = DefaultMaxAge
}
return &Poller{
feeds: valid, fetch: fetch, notes: notes, marks: marks,
embed: embed, ranker: ranker, cfg: cfg,
nextDue: map[string]time.Time{},
seen: map[string]map[string]bool{},
}
}
// Feeds returns the configured feeds (the answer path lists categories).
func (p *Poller) Feeds() []FeedConfig { return p.feeds }
// PollDue reads every feed whose interval has elapsed and returns how many
// notes were written. Errors are logged per feed, never returned: one dead feed
// must not stop the others, and there is nobody waiting on this.
func (p *Poller) PollDue(ctx context.Context, now time.Time) int {
written := 0
for _, f := range p.feeds {
if due, ok := p.nextDue[f.Name]; ok && now.Before(due) {
continue
}
interval := f.Interval
if interval <= 0 {
interval = p.cfg.DefaultInterval
}
p.nextDue[f.Name] = now.Add(interval)
n, err := p.PollFeed(ctx, f, now)
if err != nil {
// The URL is configured by him and not a secret, so it is loggable;
// item titles are not logged, only counts.
log.Printf("rss: feed %s: %v", f.Name, err)
continue
}
if n > 0 {
log.Printf("rss: feed %s: %d new item(s) noted", f.Name, n)
}
written += n
}
return written
}
// PollFeed reads one feed now, regardless of its schedule.
func (p *Poller) PollFeed(ctx context.Context, f FeedConfig, now time.Time) (int, error) {
body, err := p.fetch.Get(ctx, f.URL)
if err != nil {
return 0, err
}
feed, err := Parse(strings.NewReader(string(body.Bytes)))
if err != nil {
return 0, err
}
mark := p.mark(ctx, f.Name, now)
newest := mark
written := 0
for _, it := range feed.Items {
if written >= p.cfg.MaxItems {
break
}
if !p.fresh(f, it, mark, now) {
continue
}
if !Matches(f, it) {
continue
}
if p.ranker != nil {
ok, err := p.ranker.Relevant(ctx, it.Title+" "+it.Summary)
if err != nil {
log.Printf("rss: feed %s: relevance: %v", f.Name, err)
} else if !ok {
continue
}
}
if err := p.write(ctx, f, it, now); err != nil {
return written, err
}
written++
if it.Published.After(newest) {
newest = it.Published
}
}
if p.marks != nil && newest.After(mark) {
if err := p.marks.SetMark(ctx, f.Name, newest); err != nil {
log.Printf("rss: feed %s: save mark: %v", f.Name, err)
}
}
return written, nil
}
// mark — how far this feed was read. A feed with no mark starts MaxAge ago, so
// a first poll takes today's headlines instead of the whole archive.
func (p *Poller) mark(ctx context.Context, feed string, now time.Time) time.Time {
cold := now.Add(-p.cfg.MaxAge)
if p.marks == nil {
return cold
}
at, err := p.marks.LastMark(ctx, feed)
if err != nil || at.IsZero() {
return cold
}
return at
}
// fresh — two dedup rules, because feeds are inconsistent about dates. A dated
// item must be newer than the mark; an undated one is kept once per process by
// ID. Both are needed: dates alone re-import undated feeds forever, IDs alone
// lose their memory on restart.
func (p *Poller) fresh(f FeedConfig, it Item, mark, now time.Time) bool {
if !it.Published.IsZero() {
if !it.Published.After(mark) {
return false
}
// A feed that dates its items in the future (or a clock skew) must not
// win the mark and mute everything after it.
return !it.Published.After(now.Add(time.Hour))
}
id := it.ID
if id == "" {
id = it.Title
}
if p.seen[f.Name] == nil {
p.seen[f.Name] = map[string]bool{}
}
if p.seen[f.Name][id] {
return false
}
p.seen[f.Name][id] = true
return true
}
// write stores one item as a note. Source "rss:<feed>" is what the answer path
// filters on, and what makes a feed note distinguishable from something he said.
func (p *Poller) write(ctx context.Context, f FeedConfig, it Item, now time.Time) error {
text := NoteText(f, it)
var vec []float32
if p.embed != nil {
v, err := p.embed.Embed(ctx, text)
if err != nil {
log.Printf("rss: feed %s: embed: %v", f.Name, err)
} else {
vec = v
}
}
ts := it.Published
if ts.IsZero() {
ts = now
}
if _, err := p.notes.WriteNote(ctx, ts, text, vec, SourceFor(f.Name)); err != nil {
return fmt.Errorf("write note: %w", err)
}
return nil
}
// SourceFor is the note source for a feed.
func SourceFor(feed string) string { return "rss:" + feed }
// SourcePrefix — what the answer path matches to find feed notes.
const SourcePrefix = "rss:"
// NoteText renders an item as the note body. The category is included because
// "что нового по технологиям?" is answered by reading notes, and a note has to
// carry enough to be recognised as belonging to that category.
func NoteText(f FeedConfig, it Item) string {
var b strings.Builder
b.WriteString(it.Title)
if f.Category != "" {
fmt.Fprintf(&b, " [%s]", f.Category)
}
if it.Summary != "" {
b.WriteString("\n")
b.WriteString(trimRunes(it.Summary, 500))
}
if it.Link != "" {
b.WriteString("\n")
b.WriteString(it.Link)
}
return b.String()
}
// trimRunes cuts on a rune boundary — a note is Russian as often as English and
// half a cyrillic letter is a broken note.
func trimRunes(s string, max int) string {
r := []rune(s)
if len(r) <= max {
return s
}
return strings.TrimSpace(string(r[:max])) + "…"
}
// Matches applies the per-feed keyword filter: keep when Include is empty or one
// include matches, drop when any exclude matches. Case-insensitive substring,
// which for Russian is the honest choice — no stemmer here, so "выборы" does not
// match "выборах", and a filter he writes is a filter he can predict.
func Matches(f FeedConfig, it Item) bool {
hay := strings.ToLower(it.Title + " " + it.Summary)
for _, x := range f.Exclude {
if x = strings.ToLower(strings.TrimSpace(x)); x != "" && strings.Contains(hay, x) {
return false
}
}
if len(f.Include) == 0 {
return true
}
for _, in := range f.Include {
if in = strings.ToLower(strings.TrimSpace(in)); in != "" && strings.Contains(hay, in) {
return true
}
}
return false
}
+210
View File
@@ -0,0 +1,210 @@
package rss
import (
"context"
"errors"
"strings"
"testing"
"time"
)
type fakeFetch struct {
body string
err error
calls int
urls []string
}
func (f *fakeFetch) Get(_ context.Context, url string) (*Body, error) {
f.calls++
f.urls = append(f.urls, url)
if f.err != nil {
return nil, f.err
}
return &Body{Bytes: []byte(f.body)}, nil
}
type writtenNote struct {
ts time.Time
text string
source string
vec []float32
}
type fakeNotes struct{ notes []writtenNote }
func (n *fakeNotes) WriteNote(_ context.Context, ts time.Time, text string, vec []float32, source string) (int64, error) {
n.notes = append(n.notes, writtenNote{ts, text, source, vec})
return int64(len(n.notes)), nil
}
type fakeMarks struct{ m map[string]time.Time }
func newMarks() *fakeMarks { return &fakeMarks{m: map[string]time.Time{}} }
func (f *fakeMarks) LastMark(_ context.Context, feed string) (time.Time, error) {
return f.m[feed], nil
}
func (f *fakeMarks) SetMark(_ context.Context, feed string, at time.Time) error {
f.m[feed] = at
return nil
}
var now = time.Date(2026, 7, 28, 12, 0, 0, 0, time.UTC)
func TestPollWritesNotesWithSource(t *testing.T) {
fetch := &fakeFetch{body: rss2}
notes := &fakeNotes{}
marks := newMarks()
p := NewPoller([]FeedConfig{{Name: "habr", URL: "https://example.org/rss", Category: "технологии"}},
fetch, notes, marks, nil, nil, Config{})
if p == nil {
t.Fatal("NewPoller returned nil for a configured feed")
}
n := p.PollDue(context.Background(), now)
if n != 2 || len(notes.notes) != 2 {
t.Fatalf("wrote %d notes (returned %d), want 2", len(notes.notes), n)
}
if notes.notes[0].source != "rss:habr" {
t.Errorf("source = %q, want rss:habr", notes.notes[0].source)
}
if !strings.Contains(notes.notes[0].text, "технологии") {
t.Errorf("note does not carry its category: %q", notes.notes[0].text)
}
if !strings.Contains(notes.notes[0].text, "https://example.org/a") {
t.Errorf("note does not carry its link: %q", notes.notes[0].text)
}
// The undated item is stamped with now, the dated one with its own date.
if !notes.notes[1].ts.Equal(now) {
t.Errorf("undated item ts = %v, want now", notes.notes[1].ts)
}
}
// The whole point of a mark: polling twice must not re-note the same headlines.
func TestSecondPollIsQuiet(t *testing.T) {
fetch := &fakeFetch{body: rss2}
notes := &fakeNotes{}
p := NewPoller([]FeedConfig{{Name: "habr", URL: "u", Interval: time.Minute}}, fetch, notes, newMarks(), nil, nil, Config{})
p.PollDue(context.Background(), now)
before := len(notes.notes)
p.PollDue(context.Background(), now.Add(2*time.Minute))
if len(notes.notes) != before {
t.Fatalf("second poll wrote %d extra notes", len(notes.notes)-before)
}
}
// A mark that survives a restart is the durable half; simulate one by building a
// fresh poller over the same marks.
func TestMarkSurvivesRestart(t *testing.T) {
marks := newMarks()
fetch := &fakeFetch{body: rss2}
notes := &fakeNotes{}
feeds := []FeedConfig{{Name: "habr", URL: "u"}}
NewPoller(feeds, fetch, notes, marks, nil, nil, Config{}).PollDue(context.Background(), now)
if len(notes.notes) != 2 {
t.Fatalf("first run wrote %d", len(notes.notes))
}
notes2 := &fakeNotes{}
NewPoller(feeds, fetch, notes2, marks, nil, nil, Config{}).PollDue(context.Background(), now.Add(time.Hour))
// The dated item is behind the mark. The undated one has no date to compare,
// so it comes back — accepted and documented in fresh(): an undated feed is
// deduped per process, not forever.
for _, n := range notes2.notes {
if strings.Contains(n.text, "уязвимость") {
t.Fatalf("dated item re-noted after restart: %q", n.text)
}
}
}
func TestIntervalIsRespected(t *testing.T) {
fetch := &fakeFetch{body: rss2}
p := NewPoller([]FeedConfig{{Name: "habr", URL: "u", Interval: time.Hour}}, fetch, &fakeNotes{}, newMarks(), nil, nil, Config{})
p.PollDue(context.Background(), now)
p.PollDue(context.Background(), now.Add(time.Minute))
if fetch.calls != 1 {
t.Fatalf("fetched %d times inside one interval, want 1", fetch.calls)
}
p.PollDue(context.Background(), now.Add(2*time.Hour))
if fetch.calls != 2 {
t.Fatalf("fetched %d times, want 2 after the interval elapsed", fetch.calls)
}
}
func TestColdStartIgnoresOldItems(t *testing.T) {
old := `<rss><channel><item><title>Старое</title><link>l</link>` +
`<pubDate>Mon, 01 Jun 2026 10:00:00 +0000</pubDate></item></channel></rss>`
notes := &fakeNotes{}
p := NewPoller([]FeedConfig{{Name: "f", URL: "u"}}, &fakeFetch{body: old}, notes, newMarks(), nil, nil, Config{MaxAge: 24 * time.Hour})
if n := p.PollDue(context.Background(), now); n != 0 {
t.Fatalf("cold start imported %d old items, want 0", n)
}
}
func TestMaxItemsCap(t *testing.T) {
var b strings.Builder
b.WriteString("<rss><channel>")
for i := 0; i < 10; i++ {
b.WriteString("<item><title>t")
b.WriteByte(byte('0' + i))
b.WriteString("</title><link>https://example.org/")
b.WriteByte(byte('0' + i))
b.WriteString("</link></item>")
}
b.WriteString("</channel></rss>")
notes := &fakeNotes{}
p := NewPoller([]FeedConfig{{Name: "f", URL: "u"}}, &fakeFetch{body: b.String()}, notes, newMarks(), nil, nil, Config{MaxItems: 3})
if n := p.PollDue(context.Background(), now); n != 3 {
t.Fatalf("wrote %d notes, want the cap of 3", n)
}
}
func TestKeywordFilter(t *testing.T) {
f := FeedConfig{Include: []string{"ядр"}, Exclude: []string{"реклама"}}
if !Matches(f, Item{Title: "Новое ядро"}) {
t.Error("include did not match")
}
if Matches(f, Item{Title: "Новое ядро", Summary: "Реклама внутри"}) {
t.Error("exclude must win over include")
}
if Matches(f, Item{Title: "Погода"}) {
t.Error("non-matching item passed the include filter")
}
if !Matches(FeedConfig{}, Item{Title: "что угодно"}) {
t.Error("an unfiltered feed must keep everything")
}
}
type fakeRanker struct{ keep bool }
func (r fakeRanker) Relevant(context.Context, string) (bool, error) { return r.keep, nil }
func TestRankerCanDropEverything(t *testing.T) {
notes := &fakeNotes{}
p := NewPoller([]FeedConfig{{Name: "f", URL: "u"}}, &fakeFetch{body: rss2}, notes, newMarks(), nil, fakeRanker{false}, Config{})
if n := p.PollDue(context.Background(), now); n != 0 {
t.Fatalf("ranker rejected everything but %d notes were written", n)
}
}
func TestFetchErrorIsSurvivable(t *testing.T) {
notes := &fakeNotes{}
p := NewPoller([]FeedConfig{
{Name: "dead", URL: "u1"},
{Name: "live", URL: "u2"},
}, &fakeFetch{err: errors.New("boom")}, notes, newMarks(), nil, nil, Config{})
if n := p.PollDue(context.Background(), now); n != 0 {
t.Fatalf("n = %d", n)
}
// Both feeds were attempted: one dead feed does not abort the round.
if p.nextDue["live"].IsZero() {
t.Fatal("the second feed was never attempted")
}
}
func TestNoFeedsMeansNoPoller(t *testing.T) {
if p := NewPoller(nil, &fakeFetch{}, &fakeNotes{}, nil, nil, nil, Config{}); p != nil {
t.Fatal("NewPoller must return nil when nothing is configured")
}
if p := NewPoller([]FeedConfig{{Name: "", URL: ""}}, &fakeFetch{}, &fakeNotes{}, nil, nil, nil, Config{}); p != nil {
t.Fatal("a feed with no name or url is not a configuration")
}
}
+301
View File
@@ -0,0 +1,301 @@
// Package webfetch is the one door Maven uses to read something off the
// network, and it is a narrow one.
//
// "Never phones home" stopped being a hard constraint on 2026-07-31, but what
// replaced it is not "she may fetch anything": local sources come first (Kiwix
// on the box), external fetching is off unless configured, and only the
// utterance ever leaves — never his notes, facts or history. That policy is
// enforced by the callers. What THIS package enforces is the part that must be
// code rather than a paragraph in a plan, because it protects the homelab from
// its own assistant:
//
// - http/https only — no file://, no ftp://, no gopher;
// - no private address, ever: loopback, RFC1918 (which is what makes the
// 10.42.0.0/24 wireguard tunnel and the 192.168.1.0/24 LAN unreachable),
// link-local incl. the 169.254.169.254 cloud metadata address, CGNAT,
// unique-local v6. Checked in the dialer's Control hook, so it holds for
// every address the resolver returns AND for every hop of a redirect
// chain — a DNS name that resolves to 127.0.0.1 is refused at connect
// time, which a pre-flight lookup could not promise (rebinding);
// - an allowlist, when one is configured, and a denylist that always wins;
// - a response size cap, a total timeout, a redirect cap;
// - one request per host per interval, so a poll loop with a bug is slow
// rather than an outbound flood.
//
// Everything above is on by default with sane numbers: a zero Config is a
// usable, conservative fetcher. There is no cache and no retry — a feed poll
// or a page read that fails is simply not answered this round.
package webfetch
import (
"context"
"errors"
"fmt"
"io"
"net"
"net/http"
"net/url"
"strings"
"sync"
"syscall"
"time"
)
// Defaults. Small on purpose: this reads feeds and article pages, not ISOs.
const (
DefaultTimeout = 20 * time.Second
DefaultMaxBytes = 2 << 20 // 2 MiB
DefaultMaxRedirects = 3
DefaultHostInterval = time.Second
DefaultUserAgent = "Maven/1.0 (self-hosted personal assistant)"
)
// Errors callers distinguish. Everything else is wrapped transport error.
var (
ErrScheme = errors.New("webfetch: only http and https are allowed")
ErrBlocked = errors.New("webfetch: host is not allowed")
ErrPrivate = errors.New("webfetch: refusing to connect to a private address")
ErrTooLarge = errors.New("webfetch: response exceeds the size cap")
ErrRedirects = errors.New("webfetch: too many redirects")
ErrStatus = errors.New("webfetch: non-2xx status")
)
// Config are the limits. Every zero value means "the default above", so
// Config{} is safe; the only field that changes behaviour by being empty is
// AllowHosts (empty ⇒ any public host that is not denied).
type Config struct {
// AllowHosts — when non-empty, the ONLY hosts that may be fetched. An
// entry matches the host itself and its subdomains ("example.com" allows
// "news.example.com"). This is the knob to reach for when a capability
// should read two feeds and nothing else.
AllowHosts []string
// DenyHosts — same matching, checked first and always winning.
DenyHosts []string
Timeout time.Duration // whole request, including redirects and body read
MaxBytes int64 // response body cap
MaxRedirects int // 0 ⇒ default; negative ⇒ no redirects followed
HostInterval time.Duration // minimum spacing between requests to one host
UserAgent string
// AllowPrivate disables the private-address guard. It exists for tests
// (httptest listens on 127.0.0.1) and for an explicitly configured
// on-box mirror. Nothing in deploy/mavend.json sets it, and it should
// stay that way: with it on, any URL Maven is handed becomes an SSRF
// probe of the LAN and the wireguard range.
AllowPrivate bool
}
// Response is a fetched body, already bounded by MaxBytes.
type Response struct {
URL string // final URL after redirects
Status int
ContentType string
Body []byte
}
// Fetcher performs guarded GETs. Safe for concurrent use; the per-host rate
// limiter is shared, which is the point of sharing one Fetcher.
type Fetcher struct {
cfg Config
http *http.Client
mu sync.Mutex
last map[string]time.Time // host → when we last dialed it
}
// New builds a fetcher from cfg, filling in defaults.
func New(cfg Config) *Fetcher {
if cfg.Timeout <= 0 {
cfg.Timeout = DefaultTimeout
}
if cfg.MaxBytes <= 0 {
cfg.MaxBytes = DefaultMaxBytes
}
if cfg.MaxRedirects == 0 {
cfg.MaxRedirects = DefaultMaxRedirects
}
if cfg.HostInterval <= 0 {
cfg.HostInterval = DefaultHostInterval
}
if cfg.UserAgent == "" {
cfg.UserAgent = DefaultUserAgent
}
f := &Fetcher{cfg: cfg, last: map[string]time.Time{}}
dialer := &net.Dialer{Timeout: 10 * time.Second}
if !cfg.AllowPrivate {
// The guard lives here rather than in a pre-flight net.LookupHost so
// that it sees the address actually being connected to: every A/AAAA
// the resolver handed back, on every redirect hop, with no window in
// which the name could be re-pointed at the LAN.
dialer.Control = func(_, address string, _ syscall.RawConn) error {
host, _, err := net.SplitHostPort(address)
if err != nil {
return err
}
ip := net.ParseIP(host)
if ip == nil || IsPrivateIP(ip) {
return fmt.Errorf("%w: %s", ErrPrivate, host)
}
return nil
}
}
f.http = &http.Client{
Timeout: cfg.Timeout,
Transport: &http.Transport{DialContext: dialer.DialContext},
CheckRedirect: func(req *http.Request, via []*http.Request) error {
if len(via) > f.cfg.MaxRedirects {
return ErrRedirects
}
// A redirect is a fresh URL and gets the full check: an allowed
// host must not be able to bounce us onto a denied one.
return f.checkURL(req.URL)
},
}
return f
}
// Get fetches rawURL. The body is capped: a larger response is an error, not a
// truncation, because half an XML document is worse than none.
func (f *Fetcher) Get(ctx context.Context, rawURL string) (*Response, error) {
u, err := url.Parse(strings.TrimSpace(rawURL))
if err != nil {
return nil, fmt.Errorf("webfetch: bad url %q: %w", rawURL, err)
}
if err := f.checkURL(u); err != nil {
return nil, err
}
if err := f.waitTurn(ctx, u.Hostname()); err != nil {
return nil, err
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u.String(), nil)
if err != nil {
return nil, err
}
req.Header.Set("User-Agent", f.cfg.UserAgent)
req.Header.Set("Accept-Encoding", "identity")
resp, err := f.http.Do(req)
if err != nil {
// http.Client wraps our sentinels in *url.Error; unwrap so callers can
// still tell "blocked" from "the network is down".
for _, sentinel := range []error{ErrPrivate, ErrBlocked, ErrRedirects, ErrScheme} {
if errors.Is(err, sentinel) {
return nil, err
}
}
return nil, err
}
defer resp.Body.Close()
body, err := io.ReadAll(io.LimitReader(resp.Body, f.cfg.MaxBytes+1))
if err != nil {
return nil, err
}
if int64(len(body)) > f.cfg.MaxBytes {
return nil, fmt.Errorf("%w (%d bytes)", ErrTooLarge, f.cfg.MaxBytes)
}
if resp.StatusCode < 200 || resp.StatusCode > 299 {
return nil, fmt.Errorf("%w: %d", ErrStatus, resp.StatusCode)
}
return &Response{
URL: resp.Request.URL.String(),
Status: resp.StatusCode,
ContentType: resp.Header.Get("Content-Type"),
Body: body,
}, nil
}
// checkURL applies the scheme rule and the host lists. The address rule is the
// dialer's job (see New).
func (f *Fetcher) checkURL(u *url.URL) error {
switch u.Scheme {
case "http", "https":
default:
return fmt.Errorf("%w: %q", ErrScheme, u.Scheme)
}
host := strings.ToLower(u.Hostname())
if host == "" {
return fmt.Errorf("%w: no host", ErrBlocked)
}
if HostMatches(host, f.cfg.DenyHosts) {
return fmt.Errorf("%w: %s is denied", ErrBlocked, host)
}
if len(f.cfg.AllowHosts) > 0 && !HostMatches(host, f.cfg.AllowHosts) {
return fmt.Errorf("%w: %s is not on the allowlist", ErrBlocked, host)
}
// A literal private address is refused here as well as in the dialer, so
// the error is the specific one even when no connection is attempted.
if !f.cfg.AllowPrivate {
if ip := net.ParseIP(host); ip != nil && IsPrivateIP(ip) {
return fmt.Errorf("%w: %s", ErrPrivate, host)
}
}
return nil
}
// waitTurn blocks until this host's rate-limit interval has elapsed. It holds
// no lock while sleeping, so two hosts never wait on each other.
func (f *Fetcher) waitTurn(ctx context.Context, host string) error {
for {
f.mu.Lock()
now := time.Now()
earliest := f.last[host].Add(f.cfg.HostInterval)
if !now.Before(earliest) {
f.last[host] = now
f.mu.Unlock()
return nil
}
f.mu.Unlock()
wait := time.NewTimer(earliest.Sub(now))
select {
case <-ctx.Done():
wait.Stop()
return ctx.Err()
case <-wait.C:
}
}
}
// HostMatches reports whether host equals one of pats or is a subdomain of one.
// Exported because the crawler applies the same rule to links it decides not to
// follow, before it ever builds a request.
func HostMatches(host string, pats []string) bool {
host = strings.ToLower(strings.TrimSuffix(host, "."))
for _, p := range pats {
p = strings.ToLower(strings.TrimSpace(strings.TrimPrefix(p, "*.")))
if p == "" {
continue
}
if host == p || strings.HasSuffix(host, "."+p) {
return true
}
}
return false
}
// cgnat is 100.64.0.0/10 — carrier NAT, not covered by net.IP's helpers and not
// somewhere a personal assistant has business connecting.
var cgnat = &net.IPNet{IP: net.IPv4(100, 64, 0, 0).To4(), Mask: net.CIDRMask(10, 32)}
// IsPrivateIP reports whether ip is somewhere Maven must never reach out to:
// the box itself, the LAN, the wireguard range (10.42.0.0/24 ⊂ 10/8), the cloud
// metadata address (169.254.169.254 ⊂ link-local), or anything unroutable.
func IsPrivateIP(ip net.IP) bool {
if ip.IsLoopback() || ip.IsPrivate() || ip.IsUnspecified() ||
ip.IsLinkLocalUnicast() || ip.IsLinkLocalMulticast() ||
ip.IsInterfaceLocalMulticast() || ip.IsMulticast() {
return true
}
if v4 := ip.To4(); v4 != nil && cgnat.Contains(v4) {
return true
}
// IPv4-mapped/compatible forms of the above are handled by To4() inside the
// stdlib helpers; what is left is v6 unique-local (fc00::/7).
if len(ip) == net.IPv6len && ip.To4() == nil && ip[0]&0xfe == 0xfc {
return true
}
return false
}
+227
View File
@@ -0,0 +1,227 @@
package webfetch
import (
"context"
"errors"
"net"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
)
// The limits in this package are the reason a crawler is allowed to exist on
// this box at all, so each one has a test that fails loudly if it is removed.
func TestPrivateAddressesAreRefused(t *testing.T) {
// The wireguard range (10.42.0.0/24), the LAN (192.168.1.0/24) and the
// cloud metadata address are the three that matter here; the rest come
// along for free.
for _, s := range []string{
"127.0.0.1", "127.1.2.3", "10.42.0.7", "10.0.0.5", "192.168.1.104",
"172.16.4.4", "169.254.169.254", "100.64.1.1", "0.0.0.0",
"::1", "fc00::1", "fd12:3456::1", "fe80::1",
} {
if !IsPrivateIP(net.ParseIP(s)) {
t.Errorf("IsPrivateIP(%s) = false, want true", s)
}
}
for _, s := range []string{"8.8.8.8", "1.1.1.1", "93.184.216.34", "2606:2800:220:1::1"} {
if IsPrivateIP(net.ParseIP(s)) {
t.Errorf("IsPrivateIP(%s) = true, want false", s)
}
}
}
func TestGetRefusesPrivateLiteral(t *testing.T) {
f := New(Config{})
for _, u := range []string{
"http://127.0.0.1:8034/search",
"http://10.42.0.1/",
"http://192.168.1.104/dash",
"http://[::1]:9100/mcp",
} {
if _, err := f.Get(context.Background(), u); !errors.Is(err, ErrPrivate) {
t.Errorf("Get(%s) error = %v, want ErrPrivate", u, err)
}
}
}
// A hostname that resolves into private space must fail too — that is the
// rebinding case, and it is why the check lives in the dialer.
func TestGetRefusesPrivateResolution(t *testing.T) {
f := New(Config{})
if _, err := f.Get(context.Background(), "http://localhost:8034/"); !errors.Is(err, ErrPrivate) {
t.Fatalf("Get(localhost) error = %v, want ErrPrivate", err)
}
}
func TestGetRefusesNonHTTPSchemes(t *testing.T) {
f := New(Config{})
for _, u := range []string{"file:///etc/passwd", "ftp://example.com/x", "gopher://example.com"} {
if _, err := f.Get(context.Background(), u); !errors.Is(err, ErrScheme) {
t.Errorf("Get(%s) error = %v, want ErrScheme", u, err)
}
}
}
// testFetcher — a fetcher pointed at an httptest server, which necessarily
// listens on loopback. AllowPrivate is the test-only escape hatch.
func testFetcher(t *testing.T, cfg Config) *Fetcher {
t.Helper()
cfg.AllowPrivate = true
if cfg.HostInterval == 0 {
cfg.HostInterval = time.Nanosecond
}
return New(cfg)
}
func TestAllowAndDenyLists(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte("ok"))
}))
defer srv.Close()
f := testFetcher(t, Config{AllowHosts: []string{"example.com"}})
if _, err := f.Get(context.Background(), srv.URL); !errors.Is(err, ErrBlocked) {
t.Fatalf("off-allowlist host: error = %v, want ErrBlocked", err)
}
f = testFetcher(t, Config{DenyHosts: []string{"127.0.0.1"}})
if _, err := f.Get(context.Background(), srv.URL); !errors.Is(err, ErrBlocked) {
t.Fatalf("denied host: error = %v, want ErrBlocked", err)
}
f = testFetcher(t, Config{AllowHosts: []string{"127.0.0.1"}})
if _, err := f.Get(context.Background(), srv.URL); err != nil {
t.Fatalf("allowlisted host: %v", err)
}
}
func TestHostMatchesSubdomains(t *testing.T) {
pats := []string{"example.com", "*.news.org"}
for _, h := range []string{"example.com", "news.example.com", "a.b.example.com", "news.org", "feeds.news.org"} {
if !HostMatches(h, pats) {
t.Errorf("HostMatches(%q) = false, want true", h)
}
}
for _, h := range []string{"notexample.com", "example.com.evil.net", "org"} {
if HostMatches(h, pats) {
t.Errorf("HostMatches(%q) = true, want false", h)
}
}
}
func TestSizeCap(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte(strings.Repeat("x", 5000)))
}))
defer srv.Close()
f := testFetcher(t, Config{MaxBytes: 100})
if _, err := f.Get(context.Background(), srv.URL); !errors.Is(err, ErrTooLarge) {
t.Fatalf("error = %v, want ErrTooLarge", err)
}
f = testFetcher(t, Config{MaxBytes: 6000})
resp, err := f.Get(context.Background(), srv.URL)
if err != nil {
t.Fatalf("under the cap: %v", err)
}
if len(resp.Body) != 5000 {
t.Fatalf("body = %d bytes, want 5000", len(resp.Body))
}
}
func TestRedirectCap(t *testing.T) {
var srv *httptest.Server
srv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Redirect(w, r, srv.URL+"/again", http.StatusFound)
}))
defer srv.Close()
f := testFetcher(t, Config{MaxRedirects: 2})
if _, err := f.Get(context.Background(), srv.URL); !errors.Is(err, ErrRedirects) {
t.Fatalf("error = %v, want ErrRedirects", err)
}
}
// A redirect off the allowlist is the interesting redirect: the first hop is
// permitted, the second must not be.
func TestRedirectRecheckedAgainstDenylist(t *testing.T) {
target := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte("secret"))
}))
defer target.Close()
hop := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Redirect(w, r, target.URL, http.StatusFound)
}))
defer hop.Close()
// Reach the hop under the name "localhost" and allow only that name; the
// redirect lands on the same box under its literal address, which the
// allowlist does not cover. Without the CheckRedirect hook this fetch
// succeeds and returns "secret".
f := testFetcher(t, Config{AllowHosts: []string{"localhost"}})
viaName := strings.Replace(hop.URL, "127.0.0.1", "localhost", 1)
if _, err := f.Get(context.Background(), viaName); !errors.Is(err, ErrBlocked) {
t.Fatalf("error = %v, want ErrBlocked", err)
}
}
func TestPerHostRateLimit(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte("ok"))
}))
defer srv.Close()
f := testFetcher(t, Config{HostInterval: 60 * time.Millisecond})
start := time.Now()
for i := 0; i < 3; i++ {
if _, err := f.Get(context.Background(), srv.URL); err != nil {
t.Fatalf("request %d: %v", i, err)
}
}
if elapsed := time.Since(start); elapsed < 120*time.Millisecond {
t.Fatalf("three requests took %s, want at least 120ms of spacing", elapsed)
}
}
func TestRateLimitHonoursContext(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}))
defer srv.Close()
f := testFetcher(t, Config{HostInterval: 10 * time.Second})
if _, err := f.Get(context.Background(), srv.URL); err != nil {
t.Fatalf("first request: %v", err)
}
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond)
defer cancel()
if _, err := f.Get(ctx, srv.URL); !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("error = %v, want DeadlineExceeded", err)
}
}
func TestNon2xxIsAnError(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Error(w, "nope", http.StatusInternalServerError)
}))
defer srv.Close()
f := testFetcher(t, Config{})
if _, err := f.Get(context.Background(), srv.URL); !errors.Is(err, ErrStatus) {
t.Fatalf("error = %v, want ErrStatus", err)
}
}
func TestUserAgentIsSent(t *testing.T) {
got := make(chan string, 1)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
got <- r.Header.Get("User-Agent")
}))
defer srv.Close()
f := testFetcher(t, Config{UserAgent: "Maven/test"})
if _, err := f.Get(context.Background(), srv.URL); err != nil {
t.Fatal(err)
}
if ua := <-got; ua != "Maven/test" {
t.Fatalf("user-agent = %q", ua)
}
}