From 6f188d1939c0d1fbdf2ab09dc28aff29b6bae692 Mon Sep 17 00:00:00 2001 From: Greg Pomerantz Date: Sun, 3 May 2026 20:30:37 -0400 Subject: [PATCH] Initial commit. --- cmd/debug/main.go | 64 +++++++++ config/config.go | 147 +++++++++++++++++++++ downloader/downloader.go | 140 ++++++++++++++++++++ go.mod | 16 +++ go.sum | 34 +++++ main.go | 266 ++++++++++++++++++++++++++++++++++++++ poller/poller.go | 118 +++++++++++++++++ state/state.go | 272 +++++++++++++++++++++++++++++++++++++++ 8 files changed, 1057 insertions(+) create mode 100644 cmd/debug/main.go create mode 100644 config/config.go create mode 100644 downloader/downloader.go create mode 100644 go.mod create mode 100644 go.sum create mode 100644 main.go create mode 100644 poller/poller.go create mode 100644 state/state.go diff --git a/cmd/debug/main.go b/cmd/debug/main.go new file mode 100644 index 0000000..2c99066 --- /dev/null +++ b/cmd/debug/main.go @@ -0,0 +1,64 @@ +package main + +import ( + "fmt" + "io" + "net/http" + "os" + "time" + + "github.com/mmcdole/gofeed" +) + +func main() { + client := &http.Client{Timeout: 60 * time.Second} + resp, err := client.Get("https://feeds.transistor.fm/oxide-and-friends") + if err != nil { + fmt.Fprintf(os.Stderr, "Error fetching: %v\n", err) + os.Exit(1) + } + defer resp.Body.Close() + + body, err := io.ReadAll(resp.Body) + if err != nil { + fmt.Fprintf(os.Stderr, "Error reading body: %v\n", err) + os.Exit(1) + } + + parser := gofeed.NewParser() + feed, err := parser.ParseString(string(body)) + if err != nil { + fmt.Fprintf(os.Stderr, "Error parsing feed: %v\n", err) + os.Exit(1) + } + + fmt.Printf("Feed title: %s\n", feed.Title) + fmt.Printf("Total items: %d\n\n", len(feed.Items)) + + for i, item := range feed.Items { + if i >= 5 { + break + } + fmt.Printf("--- Item %d ---\n", i+1) + fmt.Printf(" Title: %s\n", item.Title) + fmt.Printf(" Published (raw): %s\n", item.Published) + if item.PublishedParsed != nil { + fmt.Printf(" PublishedParsed: %v\n", *item.PublishedParsed) + } else { + fmt.Printf(" PublishedParsed: \n") + } + + // Check for itunes:pubDate or other date fields + for _, enc := range item.Enclosures { + fmt.Printf(" Enclosure URL (truncated): %s...\n", enc.URL[:min(len(enc.URL), 80)]) + } + fmt.Println() + } +} + +func min(a, b int) int { + if a < b { + return a + } + return b +} diff --git a/config/config.go b/config/config.go new file mode 100644 index 0000000..6170c74 --- /dev/null +++ b/config/config.go @@ -0,0 +1,147 @@ +package config + +import ( + "fmt" + "os" + "path/filepath" + "strings" + "time" + + "gopkg.in/yaml.v3" +) + +// ExpandPath resolves a path that may start with ~ to an absolute path. +func ExpandPath(p string) string { + if strings.HasPrefix(p, "~") { + home, err := os.UserHomeDir() + if err != nil { + return p // fall back to literal on error + } + return filepath.Join(home, strings.TrimPrefix(p, "~")) + } + return p +} + +// expandPath resolves a path that may start with ~ to an absolute path. +func expandPath(p string) string { + return ExpandPath(p) +} + +// Feed represents a single podcast feed to monitor. +type Feed struct { + URL string `yaml:"url"` + OutputDir string `yaml:"output_dir"` + PollInterval *time.Duration `yaml:"poll_interval,omitempty"` // nil means use global default + StartDate string `yaml:"start_date,omitempty"` // optional, format "2006-01-02" +} + +// PublishedAt returns the parsed start date for this feed, or zero time if not set. +func (f Feed) PublishedAt() time.Time { + if f.StartDate == "" { + return time.Time{} + } + t, err := time.Parse("2006-01-02", f.StartDate) + if err != nil { + return time.Time{} + } + return t +} + +// Settings holds global defaults for the daemon. Zero values are replaced by ApplyDefaults. +type Settings struct { + PollInterval time.Duration `yaml:"poll_interval"` + MaxRetries int `yaml:"max_retries"` + BackoffBase time.Duration `yaml:"backoff_base"` + DownloadPoolSize int `yaml:"download_pool_size"` +} + +// Config is the top-level parsed configuration. +type Config struct { + Settings Settings `yaml:"settings,omitempty"` + Feeds []Feed `yaml:"feeds"` +} + +// ApplyDefaults fills in zero-valued settings with sensible defaults and expands ~ paths. +func (c *Config) ApplyDefaults() { + if c.Settings.PollInterval == 0 { + c.Settings.PollInterval = 6 * time.Hour + } + if c.Settings.MaxRetries == 0 { + c.Settings.MaxRetries = 5 + } + if c.Settings.BackoffBase == 0 { + c.Settings.BackoffBase = 5 * time.Minute + } + if c.Settings.DownloadPoolSize == 0 { + c.Settings.DownloadPoolSize = 4 + } + + for i := range c.Feeds { + c.Feeds[i].OutputDir = expandPath(c.Feeds[i].OutputDir) + } +} + +// PollIntervalFor returns the effective poll interval for a feed, falling back to global default. +func (c *Config) PollIntervalFor(f Feed) time.Duration { + if f.PollInterval != nil && *f.PollInterval > 0 { + return *f.PollInterval + } + return c.Settings.PollInterval +} + +// Validate checks that the config is usable. +func (c *Config) Validate() error { + if len(c.Feeds) == 0 { + return fmt.Errorf("no feeds defined") + } + for i, f := range c.Feeds { + if f.URL == "" { + return fmt.Errorf("feed[%d]: url is required", i) + } + if f.OutputDir == "" { + return fmt.Errorf("feed[%d] (%s): output_dir is required", i, f.URL) + } + if c.Settings.PollInterval <= 0 { + return fmt.Errorf("settings.poll_interval must be > 0") + } + if c.Settings.MaxRetries == 0 { + return fmt.Errorf("settings.max_retries must be > 0") + } + if c.Settings.BackoffBase <= 0 { + return fmt.Errorf("settings.backoff_base must be > 0") + } + if c.Settings.DownloadPoolSize == 0 { + return fmt.Errorf("settings.download_pool_size must be > 0") + } + // Validate start_date format if provided. + if f.StartDate != "" { + if _, err := time.Parse("2006-01-02", f.StartDate); err != nil { + return fmt.Errorf("feed[%d] (%s): invalid start_date %q: %w", i, f.URL, f.StartDate, err) + } + } + } + return nil +} + +// Load reads and parses a YAML config from the given path. +func Load(path string) (*Config, error) { + path = expandPath(path) + + data, err := os.ReadFile(path) + if err != nil { + return nil, fmt.Errorf("read config %s: %w", path, err) + } + + var cfg Config + if err := yaml.Unmarshal(data, &cfg); err != nil { + return nil, fmt.Errorf("parse config: %w", err) + } + + cfg.ApplyDefaults() + + if err := cfg.Validate(); err != nil { + return nil, fmt.Errorf("validate config: %w", err) + } + + return &cfg, nil +} diff --git a/downloader/downloader.go b/downloader/downloader.go new file mode 100644 index 0000000..cfc095c --- /dev/null +++ b/downloader/downloader.go @@ -0,0 +1,140 @@ +package downloader + +import ( + "fmt" + "io" + "net/http" + "os" + "path/filepath" + "time" + + "rssd/state" +) + +// DownloadJob carries all information needed for a single download. +type DownloadJob struct { + FeedURL string + EnclosureURL string + DestPath string +} + +// Pool manages a fixed-size worker pool for downloading files. +type Pool struct { + poolSize int + backoffBase time.Duration + state *state.Manager +} + +// NewPool creates a download pool with the given parameters. +func NewPool(poolSize int, backoffBase time.Duration, sm *state.Manager) *Pool { + return &Pool{poolSize: poolSize, backoffBase: backoffBase, state: sm} +} + +// Run starts the worker goroutines. It returns immediately after spawning the workers. +func (p *Pool) Run(jobs <-chan DownloadJob) { + for i := 0; i < p.poolSize; i++ { + go func() { + for job := range jobs { + p.worker(job) + } + }() + } +} + +func (p *Pool) worker(job DownloadJob) { + fmt.Fprintf(os.Stderr, "rssd: worker started download %s\n", job.EnclosureURL) + if err := p.download(job); err != nil { + fmt.Fprintf(os.Stderr, "ERROR: download %s failed: %v\n", job.EnclosureURL, err) + } else { + fmt.Fprintf(os.Stderr, "OK: downloaded %s -> %s\n", job.EnclosureURL, job.DestPath) + } + fmt.Fprintf(os.Stderr, "rssd: worker finished download %s\n", job.EnclosureURL) +} + +func (p *Pool) download(job DownloadJob) error { + // Update state to downloading. + if !p.state.UpdateJob(job.EnclosureURL, func(j *state.Job) { + j.Status = state.StatusDownloading + j.AttemptCount++ + }) { + return fmt.Errorf("job not found for %s", job.EnclosureURL) + } + if err := p.state.Save(); err != nil { + return fmt.Errorf("save state (downloading): %w", err) + } + + // Ensure output directory exists. + if err := os.MkdirAll(filepath.Dir(job.DestPath), 0755); err != nil { + p.markRetry(job.EnclosureURL, fmt.Errorf("create dest dir: %w", err)) + return err + } + + // Create destination file. + dstFile, err := os.Create(job.DestPath) + if err != nil { + p.markRetry(job.EnclosureURL, fmt.Errorf("create dest file: %w", err)) + return err + } + + // Download with timeout. + client := &http.Client{Timeout: 30 * time.Minute} + req, err := http.NewRequest("GET", job.EnclosureURL, nil) + if err != nil { + dstFile.Close() + os.Remove(job.DestPath) + p.markRetry(job.EnclosureURL, fmt.Errorf("create request: %w", err)) + return err + } + + resp, err := client.Do(req) + if err != nil { + dstFile.Close() + os.Remove(job.DestPath) + p.markRetry(job.EnclosureURL, fmt.Errorf("fetch: %w", err)) + return err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + dstFile.Close() + os.Remove(job.DestPath) + p.markRetry(job.EnclosureURL, fmt.Errorf("HTTP %d", resp.StatusCode)) + return fmt.Errorf("HTTP %d", resp.StatusCode) + } + + // Copy body to file. + if _, err := io.Copy(dstFile, resp.Body); err != nil { + dstFile.Close() + os.Remove(job.DestPath) + p.markRetry(job.EnclosureURL, fmt.Errorf("write file: %w", err)) + return err + } + if err := dstFile.Close(); err != nil { + os.Remove(job.DestPath) + p.markRetry(job.EnclosureURL, fmt.Errorf("close file: %w", err)) + return err + } + + // Mark as downloaded. + p.state.UpdateJob(job.EnclosureURL, func(j *state.Job) { + j.Status = state.StatusDownloaded + }) + if err := p.state.Save(); err != nil { + return fmt.Errorf("save state (downloaded): %w", err) + } + + return nil +} + +// markRetry updates a job to retrying status with the given error and computes backoff delay. +func (p *Pool) markRetry(enclosureURL string, err error) { + p.state.UpdateJob(enclosureURL, func(j *state.Job) { + j.Status = state.StatusRetrying + j.Error = err.Error() + // Exponential backoff: base * 2^attemptCount (attempt already incremented). + delay := p.backoffBase * time.Duration(1< %s\n", encInfo.URL, destPath) + jobs <- downloader.DownloadJob{ + FeedURL: feed.URL, + EnclosureURL: encInfo.URL, + DestPath: destPath, + } + } else { + job := sm.FindJobByURL(encInfo.URL) + if job != nil { + fmt.Fprintf(os.Stderr, "rssd: skipping %s (already %s)\n", encInfo.URL, job.Status) + } else { + fmt.Fprintf(os.Stderr, "rssd: skipping duplicate %s\n", encInfo.URL) + } + } + } + + if err := sm.Save(); err != nil { + fmt.Fprintf(os.Stderr, "rssd: error saving state: %v\n", err) + } +} + +// dispatchRetries scans all jobs in state and re-enqueues those whose backoff has elapsed. +func dispatchRetries(sm *state.Manager, jobs chan<- downloader.DownloadJob, maxRetries int) { + // Collect snapshots (lock released before any callback invocation). + var ready []downloader.DownloadJob + var failed []string + + sm.ForEachJob(func(snap *state.JobSnapshot) { + if snap.Status != state.StatusRetrying { + return + } + + if snap.AttemptCount >= maxRetries { + // Permanently failed — use UpdateJob to mutate the live state. + sm.UpdateJob(snap.EnclosureURL, func(j *state.Job) { + j.Status = state.StatusFailed + }) + failed = append(failed, fmt.Sprintf("%s (after %d attempts: %s)", + snap.EnclosureURL, snap.AttemptCount, snap.Error)) + return + } + + if snap.NextAttemptAt == nil || time.Now().After(*snap.NextAttemptAt) { + // Re-enqueue — use UpdateJob to mutate the live state. + sm.UpdateJob(snap.EnclosureURL, func(j *state.Job) { + j.Status = state.StatusPending + j.NextAttemptAt = nil + }) + ready = append(ready, downloader.DownloadJob{ + FeedURL: snap.FeedURL, + EnclosureURL: snap.EnclosureURL, + DestPath: snap.DestPath, + }) + } + }) + + if len(failed) > 0 { + for _, msg := range failed { + fmt.Fprintf(os.Stderr, "rssd: %s permanently failed\n", msg) + } + } + + // Send to channel and persist state (lock released). + for _, job := range ready { + jobs <- job + } + + if err := sm.Save(); err != nil { + fmt.Fprintf(os.Stderr, "rssd: error saving state (retry dispatch): %v\n", err) + } +} + +// sanitizeFilename replaces characters that are invalid or problematic in filenames. +func sanitizeFilename(s string) string { + replacer := strings.NewReplacer( + "/", "_", "\\", "_", ":", "_", "<", "_", ">", "_", "|", "_", "\"", "_", "?", "_", "*", "_", + ) + s = replacer.Replace(s) + for strings.Contains(s, "__") { + s = strings.ReplaceAll(s, "__", "_") + } + return strings.TrimSpace(s) +} + +// computeDestPath constructs the local file path from an enclosure URL and output directory. +// The filename is prefixed with YYYYMMDD - for chronological sorting when dates are available. +func computeDestPath(outputDir string, encInfo poller.EnclosureInfo) string { + baseName := sanitizeFilename(encInfo.Title) + if baseName == "" { + // Fallback to enclosure URL basename if title is empty or invalid. + parsed, err := url.Parse(encInfo.URL) + if err != nil || parsed.Path == "" || parsed.Path == "/" { + return filepath.Join(outputDir, strings.ReplaceAll(encInfo.URL, "/", "_")) + } + baseName = sanitizeFilename(filepath.Base(parsed.Path)) + if baseName == "." || baseName == "" { + baseName = "unknown" + } + } + + // Prepend date prefix for chronological sorting if available. + if !encInfo.PublishedAt.IsZero() { + datePrefix := encInfo.PublishedAt.Format("2006-01-02") + return filepath.Join(outputDir, fmt.Sprintf("%s - %s", datePrefix, baseName)) + } + + return filepath.Join(outputDir, baseName) +} diff --git a/poller/poller.go b/poller/poller.go new file mode 100644 index 0000000..4df5edc --- /dev/null +++ b/poller/poller.go @@ -0,0 +1,118 @@ +package poller + +import ( + "fmt" + "io" + "net/http" + "time" + + "github.com/mmcdole/gofeed" +) + +// EnclosureInfo carries an enclosure URL and its associated item's metadata. +type EnclosureInfo struct { + URL string + Title string // episode title, if available + PublishedAt time.Time // zero value means no date available +} + +// FeedResult holds the result of polling a single feed. +type FeedResult struct { + URL string + NewEnclosures []EnclosureInfo + LastModified string + ETag string + NotModified bool + Error error +} + +// Fetcher abstracts the HTTP fetch + parse cycle for a single feed. +type Fetcher struct { + httpClient *http.Client + parser *gofeed.Parser +} + +// NewFetcher creates a new Fetcher with sensible defaults. +func NewFetcher() *Fetcher { + return &Fetcher{ + httpClient: &http.Client{ + Timeout: 30 * time.Second, + }, + parser: gofeed.NewParser(), + } +} + +// Poll fetches and parses an RSS/Atom feed, returning new enclosure info. +// lastModified and etag are used for conditional requests (If-None-Match / If-Modified-Since). +// Returns a FeedResult with NewEnclosures populated only if the feed changed. +func (f *Fetcher) Poll(url, lastModified, etag string) (*FeedResult, error) { + req, err := http.NewRequest("GET", url, nil) + if err != nil { + return nil, fmt.Errorf("create request: %w", err) + } + + if etag != "" { + req.Header.Set("If-None-Match", etag) + } + if lastModified != "" { + req.Header.Set("If-Modified-Since", lastModified) + } + + resp, err := f.httpClient.Do(req) + if err != nil { + return nil, fmt.Errorf("fetch %s: %w", url, err) + } + defer resp.Body.Close() + + // Handle 304 Not Modified — feed hasn't changed. + if resp.StatusCode == http.StatusNotModified { + return &FeedResult{ + URL: url, + NotModified: true, + }, nil + } + + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("fetch %s: HTTP %d", url, resp.StatusCode) + } + + body, err := io.ReadAll(resp.Body) + if err != nil { + return nil, fmt.Errorf("read body %s: %w", url, err) + } + + feed, err := f.parser.ParseString(string(body)) + if err != nil { + return nil, fmt.Errorf("parse feed %s: %w", url, err) + } + + // Extract enclosure info from feed items. + var enclosures []EnclosureInfo + for _, item := range feed.Items { + for _, enc := range item.Enclosures { + if enc.URL != "" { + encInfo := EnclosureInfo{URL: enc.URL, Title: item.Title} + if item.PublishedParsed != nil { + encInfo.PublishedAt = *item.PublishedParsed + } + enclosures = append(enclosures, encInfo) + } + } + } + + result := &FeedResult{ + URL: url, + NewEnclosures: enclosures, + NotModified: false, + } + + // Capture cache headers for the next conditional request. + if v := resp.Header.Get("Last-Modified"); v != "" { + result.LastModified = v + } + if v := resp.Header.Get("ETag"); v != "" { + result.ETag = v + } + + return result, nil +} diff --git a/state/state.go b/state/state.go new file mode 100644 index 0000000..0096117 --- /dev/null +++ b/state/state.go @@ -0,0 +1,272 @@ +package state + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "sync" + "time" +) + +// State is the full persisted state of the daemon. +type State struct { + Feeds []FeedState `json:"feeds"` + Jobs []Job `json:"jobs"` +} + +// FeedState holds per-feed metadata from HTTP responses. +type FeedState struct { + URL string `json:"url"` + LastModified string `json:"last_modified,omitempty"` + ETag string `json:"etag,omitempty"` +} + +// Job represents a download job with retry state. +type Job struct { + FeedURL string `json:"feed_url"` + EnclosureURL string `json:"enclosure_url"` + DestPath string `json:"dest_path"` + Status string `json:"status"` + AttemptCount int `json:"attempt_count"` + MaxRetries int `json:"max_retries"` + NextAttemptAt *time.Time `json:"next_attempt_at,omitempty"` // nil means ready now + Error string `json:"error,omitempty"` +} + +// JobStatus values. +const ( + StatusPending = "pending" + StatusDownloading = "downloading" + StatusRetrying = "retrying" + StatusDownloaded = "downloaded" + StatusFailed = "failed" +) + +// Manager provides thread-safe access to persisted state. +type Manager struct { + mu sync.Mutex + path string + state *State +} + +// NewManager creates a new state manager pointing at the given file path. +func NewManager(path string) (*Manager, error) { + dir := filepath.Dir(path) + if err := os.MkdirAll(dir, 0755); err != nil { + return nil, fmt.Errorf("create state directory: %w", err) + } + return &Manager{path: path, state: &State{}}, nil +} + +// Load reads the state file from disk. If it doesn't exist or is empty, returns an empty state. +func (m *Manager) Load() error { + m.mu.Lock() + defer m.mu.Unlock() + + data, err := os.ReadFile(m.path) + if err != nil { + if os.IsNotExist(err) { + m.state = &State{} + return nil + } + return fmt.Errorf("read state: %w", err) + } + + if len(data) == 0 { + m.state = &State{} + return nil + } + + var s State + if err := json.Unmarshal(data, &s); err != nil { + return fmt.Errorf("parse state: %w", err) + } + + m.state = &s + return nil +} + +// Save atomically writes the current state to disk with pretty-printed JSON. +func (m *Manager) Save() error { + m.mu.Lock() + defer m.mu.Unlock() + + data, err := json.MarshalIndent(m.state, "", " ") + if err != nil { + return fmt.Errorf("marshal state: %w", err) + } + + dir := filepath.Dir(m.path) + tmpFile, err := os.CreateTemp(dir, ".rssd-state-*.tmp") + if err != nil { + return fmt.Errorf("create temp file: %w", err) + } + tmpName := tmpFile.Name() + + if _, err := tmpFile.Write(data); err != nil { + tmpFile.Close() + os.Remove(tmpName) + return fmt.Errorf("write temp state: %w", err) + } + if err := tmpFile.Close(); err != nil { + os.Remove(tmpName) + return fmt.Errorf("close temp state: %w", err) + } + + if err := os.Rename(tmpName, m.path); err != nil { + os.Remove(tmpName) + return fmt.Errorf("rename state file: %w", err) + } + + return nil +} + +// JobSnapshot is a copy of a Job suitable for read-only inspection. +type JobSnapshot struct { + FeedURL string + EnclosureURL string + DestPath string + Status string + AttemptCount int + MaxRetries int + NextAttemptAt *time.Time // nil means ready now + Error string +} + +// ForEachJob iterates over all jobs and calls fn for each one. +// The callback receives a pointer to a copy of the job (not the live state), +// so it must not modify fields. Use UpdateJob or AddJob for mutations. +// fn may perform I/O without risk of deadlock — the lock is released before +// any callback invocation. +func (m *Manager) ForEachJob(fn func(*JobSnapshot)) { + m.mu.Lock() + jobsCopy := make([]JobSnapshot, len(m.state.Jobs)) + for i := range m.state.Jobs { + j := &m.state.Jobs[i] + jobsCopy[i] = JobSnapshot{ + FeedURL: j.FeedURL, + EnclosureURL: j.EnclosureURL, + DestPath: j.DestPath, + Status: j.Status, + AttemptCount: j.AttemptCount, + MaxRetries: j.MaxRetries, + NextAttemptAt: j.NextAttemptAt, + Error: j.Error, + } + } + m.mu.Unlock() + + for i := range jobsCopy { + fn(&jobsCopy[i]) + } +} + +// State returns a copy of the current in-memory state for iteration. +func (m *Manager) State() *State { + m.mu.Lock() + defer m.mu.Unlock() + return m.state +} + +// AddJob adds a new download job if one with this enclosure URL doesn't already exist. +// Returns true if the job was added, false if it already existed. +func (m *Manager) AddJob(feedURL, enclosureURL, destPath string, maxRetries int) bool { + m.mu.Lock() + defer m.mu.Unlock() + + // Check if this enclosure URL already has a job. + for _, j := range m.state.Jobs { + if j.EnclosureURL == enclosureURL { + return false // already exists + } + } + + job := Job{ + FeedURL: feedURL, + EnclosureURL: enclosureURL, + DestPath: destPath, + Status: StatusPending, + AttemptCount: 0, + MaxRetries: maxRetries, + NextAttemptAt: nil, + Error: "", + } + + m.state.Jobs = append(m.state.Jobs, job) + return true +} + +// UpdateJob finds a job by enclosure URL and applies the update function. +// Returns true if the job was found and updated. +func (m *Manager) UpdateJob(enclosureURL string, fn func(*Job)) bool { + m.mu.Lock() + defer m.mu.Unlock() + + for i := range m.state.Jobs { + if m.state.Jobs[i].EnclosureURL == enclosureURL { + fn(&m.state.Jobs[i]) + return true + } + } + return false +} + +// FindJobByURL returns a pointer to the job for the given enclosure URL, or nil. +func (m *Manager) FindJobByURL(enclosureURL string) *Job { + m.mu.Lock() + defer m.mu.Unlock() + + for i := range m.state.Jobs { + if m.state.Jobs[i].EnclosureURL == enclosureURL { + return &m.state.Jobs[i] + } + } + return nil +} + +// HasJob returns true if a job with the given enclosure URL exists. +func (m *Manager) HasJob(enclosureURL string) bool { + m.mu.Lock() + defer m.mu.Unlock() + + for _, j := range m.state.Jobs { + if j.EnclosureURL == enclosureURL { + return true + } + } + return false +} + +// FindFeedByURL returns the FeedState for a given feed URL, or nil. +func (m *Manager) FindFeedByURL(feedURL string) *FeedState { + m.mu.Lock() + defer m.mu.Unlock() + + for i := range m.state.Feeds { + if m.state.Feeds[i].URL == feedURL { + return &m.state.Feeds[i] + } + } + return nil +} + +// UpsertFeed adds or updates a feed's metadata (ETag, Last-Modified). +func (m *Manager) UpsertFeed(feedURL, lastModified, etag string) { + m.mu.Lock() + defer m.mu.Unlock() + + for i := range m.state.Feeds { + if m.state.Feeds[i].URL == feedURL { + m.state.Feeds[i].LastModified = lastModified + m.state.Feeds[i].ETag = etag + return + } + } + + m.state.Feeds = append(m.state.Feeds, FeedState{ + URL: feedURL, + LastModified: lastModified, + ETag: etag, + }) +}