mirror of
https://github.com/AmanTahiliani/box-box.git
synced 2026-08-07 19:56:18 -04:00
Add paddock briefing feed ingestion
This commit is contained in:
141
internal/news/refresh.go
Normal file
141
internal/news/refresh.go
Normal file
@@ -0,0 +1,141 @@
|
||||
package news
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/AmanTahiliani/box-box/internal/store"
|
||||
)
|
||||
|
||||
const DefaultTTL = 30 * time.Minute
|
||||
|
||||
// Store is the storage surface needed by feed refreshes.
|
||||
type Store interface {
|
||||
UpsertNewsSource(store.NewsSource) error
|
||||
UpsertNewsItem(store.NewsItem) error
|
||||
}
|
||||
|
||||
// RefreshOptions configures one local news refresh run.
|
||||
type RefreshOptions struct {
|
||||
Sources []Source
|
||||
Client *http.Client
|
||||
TTL time.Duration
|
||||
DryRun bool
|
||||
Now func() time.Time
|
||||
Progress io.Writer
|
||||
}
|
||||
|
||||
// RefreshResult summarizes one local news refresh run.
|
||||
type RefreshResult struct {
|
||||
SourcesFetched int
|
||||
SourcesFailed int
|
||||
ItemsFetched int
|
||||
ItemsUpserted int
|
||||
}
|
||||
|
||||
// Refresh fetches RSS/Atom sources and stores normalized, URL-deduped items.
|
||||
func Refresh(ctx context.Context, st Store, opts RefreshOptions) (RefreshResult, error) {
|
||||
if st == nil && !opts.DryRun {
|
||||
return RefreshResult{}, errors.New("news refresh: store is required")
|
||||
}
|
||||
if len(opts.Sources) == 0 {
|
||||
opts.Sources = DefaultSources
|
||||
}
|
||||
if opts.Client == nil {
|
||||
opts.Client = &http.Client{Timeout: 10 * time.Second}
|
||||
}
|
||||
if opts.TTL <= 0 {
|
||||
opts.TTL = DefaultTTL
|
||||
}
|
||||
now := func() time.Time { return time.Now().UTC() }
|
||||
if opts.Now != nil {
|
||||
now = func() time.Time { return opts.Now().UTC() }
|
||||
}
|
||||
|
||||
var result RefreshResult
|
||||
var failures []string
|
||||
for _, source := range opts.Sources {
|
||||
fetchedAt := now()
|
||||
expiresAt := fetchedAt.Add(opts.TTL)
|
||||
if !opts.DryRun {
|
||||
if err := st.UpsertNewsSource(store.NewsSource{
|
||||
Source: source.ID,
|
||||
Name: source.Name,
|
||||
FeedURL: source.URL,
|
||||
Category: source.Category,
|
||||
Enabled: true,
|
||||
UpdatedAt: fetchedAt,
|
||||
}); err != nil {
|
||||
return result, err
|
||||
}
|
||||
}
|
||||
items, err := Fetch(ctx, opts.Client, source)
|
||||
if err != nil {
|
||||
result.SourcesFailed++
|
||||
failures = append(failures, fmt.Sprintf("%s: %v", source.ID, err))
|
||||
progressf(opts.Progress, "news: %s failed: %v\n", source.ID, err)
|
||||
continue
|
||||
}
|
||||
|
||||
result.SourcesFetched++
|
||||
result.ItemsFetched += len(items)
|
||||
progressf(opts.Progress, "news: %s fetched %d items\n", source.ID, len(items))
|
||||
if opts.DryRun {
|
||||
continue
|
||||
}
|
||||
|
||||
if err := st.UpsertNewsSource(store.NewsSource{
|
||||
Source: source.ID,
|
||||
Name: source.Name,
|
||||
FeedURL: source.URL,
|
||||
Category: source.Category,
|
||||
Enabled: true,
|
||||
FetchedAt: &fetchedAt,
|
||||
ExpiresAt: &expiresAt,
|
||||
UpdatedAt: fetchedAt,
|
||||
}); err != nil {
|
||||
return result, err
|
||||
}
|
||||
for _, item := range items {
|
||||
item.FetchedAt = fetchedAt
|
||||
publishedAt := timePtr(item.PublishedAt)
|
||||
if err := st.UpsertNewsItem(store.NewsItem{
|
||||
URL: item.URL,
|
||||
Source: item.Source,
|
||||
Title: item.Title,
|
||||
PublishedAt: publishedAt,
|
||||
Summary: item.Summary,
|
||||
Category: item.Category,
|
||||
FetchedAt: item.FetchedAt,
|
||||
}); err != nil {
|
||||
return result, err
|
||||
}
|
||||
result.ItemsUpserted++
|
||||
}
|
||||
}
|
||||
|
||||
if len(failures) > 0 {
|
||||
return result, fmt.Errorf("news refresh completed with %d source failure(s): %s", len(failures), strings.Join(failures, "; "))
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func progressf(w io.Writer, format string, args ...any) {
|
||||
if w == nil {
|
||||
return
|
||||
}
|
||||
fmt.Fprintf(w, format, args...)
|
||||
}
|
||||
|
||||
func timePtr(v time.Time) *time.Time {
|
||||
if v.IsZero() {
|
||||
return nil
|
||||
}
|
||||
t := v.UTC()
|
||||
return &t
|
||||
}
|
||||
149
internal/news/refresh_test.go
Normal file
149
internal/news/refresh_test.go
Normal file
@@ -0,0 +1,149 @@
|
||||
package news
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/AmanTahiliani/box-box/internal/store"
|
||||
)
|
||||
|
||||
func TestRefreshStoresSourcesAndItems(t *testing.T) {
|
||||
feed := `<?xml version="1.0"?>
|
||||
<rss version="2.0">
|
||||
<channel>
|
||||
<item>
|
||||
<title>Briefing one</title>
|
||||
<link>https://example.com/f1/one?utm_source=rss</link>
|
||||
<pubDate>Mon, 25 May 2026 10:00:00 GMT</pubDate>
|
||||
<description>Morning note</description>
|
||||
</item>
|
||||
<item>
|
||||
<title>Briefing duplicate newer</title>
|
||||
<link>https://example.com/f1/one</link>
|
||||
<pubDate>Mon, 25 May 2026 11:00:00 GMT</pubDate>
|
||||
<description>Updated note</description>
|
||||
</item>
|
||||
</channel>
|
||||
</rss>`
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if got := r.Header.Get("User-Agent"); got != UserAgent {
|
||||
t.Fatalf("User-Agent = %q, want %q", got, UserAgent)
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/rss+xml")
|
||||
_, _ = w.Write([]byte(feed))
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
st := openNewsTestStore(t)
|
||||
now := time.Unix(1800000000, 0).UTC()
|
||||
result, err := Refresh(context.Background(), st, RefreshOptions{
|
||||
Sources: []Source{{
|
||||
ID: "example",
|
||||
Name: "Example F1",
|
||||
URL: server.URL + "/feed.xml",
|
||||
Category: "news",
|
||||
}},
|
||||
Client: server.Client(),
|
||||
TTL: time.Hour,
|
||||
Now: func() time.Time { return now },
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Refresh() error = %v", err)
|
||||
}
|
||||
if result.SourcesFetched != 1 || result.ItemsFetched != 1 || result.ItemsUpserted != 1 || result.SourcesFailed != 0 {
|
||||
t.Fatalf("result = %+v, want one fetched/upserted item and no failures", result)
|
||||
}
|
||||
|
||||
var fetchedAt, expiresAt int64
|
||||
if err := st.DB().QueryRow(`
|
||||
SELECT fetched_at, expires_at
|
||||
FROM news_sources
|
||||
WHERE source = 'example'
|
||||
`).Scan(&fetchedAt, &expiresAt); err != nil {
|
||||
t.Fatalf("query source metadata: %v", err)
|
||||
}
|
||||
if fetchedAt != now.Unix() || expiresAt != now.Add(time.Hour).Unix() {
|
||||
t.Fatalf("source times = %d/%d, want %d/%d", fetchedAt, expiresAt, now.Unix(), now.Add(time.Hour).Unix())
|
||||
}
|
||||
|
||||
items, err := st.ListNewsItems(10, "example")
|
||||
if err != nil {
|
||||
t.Fatalf("ListNewsItems() error = %v", err)
|
||||
}
|
||||
if len(items) != 1 {
|
||||
t.Fatalf("items len = %d, want 1", len(items))
|
||||
}
|
||||
if items[0].URL != "https://example.com/f1/one" || items[0].Title != "Briefing duplicate newer" {
|
||||
t.Fatalf("stored item = %+v, want canonical newer duplicate", items[0])
|
||||
}
|
||||
if !items[0].FetchedAt.Equal(now) {
|
||||
t.Fatalf("item fetched_at = %v, want %v", items[0].FetchedAt, now)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRefreshDryRunDoesNotRequireStore(t *testing.T) {
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
_, _ = w.Write([]byte(`<?xml version="1.0"?><rss version="2.0"><channel><item><title>Dry</title><link>https://example.com/dry</link></item></channel></rss>`))
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
result, err := Refresh(context.Background(), nil, RefreshOptions{
|
||||
Sources: []Source{{ID: "dry", Name: "Dry", URL: server.URL, Category: "news"}},
|
||||
Client: server.Client(),
|
||||
DryRun: true,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Refresh() dry run error = %v", err)
|
||||
}
|
||||
if result.ItemsFetched != 1 || result.ItemsUpserted != 0 {
|
||||
t.Fatalf("result = %+v, want fetched item with no upsert", result)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRefreshContinuesAfterSourceFailure(t *testing.T) {
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if strings.Contains(r.URL.Path, "bad") {
|
||||
http.Error(w, "nope", http.StatusBadGateway)
|
||||
return
|
||||
}
|
||||
_, _ = w.Write([]byte(`<?xml version="1.0"?><rss version="2.0"><channel><item><title>Good</title><link>https://example.com/good</link></item></channel></rss>`))
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
st := openNewsTestStore(t)
|
||||
result, err := Refresh(context.Background(), st, RefreshOptions{
|
||||
Sources: []Source{
|
||||
{ID: "bad", Name: "Bad", URL: server.URL + "/bad", Category: "news"},
|
||||
{ID: "good", Name: "Good", URL: server.URL + "/good", Category: "news"},
|
||||
},
|
||||
Client: server.Client(),
|
||||
})
|
||||
if err == nil {
|
||||
t.Fatal("Refresh() error = nil, want source failure")
|
||||
}
|
||||
if result.SourcesFetched != 1 || result.SourcesFailed != 1 || result.ItemsUpserted != 1 {
|
||||
t.Fatalf("result = %+v, want one failure and one stored item", result)
|
||||
}
|
||||
items, err := st.ListNewsItems(10, "")
|
||||
if err != nil {
|
||||
t.Fatalf("ListNewsItems() error = %v", err)
|
||||
}
|
||||
if len(items) != 1 || items[0].Source != "good" {
|
||||
t.Fatalf("items = %+v, want good source item stored", items)
|
||||
}
|
||||
}
|
||||
|
||||
func openNewsTestStore(t *testing.T) *store.Store {
|
||||
t.Helper()
|
||||
st, err := store.Open(filepath.Join(t.TempDir(), "news.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("store.Open() error = %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = st.Close() })
|
||||
return st
|
||||
}
|
||||
Reference in New Issue
Block a user