diff --git a/cmd/videnc/main.go b/cmd/videnc/main.go index 89a8c8e..c3b5381 100644 --- a/cmd/videnc/main.go +++ b/cmd/videnc/main.go @@ -1,6 +1,7 @@ package main import ( + "context" "crypto/rand" "encoding/hex" "flag" @@ -61,24 +62,18 @@ func main() { metaClient := metadata.NewClient(cfg.OMDBAPIKey) - done := make(chan struct{}) - sigChan := make(chan os.Signal, 1) - signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM) - - go func() { - <-sigChan - close(done) - }() + ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer cancel() w := watcher.New(cfg.Paths.Input, 15) - w.Start(func(inputPath string) error { - return processFile(inputPath, cfg, enc, metaClient, log) - }, done) + w.Start(ctx, func(ctx context.Context, inputPath string) error { + return processFile(ctx, inputPath, cfg, enc, metaClient, log) + }) log.Info("videnc-vibe started") } -func processFile(inputPath string, cfg *types.Config, enc *encoder.Encoder, metaClient *metadata.Client, log *logger.Logger) error { +func processFile(ctx context.Context, inputPath string, cfg *types.Config, enc *encoder.Encoder, metaClient *metadata.Client, log *logger.Logger) error { log.Info(fmt.Sprintf("Processing: %s", inputPath)) filename := filepath.Base(inputPath) @@ -97,14 +92,14 @@ func processFile(inputPath string, cfg *types.Config, enc *encoder.Encoder, meta // Cleanup the work directory on every exit path, success or failure. defer os.RemoveAll(workDir) - width, height, interlaced, err := enc.GetMediaInfo(inputPath) + width, height, interlaced, err := enc.GetMediaInfo(ctx, inputPath) if err != nil { log.ErrorFile(inputPath, "Getting media info", err.Error()) failToFailed(inputPath, cfg.Paths.Failed, log) return err } - streamLangs, err := enc.GetStreamLanguages(inputPath) + streamLangs, err := enc.GetStreamLanguages(ctx, inputPath) if err != nil { log.ErrorFile(inputPath, "Getting stream languages", err.Error()) } @@ -130,12 +125,12 @@ func processFile(inputPath string, cfg *types.Config, enc *encoder.Encoder, meta var meta *types.Metadata if isSeries && tvmazeID != "" && season != "" && episode != "" { - meta, err = metaClient.FetchSeriesMetadata(tvmazeID, season, episode) + meta, err = metaClient.FetchSeriesMetadata(ctx, tvmazeID, season, episode) if err != nil { log.ErrorFile(inputPath, "Fetching series metadata", err.Error()) } } else if imdbID != "" { - meta, err = metaClient.FetchMovieMetadata(imdbID) + meta, err = metaClient.FetchMovieMetadata(ctx, imdbID) if err != nil { log.ErrorFile(inputPath, "Fetching movie metadata", err.Error()) } @@ -159,7 +154,7 @@ func processFile(inputPath string, cfg *types.Config, enc *encoder.Encoder, meta DeleteOrigin: deleteOrigin, } - if err := enc.Transcode(inputPath, workDir, job, meta, interlaced, streamLangs); err != nil { + if err := enc.Transcode(ctx, inputPath, workDir, job, meta, interlaced, streamLangs); err != nil { log.ErrorFile(inputPath, "Transcoding", err.Error()) failToFailed(inputPath, cfg.Paths.Failed, log) return err diff --git a/internal/encoder/encoder.go b/internal/encoder/encoder.go index a3bf803..0589c72 100644 --- a/internal/encoder/encoder.go +++ b/internal/encoder/encoder.go @@ -1,6 +1,7 @@ package encoder import ( + "context" "encoding/json" "fmt" "os/exec" @@ -42,8 +43,8 @@ type StreamMetadata struct { Title string `json:"title"` } -func (e *Encoder) GetStreamLanguages(path string) ([]StreamMetadata, error) { - cmd := exec.Command(e.ffprobePath, "-v", "error", "-show_streams", "-print_format", "json", path) +func (e *Encoder) GetStreamLanguages(ctx context.Context, path string) ([]StreamMetadata, error) { + cmd := exec.CommandContext(ctx, e.ffprobePath, "-v", "error", "-show_streams", "-print_format", "json", path) output, err := cmd.CombinedOutput() if err != nil { return nil, fmt.Errorf("ffprobe error: %w", err) @@ -103,8 +104,8 @@ func (e *Encoder) CheckDeps() error { return nil } -func (e *Encoder) GetMediaInfo(path string) (width, height int, interlaced bool, err error) { - cmd := exec.Command(e.ffprobePath, "-v", "error", "-show_streams", "-print_format", "json", path) +func (e *Encoder) GetMediaInfo(ctx context.Context, path string) (width, height int, interlaced bool, err error) { + cmd := exec.CommandContext(ctx, e.ffprobePath, "-v", "error", "-show_streams", "-print_format", "json", path) output, err := cmd.CombinedOutput() if err != nil { return 0, 0, false, fmt.Errorf("ffprobe error: %w", err) @@ -131,7 +132,7 @@ func (e *Encoder) GetMediaInfo(path string) (width, height int, interlaced bool, return 0, 0, false, fmt.Errorf("no video stream found") } - cmd = exec.Command(e.ffmpegPath, + cmd = exec.CommandContext(ctx, e.ffmpegPath, "-hide_banner", "-nostats", "-i", path, @@ -168,8 +169,8 @@ func detectInterlaced(idetOutput string) bool { return interlaced > prog } -func (e *Encoder) calculateZscaleWidth(path string, originalHeight int) (string, int, error) { - cmd := exec.Command(e.ffprobePath, "-v", "error", "-select_streams", "v:0", "-show_entries", "stream=width,height,sample_aspect_ratio", "-print_format", "json", path) +func (e *Encoder) calculateZscaleWidth(ctx context.Context, path string, originalHeight int) (string, int, error) { + cmd := exec.CommandContext(ctx, e.ffprobePath, "-v", "error", "-select_streams", "v:0", "-show_entries", "stream=width,height,sample_aspect_ratio", "-print_format", "json", path) output, err := cmd.CombinedOutput() if err != nil { return "", 0, fmt.Errorf("ffprobe error: %w", err) @@ -221,18 +222,18 @@ func (e *Encoder) calculateZscaleWidth(path string, originalHeight int) (string, return zscale, newWidth, nil } -func (e *Encoder) Transcode(input, workDir string, job *types.Job, metadata *types.Metadata, interlaced bool, streamLangs []StreamMetadata) error { - audioWavs, err := e.extractAudio(input, workDir, streamLangs) +func (e *Encoder) Transcode(ctx context.Context, input, workDir string, job *types.Job, metadata *types.Metadata, interlaced bool, streamLangs []StreamMetadata) error { + audioWavs, err := e.extractAudio(ctx, input, workDir, streamLangs) if err != nil { return fmt.Errorf("extracting audio: %w", err) } - opusFiles, err := e.encodeOpus(audioWavs, workDir) + opusFiles, err := e.encodeOpus(ctx, audioWavs, workDir) if err != nil { return fmt.Errorf("encoding opus: %w", err) } - if err := e.encodeVideo(input, workDir, opusFiles, job, metadata, interlaced, streamLangs); err != nil { + if err := e.encodeVideo(ctx, input, workDir, opusFiles, job, metadata, interlaced, streamLangs); err != nil { return fmt.Errorf("encoding video: %w", err) } @@ -253,7 +254,7 @@ func audioStreamsInSourceOrder(streamLangs []StreamMetadata) []StreamMetadata { return audio } -func (e *Encoder) extractAudio(input, workDir string, streamLangs []StreamMetadata) ([]string, error) { +func (e *Encoder) extractAudio(ctx context.Context, input, workDir string, streamLangs []StreamMetadata) ([]string, error) { audio := audioStreamsInSourceOrder(streamLangs) if len(audio) == 0 { return nil, fmt.Errorf("no audio streams found") @@ -262,7 +263,7 @@ func (e *Encoder) extractAudio(input, workDir string, streamLangs []StreamMetada wavs := make([]string, 0, len(audio)) for srcAudioIndex := range audio { wavPath := filepath.Join(workDir, fmt.Sprintf("audio.%d.wav", srcAudioIndex)) - cmd := exec.Command(e.ffmpegPath, "-i", input, + cmd := exec.CommandContext(ctx, e.ffmpegPath, "-i", input, "-map", fmt.Sprintf("0:a:%d", srcAudioIndex), "-vn", "-c:a", "pcm_s16le", "-ar", "48000", "-f", "wav", wavPath) @@ -274,12 +275,12 @@ func (e *Encoder) extractAudio(input, workDir string, streamLangs []StreamMetada return wavs, nil } -func (e *Encoder) encodeOpus(wavs []string, workDir string) ([]string, error) { +func (e *Encoder) encodeOpus(ctx context.Context, wavs []string, workDir string) ([]string, error) { var opusFiles []string for _, wav := range wavs { base := filepath.Base(strings.Replace(wav, ".wav", ".opus", 1)) out := filepath.Join(workDir, base) - cmd := exec.Command(e.opusencPath, "--bitrate", "128k", wav, out) + cmd := exec.CommandContext(ctx, e.opusencPath, "--bitrate", "128k", wav, out) if logOut, err := cmd.CombinedOutput(); err != nil { return nil, fmt.Errorf("opusenc: %s %w", logOut, err) } @@ -288,12 +289,12 @@ func (e *Encoder) encodeOpus(wavs []string, workDir string) ([]string, error) { return opusFiles, nil } -func (e *Encoder) encodeVideo(input, workDir string, opusFiles []string, job *types.Job, metadata *types.Metadata, interlaced bool, streamLangs []StreamMetadata) error { +func (e *Encoder) encodeVideo(ctx context.Context, input, workDir string, opusFiles []string, job *types.Job, metadata *types.Metadata, interlaced bool, streamLangs []StreamMetadata) error { outFile := filepath.Join(workDir, "output.mkv") svtParams := fmt.Sprintf("film-grain=10:film-grain-denoise=1:scd=1:qm-min=4:qm-max=15:keyint=10s") - zscaleStr, newWidth, err := e.calculateZscaleWidth(input, 0) + zscaleStr, newWidth, err := e.calculateZscaleWidth(ctx, input, 0) if err != nil { return fmt.Errorf("calculating zscale: %w", err) } @@ -386,7 +387,7 @@ func (e *Encoder) encodeVideo(input, workDir string, opusFiles []string, job *ty args = append(args, outFile) args = append([]string{"-hide_banner", "-v", "error"}, args...) - cmd := exec.Command(e.ffmpegPath, args...) + cmd := exec.CommandContext(ctx, e.ffmpegPath, args...) fmt.Printf("DEBUG FFmpeg command: ffmpeg %s\n", strings.Join(args, " ")) if out, err := cmd.CombinedOutput(); err != nil { return fmt.Errorf("ffmpeg encode: %s %w", out, err) diff --git a/internal/metadata/metadata.go b/internal/metadata/metadata.go index 42da193..7a285c5 100644 --- a/internal/metadata/metadata.go +++ b/internal/metadata/metadata.go @@ -1,6 +1,7 @@ package metadata import ( + "context" "encoding/json" "fmt" "io" @@ -109,10 +110,14 @@ func ParseMediaType(filename string) types.MediaType { return "" } -func (c *Client) FetchMovieMetadata(imdbID string) (*types.Metadata, error) { +func (c *Client) FetchMovieMetadata(ctx context.Context, imdbID string) (*types.Metadata, error) { url := fmt.Sprintf("%s?i=%s&apikey=%s", omdbAPIURL, imdbID, c.omdbAPIKey) - resp, err := c.httpClient.Get(url) + req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + return nil, fmt.Errorf("building OMDb request: %w", err) + } + resp, err := c.httpClient.Do(req) if err != nil { return nil, fmt.Errorf("fetching OMDb: %w", err) } @@ -153,7 +158,7 @@ func (c *Client) FetchMovieMetadata(imdbID string) (*types.Metadata, error) { }, nil } -func (c *Client) FetchSeriesMetadata(tvmazeID, season, episode string) (*types.Metadata, error) { +func (c *Client) FetchSeriesMetadata(ctx context.Context, tvmazeID, season, episode string) (*types.Metadata, error) { seasonNum, err := strconv.Atoi(season) if err != nil { return nil, fmt.Errorf("parsing season %q: %w", season, err) @@ -165,13 +170,13 @@ func (c *Client) FetchSeriesMetadata(tvmazeID, season, episode string) (*types.M epURL := fmt.Sprintf("%s/shows/%s/episodebynumber?season=%d&number=%d", tvmazeAPIURL, tvmazeID, seasonNum, episodeNum) var ep TVMazeEpisode - if err := c.fetchTVMazeJSON(epURL, &ep); err != nil { + if err := c.fetchTVMazeJSON(ctx, epURL, &ep); err != nil { return nil, fmt.Errorf("TVmaze episode %s S%dE%d: %w", tvmazeID, seasonNum, episodeNum, err) } showURL := fmt.Sprintf("%s/shows/%s", tvmazeAPIURL, tvmazeID) var show TVMazeShowResponse - if err := c.fetchTVMazeJSON(showURL, &show); err != nil { + if err := c.fetchTVMazeJSON(ctx, showURL, &show); err != nil { return nil, fmt.Errorf("TVmaze show %s: %w", tvmazeID, err) } @@ -192,8 +197,12 @@ func (c *Client) FetchSeriesMetadata(tvmazeID, season, episode string) (*types.M }, nil } -func (c *Client) fetchTVMazeJSON(url string, v interface{}) error { - resp, err := c.httpClient.Get(url) +func (c *Client) fetchTVMazeJSON(ctx context.Context, url string, v interface{}) error { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + return fmt.Errorf("build request: %w", err) + } + resp, err := c.httpClient.Do(req) if err != nil { return fmt.Errorf("GET: %w", err) } diff --git a/internal/watcher/watcher.go b/internal/watcher/watcher.go index 5245cf2..af358cf 100644 --- a/internal/watcher/watcher.go +++ b/internal/watcher/watcher.go @@ -1,6 +1,7 @@ package watcher import ( + "context" "fmt" "os" "path/filepath" @@ -60,23 +61,23 @@ func New(inputDir string, intervalSeconds int) *Watcher { } } -func (w *Watcher) Start(processFn func(string) error, done chan struct{}) { +func (w *Watcher) Start(ctx context.Context, processFn func(context.Context, string) error) { ticker := time.NewTicker(w.interval) defer ticker.Stop() - w.scanAndProcess(processFn) + w.scanAndProcess(ctx, processFn) for { select { - case <-done: + case <-ctx.Done(): return case <-ticker.C: - w.scanAndProcess(processFn) + w.scanAndProcess(ctx, processFn) } } } -func (w *Watcher) scanAndProcess(processFn func(string) error) { +func (w *Watcher) scanAndProcess(ctx context.Context, processFn func(context.Context, string) error) { files, err := filepath.Glob(filepath.Join(w.inputDir, "*.mkv")) if err != nil { fmt.Printf("Error scanning input directory: %v\n", err) @@ -118,7 +119,7 @@ func (w *Watcher) scanAndProcess(processFn func(string) error) { continue } - if err := processFn(file); err != nil { + if err := processFn(ctx, file); err != nil { fmt.Printf("Error processing %s: %v\n", file, err) nextFailed[file] = failureRecord{ failedAt: now,