diff --git a/cmd/podsync/main.go b/cmd/podsync/main.go index 8cb30c7..00e80cc 100644 --- a/cmd/podsync/main.go +++ b/cmd/podsync/main.go @@ -12,6 +12,7 @@ import ( "golang.org/x/sync/errgroup" "github.com/mxpv/podsync/pkg/config" + "github.com/mxpv/podsync/pkg/storage" "github.com/mxpv/podsync/pkg/ytdl" ) @@ -85,13 +86,18 @@ func main() { log.WithError(err).Fatal("failed to load configuration file") } + db, err := storage.NewBadger(opts.ConfigPath) + if err != nil { + log.WithError(err).Fatal("failed to open database") + } + // Queue of feeds to update updates := make(chan *config.Feed, 16) defer close(updates) // Run updater thread log.Debug("creating updater") - updater, err := NewUpdater(cfg, downloader) + updater, err := NewUpdater(cfg, downloader, db) if err != nil { log.WithError(err).Fatal("failed to create updater") } @@ -165,5 +171,9 @@ func main() { log.WithError(err).Error("wait error") } + if err := db.Close(); err != nil { + log.WithError(err).Error("failed to close database") + } + log.Info("gracefully stopped") } diff --git a/cmd/podsync/updater.go b/cmd/podsync/updater.go index 1223a11..fbcae62 100644 --- a/cmd/podsync/updater.go +++ b/cmd/podsync/updater.go @@ -18,6 +18,7 @@ import ( "github.com/mxpv/podsync/pkg/feed" "github.com/mxpv/podsync/pkg/link" "github.com/mxpv/podsync/pkg/model" + "github.com/mxpv/podsync/pkg/storage" ) type Downloader interface { @@ -27,10 +28,15 @@ type Downloader interface { type Updater struct { config *config.Config downloader Downloader + db storage.Storage } -func NewUpdater(config *config.Config, downloader Downloader) (*Updater, error) { - return &Updater{config: config, downloader: downloader}, nil +func NewUpdater(config *config.Config, downloader Downloader, db storage.Storage) (*Updater, error) { + return &Updater{ + config: config, + downloader: downloader, + db: db, + }, nil } func (u *Updater) Update(ctx context.Context, feedConfig *config.Feed) error { @@ -48,6 +54,26 @@ func (u *Updater) Update(ctx context.Context, feedConfig *config.Feed) error { return errors.Wrapf(err, "failed to create directory for feed %q", feedConfig.ID) } + if err := u.updateFeed(ctx, feedConfig); err != nil { + return err + } + + if err := u.downloadEpisodes(ctx, feedConfig, feedPath); err != nil { + return err + } + + if err := u.buildXML(ctx, feedConfig); err != nil { + return err + } + + elapsed := time.Since(started) + nextUpdate := time.Now().Add(feedConfig.UpdatePeriod.Duration) + log.Infof("successfully updated feed in %s, next update at %s", elapsed, nextUpdate.Format(time.Kitchen)) + return nil +} + +// updateFeed pulls API for new episodes and saves them to database +func (u *Updater) updateFeed(ctx context.Context, feedConfig *config.Feed) error { // Create an updater for this feed type provider, err := u.makeBuilder(ctx, feedConfig) if err != nil { @@ -63,60 +89,126 @@ func (u *Updater) Update(ctx context.Context, feedConfig *config.Feed) error { log.Debugf("received %d episode(s) for %q", len(result.Episodes), result.Title) - // Since there is no way to detect the size of an episode after download and encoding via API, - // we'll patch XML feed with values from this map - sizes := map[string]int64{} + if err := u.db.AddFeed(ctx, result); err != nil { + return err + } - // The number of episodes downloaded during this update - downloaded := 0 + log.Debug("successfully saved updates to storage") + return nil +} - // Download and encode episodes - for idx, episode := range result.Episodes { +func (u *Updater) downloadEpisodes(ctx context.Context, feedConfig *config.Feed, targetDir string) error { + var ( + feedID = feedConfig.ID + updateList []*model.Episode + ) + + // Build the list of files to download + if err := u.db.WalkFiles(ctx, feedID, func(file *model.File) error { + if file.Status != model.EpisodeNew && file.Status != model.EpisodeError { + // File already downloaded + return nil + } + + episode, err := u.db.GetEpisode(ctx, feedID, file.EpisodeID) + if err != nil { + return errors.Wrapf(err, "failed to query episode %q from database", file.EpisodeID) + } + + updateList = append(updateList, episode) + return nil + }); err != nil { + return errors.Wrapf(err, "failed to build update list") + } + + log.Debugf("update list: %+v", updateList) + + // Download pending episodes + for idx, episode := range updateList { logger := log.WithFields(log.Fields{ "index": idx, "episode_id": episode.ID, }) - episodePath := filepath.Join(feedPath, u.episodeName(feedConfig, episode)) - _, err := os.Stat(episodePath) - if err != nil && !os.IsNotExist(err) { - return errors.Wrap(err, "failed to check whether episode exists") - } + // Check whether episode exists on disk + + episodePath := filepath.Join(targetDir, u.episodeName(feedConfig, episode)) + stat, err := os.Stat(episodePath) + if err == nil { + logger.Infof("episode %q already exists on disk (%s)", episode.ID, episodePath) - if os.IsNotExist(err) { - // There is no file on disk, download episode - logger.Infof("! downloading episode %s", episode.VideoURL) - if output, err := u.downloader.Download(ctx, feedConfig, episode, feedPath); err == nil { - downloaded++ - } else { - // YouTube might block host with HTTP Error 429: Too Many Requests - // We still need to generate XML, so just stop sending download requests and - // retry next time - if strings.Contains(output, "HTTP Error 429") { - logger.WithError(err).Warnf("got too many requests error, will retry download next time") - break - } - - logger.WithError(err).Errorf("youtube-dl error: %s", output) + // File already exists, update file status and disk size + if err := u.db.UpdateFile(feedID, episode.ID, func(file *model.File) error { + file.Size = stat.Size() + file.Status = model.EpisodeDownloaded + return nil + }); err != nil { + logger.WithError(err).Error("failed to update file info") + return err } + + return nil + } else if os.IsNotExist(err) { + // Will download, do nothing here } else { - // Episode already downloaded - logger.Debug("skipping download of episode") + logger.WithError(err).Error("failed to stat file") + return err } - // Record file size - if size, err := u.fileSize(episodePath); err != nil { - // Don't return on error, use estimated file size provided by builders - logger.WithError(err).Error("failed to get episode file size") - } else { //nolint - logger.Debugf("file size %d", size) - sizes[episode.ID] = size + // Download episode to disk + + logger.Infof("! downloading episode %s", episode.VideoURL) + output, err := u.downloader.Download(ctx, feedConfig, episode, episodePath) + if err != nil { + logger.WithError(err).Errorf("youtube-dl error: %s", output) + + // YouTube might block host with HTTP Error 429: Too Many Requests + // We still need to generate XML, so just stop sending download requests and + // retry next time + if strings.Contains(output, "HTTP Error 429") { + break + } + + if err := u.db.UpdateFile(feedID, episode.ID, func(file *model.File) error { + file.Status = model.EpisodeError + return nil + }); err != nil { + return err + } + + continue } + + // Update file status in database + + if err := u.db.UpdateFile(feedID, episode.ID, func(file *model.File) error { + // Record file size of newly downloaded file + size, err := u.fileSize(episodePath) + if err != nil { + logger.WithError(err).Error("failed to get episode file size") + } else { + file.Size = size + } + + file.Status = model.EpisodeDownloaded + return nil + }); err != nil { + return err + } + } + + return nil +} + +func (u *Updater) buildXML(ctx context.Context, feedConfig *config.Feed) error { + feed, err := u.db.GetFeed(ctx, feedConfig.ID) + if err != nil { + return err } // Build iTunes XML feed with data received from builder log.Debug("building iTunes podcast feed") - podcast, err := u.buildPodcast(result, feedConfig, sizes) + podcast, err := u.buildPodcast(feed, feedConfig) if err != nil { return err } @@ -129,18 +221,10 @@ func (u *Updater) Update(ctx context.Context, feedConfig *config.Feed) error { return errors.Wrapf(err, "failed to write XML feed to disk") } - elapsed := time.Since(started) - nextUpdate := time.Now().Add(feedConfig.UpdatePeriod.Duration) - log.Infof( - "successfully updated feed in %s, downloaded: %d episode(s), next update at %s", - elapsed, - downloaded, - nextUpdate.Format(time.Kitchen), - ) return nil } -func (u *Updater) buildPodcast(feed *model.Feed, cfg *config.Feed, sizes map[string]int64) (*itunes.Podcast, error) { +func (u *Updater) buildPodcast(feed *model.Feed, cfg *config.Feed) (*itunes.Podcast, error) { const ( podsyncGenerator = "Podsync generator (support us at https://github.com/mxpv/podsync)" defaultCategory = "TV & Film" @@ -167,11 +251,6 @@ func (u *Updater) buildPodcast(feed *model.Feed, cfg *config.Feed, sizes map[str } for i, episode := range feed.Episodes { - // Fixup episode size after downloading and encoding - if size, ok := sizes[episode.ID]; ok { - episode.Size = size - } - item := itunes.Item{ GUID: episode.ID, Link: episode.VideoURL, diff --git a/pkg/model/feed.go b/pkg/model/feed.go index 44d5460..cf5c739 100644 --- a/pkg/model/feed.go +++ b/pkg/model/feed.go @@ -63,6 +63,7 @@ type EpisodeStatus string const ( EpisodeNew = EpisodeStatus("new") // New episode received via API EpisodeDownloaded = EpisodeStatus("downloaded") // Downloaded, encoded and available for download + EpisodeError = EpisodeStatus("error") // Could not download, will retry EpisodeCleaned = EpisodeStatus("cleaned") // Downloaded and later removed from disk due to update strategy ) diff --git a/pkg/storage/badger.go b/pkg/storage/badger.go index 0a77569..6e2f52f 100644 --- a/pkg/storage/badger.go +++ b/pkg/storage/badger.go @@ -105,6 +105,42 @@ func (b *Badger) AddFeed(_ context.Context, feed *model.Feed) error { }) } +func (b *Badger) GetFeed(_ context.Context, feedID string) (*model.Feed, error) { + var ( + feed = model.Feed{} + feedKey = b.getKey(feedPath, feedID) + ) + + if err := b.db.View(func(txn *badger.Txn) error { + // Query feed + if err := b.getObj(txn, feedKey, &feed); err != nil { + return err + } + + // Query episodes + opts := badger.DefaultIteratorOptions + opts.Prefix = b.getKey(episodePrefix, feedID) + opts.PrefetchValues = true + if err := b.iterator(txn, opts, func(item *badger.Item) error { + episode := &model.Episode{} + if err := b.getObj(txn, item.Key(), &episode); err != nil { + return err + } + + feed.Episodes = append(feed.Episodes, episode) + return nil + }); err != nil { + return err + } + + return nil + }); err != nil { + return nil, err + } + + return &feed, nil +} + func (b *Badger) WalkFeeds(_ context.Context, cb func(feed *model.Feed) error) error { return b.db.View(func(txn *badger.Txn) error { opts := badger.DefaultIteratorOptions @@ -153,11 +189,26 @@ func (b *Badger) DeleteFeed(_ context.Context, feedID string) error { }) } +func (b *Badger) GetEpisode(_ context.Context, feedID string, episodeID string) (*model.Episode, error) { + var ( + episode model.Episode + err error + key = b.getKey(episodePath, feedID, episodeID) + ) + + err = b.db.View(func(txn *badger.Txn) error { + return b.getObj(txn, key, &episode) + }) + + return &episode, err +} + func (b *Badger) WalkFiles(_ context.Context, feedID string, cb func(file *model.File) error) error { + opts := badger.DefaultIteratorOptions + opts.Prefix = b.getKey(filePrefix, feedID) + opts.PrefetchValues = true + return b.db.View(func(txn *badger.Txn) error { - opts := badger.DefaultIteratorOptions - opts.Prefix = b.getKey(filePrefix, feedID) - opts.PrefetchValues = true return b.iterator(txn, opts, func(item *badger.Item) error { file := &model.File{} if err := b.unmarshalObj(item, file); err != nil { @@ -169,30 +220,30 @@ func (b *Badger) WalkFiles(_ context.Context, feedID string, cb func(file *model }) } -func (b *Badger) UpdateFile(file *model.File, cb func() error) error { +func (b *Badger) UpdateFile(feedID string, episodeID string, cb func(file *model.File) error) error { + var ( + key = b.getKey(filePath, feedID, episodeID) + file = &model.File{} + ) + return b.db.Update(func(txn *badger.Txn) error { - fileKey := b.getKey(filePath, file.FeedID, file.EpisodeID) - - fileOld := &model.File{} - if err := b.getObj(txn, fileKey, fileOld); err != nil { + if err := b.getObj(txn, key, file); err != nil { return err } - if file.Size > 0 { - fileOld.Size = file.Size - } - - fileOld.Status = file.Status - - if err := b.setObj(txn, fileKey, fileOld, true); err != nil { + if err := cb(file); err != nil { return err } - if cb != nil { - return cb() + if file.FeedID != feedID { + return errors.New("can't change feed ID") } - return nil + if file.EpisodeID != episodeID { + return errors.New("can't change episode ID") + } + + return b.setObj(txn, key, file, true) }) } diff --git a/pkg/storage/badger_test.go b/pkg/storage/badger_test.go index 01134ac..774f39f 100644 --- a/pkg/storage/badger_test.go +++ b/pkg/storage/badger_test.go @@ -53,10 +53,30 @@ func TestBadger_AddFeed(t *testing.T) { require.NoError(t, err) defer db.Close() - err = db.AddFeed(context.Background(), getFeed()) + err = db.AddFeed(testCtx, getFeed()) assert.NoError(t, err) } +func TestBadger_GetFeed(t *testing.T) { + dir, err := ioutil.TempDir("", "podsync-badger-") + assert.NoError(t, err) + defer os.RemoveAll(dir) + + db, err := NewBadger(dir) + require.NoError(t, err) + defer db.Close() + + feed := getFeed() + feed.Episodes = nil + + err = db.AddFeed(testCtx, feed) + require.NoError(t, err) + + actual, err := db.GetFeed(testCtx, feed.ID) + assert.NoError(t, err) + assert.Equal(t, feed, actual) +} + func TestBadger_WalkFeeds(t *testing.T) { dir, err := ioutil.TempDir("", "podsync-badger-") assert.NoError(t, err) @@ -145,7 +165,6 @@ func TestBadger_WalkFiles(t *testing.T) { assert.Equal(t, 2, called) } - func TestBadger_UpdateFile(t *testing.T) { dir, err := ioutil.TempDir("", "podsync-badger-") assert.NoError(t, err) @@ -159,14 +178,9 @@ func TestBadger_UpdateFile(t *testing.T) { err = db.AddFeed(testCtx, feed) assert.NoError(t, err) - update := &model.File{ - EpisodeID: feed.Episodes[0].ID, - FeedID: feed.ID, - Size: 333, - Status: model.EpisodeDownloaded, - } - - err = db.UpdateFile(update, func() error { + err = db.UpdateFile(feed.ID, feed.Episodes[0].ID, func(file *model.File) error { + file.Size = 333 + file.Status = model.EpisodeDownloaded return nil }) assert.NoError(t, err) @@ -175,7 +189,10 @@ func TestBadger_UpdateFile(t *testing.T) { err = db.WalkFiles(testCtx, feed.ID, func(file *model.File) error { if first { - assert.Equal(t, update, file) + assert.Equal(t, feed.ID, file.FeedID) + assert.Equal(t, feed.Episodes[0].ID, file.EpisodeID) + assert.EqualValues(t, 333, file.Size) + assert.Equal(t, model.EpisodeDownloaded, file.Status) first = false } diff --git a/pkg/storage/storage.go b/pkg/storage/storage.go index 4f1982c..561f862 100644 --- a/pkg/storage/storage.go +++ b/pkg/storage/storage.go @@ -27,16 +27,21 @@ type Storage interface { // - Insert File model for each new episode AddFeed(ctx context.Context, feed *model.Feed) error + // GetFeed gets a feed by ID + GetFeed(ctx context.Context, feedID string) (*model.Feed, error) + // WalkFeeds iterates over feeds saved to database WalkFeeds(ctx context.Context, cb func(feed *model.Feed) error) error // DeleteFeed deletes feed and all related data from database DeleteFeed(ctx context.Context, feedID string) error + // GetEpisode gets episode by identifier + GetEpisode(ctx context.Context, feedID string, episodeID string) (*model.Episode, error) + // WalkFiles walks all files for the given feed ID WalkFiles(ctx context.Context, feedID string, cb func(file *model.File) error) error - // UpdateFile updates file's status and (optionally) size. - // Callback can be used to rollback update transaction. - UpdateFile(file *model.File, cb func() error) error + // UpdateFile updates file via callback function. + UpdateFile(feedID string, episodeID string, cb func(file *model.File) error) error }