Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| cb3641e7bb |
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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**
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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), " "))
|
||||
}
|
||||
@@ -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><p>Патч уже <b>вышел</b>.</p></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 & 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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user