Read RSS and Atom feeds, and speak about them only when asked (#258) #66

Closed
claude wants to merge 1 commits from overnight/rss-feeds into overnight/email-poller
17 changed files with 2069 additions and 0 deletions
+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,
+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
@@ -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
+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)
}
}