267 lines
8.0 KiB
Go
267 lines
8.0 KiB
Go
package main
|
|
|
|
import (
|
|
"fmt"
|
|
"net/url"
|
|
"os"
|
|
"os/signal"
|
|
"path/filepath"
|
|
"strings"
|
|
"syscall"
|
|
"time"
|
|
|
|
"rssd/config"
|
|
"rssd/downloader"
|
|
"rssd/poller"
|
|
"rssd/state"
|
|
)
|
|
|
|
const (
|
|
stateFile = "state.json"
|
|
)
|
|
|
|
func main() {
|
|
fmt.Fprintf(os.Stderr, "rssd: starting\n")
|
|
|
|
// Load configuration.
|
|
cfgPath := os.Getenv("RSSD_CONFIG")
|
|
if cfgPath == "" {
|
|
cfgPath = "~/.config/rssd/config.yaml"
|
|
}
|
|
cfgPath = config.ExpandPath(cfgPath) // resolve ~ early for consistent path usage
|
|
cfg, err := config.Load(cfgPath)
|
|
if err != nil {
|
|
fmt.Fprintf(os.Stderr, "rssd: fatal: %v\n", err)
|
|
os.Exit(1)
|
|
}
|
|
|
|
// Load state — store alongside the config file.
|
|
stateDir := filepath.Dir(cfgPath)
|
|
sm, err := state.NewManager(filepath.Join(stateDir, stateFile))
|
|
if err != nil {
|
|
fmt.Fprintf(os.Stderr, "rssd: fatal: %v\n", err)
|
|
os.Exit(1)
|
|
}
|
|
if err := sm.Load(); err != nil {
|
|
fmt.Fprintf(os.Stderr, "rssd: warning: could not load state (%v), starting fresh\n", err)
|
|
}
|
|
|
|
// Start download worker pool.
|
|
jobs := make(chan downloader.DownloadJob)
|
|
pool := downloader.NewPool(cfg.Settings.DownloadPoolSize, cfg.Settings.BackoffBase, sm)
|
|
go pool.Run(jobs)
|
|
|
|
// Re-enqueue any "pending" jobs from a previous run that were never picked up by workers.
|
|
sm.ForEachJob(func(j *state.JobSnapshot) {
|
|
if j.Status == state.StatusPending {
|
|
fmt.Fprintf(os.Stderr, "rssd: re-enqueued pending job %s\n", j.EnclosureURL)
|
|
jobs <- downloader.DownloadJob{
|
|
FeedURL: j.FeedURL,
|
|
EnclosureURL: j.EnclosureURL,
|
|
DestPath: j.DestPath,
|
|
}
|
|
}
|
|
})
|
|
|
|
fmt.Fprintf(os.Stderr, "rssd: workers started, polling feeds\n")
|
|
sigCh := make(chan os.Signal, 1)
|
|
signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT)
|
|
|
|
fmt.Fprintf(os.Stderr, "rssd: %d feeds configured, polling started\n", len(cfg.Feeds))
|
|
|
|
// Launch a poller goroutine per feed.
|
|
for _, feed := range cfg.Feeds {
|
|
go pollFeed(feed, sm, jobs, cfg)
|
|
}
|
|
|
|
// Periodically dispatch retrying jobs whose backoff has elapsed.
|
|
retryTicker := time.NewTicker(1 * time.Minute)
|
|
defer retryTicker.Stop()
|
|
go func() {
|
|
for range retryTicker.C {
|
|
dispatchRetries(sm, jobs, cfg.Settings.MaxRetries)
|
|
}
|
|
}()
|
|
|
|
// Wait for shutdown signal.
|
|
sig := <-sigCh
|
|
fmt.Fprintf(os.Stderr, "rssd: received %v, shutting down...\n", sig)
|
|
|
|
// Close jobs channel to stop accepting new work. Workers will finish in-flight downloads
|
|
// but we don't wait — just exit and let the next run handle any remaining state.
|
|
close(jobs)
|
|
|
|
// Save final state before exiting so no progress is lost.
|
|
if err := sm.Save(); err != nil {
|
|
fmt.Fprintf(os.Stderr, "rssd: error saving final state: %v\n", err)
|
|
}
|
|
fmt.Fprintf(os.Stderr, "rssd: stopped\n")
|
|
}
|
|
|
|
func pollFeed(feed config.Feed, sm *state.Manager, jobs chan<- downloader.DownloadJob, cfg *config.Config) {
|
|
fetcher := poller.NewFetcher()
|
|
interval := cfg.PollIntervalFor(feed)
|
|
|
|
fmt.Fprintf(os.Stderr, "rssd: polling %s every %s\n", feed.URL, interval)
|
|
|
|
// Poll immediately on startup.
|
|
pollAndEnqueue(feed, sm, jobs, fetcher, cfg)
|
|
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
|
|
for range ticker.C {
|
|
pollAndEnqueue(feed, sm, jobs, fetcher, cfg)
|
|
}
|
|
}
|
|
|
|
func pollAndEnqueue(feed config.Feed, sm *state.Manager, jobs chan<- downloader.DownloadJob, fetcher *poller.Fetcher, cfg *config.Config) {
|
|
feedState := sm.FindFeedByURL(feed.URL)
|
|
var lastModified, etag string
|
|
if feedState != nil {
|
|
lastModified = feedState.LastModified
|
|
etag = feedState.ETag
|
|
}
|
|
|
|
result, err := fetcher.Poll(feed.URL, lastModified, etag)
|
|
if err != nil {
|
|
fmt.Fprintf(os.Stderr, "rssd: error polling %s: %v\n", feed.URL, err)
|
|
return
|
|
}
|
|
if result.NotModified {
|
|
fmt.Fprintf(os.Stderr, "rssd: %s unchanged (skipping)\n", feed.URL)
|
|
return
|
|
}
|
|
|
|
fmt.Fprintf(os.Stderr, "rssd: %s polled — %d enclosures found\n", feed.URL, len(result.NewEnclosures))
|
|
|
|
if len(result.NewEnclosures) == 0 {
|
|
fmt.Fprintf(os.Stderr, "rssd: warning: %s has no enclosures in latest items\n", feed.URL)
|
|
}
|
|
|
|
// Save updated feed metadata for next conditional request.
|
|
sm.UpsertFeed(feed.URL, result.LastModified, result.ETag)
|
|
|
|
// Enqueue new enclosures (filtered by start_date).
|
|
startDate := feed.PublishedAt()
|
|
for _, encInfo := range result.NewEnclosures {
|
|
// Skip items older than the configured start date.
|
|
if !startDate.IsZero() && (!encInfo.PublishedAt.IsZero() && encInfo.PublishedAt.Before(startDate)) {
|
|
fmt.Fprintf(os.Stderr, "rssd: skipping %s (published %s, before start_date %s)\n",
|
|
encInfo.URL,
|
|
encInfo.PublishedAt.Format("2006-01-02"),
|
|
feed.StartDate)
|
|
continue
|
|
}
|
|
|
|
destPath := computeDestPath(feed.OutputDir, encInfo)
|
|
if sm.AddJob(feed.URL, encInfo.URL, destPath, cfg.Settings.MaxRetries) {
|
|
fmt.Fprintf(os.Stderr, "rssd: enqueued %s -> %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)
|
|
}
|