From 412defb939ff19e3afc10cf5dcf55dba32a252e2 Mon Sep 17 00:00:00 2001 From: RXWatcher <14085001+RXWatcher@users.noreply.github.com> Date: Sun, 24 May 2026 16:06:27 +0200 Subject: [PATCH] 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) --- cmd/silo/main.go | 2 + go.mod | 7 + go.sum | 22 ++ internal/audiobooks/podcastfeed/refresher.go | 355 ++++++++++++++++++ .../audiobooks/podcastfeed/refresher_test.go | 285 ++++++++++++++ internal/audiobooks/podcastfeed/store.go | 144 +++++++ .../taskmanager/tasks/sync_podcast_feeds.go | 54 +++ migrations/144_podcast_episode_guid.down.sql | 5 + migrations/144_podcast_episode_guid.up.sql | 19 + 9 files changed, 893 insertions(+) create mode 100644 internal/audiobooks/podcastfeed/refresher.go create mode 100644 internal/audiobooks/podcastfeed/refresher_test.go create mode 100644 internal/audiobooks/podcastfeed/store.go create mode 100644 internal/taskmanager/tasks/sync_podcast_feeds.go create mode 100644 migrations/144_podcast_episode_guid.down.sql create mode 100644 migrations/144_podcast_episode_guid.up.sql diff --git a/cmd/silo/main.go b/cmd/silo/main.go index c088880f..3754e194 100644 --- a/cmd/silo/main.go +++ b/cmd/silo/main.go @@ -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 { diff --git a/go.mod b/go.mod index 5863e860..77596d3f 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/go.sum b/go.sum index 05fca771..9bcf8db5 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/internal/audiobooks/podcastfeed/refresher.go b/internal/audiobooks/podcastfeed/refresher.go new file mode 100644 index 00000000..a9fd8d5e --- /dev/null +++ b/internal/audiobooks/podcastfeed/refresher.go @@ -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 , 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 "" +} diff --git a/internal/audiobooks/podcastfeed/refresher_test.go b/internal/audiobooks/podcastfeed/refresher_test.go new file mode 100644 index 00000000..9ad6ded1 --- /dev/null +++ b/internal/audiobooks/podcastfeed/refresher_test.go @@ -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 = ` + + + Test Show + Host McHostface + + Episode 1 + First one. + Mon, 01 Apr 2026 12:00:00 GMT + episode-guid-1 + + 00:30:00 + 1 + 1 + + + Episode 2 + Second one. + Mon, 08 Apr 2026 12:00:00 GMT + episode-guid-2 + + 45:30 + 2 + + +` + +// 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 = ` +Mixed +Text-onlyt-1 +Has audioa-1 + +` + 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) + } +} diff --git a/internal/audiobooks/podcastfeed/store.go b/internal/audiobooks/podcastfeed/store.go new file mode 100644 index 00000000..45d76678 --- /dev/null +++ b/internal/audiobooks/podcastfeed/store.go @@ -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 +} diff --git a/internal/taskmanager/tasks/sync_podcast_feeds.go b/internal/taskmanager/tasks/sync_podcast_feeds.go new file mode 100644 index 00000000..2589c6bd --- /dev/null +++ b/internal/taskmanager/tasks/sync_podcast_feeds.go @@ -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 +} diff --git a/migrations/144_podcast_episode_guid.down.sql b/migrations/144_podcast_episode_guid.down.sql new file mode 100644 index 00000000..741249f8 --- /dev/null +++ b/migrations/144_podcast_episode_guid.down.sql @@ -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; diff --git a/migrations/144_podcast_episode_guid.up.sql b/migrations/144_podcast_episode_guid.up.sql new file mode 100644 index 00000000..229338ec --- /dev/null +++ b/migrations/144_podcast_episode_guid.up.sql @@ -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 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;