diff --git a/cmd/mavend/actions_query.go b/cmd/mavend/actions_query.go index 69216ff..4cfa1b7 100644 --- a/cmd/mavend/actions_query.go +++ b/cmd/mavend/actions_query.go @@ -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. diff --git a/cmd/mavend/feeds.go b/cmd/mavend/feeds.go new file mode 100644 index 0000000..7797697 --- /dev/null +++ b/cmd/mavend/feeds.go @@ -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:", 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) +} diff --git a/cmd/mavend/feeds_test.go b/cmd/mavend/feeds_test.go new file mode 100644 index 0000000..50783b0 --- /dev/null +++ b/cmd/mavend/feeds_test.go @@ -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") + } +} diff --git a/cmd/mavend/main.go b/cmd/mavend/main.go index 5d1cf2c..4053be3 100644 --- a/cmd/mavend/main.go +++ b/cmd/mavend/main.go @@ -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() diff --git a/cmd/mavend/voice.go b/cmd/mavend/voice.go index 79327d0..4ef572c 100644 --- a/cmd/mavend/voice.go +++ b/cmd/mavend/voice.go @@ -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 diff --git a/cmd/mavend/voicewire.go b/cmd/mavend/voicewire.go index 0e5ce5b..21032df 100644 --- a/cmd/mavend/voicewire.go +++ b/cmd/mavend/voicewire.go @@ -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, diff --git a/deploy/README.md b/deploy/README.md index afdeba2..0cb1a24 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -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:`, 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:`, 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** diff --git a/docs/plans/13-rss-news-feeds.md b/docs/plans/13-rss-news-feeds.md index 9edd1b5..5a4289e 100644 --- a/docs/plans/13-rss-news-feeds.md +++ b/docs/plans/13-rss-news-feeds.md @@ -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:` 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. diff --git a/internal/config/config.go b/internal/config/config.go index 5becbb8..47d476c 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -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:" 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:" + 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 diff --git a/internal/router/feeds.go b/internal/router/feeds.go new file mode 100644 index 0000000..91a6dac --- /dev/null +++ b/internal/router/feeds.go @@ -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) +} diff --git a/internal/router/feeds_test.go b/internal/router/feeds_test.go new file mode 100644 index 0000000..72ffdd1 --- /dev/null +++ b/internal/router/feeds_test.go @@ -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") + } +} diff --git a/internal/rss/feed.go b/internal/rss/feed.go new file mode 100644 index 0000000..71d5580 --- /dev/null +++ b/internal/rss/feed.go @@ -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)]*>.*?|]*>.*?`) + 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), " ")) +} diff --git a/internal/rss/feed_test.go b/internal/rss/feed_test.go new file mode 100644 index 0000000..04a87f3 --- /dev/null +++ b/internal/rss/feed_test.go @@ -0,0 +1,112 @@ +package rss + +import ( + "strings" + "testing" + "time" +) + +const rss2 = ` + + + Хабр + + Новая уязвимость в ядре + https://example.org/a + <p>Патч уже <b>вышел</b>.</p> + tag:example.org,a + Mon, 28 Jul 2026 10:00:00 +0000 + + + Без даты + https://example.org/b + + +` + +const atom = ` + + Example Atom + + Release 2.0 + + + urn:uuid:1 + 2026-07-30T12:30:00Z + Ships & works + +` + +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("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(`x`)) + 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(`

hi

there`) + if got != "hi there" { + t.Fatalf("got %q", got) + } +} diff --git a/internal/rss/poller.go b/internal/rss/poller.go new file mode 100644 index 0000000..073f774 --- /dev/null +++ b/internal/rss/poller.go @@ -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:" + 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:"). +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:" 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 +} diff --git a/internal/rss/poller_test.go b/internal/rss/poller_test.go new file mode 100644 index 0000000..f071399 --- /dev/null +++ b/internal/rss/poller_test.go @@ -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 := `Староеl` + + `Mon, 01 Jun 2026 10:00:00 +0000` + 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("") + for i := 0; i < 10; i++ { + b.WriteString("t") + b.WriteByte(byte('0' + i)) + b.WriteString("https://example.org/") + b.WriteByte(byte('0' + i)) + b.WriteString("") + } + b.WriteString("") + 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") + } +} diff --git a/internal/webfetch/webfetch.go b/internal/webfetch/webfetch.go new file mode 100644 index 0000000..95ab273 --- /dev/null +++ b/internal/webfetch/webfetch.go @@ -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 +} diff --git a/internal/webfetch/webfetch_test.go b/internal/webfetch/webfetch_test.go new file mode 100644 index 0000000..667190d --- /dev/null +++ b/internal/webfetch/webfetch_test.go @@ -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) + } +}