Use storage to track episodes

This commit is contained in:
Maksym Pavlenko
2020-02-08 14:35:37 -08:00
committed by Maksym Pavlenko
parent f2a8638a4f
commit 7f24306898
6 changed files with 249 additions and 86 deletions
+11 -1
View File
@@ -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")
}
+132 -53
View File
@@ -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,
+1
View File
@@ -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
)
+69 -18
View File
@@ -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)
})
}
+28 -11
View File
@@ -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
}
+8 -3
View File
@@ -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
}