mirror of
https://github.com/hcengineering/platform.git
synced 2026-09-10 11:47:42 +02:00
Fix incomplete upload (#14)
This commit is contained in:
+2
-2
@@ -57,11 +57,11 @@ linters-settings:
|
||||
goimports:
|
||||
local-prefixes: github.com/networkservicemesh/sdk
|
||||
gocyclo:
|
||||
min-complexity: 20
|
||||
min-complexity: 30
|
||||
dupl:
|
||||
threshold: 150
|
||||
funlen:
|
||||
lines: 180
|
||||
lines: 200
|
||||
statements: 100
|
||||
goconst:
|
||||
min-len: 2
|
||||
|
||||
+1
-1
@@ -11,7 +11,7 @@
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
FROM golang:1.24.1 AS builder
|
||||
FROM golang:1.24.4 AS builder
|
||||
ENV GO111MODULE=on
|
||||
ENV CGO_ENABLED=0
|
||||
ENV GOBIN=/bin
|
||||
|
||||
@@ -8,7 +8,7 @@ require (
|
||||
github.com/aws/aws-sdk-go-v2/credentials v1.17.59
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.77.0
|
||||
github.com/getsentry/sentry-go v0.31.1
|
||||
github.com/golang-jwt/jwt/v5 v5.2.1
|
||||
github.com/golang-jwt/jwt/v5 v5.2.2
|
||||
github.com/google/uuid v1.6.0
|
||||
github.com/kelseyhightower/envconfig v1.4.0
|
||||
github.com/pkg/errors v0.9.1
|
||||
|
||||
@@ -45,8 +45,8 @@ github.com/getsentry/sentry-go v0.31.1 h1:ELVc0h7gwyhnXHDouXkhqTFSO5oslsRDk0++ey
|
||||
github.com/getsentry/sentry-go v0.31.1/go.mod h1:CYNcMMz73YigoHljQRG+qPF+eMq8gG72XcGN/p71BAY=
|
||||
github.com/go-errors/errors v1.4.2 h1:J6MZopCL4uSllY1OfXM374weqZFFItUbrImctkmUxIA=
|
||||
github.com/go-errors/errors v1.4.2/go.mod h1:sIVyrIiJhuEF+Pj9Ebtd6P/rEYROXFi3BopGUQ5a5Og=
|
||||
github.com/golang-jwt/jwt/v5 v5.2.1 h1:OuVbFODueb089Lh128TAcimifWaLhJwVflnrgM17wHk=
|
||||
github.com/golang-jwt/jwt/v5 v5.2.1/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk=
|
||||
github.com/golang-jwt/jwt/v5 v5.2.2 h1:Rl4B7itRWVtYIHFrSNd7vhTiz9UpLdi6gZhZ3wEeDy8=
|
||||
github.com/golang-jwt/jwt/v5 v5.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk=
|
||||
github.com/golang/mock v1.6.0 h1:ErTB+efbowRARo13NNdxyJji2egdxLGQhRaY+DUumQc=
|
||||
github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs=
|
||||
github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI=
|
||||
|
||||
@@ -19,7 +19,6 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
@@ -30,12 +29,37 @@ import (
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// LogLevel is ffmpeg log level
|
||||
type LogLevel string
|
||||
|
||||
const (
|
||||
// LogLevelQuiet is quiet log level
|
||||
LogLevelQuiet LogLevel = "quiet"
|
||||
// LogLevelPanic is panic log level
|
||||
LogLevelPanic LogLevel = "panic"
|
||||
// LogLevelFatal is fatal log level
|
||||
LogLevelFatal LogLevel = "fatal"
|
||||
// LogLevelError is error log level
|
||||
LogLevelError LogLevel = "error"
|
||||
// LogLevelWarning is warning log level
|
||||
LogLevelWarning LogLevel = "warning"
|
||||
// LogLevelInfo is info log level
|
||||
LogLevelInfo LogLevel = "info"
|
||||
// LogLevelVerbose is verbose log level
|
||||
LogLevelVerbose LogLevel = "verbose"
|
||||
// LogLevelDebug is debug log level
|
||||
LogLevelDebug LogLevel = "debug"
|
||||
// LogLevelTrace is trace log level
|
||||
LogLevelTrace LogLevel = "trace"
|
||||
)
|
||||
|
||||
// Options represents configuration for the ffmpeg command
|
||||
type Options struct {
|
||||
Input string
|
||||
OutputDir string
|
||||
ScalingLevels []string
|
||||
Level string
|
||||
LogLevel LogLevel
|
||||
Transcode bool
|
||||
Threads int
|
||||
UploadID string
|
||||
@@ -51,8 +75,6 @@ func newFfmpegCommand(ctx context.Context, in io.Reader, args []string) (*exec.C
|
||||
logger.Debug("prepared command: ", zap.Strings("args", args))
|
||||
|
||||
var result = exec.CommandContext(ctx, "ffmpeg", args...)
|
||||
result.Stderr = os.Stdout
|
||||
result.Stdout = os.Stdout
|
||||
result.Stdin = in
|
||||
|
||||
return result, nil
|
||||
@@ -60,6 +82,8 @@ func newFfmpegCommand(ctx context.Context, in io.Reader, args []string) (*exec.C
|
||||
|
||||
func buildCommonCommand(opts *Options) []string {
|
||||
var result = []string{
|
||||
"-y", // Overwrite output files without asking.
|
||||
"-v", string(opts.LogLevel),
|
||||
"-threads", fmt.Sprint(opts.Threads),
|
||||
"-i", opts.Input,
|
||||
}
|
||||
@@ -96,7 +120,7 @@ func BuildRawVideoCommand(opts *Options) []string {
|
||||
"-g", "60",
|
||||
"-f", "hls",
|
||||
"-hls_time", "5",
|
||||
"-hls_flags", "split_by_time",
|
||||
"-hls_flags", "split_by_time+temp_file",
|
||||
"-hls_list_size", "0",
|
||||
"-hls_segment_filename", filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", opts.Level)),
|
||||
filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, opts.Level)))
|
||||
@@ -107,7 +131,7 @@ func BuildRawVideoCommand(opts *Options) []string {
|
||||
"-c:v", "copy", // Copy video stream
|
||||
"-f", "hls",
|
||||
"-hls_time", "5",
|
||||
"-hls_flags", "split_by_time",
|
||||
"-hls_flags", "split_by_time+temp_file",
|
||||
"-hls_list_size", "0",
|
||||
"-hls_segment_filename", filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", opts.Level)),
|
||||
filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, opts.Level)))
|
||||
@@ -142,7 +166,15 @@ func BuildScalingVideoCommand(opts *Options) []string {
|
||||
"-g", "60",
|
||||
"-f", "hls",
|
||||
"-hls_time", "5",
|
||||
"-hls_flags", "split_by_time",
|
||||
// Use HLS flags
|
||||
// - split_by_time
|
||||
// Allow segments to start on frames other than key frames.
|
||||
// This improves behavior on some players when the time between key frames is inconsistent,
|
||||
// but may make things worse on others, and can cause some oddities during seeking.
|
||||
// This flag should be used with the hls_time option.
|
||||
// - temp_file
|
||||
// Write segment data to filename.tmp and rename to filename only once the segment is complete.
|
||||
"-hls_flags", "split_by_time+temp_file",
|
||||
"-hls_list_size", "0",
|
||||
"-hls_segment_filename", filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", level)),
|
||||
filepath.Join(opts.OutputDir, opts.UploadID, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, level)))
|
||||
|
||||
@@ -28,10 +28,11 @@ func Test_BuildVideoCommand_Scaling(t *testing.T) {
|
||||
Input: "pipe:0",
|
||||
UploadID: "1",
|
||||
Threads: 4,
|
||||
LogLevel: mediaconvert.LogLevelDebug,
|
||||
ScalingLevels: []string{"720p", "480p"},
|
||||
})
|
||||
|
||||
const expected = `-threads 4 -i pipe:0 -map 0:v -vf scale=-2:720 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time -hls_list_size 0 -hls_segment_filename test/1/1_%03d_720p.ts test/1/1_720p_master.m3u8 -map 0:v -vf scale=-2:480 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
|
||||
const expected = `-y -v debug -threads 4 -i pipe:0 -map 0:v -vf scale=-2:720 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_720p.ts test/1/1_720p_master.m3u8 -map 0:v -vf scale=-2:480 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
|
||||
|
||||
require.Contains(t, expected, strings.Join(scaleCommand, " "))
|
||||
}
|
||||
@@ -42,11 +43,12 @@ func Test_BuildVideoCommand_Scaling_NoRaw(t *testing.T) {
|
||||
Input: "pipe:0",
|
||||
UploadID: "1",
|
||||
Threads: 4,
|
||||
LogLevel: mediaconvert.LogLevelDebug,
|
||||
Level: "720p",
|
||||
ScalingLevels: []string{"720p", "480p"},
|
||||
})
|
||||
|
||||
const expected = `-threads 4 -i pipe:0 -map 0:v -vf scale=-2:480 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
|
||||
const expected = `-y -v debug -threads 4 -i pipe:0 -map 0:v -vf scale=-2:480 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
|
||||
|
||||
require.Contains(t, expected, strings.Join(scaleCommand, " "))
|
||||
}
|
||||
@@ -57,11 +59,12 @@ func Test_BuildVideoCommand_Raw_NoTranscode(t *testing.T) {
|
||||
Input: "pipe:0",
|
||||
UploadID: "1",
|
||||
Threads: 4,
|
||||
LogLevel: mediaconvert.LogLevelDebug,
|
||||
Level: resconv.Level("651:490"),
|
||||
Transcode: false,
|
||||
})
|
||||
|
||||
const expected = `"-threads 4 -i pipe:0 -c:a copy -c:v copy -f hls -hls_time 5 -hls_flags split_by_time -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
|
||||
const expected = `"-y -v debug -threads 4 -i pipe:0 -c:a copy -c:v copy -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
|
||||
|
||||
require.Contains(t, expected, strings.Join(rawCommand, " "))
|
||||
}
|
||||
@@ -72,11 +75,12 @@ func Test_BuildVideoCommand_Raw_Transcode(t *testing.T) {
|
||||
Input: "pipe:0",
|
||||
UploadID: "1",
|
||||
Threads: 4,
|
||||
LogLevel: mediaconvert.LogLevelDebug,
|
||||
Level: resconv.Level("651:490"),
|
||||
Transcode: true,
|
||||
})
|
||||
|
||||
const expected = `-threads 4 -i pipe:0 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
|
||||
const expected = `-y -v debug -threads 4 -i pipe:0 -c:a aac -c:v libx264 -preset veryfast -crf 23 -g 60 -f hls -hls_time 5 -hls_flags split_by_time+temp_file -hls_list_size 0 -hls_segment_filename test/1/1_%03d_480p.ts test/1/1_480p_master.m3u8`
|
||||
|
||||
require.Contains(t, expected, strings.Join(rawCommand, " "))
|
||||
}
|
||||
|
||||
@@ -299,6 +299,8 @@ func IsSupportedMediaType(mediaType string) bool {
|
||||
return true
|
||||
case "video/webm":
|
||||
return true
|
||||
case "video/quicktime":
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -16,8 +16,10 @@
|
||||
package mediaconvert
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
@@ -42,6 +44,13 @@ type Transcoder struct {
|
||||
logger *zap.Logger
|
||||
}
|
||||
|
||||
// Command represents a ffmpeg command
|
||||
type Command struct {
|
||||
cmd *exec.Cmd
|
||||
stdoutBuf bytes.Buffer
|
||||
stderrBuf bytes.Buffer
|
||||
}
|
||||
|
||||
// NewTranscoder creates a new instance of task transcoder
|
||||
func NewTranscoder(ctx context.Context, cfg *config.Config) *Transcoder {
|
||||
var p = &Transcoder{
|
||||
@@ -54,7 +63,7 @@ func NewTranscoder(ctx context.Context, cfg *config.Config) *Transcoder {
|
||||
}
|
||||
|
||||
// Transcode handles one transcoding task
|
||||
func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, error) {
|
||||
func (p *Transcoder) Transcode(ctx context.Context, task *Task) (*TaskResult, error) {
|
||||
var logger = p.logger.With(zap.String("task-id", task.ID))
|
||||
|
||||
logger.Debug("start")
|
||||
@@ -64,7 +73,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err
|
||||
var tokenString, err = token.NewToken(p.cfg.ServerSecret, task.Workspace, "stream", "datalake")
|
||||
if err != nil {
|
||||
logger.Error("can not create token", zap.Error(err))
|
||||
return TaskResult{}, errors.Wrapf(err, "can not create token")
|
||||
return nil, errors.Wrapf(err, "can not create token")
|
||||
}
|
||||
|
||||
logger.Debug("phase 2: preparing fs")
|
||||
@@ -73,10 +82,11 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err
|
||||
err = os.MkdirAll(destinationFolder, os.ModePerm)
|
||||
if err != nil {
|
||||
logger.Error("can not create temporary folder", zap.Error(err))
|
||||
return TaskResult{}, errors.Wrapf(err, "can not create temporary folder")
|
||||
return nil, errors.Wrapf(err, "can not create temporary folder")
|
||||
}
|
||||
|
||||
defer func() {
|
||||
logger.Debug("remove temporary folder")
|
||||
if err = os.RemoveAll(destinationFolder); err != nil {
|
||||
logger.Error("failed to cleanup temporary folder", zap.Error(err))
|
||||
}
|
||||
@@ -87,38 +97,38 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err
|
||||
remoteStorage, err := storage.NewStorageByURL(ctx, p.cfg.Endpoint(), p.cfg.EndpointURL.Scheme, tokenString, task.Workspace)
|
||||
if err != nil {
|
||||
logger.Error("can not create storage by url", zap.Error(err), zap.String("url", p.cfg.EndpointURL.String()))
|
||||
return TaskResult{}, errors.Wrapf(err, "can not create storage by url")
|
||||
return nil, errors.Wrapf(err, "can not create storage by url")
|
||||
}
|
||||
|
||||
stat, err := remoteStorage.StatFile(ctx, task.Source)
|
||||
if err != nil {
|
||||
logger.Error("can not stat file", zap.Error(err), zap.String("filepath", task.Source))
|
||||
return TaskResult{}, errors.Wrapf(err, "can not stat file")
|
||||
return nil, errors.Wrapf(err, "can not stat file")
|
||||
}
|
||||
|
||||
if !IsSupportedMediaType(stat.Type) {
|
||||
logger.Info("unsupported media type", zap.String("type", stat.Type))
|
||||
return TaskResult{}, fmt.Errorf("unsupported media type: %s", stat.Type)
|
||||
return nil, fmt.Errorf("unsupported media type: %s", stat.Type)
|
||||
}
|
||||
|
||||
sourceFilePath := filepath.Join(destinationFolder, filename)
|
||||
if err = remoteStorage.GetFile(ctx, task.Source, sourceFilePath); err != nil {
|
||||
logger.Error("can not download source file", zap.Error(err), zap.String("filepath", task.Source))
|
||||
// TODO: reschedule
|
||||
return TaskResult{}, errors.Wrapf(err, "can not download source file")
|
||||
return nil, errors.Wrapf(err, "can not download source file")
|
||||
}
|
||||
|
||||
logger.Debug("phase 4: prepare to transcode")
|
||||
probe, err := ffprobe.ProbeURL(ctx, sourceFilePath)
|
||||
if err != nil {
|
||||
logger.Error("can not get ffprobe", zap.Error(err), zap.String("filepath", sourceFilePath))
|
||||
return TaskResult{}, errors.Wrapf(err, "can not get ffprobe")
|
||||
return nil, errors.Wrapf(err, "can not get ffprobe")
|
||||
}
|
||||
|
||||
videoStream := probe.FirstVideoStream()
|
||||
if videoStream == nil {
|
||||
logger.Error("no video stream found", zap.String("filepath", sourceFilePath))
|
||||
return TaskResult{}, errors.Wrapf(err, "no video stream found")
|
||||
return nil, errors.Wrapf(err, "no video stream found")
|
||||
}
|
||||
|
||||
logger.Debug("video stream found", zap.String("codec", videoStream.CodecName), zap.Int("width", videoStream.Width), zap.Int("height", videoStream.Height))
|
||||
@@ -136,6 +146,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err
|
||||
Input: sourceFilePath,
|
||||
OutputDir: p.cfg.OutputDir,
|
||||
Level: level,
|
||||
LogLevel: LogLevel(p.cfg.LogLevel),
|
||||
Transcode: !IsHLSSupportedVideoCodec(codec),
|
||||
ScalingLevels: append(sublevels, level),
|
||||
UploadID: task.ID,
|
||||
@@ -157,7 +168,7 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err
|
||||
err = manifest.GenerateHLSPlaylist(opts.ScalingLevels, p.cfg.OutputDir, opts.UploadID)
|
||||
if err != nil {
|
||||
logger.Error("can not generate hls playlist", zap.String("out", p.cfg.OutputDir), zap.String("uploadID", opts.UploadID))
|
||||
return TaskResult{}, errors.Wrapf(err, "can not generate hls playlist")
|
||||
return nil, errors.Wrapf(err, "can not generate hls playlist")
|
||||
}
|
||||
|
||||
go uploader.Start()
|
||||
@@ -169,34 +180,53 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err
|
||||
BuildRawVideoCommand(&opts),
|
||||
BuildScalingVideoCommand(&opts),
|
||||
}
|
||||
var cmds []*exec.Cmd
|
||||
var cmds []Command
|
||||
|
||||
for _, args := range argsSlice {
|
||||
cmd, cmdErr := newFfmpegCommand(ctx, nil, args)
|
||||
if cmdErr != nil {
|
||||
logger.Error("can not create a new command", zap.Error(cmdErr), zap.Strings("args", args))
|
||||
go uploader.Cancel()
|
||||
return TaskResult{}, errors.Wrapf(err, "can not create a new command")
|
||||
return nil, errors.Wrapf(cmdErr, "can not create a new command")
|
||||
}
|
||||
cmds = append(cmds, cmd)
|
||||
|
||||
var command = Command{
|
||||
cmd: cmd,
|
||||
stdoutBuf: bytes.Buffer{},
|
||||
stderrBuf: bytes.Buffer{},
|
||||
}
|
||||
|
||||
cmd.Stdout = io.MultiWriter(os.Stdout, &command.stdoutBuf)
|
||||
cmd.Stderr = io.MultiWriter(os.Stderr, &command.stderrBuf)
|
||||
|
||||
cmds = append(cmds, command)
|
||||
if err = cmd.Start(); err != nil {
|
||||
logger.Error("can not start a command", zap.Error(err), zap.Strings("args", args))
|
||||
go uploader.Cancel()
|
||||
return TaskResult{}, errors.Wrapf(err, "can not start a command")
|
||||
return nil, errors.Wrapf(err, "can not start a command")
|
||||
}
|
||||
}
|
||||
|
||||
logger.Debug("phase 7: wait for the result")
|
||||
|
||||
for _, cmd := range cmds {
|
||||
if err = cmd.Wait(); err != nil {
|
||||
logger.Error("can not wait for command end ", zap.Error(err))
|
||||
go uploader.Cancel()
|
||||
return TaskResult{}, errors.Wrapf(err, "can not wait for command end")
|
||||
if err = cmd.cmd.Wait(); err == nil {
|
||||
continue
|
||||
}
|
||||
|
||||
logger.Error("can not wait for command end ", zap.Error(err))
|
||||
if _, err = os.Stdout.Write(cmd.stdoutBuf.Bytes()); err != nil {
|
||||
logger.Error("can not write stdout ", zap.Error(err))
|
||||
}
|
||||
if _, err = os.Stderr.Write(cmd.stderrBuf.Bytes()); err != nil {
|
||||
logger.Error("can not write stderr", zap.Error(err))
|
||||
}
|
||||
go uploader.Cancel()
|
||||
return nil, errors.Wrapf(err, "can not wait for command end")
|
||||
}
|
||||
|
||||
logger.Debug("phase 8: schedule cleanup")
|
||||
go uploader.Stop()
|
||||
uploader.Stop()
|
||||
|
||||
logger.Debug("phase 9: try to set metadata")
|
||||
|
||||
@@ -226,5 +256,5 @@ func (p *Transcoder) Transcode(ctx context.Context, task *Task) (TaskResult, err
|
||||
}
|
||||
}
|
||||
|
||||
return result, nil
|
||||
return &result, nil
|
||||
}
|
||||
|
||||
@@ -170,6 +170,11 @@ func (d *DatalakeStorage) DeleteFile(ctx context.Context, fileName string) error
|
||||
return errors.Wrapf(err, "delete failed")
|
||||
}
|
||||
|
||||
if err := okResponse(res); err != nil {
|
||||
logRequestError(logger, err, "bad status code", res)
|
||||
return err
|
||||
}
|
||||
|
||||
logger.Debug("deleted")
|
||||
|
||||
return nil
|
||||
@@ -205,14 +210,11 @@ func (d *DatalakeStorage) PatchMeta(ctx context.Context, filename string, md *Me
|
||||
return err
|
||||
}
|
||||
|
||||
if resp.StatusCode() != fasthttp.StatusOK {
|
||||
var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode())
|
||||
logger.Debug("bad status code", zap.Error(err))
|
||||
if err := okResponse(resp); err != nil {
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return err
|
||||
}
|
||||
|
||||
fmt.Println(string(resp.Body()))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -237,9 +239,8 @@ func (d *DatalakeStorage) GetMeta(ctx context.Context, filename string) (*Metada
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if resp.StatusCode() != fasthttp.StatusOK {
|
||||
var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode())
|
||||
logger.Debug("bad status code", zap.Error(err))
|
||||
if err := okResponse(resp); err != nil {
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -270,10 +271,8 @@ func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination str
|
||||
return err
|
||||
}
|
||||
|
||||
// Check the response status code
|
||||
if resp.StatusCode() != fasthttp.StatusOK {
|
||||
var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode())
|
||||
logger.Debug("bad status code", zap.Error(err))
|
||||
if err := okResponse(resp); err != nil {
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -286,12 +285,18 @@ func (d *DatalakeStorage) GetFile(ctx context.Context, filename, destination str
|
||||
defer func() {
|
||||
_ = file.Close()
|
||||
}()
|
||||
if err := resp.BodyWriteTo(file); err != nil {
|
||||
if err = resp.BodyWriteTo(file); err != nil {
|
||||
logger.Debug("can't write to file", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
|
||||
logger.Debug("file downloaded successfully")
|
||||
stat, err := os.Stat(destination)
|
||||
if err != nil {
|
||||
logger.Error("can't stat the file", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
|
||||
logger.Info("file downloaded successfully", zap.Int64("size", stat.Size()))
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -316,9 +321,7 @@ func (d *DatalakeStorage) StatFile(ctx context.Context, filename string) (*BlobI
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Check the response status code
|
||||
if resp.StatusCode() != fasthttp.StatusOK {
|
||||
var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode())
|
||||
if err := okResponse(resp); err != nil {
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return nil, err
|
||||
}
|
||||
@@ -346,7 +349,7 @@ func (d *DatalakeStorage) SetParent(ctx context.Context, filename, parent string
|
||||
req.SetRequestURI(d.baseURL + "/blob/" + d.workspace + "/" + objectKey + "/parent")
|
||||
req.Header.SetMethod(fasthttp.MethodPatch)
|
||||
req.Header.Add("Authorization", "Bearer "+d.token)
|
||||
req.Header.Add("Content-Type", "application/json")
|
||||
req.Header.SetContentType("application/json")
|
||||
|
||||
body := map[string]any{
|
||||
"parent": parentKey,
|
||||
@@ -365,15 +368,20 @@ func (d *DatalakeStorage) SetParent(ctx context.Context, filename, parent string
|
||||
return err
|
||||
}
|
||||
|
||||
// Check the response status code
|
||||
var statusOK = resp.StatusCode() >= 200 && resp.StatusCode() < 300
|
||||
if !statusOK {
|
||||
var err = fmt.Errorf("unexpected status code: %d", resp.StatusCode())
|
||||
logger.Debug("bad status code", zap.Error(err), zap.Int("status", resp.StatusCode()), zap.String("response", resp.String()))
|
||||
if err := okResponse(resp); err != nil {
|
||||
logRequestError(logger, err, "bad status code", resp)
|
||||
return err
|
||||
}
|
||||
|
||||
logger.Debug("finished")
|
||||
return nil
|
||||
}
|
||||
|
||||
func okResponse(res *fasthttp.Response) error {
|
||||
var statusOK = res.StatusCode() >= 200 && res.StatusCode() < 300
|
||||
|
||||
if !statusOK {
|
||||
return fmt.Errorf("unexpected status code: %d", res.StatusCode())
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -34,7 +34,9 @@ import (
|
||||
|
||||
// See at https://man7.org/linux/man-pages/man7/inotify.7.html
|
||||
const inotifyCloseWrite uint32 = 0x8 // IN_CLOSE_WRITE
|
||||
const inotifyMovedFrom uint32 = 0x40 // IN_MOVED_FROM
|
||||
const inotifyMovedTo uint32 = 0x80 // IN_MOVED_TO
|
||||
const inotifyDelete uint32 = 0x200 // IN_DELETE
|
||||
|
||||
// Uploader represents file uploader
|
||||
type Uploader interface {
|
||||
@@ -102,50 +104,66 @@ func New(ctx context.Context, s storage.Storage, opts Options) Uploader {
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) Stop() {
|
||||
u.logger.Info("stopping upload")
|
||||
u.stop(false)
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) Cancel() {
|
||||
u.logger.Info("canceling upload")
|
||||
u.stop(true)
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) scanInitialFiles() {
|
||||
u.workerWaitGroup.Add(1)
|
||||
func (u *uploaderImpl) scanFiles() {
|
||||
logger := u.logger.With(zap.String("dir", u.options.Dir))
|
||||
|
||||
go func() {
|
||||
defer u.workerWaitGroup.Done()
|
||||
logger.Info("scan files")
|
||||
files, err := os.ReadDir(u.options.Dir)
|
||||
if err != nil {
|
||||
logger.Error("failed to read files", zap.Error(err))
|
||||
return
|
||||
}
|
||||
|
||||
logger := u.logger.With(zap.String("dir", u.options.Dir))
|
||||
|
||||
logger.Info("initial file scan")
|
||||
initFiles, err := os.ReadDir(u.options.Dir)
|
||||
if err != nil {
|
||||
logger.Error("failed to read initial files", zap.Error(err))
|
||||
return
|
||||
count := 0
|
||||
for _, f := range files {
|
||||
if f.IsDir() {
|
||||
continue
|
||||
}
|
||||
|
||||
for _, f := range initFiles {
|
||||
if f.IsDir() {
|
||||
continue
|
||||
}
|
||||
var filePath = filepath.Join(u.options.Dir, f.Name())
|
||||
|
||||
// Ignore source file
|
||||
var filePath = filepath.Join(u.options.Dir, f.Name())
|
||||
if filePath == u.options.SourceFile {
|
||||
continue
|
||||
}
|
||||
u.filesCh <- filePath
|
||||
// Ignore source file
|
||||
if filePath == u.options.SourceFile {
|
||||
continue
|
||||
}
|
||||
|
||||
logger.Info("initial file scan complete", zap.Int("count", len(initFiles)))
|
||||
}()
|
||||
if _, uploaded := u.sentFiles.Load(filePath); uploaded {
|
||||
logger.Debug("file already uploaded", zap.String("file", filePath))
|
||||
continue
|
||||
}
|
||||
|
||||
u.filesCh <- filePath
|
||||
count++
|
||||
}
|
||||
|
||||
logger.Info("scan complete", zap.Int("count", count))
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) stop(rollback bool) {
|
||||
// Stop watching for new files
|
||||
close(u.watcherStopCh)
|
||||
<-u.watcherDoneCh
|
||||
u.logger.Debug("file watch stopped")
|
||||
|
||||
// Scan remaining files in the directory
|
||||
u.scanFiles()
|
||||
|
||||
// Close filesCh so no new files added
|
||||
close(u.filesCh)
|
||||
|
||||
// Wait for all workers to finish processing
|
||||
u.workerWaitGroup.Wait()
|
||||
u.logger.Debug("workers done")
|
||||
|
||||
// Perform rollback
|
||||
if rollback {
|
||||
u.logger.Debug("starting rollback...")
|
||||
var i uint32
|
||||
@@ -159,15 +177,25 @@ func (u *uploaderImpl) stop(rollback bool) {
|
||||
})
|
||||
u.logger.Debug("rollback done")
|
||||
}
|
||||
close(u.filesCh)
|
||||
u.workerWaitGroup.Wait()
|
||||
u.logger.Debug("workers done")
|
||||
|
||||
u.uploadCancel()
|
||||
_ = os.RemoveAll(u.options.Dir)
|
||||
|
||||
remainingFiles, err := os.ReadDir(u.options.Dir)
|
||||
if err != nil && !os.IsNotExist(err) {
|
||||
u.logger.Error("failed to read dir", zap.Error(err))
|
||||
}
|
||||
// log remaining files
|
||||
if len(remainingFiles) > 0 {
|
||||
files := make([]string, 0, len(remainingFiles))
|
||||
for _, entry := range remainingFiles {
|
||||
files = append(files, entry.Name())
|
||||
}
|
||||
u.logger.Info("remaining files", zap.Int("count", len(files)), zap.Any("files", files))
|
||||
}
|
||||
|
||||
u.sentFiles.Clear()
|
||||
|
||||
u.logger.Debug("finish done", zap.Bool("cancel", rollback))
|
||||
u.logger.Debug("stopped", zap.Bool("rollback", rollback))
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) Start() {
|
||||
@@ -178,7 +206,7 @@ func (u *uploaderImpl) Start() {
|
||||
|
||||
<-watcherReady
|
||||
|
||||
u.scanInitialFiles()
|
||||
u.scanFiles()
|
||||
}
|
||||
|
||||
func (u *uploaderImpl) startWorkers() {
|
||||
@@ -249,10 +277,12 @@ func (u *uploaderImpl) uploadAndDelete(f string) {
|
||||
return
|
||||
}
|
||||
|
||||
// Check if the file exists
|
||||
_, err := os.Stat(f)
|
||||
if err != nil {
|
||||
logger.Debug("file does not exist")
|
||||
if err := waitFileExists(f); err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
logger.Debug("file does not exist", zap.Error(err))
|
||||
} else {
|
||||
logger.Error("failed to stat file", zap.Error(err))
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
@@ -300,6 +330,7 @@ func (u *uploaderImpl) uploadAndDelete(f string) {
|
||||
}
|
||||
}
|
||||
|
||||
// startWatch watches for changes in the directory and uploads created files
|
||||
func (u *uploaderImpl) startWatch(ready chan<- struct{}) {
|
||||
defer close(u.watcherDoneCh)
|
||||
|
||||
@@ -317,7 +348,7 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) {
|
||||
}
|
||||
}()
|
||||
|
||||
if err := watcher.AddWatch(u.options.Dir, inotifyCloseWrite|inotifyMovedTo); err != nil {
|
||||
if err := watcher.AddWatch(u.options.Dir, inotifyCloseWrite|inotifyMovedTo|inotifyDelete|inotifyMovedFrom); err != nil {
|
||||
logger.Error("can not start watching", zap.Error(err))
|
||||
close(ready)
|
||||
return
|
||||
@@ -336,18 +367,17 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) {
|
||||
logger.Error("file channel was closed")
|
||||
return
|
||||
}
|
||||
if !strings.Contains(event.Name, u.options.Dir) {
|
||||
if event.Name == u.options.Dir ||
|
||||
event.Name == u.options.SourceFile ||
|
||||
strings.HasSuffix(event.Name, ".tmp") {
|
||||
continue
|
||||
}
|
||||
if event.Name == u.options.Dir {
|
||||
continue
|
||||
}
|
||||
if event.Name == u.options.SourceFile {
|
||||
continue
|
||||
}
|
||||
if strings.HasSuffix(event.Name, ".tmp") {
|
||||
|
||||
if event.Mask&(inotifyDelete|inotifyMovedFrom) != 0 {
|
||||
logger.Debug("file deleted or moved away", zap.String("event", event.Name), zap.Uint32("mask", event.Mask))
|
||||
continue
|
||||
}
|
||||
|
||||
logger.Debug("received an event", zap.String("event", event.Name), zap.Uint32("mask", event.Mask))
|
||||
|
||||
u.filesCh <- event.Name
|
||||
@@ -359,3 +389,19 @@ func (u *uploaderImpl) startWatch(ready chan<- struct{}) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func waitFileExists(file string) error {
|
||||
var err error
|
||||
var stat os.FileInfo
|
||||
|
||||
for range 10 {
|
||||
stat, err = os.Stat(file)
|
||||
if err == nil && stat.Size() > 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user