feat(audiobooks): port podcast RSS feed refresher

Sub-plan 5 of the audiobooks absorption. Ports the plugin's
podcastfeed.Refresher into internal/audiobooks/podcastfeed and
registers it as a 10-minute scheduled task in silo's task manager.

For each row in podcast_feeds (sub-plan 1 mig 140), the refresher:
  - Fetches the RSS feed with a 30-second HTTP timeout
  - Parses entries via gofeed (github.com/mmcdole/gofeed)
  - UPSERTs new entries into episodes (joined to the podcast's
    media_items row via podcast_feeds.media_item_id = series_id)
  - Updates podcast_feeds.last_refreshed_at + last_refresh_error

Migration 144 adds podcast_guid and podcast_audio_url columns to
episodes, plus a partial unique index on (series_id, podcast_guid)
WHERE podcast_guid IS NOT NULL so the refresher can ON CONFLICT upsert
without touching non-podcast episode rows.

Enclosure file download is deferred; this stage upserts episode
metadata only. Existing content_ids are reused when a feed re-emits a
known guid so per-user progress rows stay bound to the same FK target.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
RXWatcher
2026-05-24 16:06:27 +02:00
co-authored by Claude Opus 4.7
parent 8de72bf424
commit 412defb939
9 changed files with 893 additions and 0 deletions
+2
View File
@@ -31,6 +31,7 @@ import (
"github.com/Silo-Server/silo-server/internal/activitylog"
"github.com/Silo-Server/silo-server/internal/adminjob"
"github.com/Silo-Server/silo-server/internal/audiobooks"
"github.com/Silo-Server/silo-server/internal/audiobooks/podcastfeed"
"github.com/Silo-Server/silo-server/internal/api"
"github.com/Silo-Server/silo-server/internal/api/handlers"
"github.com/Silo-Server/silo-server/internal/auth"
@@ -1271,6 +1272,7 @@ func main() {
historyReconciler := watchstate.NewHistoryReconciler(deps.DB, historyResolver)
taskMgr.Register(tasks.NewRepairProviderIDIntegrityTask(metadata.NewProviderIDIntegrityRepairer(deps.DB), historyReconciler))
taskMgr.Register(tasks.NewReconcileWatchHistoryTask(historyReconciler))
taskMgr.Register(tasks.NewSyncPodcastFeedsTask(podcastfeed.New(), podcastfeed.NewDBStore(deps.DB)))
if pluginInstallationStore != nil && pluginRuntimeConfigStore != nil && pluginService != nil {
pluginTasks, err := plugins.NewTaskRegistryWithTypedResolver(pluginInstallationStore, pluginRuntimeConfigStore, pluginService).Tasks(appCtx)
if err != nil {
+7
View File
@@ -34,14 +34,21 @@ require (
)
require (
github.com/PuerkitoBio/goquery v1.8.0 // indirect
github.com/andybalholm/brotli v1.2.0 // indirect
github.com/andybalholm/cascadia v1.3.1 // indirect
github.com/fatih/color v1.13.0 // indirect
github.com/golang/protobuf v1.5.4 // indirect
github.com/gookit/color v1.5.4 // indirect
github.com/hashicorp/yamux v0.1.2 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/klauspost/compress v1.18.0 // indirect
github.com/mattn/go-colorable v0.1.12 // indirect
github.com/mattn/go-isatty v0.0.17 // indirect
github.com/mmcdole/gofeed v1.3.0 // indirect
github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.2 // indirect
github.com/oklog/run v1.1.0 // indirect
github.com/quic-go/qpack v0.5.1 // indirect
github.com/quic-go/quic-go v0.53.0 // indirect
+22
View File
@@ -1,9 +1,13 @@
entgo.io/ent v0.14.3 h1:wokAV/kIlH9TeklJWGGS7AYJdVckr0DloWjIcO9iIIQ=
entgo.io/ent v0.14.3/go.mod h1:aDPE/OziPEu8+OWbzy4UlvWmD2/kbRuWfK2A40hcxJM=
github.com/PuerkitoBio/goquery v1.8.0 h1:PJTF7AmFCFKk1N6V6jmKfrNH9tV5pNE6lZMkG0gta/U=
github.com/PuerkitoBio/goquery v1.8.0/go.mod h1:ypIiRMtY7COPGk+I/YbZLbxsxn9g5ejnI2HSMtkjZvI=
github.com/Silo-Server/silo-plugin-sdk v0.4.0 h1:DJkRROQfr/kfwnF5dUdkhmbCso1KDLZc7uK7YfJnaO0=
github.com/Silo-Server/silo-plugin-sdk v0.4.0/go.mod h1:etqmxLTwjxpFH9goAjBDfNDoqHMv2/sqUXu8yx3hNfA=
github.com/andybalholm/brotli v1.2.0 h1:ukwgCxwYrmACq68yiUqwIWnGY0cTPox/M94sVwToPjQ=
github.com/andybalholm/brotli v1.2.0/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY=
github.com/andybalholm/cascadia v1.3.1 h1:nhxRkql1kdYCc8Snf7D5/D3spOX+dBgjA6u8x004T2c=
github.com/andybalholm/cascadia v1.3.1/go.mod h1:R4bJ1UQfqADjvDa4P6HZHLh/3OxWWEqc0Sk8XGwHqvA=
github.com/aws/aws-sdk-go-v2 v1.41.5 h1:dj5kopbwUsVUVFgO4Fi5BIT3t4WyqIDjGKCangnV/yY=
github.com/aws/aws-sdk-go-v2 v1.41.5/go.mod h1:mwsPRE8ceUUpiTgF7QmQIJ7lgsKUPQOUl3o72QBrE1o=
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.8 h1:eBMB84YGghSocM7PsjmmPffTa+1FBUeNvGvFou6V/4o=
@@ -67,6 +71,7 @@ github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek
github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps=
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/gookit/color v1.5.4 h1:FZmqs7XOyGgCAxmWyPslpiok1k05wmY3SJTytgvYFs0=
@@ -99,6 +104,8 @@ github.com/jmoiron/sqlx v1.3.5 h1:vFFPA71p1o5gAeqtEAwLU4dnX2napprKtHr7PYIcN3g=
github.com/jmoiron/sqlx v1.3.5/go.mod h1:nRVWtLre0KfCLJvgxzCsLVMogSvQ1zNJtpYr2Ccp0mQ=
github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0=
github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4=
github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM=
github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo=
github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo=
github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ=
github.com/klauspost/cpuid/v2 v2.0.9 h1:lgaqFMSdTdQYdZ04uHyN2d/eKdOMyi2YLSvlQIBFYa4=
@@ -120,6 +127,15 @@ github.com/mattn/go-isatty v0.0.17 h1:BTarxUcIeDqL27Mc+vyvdWYSL28zpIhv3RoTdsLMPn
github.com/mattn/go-isatty v0.0.17/go.mod h1:kYGgaQfpe5nmfYZH+SKPsOc2e4SrIfOl2e/yFXSvRLM=
github.com/mattn/go-sqlite3 v1.14.34 h1:3NtcvcUnFBPsuRcno8pUtupspG/GM+9nZ88zgJcp6Zk=
github.com/mattn/go-sqlite3 v1.14.34/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y=
github.com/mmcdole/gofeed v1.3.0 h1:5yn+HeqlcvjMeAI4gu6T+crm7d0anY85+M+v6fIFNG4=
github.com/mmcdole/gofeed v1.3.0/go.mod h1:9TGv2LcJhdXePDzxiuMnukhV2/zb6VtnZt1mS+SjkLE=
github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23 h1:Zr92CAlFhy2gL+V1F+EyIuzbQNbSgP4xhTODZtrXUtk=
github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23/go.mod h1:v+25+lT2ViuQ7mVxcncQ8ch1URund48oH+jhjiwEgS8=
github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg=
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M=
github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
github.com/oklog/run v1.1.0 h1:GEenZ1cK0+q0+wsJew9qUg/DyD8k3JzYsZAi5gYi2mA=
@@ -221,22 +237,28 @@ golang.org/x/image v0.39.0 h1:skVYidAEVKgn8lZ602XO75asgXBgLj9G/FE3RbuPFww=
golang.org/x/image v0.39.0/go.mod h1:sIbmppfU+xFLPIG0FoVUTvyBMmgng1/XAMhQ2ft0hpA=
golang.org/x/mod v0.34.0 h1:xIHgNUUnW6sYkcM5Jleh05DvLOtwc6RitGHbDk4akRI=
golang.org/x/mod v0.34.0/go.mod h1:ykgH52iCZe79kzLLMhyCUzhMci+nQj+0XkbXpNYtVjY=
golang.org/x/net v0.0.0-20210916014120-12bc252f5db8/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y=
golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA=
golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs=
golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4=
golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.0.0-20200116001909-b77594299b42/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20200223170610-d5e6a3e2c0ae/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210630005230-0f9fa26af87c/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20210927094055-39ccf1dd6fa6/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220503163025-988cb79eb6c6/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI=
golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/text v0.36.0 h1:JfKh3XmcRPqZPKevfXVpI1wXPTqbkE5f7JA92a55Yxg=
golang.org/x/text v0.36.0/go.mod h1:NIdBknypM8iqVmPiuco0Dh6P5Jcdk8lJL0CUebqK164=
golang.org/x/time v0.14.0 h1:MRx4UaLrDotUKUdCIqzPC48t1Y9hANFKIRpNx+Te8PI=
golang.org/x/time v0.14.0/go.mod h1:eL/Oa2bBBK0TkX57Fyni+NgnyQQN4LitPmob2Hjnqw4=
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
golang.org/x/tools v0.43.0 h1:12BdW9CeB3Z+J/I/wj34VMl8X+fEXBxVR90JeMX5E7s=
golang.org/x/tools v0.43.0/go.mod h1:uHkMso649BX2cZK6+RpuIPXS3ho2hZo4FVwfoy1vIk0=
gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk=
@@ -0,0 +1,355 @@
// Package podcastfeed fetches and parses RSS / Atom podcast feeds and
// upserts the resulting episodes into silo's episodes table. The Refresher
// is invoked from two places: the task scheduler (periodic background
// refresh for every podcast whose refresh window has elapsed) and an admin
// POST endpoint (force-refresh for troubleshooting / manual triggers after
// seeding a feed URL).
//
// Episode identity is keyed by the feed's <guid>, stored in
// episodes.podcast_guid. Upserts via Store.UpsertPodcastEpisode are
// idempotent: re-emitted feed items update existing rows without producing
// duplicates, so listener progress rows survive feed re-emits.
package podcastfeed
import (
"context"
"errors"
"fmt"
"log/slog"
"net/http"
"strings"
"time"
"github.com/mmcdole/gofeed"
"github.com/oklog/ulid/v2"
)
// PodcastFeed is one row from podcast_feeds joined to its media_items title.
// Only the fields the refresher needs are populated.
type PodcastFeed struct {
// MediaItemID is podcast_feeds.media_item_id — the FK to media_items.content_id.
// Used as the series_id when upserting episode rows.
MediaItemID string
// FeedURL is the RSS / Atom subscription URL.
FeedURL string
// RefreshIntervalSeconds controls how often the feed should be polled.
// Zero or negative defaults to 6 hours.
RefreshIntervalSeconds int
// LastRefreshedAt is nil for feeds that have never been refreshed.
LastRefreshedAt *time.Time
}
// PodcastEpisode is the data written to the episodes table for one RSS item.
type PodcastEpisode struct {
// ContentID is the episodes.content_id (ULID). New episodes get a
// fresh ULID; re-emitted episodes reuse the stored ID so foreign keys
// from progress tables survive the refresh.
ContentID string
SeriesID string // = PodcastFeed.MediaItemID
GUID string // stored in episodes.podcast_guid
Title string
Overview string
AudioURL string // stored in episodes.podcast_audio_url
DurationSeconds int // stored in episodes.runtime (minutes are fine; we store seconds via adapter)
EpisodeNumber int // 0 when not present in the feed
SeasonNumber int // 0 when not present in the feed
PublishedAt *time.Time
StillPath string // episode cover image URL (remote)
}
// Store is the narrow database surface the Refresher needs. Implemented by
// *DBStore (this package); surfaced as an interface so tests inject a stub
// without Postgres.
type Store interface {
// ListPodcastFeeds returns all rows from podcast_feeds. A feed with
// an empty FeedURL must still be returned — RefreshDue will skip it.
ListPodcastFeeds(ctx context.Context) ([]PodcastFeed, error)
// GetEpisodeIDsByGUID returns a map of guid → content_id for every
// episode in the given series whose podcast_guid matches any element
// of guids. Used to reuse stored IDs across feed refreshes.
GetEpisodeIDsByGUID(ctx context.Context, seriesID string, guids []string) (map[string]string, error)
// UpsertPodcastEpisode inserts or updates an episode keyed by
// (series_id, podcast_guid). On conflict the mutable fields are
// updated; the content_id is preserved.
UpsertPodcastEpisode(ctx context.Context, e PodcastEpisode) error
// MarkFeedRefreshed records a refresh attempt on podcast_feeds,
// setting last_refreshed_at = now() and last_refresh_error to the
// supplied string (empty string clears the error).
MarkFeedRefreshed(ctx context.Context, mediaItemID string, lastError string) error
}
// Refresher is the long-lived feed-refresh worker. One instance is created
// by the task scheduler and another can be created on demand by the admin
// force-refresh endpoint — both share the same HTTP client + parser.
type Refresher struct {
hc *http.Client
parser *gofeed.Parser
}
// New builds a Refresher with a 30-second HTTP timeout. Podcast feeds are
// typically <1 MB, but the parser reads the whole document so a hard
// timeout is the safety net.
func New() *Refresher {
return &Refresher{
hc: &http.Client{Timeout: 30 * time.Second},
parser: gofeed.NewParser(),
}
}
// WithHTTPClient overrides the default HTTP client. Used in tests to point
// at httptest.NewServer fixtures.
func (r *Refresher) WithHTTPClient(hc *http.Client) *Refresher {
r.hc = hc
return r
}
// RefreshDue walks every podcast_feeds row and refreshes those whose
// last_refreshed_at + refresh_interval has elapsed. Per-feed failures are
// logged at Warn and do not abort the walk. Returns the count of feeds
// that were attempted.
func (r *Refresher) RefreshDue(ctx context.Context, s Store) (int, error) {
feeds, err := s.ListPodcastFeeds(ctx)
if err != nil {
return 0, fmt.Errorf("list podcast feeds: %w", err)
}
now := time.Now()
attempted := 0
for _, f := range feeds {
if f.FeedURL == "" {
continue
}
if !isDue(f, now) {
continue
}
attempted++
if err := r.RefreshOne(ctx, s, f); err != nil {
slog.Warn("podcast feed refresh failed",
"media_item_id", f.MediaItemID,
"feed_url", f.FeedURL,
"err", err.Error(),
)
}
}
return attempted, nil
}
// RefreshOne refreshes a single podcast feed. Public so the admin
// force-refresh endpoint can call it directly. Records success or failure
// in podcast_feeds.last_refresh_error regardless of outcome.
func (r *Refresher) RefreshOne(ctx context.Context, s Store, f PodcastFeed) error {
if f.FeedURL == "" {
_ = s.MarkFeedRefreshed(ctx, f.MediaItemID, "no feed_url configured")
return errors.New("no feed_url configured")
}
feed, err := r.fetchAndParse(ctx, f.FeedURL)
if err != nil {
_ = s.MarkFeedRefreshed(ctx, f.MediaItemID, err.Error())
return err
}
count, err := r.upsertItems(ctx, s, f, feed)
if err != nil {
_ = s.MarkFeedRefreshed(ctx, f.MediaItemID, err.Error())
return err
}
_ = s.MarkFeedRefreshed(ctx, f.MediaItemID, "")
slog.Debug("podcast feed refreshed",
"media_item_id", f.MediaItemID,
"episodes", count,
)
return nil
}
func (r *Refresher) fetchAndParse(ctx context.Context, feedURL string) (*gofeed.Feed, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, feedURL, nil)
if err != nil {
return nil, fmt.Errorf("new request: %w", err)
}
// Some feed hosts gate on a recognisable User-Agent.
req.Header.Set("User-Agent", "silo/podcast-refresher (+https://siloapp.com)")
req.Header.Set("Accept", "application/rss+xml, application/atom+xml, application/xml;q=0.9, */*;q=0.5")
resp, err := r.hc.Do(req)
if err != nil {
return nil, fmt.Errorf("fetch %q: %w", feedURL, err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("fetch %q: status %d", feedURL, resp.StatusCode)
}
feed, err := r.parser.Parse(resp.Body)
if err != nil {
return nil, fmt.Errorf("parse %q: %w", feedURL, err)
}
return feed, nil
}
func (r *Refresher) upsertItems(ctx context.Context, s Store, f PodcastFeed, feed *gofeed.Feed) (int, error) {
if feed == nil || len(feed.Items) == 0 {
return 0, nil
}
// Collect GUIDs so we can look up existing content IDs in bulk and
// reuse them — preserving FK targets for progress rows.
guids := make([]string, 0, len(feed.Items))
for _, item := range feed.Items {
g := strings.TrimSpace(item.GUID)
if g == "" {
g = strings.TrimSpace(item.Link)
}
if g == "" {
continue
}
guids = append(guids, g)
}
existing, err := s.GetEpisodeIDsByGUID(ctx, f.MediaItemID, guids)
if err != nil {
return 0, fmt.Errorf("guid lookup: %w", err)
}
count := 0
for _, item := range feed.Items {
guid := strings.TrimSpace(item.GUID)
if guid == "" {
guid = strings.TrimSpace(item.Link)
}
if guid == "" || item.Title == "" {
continue
}
audioURL, _, _ := pickEnclosure(item)
if audioURL == "" {
// Text-only post or video — skip.
continue
}
contentID, ok := existing[guid]
if !ok {
contentID = ulid.Make().String()
}
ep := PodcastEpisode{
ContentID: contentID,
SeriesID: f.MediaItemID,
GUID: guid,
Title: item.Title,
Overview: item.Description,
AudioURL: audioURL,
DurationSeconds: durationFromItem(item),
EpisodeNumber: intFromItem(item, "episode"),
SeasonNumber: intFromItem(item, "season"),
PublishedAt: item.PublishedParsed,
StillPath: coverFromItem(item),
}
if err := s.UpsertPodcastEpisode(ctx, ep); err != nil {
slog.Warn("podcast episode upsert failed",
"media_item_id", f.MediaItemID,
"guid", guid,
"err", err.Error(),
)
continue
}
count++
}
return count, nil
}
// isDue reports whether a feed's refresh window has elapsed.
func isDue(f PodcastFeed, now time.Time) bool {
if f.LastRefreshedAt == nil {
return true
}
interval := time.Duration(f.RefreshIntervalSeconds) * time.Second
if interval <= 0 {
interval = 6 * time.Hour
}
return now.Sub(*f.LastRefreshedAt) >= interval
}
// pickEnclosure returns the first audio enclosure URL + MIME type for the
// given feed item. Returns empty strings when no audio enclosure is found.
func pickEnclosure(item *gofeed.Item) (audioURL, mimeType string, audioBytes int64) {
for _, enc := range item.Enclosures {
mt := strings.ToLower(enc.Type)
if mt != "" && !strings.HasPrefix(mt, "audio/") {
continue
}
audioBytes = 0
if enc.Length != "" {
var n int64
_, _ = fmt.Sscan(enc.Length, &n)
audioBytes = n
}
return enc.URL, enc.Type, audioBytes
}
return "", "", 0
}
// durationFromItem reads the iTunes-namespaced duration field from a gofeed
// item, parsing "HH:MM:SS" / "MM:SS" / "SSS" forms. Returns 0 on missing or
// unparseable input.
func durationFromItem(item *gofeed.Item) int {
if item.ITunesExt == nil {
return 0
}
raw := strings.TrimSpace(item.ITunesExt.Duration)
if raw == "" {
return 0
}
return parseDuration(raw)
}
func parseDuration(raw string) int {
parts := strings.Split(raw, ":")
var h, m, s int
switch len(parts) {
case 1:
_, _ = fmt.Sscan(parts[0], &s)
case 2:
_, _ = fmt.Sscan(parts[0], &m)
_, _ = fmt.Sscan(parts[1], &s)
case 3:
_, _ = fmt.Sscan(parts[0], &h)
_, _ = fmt.Sscan(parts[1], &m)
_, _ = fmt.Sscan(parts[2], &s)
default:
return 0
}
return h*3600 + m*60 + s
}
// intFromItem returns the iTunes-namespaced season or episode number as an
// int, or 0 if absent or unparseable.
func intFromItem(item *gofeed.Item, field string) int {
if item.ITunesExt == nil {
return 0
}
var raw string
switch field {
case "episode":
raw = item.ITunesExt.Episode
case "season":
raw = item.ITunesExt.Season
default:
return 0
}
raw = strings.TrimSpace(raw)
if raw == "" {
return 0
}
var n int
if _, err := fmt.Sscan(raw, &n); err != nil {
return 0
}
return n
}
// coverFromItem extracts an episode-level cover image URL, falling back to
// the iTunes image extension when present.
func coverFromItem(item *gofeed.Item) string {
if item.Image != nil && item.Image.URL != "" {
return item.Image.URL
}
if item.ITunesExt != nil && item.ITunesExt.Image != "" {
return item.ITunesExt.Image
}
return ""
}
@@ -0,0 +1,285 @@
package podcastfeed_test
import (
"context"
"net/http"
"net/http/httptest"
"sync"
"testing"
"time"
"github.com/Silo-Server/silo-server/internal/audiobooks/podcastfeed"
)
// fakeStore implements podcastfeed.Store for unit tests. It records every
// UpsertPodcastEpisode call so assertions can inspect what the refresher
// would have written without standing up Postgres.
type fakeStore struct {
mu sync.Mutex
feeds []podcastfeed.PodcastFeed
existingByGUID map[string]string
upsertedEpisodes []podcastfeed.PodcastEpisode
refreshed map[string]string // media_item_id → last_error
}
func newFakeStore(feeds ...podcastfeed.PodcastFeed) *fakeStore {
return &fakeStore{
feeds: feeds,
existingByGUID: map[string]string{},
refreshed: map[string]string{},
}
}
func (f *fakeStore) ListPodcastFeeds(_ context.Context) ([]podcastfeed.PodcastFeed, error) {
f.mu.Lock()
defer f.mu.Unlock()
return append([]podcastfeed.PodcastFeed(nil), f.feeds...), nil
}
func (f *fakeStore) GetEpisodeIDsByGUID(_ context.Context, _ string, guids []string) (map[string]string, error) {
f.mu.Lock()
defer f.mu.Unlock()
out := map[string]string{}
for _, g := range guids {
if id, ok := f.existingByGUID[g]; ok {
out[g] = id
}
}
return out, nil
}
func (f *fakeStore) UpsertPodcastEpisode(_ context.Context, e podcastfeed.PodcastEpisode) error {
f.mu.Lock()
defer f.mu.Unlock()
f.upsertedEpisodes = append(f.upsertedEpisodes, e)
return nil
}
func (f *fakeStore) MarkFeedRefreshed(_ context.Context, mediaItemID string, lastError string) error {
f.mu.Lock()
defer f.mu.Unlock()
f.refreshed[mediaItemID] = lastError
return nil
}
// rssFixture is a small but realistic RSS 2.0 + iTunes-namespaced feed
// covering the fields the refresher actually reads.
const rssFixture = `<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:itunes="http://www.itunes.com/dtds/podcast-1.0.dtd">
<channel>
<title>Test Show</title>
<itunes:author>Host McHostface</itunes:author>
<item>
<title>Episode 1</title>
<description>First one.</description>
<pubDate>Mon, 01 Apr 2026 12:00:00 GMT</pubDate>
<guid>episode-guid-1</guid>
<enclosure url="https://cdn.example.com/ep1.mp3" length="12345" type="audio/mpeg"/>
<itunes:duration>00:30:00</itunes:duration>
<itunes:episode>1</itunes:episode>
<itunes:season>1</itunes:season>
</item>
<item>
<title>Episode 2</title>
<description>Second one.</description>
<pubDate>Mon, 08 Apr 2026 12:00:00 GMT</pubDate>
<guid>episode-guid-2</guid>
<enclosure url="https://cdn.example.com/ep2.mp3" length="23456" type="audio/mpeg"/>
<itunes:duration>45:30</itunes:duration>
<itunes:episode>2</itunes:episode>
</item>
</channel>
</rss>`
// TestRefreshOne_InsertsNewEpisodes drives the refresher against a fake
// feed and confirms it upserts every item with the expected fields.
func TestRefreshOne_InsertsNewEpisodes(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/rss+xml")
_, _ = w.Write([]byte(rssFixture))
}))
defer srv.Close()
feed := podcastfeed.PodcastFeed{MediaItemID: "mi-1", FeedURL: srv.URL}
fs := newFakeStore(feed)
r := podcastfeed.New().WithHTTPClient(srv.Client())
if err := r.RefreshOne(context.Background(), fs, feed); err != nil {
t.Fatalf("RefreshOne: %v", err)
}
if len(fs.upsertedEpisodes) != 2 {
t.Fatalf("upserts = %d, want 2", len(fs.upsertedEpisodes))
}
first := fs.upsertedEpisodes[0]
if first.Title != "Episode 1" {
t.Errorf("first.Title = %q", first.Title)
}
if first.AudioURL != "https://cdn.example.com/ep1.mp3" {
t.Errorf("first.AudioURL = %q", first.AudioURL)
}
if first.GUID != "episode-guid-1" {
t.Errorf("first.GUID = %q", first.GUID)
}
if first.DurationSeconds != 1800 {
t.Errorf("first.DurationSeconds = %d, want 1800 (00:30:00)", first.DurationSeconds)
}
if first.EpisodeNumber != 1 {
t.Errorf("first.EpisodeNumber = %d, want 1", first.EpisodeNumber)
}
if first.SeasonNumber != 1 {
t.Errorf("first.SeasonNumber = %d, want 1", first.SeasonNumber)
}
if first.PublishedAt == nil {
t.Errorf("first.PublishedAt must be parsed")
}
// 45:30 mm:ss must parse to 2730s on the second item.
if fs.upsertedEpisodes[1].DurationSeconds != 2730 {
t.Errorf("second.DurationSeconds = %d, want 2730 (45:30)", fs.upsertedEpisodes[1].DurationSeconds)
}
// Mark-refreshed bookkeeping must record success (empty last_error).
if got := fs.refreshed["mi-1"]; got != "" {
t.Errorf("refreshed[mi-1] = %q, want empty (success)", got)
}
}
// TestRefreshOne_ReusesExistingEpisodeID confirms the idempotent-upsert
// contract: when a feed re-emits an item we've already stored, we keep
// the existing content_id so per-user progress rows don't lose their FK.
func TestRefreshOne_ReusesExistingEpisodeID(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/rss+xml")
_, _ = w.Write([]byte(rssFixture))
}))
defer srv.Close()
feed := podcastfeed.PodcastFeed{MediaItemID: "mi-1", FeedURL: srv.URL}
fs := newFakeStore(feed)
// Pretend episode 1 already exists with a stored id.
fs.existingByGUID["episode-guid-1"] = "stored-ulid-for-ep1"
r := podcastfeed.New().WithHTTPClient(srv.Client())
if err := r.RefreshOne(context.Background(), fs, feed); err != nil {
t.Fatalf("RefreshOne: %v", err)
}
if len(fs.upsertedEpisodes) != 2 {
t.Fatalf("upserts = %d, want 2", len(fs.upsertedEpisodes))
}
if fs.upsertedEpisodes[0].ContentID != "stored-ulid-for-ep1" {
t.Errorf("existing episode id was rotated: got %q, want stored-ulid-for-ep1",
fs.upsertedEpisodes[0].ContentID)
}
// Episode 2 is new — must mint a non-empty id, but not the stored one.
if fs.upsertedEpisodes[1].ContentID == "" || fs.upsertedEpisodes[1].ContentID == "stored-ulid-for-ep1" {
t.Errorf("new episode id = %q (must be a fresh ULID)", fs.upsertedEpisodes[1].ContentID)
}
}
// TestRefreshOne_UpstreamFailure records the error in the feed's
// last_refresh_error column so operators see the cause without reading logs.
func TestRefreshOne_UpstreamFailure(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
http.Error(w, "feed gone", http.StatusNotFound)
}))
defer srv.Close()
feed := podcastfeed.PodcastFeed{MediaItemID: "mi-1", FeedURL: srv.URL}
fs := newFakeStore(feed)
r := podcastfeed.New().WithHTTPClient(srv.Client())
err := r.RefreshOne(context.Background(), fs, feed)
if err == nil {
t.Fatal("RefreshOne must surface upstream 404 as an error")
}
if fs.refreshed["mi-1"] == "" {
t.Errorf("MarkFeedRefreshed not called with the error message")
}
}
// TestRefreshDue_OnlyWalksDuePodcasts confirms the refresh-interval gate:
// a feed refreshed two minutes ago with a 6-hour interval must not be
// re-fetched.
func TestRefreshDue_OnlyWalksDuePodcasts(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/rss+xml")
_, _ = w.Write([]byte(rssFixture))
}))
defer srv.Close()
twoMinAgo := time.Now().Add(-2 * time.Minute)
hoursAgo := time.Now().Add(-12 * time.Hour)
fs := newFakeStore(
podcastfeed.PodcastFeed{
MediaItemID: "due",
FeedURL: srv.URL,
RefreshIntervalSeconds: 21600, // 6 hours
LastRefreshedAt: &hoursAgo,
},
podcastfeed.PodcastFeed{
MediaItemID: "fresh",
FeedURL: srv.URL,
RefreshIntervalSeconds: 21600,
LastRefreshedAt: &twoMinAgo,
},
podcastfeed.PodcastFeed{
MediaItemID: "never",
FeedURL: srv.URL,
// LastRefreshedAt is nil → always due
},
podcastfeed.PodcastFeed{
MediaItemID: "noFeed",
// FeedURL is empty → skipped without attempting
},
)
r := podcastfeed.New().WithHTTPClient(srv.Client())
attempted, err := r.RefreshDue(context.Background(), fs)
if err != nil {
t.Fatalf("RefreshDue: %v", err)
}
// "due" + "never" → 2 attempts. "fresh" + "noFeed" → skipped.
if attempted != 2 {
t.Errorf("attempted = %d, want 2", attempted)
}
if _, ok := fs.refreshed["fresh"]; ok {
t.Errorf("fresh feed was refreshed when it should not have been")
}
if _, ok := fs.refreshed["due"]; !ok {
t.Errorf("due feed was not refreshed")
}
if _, ok := fs.refreshed["never"]; !ok {
t.Errorf("never-refreshed feed was not refreshed")
}
}
// TestRefreshOne_SkipsItemsWithoutAudio drops feed items that have no
// audio enclosure (text-only posts that some podcasts mix in).
func TestRefreshOne_SkipsItemsWithoutAudio(t *testing.T) {
const mixedFixture = `<?xml version="1.0"?>
<rss version="2.0"><channel><title>Mixed</title>
<item><title>Text-only</title><guid>t-1</guid></item>
<item><title>Has audio</title><guid>a-1</guid>
<enclosure url="https://cdn.example.com/a.mp3" type="audio/mpeg"/>
</item></channel></rss>`
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
_, _ = w.Write([]byte(mixedFixture))
}))
defer srv.Close()
feed := podcastfeed.PodcastFeed{MediaItemID: "mi-p", FeedURL: srv.URL}
fs := newFakeStore(feed)
r := podcastfeed.New().WithHTTPClient(srv.Client())
if err := r.RefreshOne(context.Background(), fs, feed); err != nil {
t.Fatalf("RefreshOne: %v", err)
}
if len(fs.upsertedEpisodes) != 1 {
t.Fatalf("upserts = %d, want 1 (text-only item skipped)", len(fs.upsertedEpisodes))
}
if fs.upsertedEpisodes[0].Title != "Has audio" {
t.Errorf("wrong item kept: %q", fs.upsertedEpisodes[0].Title)
}
}
+144
View File
@@ -0,0 +1,144 @@
package podcastfeed
import (
"context"
"fmt"
"github.com/jackc/pgx/v5/pgxpool"
)
// DBStore implements Store backed by silo's Postgres pool.
type DBStore struct {
pool *pgxpool.Pool
}
// NewDBStore constructs a production Store.
func NewDBStore(pool *pgxpool.Pool) *DBStore {
return &DBStore{pool: pool}
}
// ListPodcastFeeds returns all rows from podcast_feeds.
func (s *DBStore) ListPodcastFeeds(ctx context.Context) ([]PodcastFeed, error) {
rows, err := s.pool.Query(ctx, `
SELECT media_item_id,
feed_url,
refresh_interval_seconds,
last_refreshed_at
FROM podcast_feeds
ORDER BY last_refreshed_at ASC NULLS FIRST
`)
if err != nil {
return nil, fmt.Errorf("list podcast feeds: %w", err)
}
defer rows.Close()
var out []PodcastFeed
for rows.Next() {
var f PodcastFeed
if err := rows.Scan(
&f.MediaItemID,
&f.FeedURL,
&f.RefreshIntervalSeconds,
&f.LastRefreshedAt,
); err != nil {
return nil, fmt.Errorf("scan podcast feed: %w", err)
}
out = append(out, f)
}
return out, rows.Err()
}
// GetEpisodeIDsByGUID returns a map of podcast_guid → content_id for
// episodes in the given series that match any of the supplied GUIDs.
func (s *DBStore) GetEpisodeIDsByGUID(ctx context.Context, seriesID string, guids []string) (map[string]string, error) {
out := make(map[string]string, len(guids))
if seriesID == "" || len(guids) == 0 {
return out, nil
}
rows, err := s.pool.Query(ctx, `
SELECT content_id, podcast_guid
FROM episodes
WHERE series_id = $1
AND podcast_guid = ANY($2::text[])
`, seriesID, guids)
if err != nil {
return nil, fmt.Errorf("guid lookup: %w", err)
}
defer rows.Close()
for rows.Next() {
var id, guid string
if err := rows.Scan(&id, &guid); err != nil {
return nil, fmt.Errorf("scan episode guid: %w", err)
}
out[guid] = id
}
return out, rows.Err()
}
// UpsertPodcastEpisode inserts or updates an episode keyed by
// (series_id, podcast_guid). The season / episode number columns use 0
// for podcast episodes that lack a structured number in the feed.
//
// runtime is stored in seconds (the column name comes from TV; for
// podcasts we store raw seconds to preserve fidelity — ABS and the
// native API read it back as seconds).
func (s *DBStore) UpsertPodcastEpisode(ctx context.Context, e PodcastEpisode) error {
if e.ContentID == "" || e.SeriesID == "" || e.GUID == "" || e.Title == "" {
return fmt.Errorf("content_id, series_id, guid, title required")
}
_, err := s.pool.Exec(ctx, `
INSERT INTO episodes (
content_id, series_id,
season_number, episode_number,
title, overview, air_date, runtime,
still_path,
podcast_guid, podcast_audio_url
) VALUES (
$1, $2,
$3, $4,
$5, NULLIF($6,''), $7, $8,
NULLIF($9,''),
$10, NULLIF($11,'')
)
ON CONFLICT (series_id, podcast_guid) WHERE podcast_guid IS NOT NULL
DO UPDATE SET
title = EXCLUDED.title,
overview = EXCLUDED.overview,
air_date = EXCLUDED.air_date,
runtime = EXCLUDED.runtime,
still_path = EXCLUDED.still_path,
podcast_audio_url = EXCLUDED.podcast_audio_url,
updated_at = NOW()
`,
e.ContentID,
e.SeriesID,
e.SeasonNumber,
e.EpisodeNumber,
e.Title,
e.Overview,
e.PublishedAt,
e.DurationSeconds,
e.StillPath,
e.GUID,
e.AudioURL,
)
if err != nil {
return fmt.Errorf("upsert podcast episode: %w", err)
}
return nil
}
// MarkFeedRefreshed sets last_refreshed_at = now() and last_refresh_error
// on the podcast_feeds row. An empty lastError clears the error field.
func (s *DBStore) MarkFeedRefreshed(ctx context.Context, mediaItemID string, lastError string) error {
_, err := s.pool.Exec(ctx, `
UPDATE podcast_feeds
SET last_refreshed_at = now(),
last_refresh_error = NULLIF($2, ''),
updated_at = now()
WHERE media_item_id = $1
`, mediaItemID, lastError)
if err != nil {
return fmt.Errorf("mark feed refreshed: %w", err)
}
return nil
}
@@ -0,0 +1,54 @@
package tasks
import (
"context"
"encoding/json"
"fmt"
"github.com/Silo-Server/silo-server/internal/audiobooks/podcastfeed"
"github.com/Silo-Server/silo-server/internal/taskmanager"
)
// SyncPodcastFeedsTask refreshes RSS podcast feeds that are due for a poll.
// It wraps podcastfeed.Refresher so the task manager can invoke it on a
// schedule, report progress, and display it in the admin task panel.
type SyncPodcastFeedsTask struct {
refresher *podcastfeed.Refresher
store podcastfeed.Store
}
// NewSyncPodcastFeedsTask constructs the task. refresher and store come
// from the audiobooks subsystem wired in cmd/silo/main.go.
func NewSyncPodcastFeedsTask(refresher *podcastfeed.Refresher, store podcastfeed.Store) *SyncPodcastFeedsTask {
return &SyncPodcastFeedsTask{refresher: refresher, store: store}
}
func (t *SyncPodcastFeedsTask) Key() string { return "sync_podcast_feeds" }
func (t *SyncPodcastFeedsTask) Name() string { return "Sync Podcast Feeds" }
func (t *SyncPodcastFeedsTask) Description() string {
return "Refreshes RSS podcast feeds that are due for polling and upserts new episodes"
}
func (t *SyncPodcastFeedsTask) Category() taskmanager.TaskCategory {
return taskmanager.TaskCategoryLibrary
}
func (t *SyncPodcastFeedsTask) IsHidden() bool { return false }
func (t *SyncPodcastFeedsTask) DefaultTriggers() []taskmanager.TriggerConfig {
return []taskmanager.TriggerConfig{
{Type: taskmanager.TriggerTypeInterval, IntervalMs: 10 * 60 * 1000}, // every 10 minutes
}
}
func (t *SyncPodcastFeedsTask) Execute(ctx context.Context, progress taskmanager.ProgressReporter) error {
progress.Report(0, "Checking for due podcast feeds")
attempted, err := t.refresher.RefreshDue(ctx, t.store)
if err != nil {
return fmt.Errorf("podcast feed refresh: %w", err)
}
result, _ := json.Marshal(map[string]int{"feeds_attempted": attempted})
progress.SetResultData(result)
progress.Report(100, fmt.Sprintf("Podcast feed sync complete (%d feeds attempted)", attempted))
return nil
}
@@ -0,0 +1,5 @@
DROP INDEX IF EXISTS public.idx_episodes_podcast_guid;
ALTER TABLE public.episodes
DROP COLUMN IF EXISTS podcast_guid,
DROP COLUMN IF EXISTS podcast_audio_url;
@@ -0,0 +1,19 @@
-- Podcast-specific episode metadata. RSS episodes are identified by an
-- opaque GUID from the feed (not a sequential episode number), and they
-- carry a remote audio URL rather than a local file path.
--
-- podcast_guid — the RSS <guid> value; unique per show so the feed
-- refresher can upsert without duplicating.
-- podcast_audio_url — the remote enclosure URL for streaming / download.
--
-- Both columns are NULL for non-podcast episodes (movies/TV/audiobooks).
ALTER TABLE public.episodes
ADD COLUMN IF NOT EXISTS podcast_guid text,
ADD COLUMN IF NOT EXISTS podcast_audio_url text;
-- Unique constraint: one row per (series_id, podcast_guid) so the feed
-- refresher can ON CONFLICT (series_id, podcast_guid) DO UPDATE.
CREATE UNIQUE INDEX IF NOT EXISTS idx_episodes_podcast_guid
ON public.episodes (series_id, podcast_guid)
WHERE podcast_guid IS NOT NULL;