diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index e4cc858..7b1476b 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -29,7 +29,7 @@ jobs: username: ${{ github.repository_owner }} password: ${{ secrets.GITHUB_TOKEN }} - - uses: goreleaser/goreleaser-action@v2.9.1 + - uses: goreleaser/goreleaser-action@v3.0.0 if: startsWith(github.ref, 'refs/tags/') with: version: latest diff --git a/README.md b/README.md index fc57d17..4e41c91 100644 --- a/README.md +++ b/README.md @@ -67,7 +67,10 @@ Minimal configuration would look like this: ```toml [server] port = 8080 -data_dir = "/data/podsync/" + +[storage] + [storage.local] + data_dir = "/data/podsync/" [tokens] youtube = "PASTE YOUR API KEY HERE" diff --git a/cloud_formation.yml b/cloud_formation.yml index d9de555..620af73 100644 --- a/cloud_formation.yml +++ b/cloud_formation.yml @@ -169,8 +169,11 @@ Resources: tee /home/ec2-user/podsync/config.toml < len(ids) { + end = len(ids) + } + // Save each slice as comma-delimited string + idsList = append(idsList, strings.Join(ids[i:end], ",")) } - for _, video := range req.Items { - var ( - snippet = video.Snippet - videoID = video.Id - videoURL = fmt.Sprintf("https://youtube.com/watch?v=%s", video.Id) - image = yt.selectThumbnail(snippet.Thumbnails, feed.Quality, videoID) - ) + // Show how many API calls will be required + log.Debugf("Expected to make %d API calls to get the descriptions for %d episode(s).", len(idsList), len(ids)) - // Parse date added to playlist / publication date - dateStr := "" - playlistItem, ok := playlist[video.Id] - if ok { - dateStr = playlistItem.PublishedAt - } else { - dateStr = snippet.PublishedAt - } - - pubDate, err := yt.parseDate(dateStr) + // Loop in each slices of 50 (or less) IDs and query their description + for _, idsI := range idsList { + req, err := yt.client.Videos.List("id,snippet,contentDetails").Id(idsI).Context(ctx).Do(yt.key) if err != nil { - return errors.Wrapf(err, "failed to parse video publish date: %s", dateStr) + return errors.Wrap(err, "failed to query video descriptions") } - // Sometimes YouTube retrun empty content defailt, use arbitrary one - var seconds int64 = 1 - if video.ContentDetails != nil { - // Parse duration - d, err := duration.FromString(video.ContentDetails.Duration) - if err != nil { - return errors.Wrapf(err, "failed to parse duration %s", video.ContentDetails.Duration) + for _, video := range req.Items { + var ( + snippet = video.Snippet + videoID = video.Id + videoURL = fmt.Sprintf("https://youtube.com/watch?v=%s", video.Id) + image = yt.selectThumbnail(snippet.Thumbnails, feed.Quality, videoID) + ) + + // Parse date added to playlist / publication date + dateStr := "" + playlistItem, ok := playlist[video.Id] + if ok { + dateStr = playlistItem.PublishedAt + } else { + dateStr = snippet.PublishedAt } - seconds = int64(d.ToDuration().Seconds()) + pubDate, err := yt.parseDate(dateStr) + if err != nil { + return errors.Wrapf(err, "failed to parse video publish date: %s", dateStr) + } + + // Sometimes YouTube retrun empty content defailt, use arbitrary one + var seconds int64 = 1 + if video.ContentDetails != nil { + // Parse duration + d, err := duration.FromString(video.ContentDetails.Duration) + if err != nil { + return errors.Wrapf(err, "failed to parse duration %s", video.ContentDetails.Duration) + } + + seconds = int64(d.ToDuration().Seconds()) + } + + var ( + order = strconv.FormatInt(playlistItem.Position, 10) + size = yt.getSize(seconds, feed) + ) + + feed.Episodes = append(feed.Episodes, &model.Episode{ + ID: video.Id, + Title: snippet.Title, + Description: snippet.Description, + Thumbnail: image, + Duration: seconds, + Size: size, + VideoURL: videoURL, + PubDate: pubDate, + Order: order, + Status: model.EpisodeNew, + }) } - - var ( - order = strconv.FormatInt(playlistItem.Position, 10) - size = yt.getSize(seconds, feed) - ) - - feed.Episodes = append(feed.Episodes, &model.Episode{ - ID: video.Id, - Title: snippet.Title, - Description: snippet.Description, - Thumbnail: image, - Duration: seconds, - Size: size, - VideoURL: videoURL, - PubDate: pubDate, - Order: order, - Status: model.EpisodeNew, - }) } return nil diff --git a/pkg/fs/local.go b/pkg/fs/local.go index 8a7734f..687ba54 100644 --- a/pkg/fs/local.go +++ b/pkg/fs/local.go @@ -11,6 +11,11 @@ import ( log "github.com/sirupsen/logrus" ) +// LocalConfig is the storage configuration for local file system +type LocalConfig struct { + DataDir string `yaml:"data_dir"` +} + // Local implements local file storage type Local struct { rootDir string @@ -65,3 +70,18 @@ func (l *Local) copyFile(source io.Reader, destinationPath string) (int64, error return written, nil } + +func (l *Local) Size(_ctx context.Context, name string) (int64, error) { + file, err := l.Open(name) + if err != nil { + return 0, err + } + defer file.Close() + + stat, err := file.Stat() + if err != nil { + return 0, err + } + + return stat.Size(), nil +} diff --git a/pkg/fs/local_test.go b/pkg/fs/local_test.go index 67bcb09..b20ca01 100644 --- a/pkg/fs/local_test.go +++ b/pkg/fs/local_test.go @@ -52,7 +52,7 @@ func TestLocal_Size(t *testing.T) { _, err = stor.Create(testCtx, "1/test", bytes.NewBuffer([]byte{1, 5, 7, 8, 3})) assert.NoError(t, err) - sz, err := Size(stor, "1/test") + sz, err := stor.Size(testCtx, "1/test") assert.NoError(t, err) assert.EqualValues(t, 5, sz) } @@ -61,7 +61,7 @@ func TestLocal_NoSize(t *testing.T) { stor, err := NewLocal("") assert.NoError(t, err) - _, err = Size(stor, "1/test") + _, err = stor.Size(testCtx, "1/test") assert.True(t, os.IsNotExist(err)) } @@ -80,7 +80,7 @@ func TestLocal_Delete(t *testing.T) { err = stor.Delete(testCtx, "1/test") assert.NoError(t, err) - _, err = Size(stor, "1/test") + _, err = stor.Size(testCtx, "1/test") assert.True(t, os.IsNotExist(err)) _, err = os.Stat(filepath.Join(tmpDir, "1", "test")) diff --git a/pkg/fs/s3.go b/pkg/fs/s3.go new file mode 100644 index 0000000..688d3a0 --- /dev/null +++ b/pkg/fs/s3.go @@ -0,0 +1,118 @@ +package fs + +import ( + "context" + "io" + "net/http" + "os" + + "github.com/aws/aws-sdk-go/aws" + "github.com/aws/aws-sdk-go/aws/awserr" + "github.com/aws/aws-sdk-go/aws/session" + "github.com/aws/aws-sdk-go/service/s3" + "github.com/aws/aws-sdk-go/service/s3/s3iface" + "github.com/aws/aws-sdk-go/service/s3/s3manager" + "github.com/pkg/errors" + log "github.com/sirupsen/logrus" +) + +// S3Config is the configuration for a S3-compatible storage provider +type S3Config struct { + // S3 Bucket to store files + Bucket string `toml:"bucket"` + // Region of the S3 service + Region string `toml:"region"` + // EndpointURL is an HTTP endpoint of the S3 API + EndpointURL string `toml:"endpoint_url"` +} + +// S3 implements file storage for S3-compatible providers. +type S3 struct { + api s3iface.S3API + uploader *s3manager.Uploader + bucket string +} + +func NewS3(c S3Config) (*S3, error) { + cfg := aws.NewConfig(). + WithEndpoint(c.EndpointURL). + WithRegion(c.Region). + WithLogger(s3logger{}). + WithLogLevel(aws.LogDebug) + sess, err := session.NewSessionWithOptions(session.Options{Config: *cfg}) + if err != nil { + return nil, errors.Wrap(err, "failed to initialize S3 session") + } + return &S3{ + api: s3.New(sess), + uploader: s3manager.NewUploader(sess), + bucket: c.Bucket, + }, nil +} + +func (s *S3) Open(_name string) (http.File, error) { + return nil, errors.New("serving files from S3 is not supported") +} + +func (s *S3) Delete(ctx context.Context, name string) error { + _, err := s.api.DeleteObjectWithContext(ctx, &s3.DeleteObjectInput{ + Bucket: &s.bucket, + Key: &name, + }) + return err +} + +func (s *S3) Create(ctx context.Context, name string, reader io.Reader) (int64, error) { + logger := log.WithField("name", name) + + logger.Infof("uploading file to %s", s.bucket) + r := &readerWithN{Reader: reader} + _, err := s.uploader.UploadWithContext(ctx, &s3manager.UploadInput{ + Bucket: &s.bucket, + Key: &name, + Body: r, + }) + if err != nil { + return 0, errors.Wrap(err, "failed to upload file") + } + + logger.Debugf("written %d bytes", r.n) + return int64(r.n), nil +} + +func (s *S3) Size(ctx context.Context, name string) (int64, error) { + logger := log.WithField("name", name) + + logger.Debugf("getting file size from %s", s.bucket) + resp, err := s.api.HeadObjectWithContext(ctx, &s3.HeadObjectInput{ + Bucket: &s.bucket, + Key: &name, + }) + if err != nil { + if awsErr, ok := err.(awserr.Error); ok { + if awsErr.Code() == "NotFound" { + return 0, os.ErrNotExist + } + } + return 0, errors.Wrap(err, "failed to get file size") + } + + return *resp.ContentLength, nil +} + +type readerWithN struct { + io.Reader + n int +} + +func (r *readerWithN) Read(p []byte) (n int, err error) { + n, err = r.Reader.Read(p) + r.n += n + return +} + +type s3logger struct{} + +func (s s3logger) Log(args ...interface{}) { + log.Debug(args...) +} diff --git a/pkg/fs/s3_test.go b/pkg/fs/s3_test.go new file mode 100644 index 0000000..ecdb65b --- /dev/null +++ b/pkg/fs/s3_test.go @@ -0,0 +1,107 @@ +package fs + +import ( + "bytes" + "io/ioutil" + "os" + "testing" + + "github.com/aws/aws-sdk-go/aws" + "github.com/aws/aws-sdk-go/aws/awserr" + "github.com/aws/aws-sdk-go/aws/client/metadata" + "github.com/aws/aws-sdk-go/aws/request" + "github.com/aws/aws-sdk-go/service/s3" + "github.com/aws/aws-sdk-go/service/s3/s3iface" + "github.com/aws/aws-sdk-go/service/s3/s3manager" + "github.com/stretchr/testify/assert" +) + +func TestS3_Create(t *testing.T) { + files := make(map[string][]byte) + stor, err := newMockS3(files) + assert.NoError(t, err) + + written, err := stor.Create(testCtx, "1/test", bytes.NewBuffer([]byte{1, 5, 7, 8, 3})) + assert.NoError(t, err) + assert.EqualValues(t, 5, written) + + d, ok := files["1/test"] + assert.True(t, ok) + assert.EqualValues(t, 5, len(d)) +} + +func TestS3_Size(t *testing.T) { + files := make(map[string][]byte) + stor, err := newMockS3(files) + assert.NoError(t, err) + + _, err = stor.Create(testCtx, "1/test", bytes.NewBuffer([]byte{1, 5, 7, 8, 3})) + assert.NoError(t, err) + + sz, err := stor.Size(testCtx, "1/test") + assert.NoError(t, err) + assert.EqualValues(t, 5, sz) +} + +func TestS3_NoSize(t *testing.T) { + files := make(map[string][]byte) + stor, err := newMockS3(files) + assert.NoError(t, err) + + _, err = stor.Size(testCtx, "1/test") + assert.True(t, os.IsNotExist(err)) +} + +func TestS3_Delete(t *testing.T) { + files := make(map[string][]byte) + stor, err := newMockS3(files) + assert.NoError(t, err) + + _, err = stor.Create(testCtx, "1/test", bytes.NewBuffer([]byte{1, 5, 7, 8, 3})) + assert.NoError(t, err) + + err = stor.Delete(testCtx, "1/test") + assert.NoError(t, err) + + _, err = stor.Size(testCtx, "1/test") + assert.True(t, os.IsNotExist(err)) + + _, ok := files["1/test"] + assert.False(t, ok) +} + +type mockS3API struct { + s3iface.S3API + files map[string][]byte +} + +func newMockS3(files map[string][]byte) (*S3, error) { + api := &mockS3API{files: files} + return &S3{ + api: api, + uploader: s3manager.NewUploaderWithClient(api), + bucket: "mock-bucket", + }, nil +} + +func (m *mockS3API) PutObjectRequest(input *s3.PutObjectInput) (*request.Request, *s3.PutObjectOutput) { + content, _ := ioutil.ReadAll(input.Body) + req := request.New(aws.Config{}, metadata.ClientInfo{}, request.Handlers{}, nil, &request.Operation{}, nil, nil) + m.files[*input.Key] = content + return req, &s3.PutObjectOutput{} +} + +func (m *mockS3API) HeadObjectWithContext(ctx aws.Context, input *s3.HeadObjectInput, opts ...request.Option) (*s3.HeadObjectOutput, error) { + if _, ok := m.files[*input.Key]; ok { + return &s3.HeadObjectOutput{ContentLength: aws.Int64(int64(len(m.files[*input.Key])))}, nil + } + return nil, awserr.New("NotFound", "", nil) +} + +func (m *mockS3API) DeleteObjectWithContext(ctx aws.Context, input *s3.DeleteObjectInput, opts ...request.Option) (*s3.DeleteObjectOutput, error) { + if _, ok := m.files[*input.Key]; ok { + delete(m.files, *input.Key) + return &s3.DeleteObjectOutput{}, nil + } + return nil, awserr.New("NotFound", "", nil) +} diff --git a/pkg/fs/storage.go b/pkg/fs/storage.go index 6540294..efdf190 100644 --- a/pkg/fs/storage.go +++ b/pkg/fs/storage.go @@ -16,20 +16,15 @@ type Storage interface { // Delete deletes the file Delete(ctx context.Context, name string) error + + // Size returns a storage object's size in bytes + Size(ctx context.Context, name string) (int64, error) } -// Size returns storage object's size in bytes. -func Size(storage http.FileSystem, name string) (int64, error) { - file, err := storage.Open(name) - if err != nil { - return 0, err - } - defer file.Close() - - stat, err := file.Stat() - if err != nil { - return 0, err - } - - return stat.Size(), nil +// Config is a configuration for the file storage backend +type Config struct { + // Type is the type of file system to use + Type string `toml:"type"` + Local LocalConfig `toml:"local"` + S3 S3Config `toml:"s3"` } diff --git a/services/update/updater.go b/services/update/updater.go index df517c8..a26e360 100644 --- a/services/update/updater.go +++ b/services/update/updater.go @@ -157,8 +157,12 @@ func (u *Manager) downloadEpisodes(ctx context.Context, feedConfig *feed.Config) // Build the list of files to download if err := u.db.WalkEpisodes(ctx, feedID, func(episode *model.Episode) error { + var ( + logger = log.WithFields(log.Fields{"episode_id": episode.ID}) + ) if episode.Status != model.EpisodeNew && episode.Status != model.EpisodeError { // File already downloaded + logger.Infof("skipping due to already downloaded") return nil } @@ -200,7 +204,7 @@ func (u *Manager) downloadEpisodes(ctx context.Context, feedConfig *feed.Config) ) // Check whether episode already exists - size, err := fs.Size(u.fs, fmt.Sprintf("%s/%s", feedID, episodeName)) + size, err := u.fs.Size(ctx, fmt.Sprintf("%s/%s", feedID, episodeName)) if err == nil { logger.Infof("episode %q already exists on disk", episode.ID)