diff --git a/frontend/public/locales/en-US/dashboard.json b/frontend/public/locales/en-US/dashboard.json index c353f1af..60867d21 100644 --- a/frontend/public/locales/en-US/dashboard.json +++ b/frontend/public/locales/en-US/dashboard.json @@ -1218,11 +1218,15 @@ "downloader": "Downloader", "aria2Des": "Start Aria2 as the same user/access level running Cloudreve on the target node server, enable the RPC service in the Aria2 config file, for more information and guidelines, refer the \"Remote download\" section of the documentation.", "qbittorrentDes": "Start qBittorrent as the same user running Cloudreve on the target node server, enable the Web UI service in the qBittorrent settings, for more information and guidelines, refer the \"Remote download\" section of the documentation.", + "ytdlpDes": "Run yt-dlp on the target node server with the same user/access level as Cloudreve. yt-dlp is invoked as a CLI process per task - no daemon required.", "rpcServer": "RPC Server", "rpcServerHelpDes": "RPC server address contain full port number, e.g. <0>http://127.0.0.1:6800/.", "rpcToken": "RPC Token", "rpcTokenDes": "Consistent with <0>rpc-secret in the Aria2 configuration file; leave blank if not set.", "downloaderOptionDes": "Additional downloader configuration when creating a download task, written in JSON key-value format, see the <0>downloader official documentation for available parameters.", + "ytdlpOptionDes": "Additional yt-dlp CLI flags when creating a download task, written in JSON key-value format - keys become --flag arguments verbatim, see the <0>yt-dlp options reference.", + "ytdlpBinary": "yt-dlp binary", + "ytdlpBinaryDes": "Path or name of the yt-dlp executable on the node server. Defaults to yt-dlp resolved from PATH.", "refreshInterval": "Status refresh interval (seconds)", "refreshIntervalDes": "The interval at which Cloudreve requests a refresh of the task state from the downloader. The actual refresh interval also depends on the configuration of the \"Remote download\" queue and the busyness of the downloader.", "waitForSeeding": "Wait for seeding", diff --git a/frontend/public/locales/zh-CN/dashboard.json b/frontend/public/locales/zh-CN/dashboard.json index ef4db013..ddfc2272 100644 --- a/frontend/public/locales/zh-CN/dashboard.json +++ b/frontend/public/locales/zh-CN/dashboard.json @@ -1218,11 +1218,15 @@ "downloader": "下载器", "aria2Des": "请在目标节点服务器上以和运行 Cloudreve 相同的用户/权限启动 Aria2,并在 Aria2 的配置文件中开启 RPC 服务,更多信息及指引请参考文档的 “离线下载” 章节。", "qbittorrentDes": "请在目标节点服务器上以和运行 Cloudreve 相同的用户/权限启动 qBittorrent, 并在 qBittorrent 的设置中开启 “Web UI” 服务,更多信息及指引请参考文档的 “离线下载” 章节。", + "ytdlpDes": "在目标节点服务器上安装 yt-dlp,并以和运行 Cloudreve 相同的用户/权限调用。yt-dlp 会按任务以命令行进程方式运行,无需常驻服务。", "rpcServer": "RPC 服务器地址", "rpcServerHelpDes": "包含端口的完整 RPC 服务器地址,例如:<0>http://127.0.0.1:6800/。", "rpcToken": "RPC 授权令牌", "rpcTokenDes": "与 Aria2 配置文件中 <0>rpc-secret 保持一致,未设置请留空。", "downloaderOptionDes": "在创建下载任务时额外携带的下载器配置,以 JSON 键值对格式书写,具体可参考 <0>下载器官方文档。", + "ytdlpOptionDes": "创建下载任务时附加的 yt-dlp 命令行参数,以 JSON 键值对格式书写,键名将原样转换为 --flag 参数,具体可参考 <0>yt-dlp 选项文档。", + "ytdlpBinary": "yt-dlp 可执行文件", + "ytdlpBinaryDes": "节点服务器上 yt-dlp 可执行文件的路径或名称,留空时从 PATH 解析 yt-dlp。", "refreshInterval": "状态刷新间隔 (秒)", "refreshIntervalDes": "Cloudreve 向下载器请求刷新任务状态的间隔,实际刷新间隔也取决于 “离线下载” 队列的配置和繁忙程度。", "waitForSeeding": "等待做种完成", diff --git a/frontend/src/api/dashboard.ts b/frontend/src/api/dashboard.ts index 6f8965fb..c7ca98ce 100644 --- a/frontend/src/api/dashboard.ts +++ b/frontend/src/api/dashboard.ts @@ -183,6 +183,7 @@ export interface Node extends CommonMixin { export enum DownloaderProvider { qbittorrent = "qbittorrent", aria2 = "aria2", + ytdlp = "ytdlp", } export interface QBittorrentSetting { @@ -200,6 +201,12 @@ export interface Aria2Setting { temp_path?: string; } +export interface YtdlpSetting { + binary?: string; + options?: Record; + temp_path?: string; +} + export interface URLValidationSetting { disabled?: boolean; allowed_hosts?: string[]; @@ -210,6 +217,7 @@ export interface NodeSetting { provider?: DownloaderProvider; qbittorrent?: QBittorrentSetting; aria2?: Aria2Setting; + ytdlp?: YtdlpSetting; interval?: number; wait_for_seeding?: boolean; url_validation?: URLValidationSetting; diff --git a/frontend/src/component/Admin/Node/EditNode/CapabilitiesSection.tsx b/frontend/src/component/Admin/Node/EditNode/CapabilitiesSection.tsx index e9966652..f9bdb07f 100644 --- a/frontend/src/component/Admin/Node/EditNode/CapabilitiesSection.tsx +++ b/frontend/src/component/Admin/Node/EditNode/CapabilitiesSection.tsx @@ -41,6 +41,7 @@ const CapabilitiesSection = () => { const [editedConfigAria2, setEditedConfigAria2] = useState(""); const [editedConfigQbittorrent, setEditedConfigQbittorrent] = useState(""); + const [editedConfigYtdlp, setEditedConfigYtdlp] = useState(""); const [testDownloaderLoading, setTestDownloaderLoading] = useState(false); const [storeFilesHintDialogOpen, setStoreFilesHintDialogOpen] = useState(false); @@ -64,6 +65,12 @@ const CapabilitiesSection = () => { ); }, [values.settings?.qbittorrent?.options]); + useEffect(() => { + setEditedConfigYtdlp( + values.settings?.ytdlp?.options ? JSON.stringify(values.settings?.ytdlp?.options, null, 2) : "", + ); + }, [values.settings?.ytdlp?.options]); + const onCapabilityChange = useCallback( (capability: number) => (e: React.ChangeEvent) => { setNode((p: Node) => ({ @@ -136,6 +143,56 @@ const CapabilitiesSection = () => { [setNode], ); + const onYtdlpBinaryChange = useCallback( + (e: React.ChangeEvent) => { + setNode((p: Node) => ({ + ...p, + settings: { + ...p.settings, + ytdlp: { + ...p.settings?.ytdlp, + binary: e.target.value ? e.target.value : undefined, + }, + }, + })); + }, + [setNode], + ); + + const onYtdlpTempPathChange = useCallback( + (e: React.ChangeEvent) => { + setNode((p: Node) => ({ + ...p, + settings: { + ...p.settings, + ytdlp: { + ...p.settings?.ytdlp, + temp_path: e.target.value ? e.target.value : undefined, + }, + }, + })); + }, + [setNode], + ); + + const onEditedConfigYtdlpBlur = useCallback( + (value: string) => { + var res: Record | undefined = undefined; + if (value) { + try { + res = JSON.parse(value); + } catch (e) { + console.error(e); + } + } + setNode((p: Node) => ({ + ...p, + settings: { ...p.settings, ytdlp: { ...p.settings?.ytdlp, options: res } }, + })); + }, + [editedConfigYtdlp, setNode], + ); + const onQBittorrentServerChange = useCallback( (e: React.ChangeEvent) => { setNode((p: Node) => ({ @@ -445,11 +502,16 @@ const CapabilitiesSection = () => { + + + {values.settings?.provider === DownloaderProvider.qbittorrent ? t("node.qbittorrentDes") - : t("node.aria2Des")} + : values.settings?.provider === DownloaderProvider.ytdlp + ? t("node.ytdlpDes") + : t("node.aria2Des")} @@ -596,6 +658,62 @@ const CapabilitiesSection = () => { )} + {values.settings?.provider === DownloaderProvider.ytdlp && ( + <> + + + + {t("node.ytdlpBinaryDes")} + + + + + }> + setEditedConfigYtdlp(value || "")} + onBlur={onEditedConfigYtdlpBlur} + height="200px" + minHeight="200px" + options={{ + wordWrap: "on", + minimap: { enabled: false }, + scrollBeyondLastLine: false, + }} + /> + + + , + ]} + /> + + + + + + + {t("node.tempPathDes")} + + + + )} + " CLI flags verbatim (bool true -> flag only). + YtdlpSetting struct { + Binary string `json:"binary,omitempty"` + Options map[string]any `json:"options,omitempty"` + TempPath string `json:"temp_path,omitempty"` + } + TaskPublicState struct { Error string `json:"error,omitempty"` ErrorHistory []string `json:"error_history,omitempty"` @@ -391,6 +400,7 @@ const ( const ( DownloaderProviderAria2 = DownloaderProvider("aria2") DownloaderProviderQBittorrent = DownloaderProvider("qbittorrent") + DownloaderProviderYtDlp = DownloaderProvider("ytdlp") ) type ( diff --git a/pkg/cluster/node.go b/pkg/cluster/node.go index 831259e4..dfd7ab29 100644 --- a/pkg/cluster/node.go +++ b/pkg/cluster/node.go @@ -18,6 +18,7 @@ import ( "github.com/cloudreve/Cloudreve/v4/pkg/downloader/aria2" "github.com/cloudreve/Cloudreve/v4/pkg/downloader/qbittorrent" "github.com/cloudreve/Cloudreve/v4/pkg/downloader/slave" + "github.com/cloudreve/Cloudreve/v4/pkg/downloader/ytdlp" "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs" "github.com/cloudreve/Cloudreve/v4/pkg/logging" "github.com/cloudreve/Cloudreve/v4/pkg/queue" @@ -219,6 +220,8 @@ func (b *masterNode) CreateDownloader(ctx context.Context, c request.Client, set func NewDownloader(ctx context.Context, c request.Client, settings setting.Provider, options *types.NodeSetting) (downloader.Downloader, error) { if options.Provider == types.DownloaderProviderQBittorrent { return qbittorrent.NewClient(logging.FromContext(ctx), c, settings, options.QBittorrentSetting) + } else if options.Provider == types.DownloaderProviderYtDlp { + return ytdlp.New(logging.FromContext(ctx), settings, options.YtdlpSetting), nil } else if options.Provider == types.DownloaderProviderAria2 { return aria2.New(logging.FromContext(ctx), settings, options.Aria2Setting), nil } else if options.Provider == "" { diff --git a/pkg/downloader/ytdlp/ytdlp.go b/pkg/downloader/ytdlp/ytdlp.go new file mode 100644 index 00000000..77c148e1 --- /dev/null +++ b/pkg/downloader/ytdlp/ytdlp.go @@ -0,0 +1,387 @@ +package ytdlp + +import ( + "bufio" + "context" + "fmt" + "io" + "io/fs" + "os" + "os/exec" + "path/filepath" + "regexp" + "sort" + "strconv" + "strings" + "sync" + + "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/cloudreve/Cloudreve/v4/pkg/downloader" + "github.com/cloudreve/Cloudreve/v4/pkg/logging" + "github.com/cloudreve/Cloudreve/v4/pkg/setting" + "github.com/cloudreve/Cloudreve/v4/pkg/util" + "github.com/gofrs/uuid" + "github.com/samber/lo" +) + +const ( + // YtDlpTempFolder is the sub-directory under the node's temp path holding + // per-task download directories. + YtDlpTempFolder = "ytdlp" + defaultBinary = "yt-dlp" +) + +var ( + // [download] 12.3% of ~10.50MiB at 1.23MiB/s ETA 00:08 + progressRe = regexp.MustCompile(`\[download\]\s+([\d.]+)%\s+of\s+~?([\d.]+\s*(?:[KMGTP]i?B|B))\s+at\s+([\d.]+\s*(?:[KMGTP]i?B|B))/s`) + // [download] 100% of 10.50MiB in 00:08 + completeRe = regexp.MustCompile(`\[download\]\s+100%`) + sizeRe = regexp.MustCompile(`^([\d.]+)\s*([KMGTP]i?B|B)$`) +) + +type ( + ytdlpClient struct { + l logging.Logger + settings setting.Provider + options *types.YtdlpSetting + + mu sync.Mutex + procs map[string]*ytdlpProcess + } + + ytdlpProcess struct { + cmd *exec.Cmd + dir string + status *downloader.TaskStatus + + mu sync.Mutex + done chan struct{} + errTail []string + } +) + +// New creates a yt-dlp based Downloader. yt-dlp is a CLI tool, so the client +// manages spawned processes and parses their progress output. +func New(l logging.Logger, settings setting.Provider, options *types.YtdlpSetting) downloader.Downloader { + return &ytdlpClient{ + l: l, + settings: settings, + options: options, + procs: map[string]*ytdlpProcess{}, + } +} + +func (c *ytdlpClient) binary() string { + if c.options != nil && c.options.Binary != "" { + return c.options.Binary + } + return defaultBinary +} + +func (c *ytdlpClient) tempPath(ctx context.Context) string { + base := "" + if c.options != nil { + base = util.RelativePath(c.options.TempPath) + } + if base == "" && c.settings != nil { + base = util.DataPath(c.settings.TempPath(ctx)) + } + guid, _ := uuid.NewV4() + return filepath.Join(base, YtDlpTempFolder, guid.String()) +} + +// optionArgs converts the settings+per-task option maps into yt-dlp CLI flags. +// Keys are used verbatim: "format" -> "--format ", "no_playlist" (bool) -> +// "--no-playlist". This mirrors how aria2 options map to RPC fields. +func optionArgs(options map[string]any) []string { + keys := lo.Keys(options) + sort.Strings(keys) + + args := make([]string, 0, len(options)*2) + for _, k := range keys { + v := options[k] + flag := "--" + strings.ReplaceAll(k, "_", "-") + switch val := v.(type) { + case bool: + if val { + args = append(args, flag) + } + case nil: + case []string: + for _, item := range val { + args = append(args, flag, item) + } + case []any: + for _, item := range val { + args = append(args, flag, fmt.Sprintf("%v", item)) + } + default: + args = append(args, flag, fmt.Sprintf("%v", v)) + } + } + return args +} + +func (c *ytdlpClient) CreateTask(ctx context.Context, url string, options map[string]interface{}) (*downloader.TaskHandle, error) { + dir := c.tempPath(ctx) + if err := os.MkdirAll(dir, 0o755); err != nil { + return nil, fmt.Errorf("failed to create temp dir: %w", err) + } + + merged := map[string]any{} + if c.options != nil { + for k, v := range c.options.Options { + merged[k] = v + } + } + for k, v := range options { + merged[k] = v + } + + // "output" is a pseudo-key: it replaces the -o template rather than + // producing a "--output" flag (which yt-dlp does not have). + output := "%(title)s.%(ext)s" + if v, ok := merged["output"]; ok { + output = fmt.Sprintf("%v", v) + delete(merged, "output") + } + + args := []string{ + "--newline", "--no-colors", + "-P", dir, + "-o", output, + } + args = append(args, optionArgs(merged)...) + // "--" guards against URLs that could be mistaken for flags. + args = append(args, "--", url) + + c.l.Info("Creating yt-dlp task: %s %s", c.binary(), strings.Join(args, " ")) + // Deliberately not CommandContext: the queue ctx is canceled when each + // task iteration returns, which would kill the process as soon as the + // create phase suspends. Lifetime is managed via Cancel/exit instead. + cmd := exec.Command(c.binary(), args...) + stdout, err := cmd.StdoutPipe() + if err != nil { + return nil, fmt.Errorf("failed to open yt-dlp stdout: %w", err) + } + stderr, err := cmd.StderrPipe() + if err != nil { + return nil, fmt.Errorf("failed to open yt-dlp stderr: %w", err) + } + + if err := cmd.Start(); err != nil { + return nil, fmt.Errorf("failed to start yt-dlp: %w", err) + } + + p := &ytdlpProcess{ + cmd: cmd, + dir: dir, + done: make(chan struct{}), + status: &downloader.TaskStatus{ + Name: url, + State: downloader.StatusDownloading, + SavePath: filepath.ToSlash(dir), + }, + } + go p.run(stdout, stderr) + + c.mu.Lock() + c.procs[p.id()] = p + c.mu.Unlock() + + return &downloader.TaskHandle{ID: p.id(), Hash: dir}, nil +} + +func (p *ytdlpProcess) id() string { return fmt.Sprintf("%d|%s", p.cmd.Process.Pid, p.dir) } + +// run consumes the process output until exit and records terminal state. +func (p *ytdlpProcess) run(stdout, stderr io.Reader) { + defer close(p.done) + var wg sync.WaitGroup + wg.Add(2) + go func() { + defer wg.Done() + p.scanProgress(stdout) + }() + go func() { + defer wg.Done() + p.scanErrors(stderr) + }() + + err := p.cmd.Wait() + wg.Wait() + + p.mu.Lock() + defer p.mu.Unlock() + if err != nil { + p.status.State = downloader.StatusError + if len(p.errTail) > 0 { + p.status.ErrorMessage = p.errTail[len(p.errTail)-1] + } else { + p.status.ErrorMessage = err.Error() + } + return + } + p.status.State = downloader.StatusCompleted + if p.status.Total == 0 { + p.status.Total = p.status.Downloaded + } +} + +func (p *ytdlpProcess) scanProgress(r io.Reader) { + scanner := bufio.NewScanner(r) + scanner.Buffer(make([]byte, 64*1024), 1024*1024) + for scanner.Scan() { + line := scanner.Text() + if m := progressRe.FindStringSubmatch(line); m != nil { + p.mu.Lock() + p.status.Total = parseSize(m[2]) + p.status.DownloadSpeed = parseSize(m[3]) + if p.status.Total > 0 { + pct, _ := strconv.ParseFloat(m[1], 64) + p.status.Downloaded = int64(pct / 100 * float64(p.status.Total)) + } + p.mu.Unlock() + continue + } + if completeRe.MatchString(line) { + p.mu.Lock() + p.status.Downloaded = p.status.Total + p.mu.Unlock() + } + } +} + +func (p *ytdlpProcess) scanErrors(r io.Reader) { + scanner := bufio.NewScanner(r) + scanner.Buffer(make([]byte, 64*1024), 1024*1024) + for scanner.Scan() { + line := strings.TrimSpace(scanner.Text()) + if line == "" { + continue + } + p.mu.Lock() + p.errTail = append(p.errTail, line) + if len(p.errTail) > 8 { + p.errTail = p.errTail[len(p.errTail)-8:] + } + p.mu.Unlock() + } +} + +// parseSize converts yt-dlp human sizes ("10.5MiB", "2 GiB") to bytes. +func parseSize(s string) int64 { + m := sizeRe.FindStringSubmatch(strings.TrimSpace(s)) + if m == nil { + return 0 + } + num, err := strconv.ParseFloat(m[1], 64) + if err != nil { + return 0 + } + units := map[string]float64{ + "B": 1, "KB": 1e3, "MB": 1e6, "GB": 1e9, "TB": 1e12, "PB": 1e15, + "KiB": 1 << 10, "MiB": 1 << 20, "GiB": 1 << 30, "TiB": 1 << 40, "PiB": 1 << 50, + } + mult, ok := units[m[2]] + if !ok { + return 0 + } + return int64(num * mult) +} + +func (c *ytdlpClient) getProcess(id string) (*ytdlpProcess, error) { + c.mu.Lock() + defer c.mu.Unlock() + if p, ok := c.procs[id]; ok { + return p, nil + } + return nil, downloader.ErrTaskNotFount +} + +func (c *ytdlpClient) Info(ctx context.Context, handle *downloader.TaskHandle) (*downloader.TaskStatus, error) { + p, err := c.getProcess(handle.ID) + if err != nil { + return nil, err + } + + p.mu.Lock() + status := *p.status + terminal := status.State == downloader.StatusCompleted || status.State == downloader.StatusError + p.mu.Unlock() + + // Enumerate downloaded files for the transfer phase. + files, err := listDirFiles(p.dir) + if err != nil { + return nil, err + } + status.Files = files + + if terminal { + // Terminal status is never queried again (monitor moves to the + // transfer phase) - release the process entry. + c.mu.Lock() + delete(c.procs, handle.ID) + c.mu.Unlock() + } + return &status, nil +} + +// listDirFiles returns downloaded files relative to dir, mimicking the aria2 +// file list shape the transfer phase consumes. +func listDirFiles(dir string) ([]downloader.TaskFile, error) { + files := []downloader.TaskFile{} + err := filepath.WalkDir(dir, func(path string, d fs.DirEntry, err error) error { + if err != nil || d.IsDir() { + return err + } + info, err := d.Info() + if err != nil { + return err + } + rel, err := filepath.Rel(dir, path) + if err != nil { + return err + } + files = append(files, downloader.TaskFile{ + Index: len(files), + Name: filepath.ToSlash(rel), + Size: info.Size(), + Progress: 1, + Selected: true, + }) + return nil + }) + return files, err +} + +func (c *ytdlpClient) Cancel(ctx context.Context, handle *downloader.TaskHandle) error { + p, err := c.getProcess(handle.ID) + if err != nil { + return nil // already gone + } + if p.cmd.Process != nil { + _ = p.cmd.Process.Kill() + } + c.mu.Lock() + for k, v := range c.procs { + if v == p { + delete(c.procs, k) + } + } + c.mu.Unlock() + return nil +} + +// SetFilesToDownload is a no-op: yt-dlp resolves its own file list from the URL. +func (c *ytdlpClient) SetFilesToDownload(ctx context.Context, handle *downloader.TaskHandle, args ...*downloader.SetFileToDownloadArgs) error { + return nil +} + +func (c *ytdlpClient) Test(ctx context.Context) (string, error) { + out, err := exec.CommandContext(ctx, c.binary(), "--version").Output() + if err != nil { + return "", fmt.Errorf("yt-dlp test failed: %w", err) + } + return fmt.Sprintf("yt-dlp %s", strings.TrimSpace(string(out))), nil +} diff --git a/pkg/downloader/ytdlp/ytdlp_test.go b/pkg/downloader/ytdlp/ytdlp_test.go new file mode 100644 index 00000000..538907af --- /dev/null +++ b/pkg/downloader/ytdlp/ytdlp_test.go @@ -0,0 +1,171 @@ +package ytdlp + +import ( + "context" + "os" + "path/filepath" + "runtime" + "testing" + "time" + + "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/cloudreve/Cloudreve/v4/pkg/downloader" + "github.com/cloudreve/Cloudreve/v4/pkg/logging" + "github.com/cloudreve/Cloudreve/v4/pkg/setting" + "github.com/stretchr/testify/require" +) + +type stubSettingProvider struct { + setting.Provider +} + +func TestParseSize(t *testing.T) { + cases := map[string]int64{ + "10.50MiB": 10.5 * (1 << 20), + "2 GiB": 2 << 30, + "500KB": 500 * 1e3, + "1.5B": 1, + "": 0, + "10.50ZiB": 0, + "notasize": 0, + } + for in, want := range cases { + require.Equal(t, want, parseSize(in), "input %q", in) + } +} + +func TestOptionArgs(t *testing.T) { + args := optionArgs(map[string]any{ + "no_playlist": true, + "ignore_errors": false, + "format": "best", + "retries": 3, + "nil_key": nil, + "add_headers": []string{"A: 1", "B: 2"}, + }) + require.Equal(t, []string{ + "--add-headers", "A: 1", "--add-headers", "B: 2", + "--format", "best", + "--no-playlist", + "--retries", "3", + }, args) +} + +// fakeYtdlp writes a shell script emulating yt-dlp's progress output and file +// creation under the -P directory, then returns its path. +func fakeYtdlp(t *testing.T, script string) string { + t.Helper() + if runtime.GOOS == "windows" { + t.Skip("shell-script stub requires a POSIX shell") + } + path := filepath.Join(t.TempDir(), "yt-dlp") + require.NoError(t, os.WriteFile(path, []byte("#!/bin/sh\n"+script), 0o755)) + return path +} + +func newTestClient(binary, temp string) downloader.Downloader { + return New( + logging.NewConsoleLogger(logging.LevelError), + &stubSettingProvider{}, + &types.YtdlpSetting{Binary: binary, TempPath: temp}, + ) +} + +func awaitInfo(t *testing.T, d downloader.Downloader, h *downloader.TaskHandle, want downloader.Status) *downloader.TaskStatus { + t.Helper() + deadline := time.Now().Add(5 * time.Second) + var last *downloader.TaskStatus + for time.Now().Before(deadline) { + s, err := d.Info(context.Background(), h) + if err != nil { + // Terminal entries are released after being reported once. + if want == downloader.StatusCompleted || want == downloader.StatusError { + return last + } + t.Fatalf("info failed: %v", err) + } + last = s + if s.State == want { + return s + } + time.Sleep(20 * time.Millisecond) + } + t.Fatalf("state %q not reached, last: %+v", want, last) + return nil +} + +func TestCreateTaskCompletes(t *testing.T) { + binary := fakeYtdlp(t, ` +dir="" +prev="" +for a in "$@"; do + if [ "$prev" = "-P" ]; then dir="$a"; fi + prev="$a" +done +echo "[download] 50.0% of ~10.00MiB at 1.00MiB/s ETA 00:05" +echo "[download] 100% of 10.00MiB in 00:01" +echo data > "$dir/video.mp4" +`) + d := newTestClient(binary, t.TempDir()) + + handle, err := d.CreateTask(context.Background(), "https://example.com/v", map[string]interface{}{ + "no_playlist": true, + }) + require.NoError(t, err) + require.NotEmpty(t, handle.ID) + + status := awaitInfo(t, d, handle, downloader.StatusCompleted) + require.NotNil(t, status) + require.Equal(t, int64(10<<20), status.Total) + require.Len(t, status.Files, 1) + require.Equal(t, "video.mp4", status.Files[0].Name) + require.True(t, status.Files[0].Selected) + require.Equal(t, filepath.ToSlash(handle.Hash), status.SavePath) +} + +func TestCreateTaskError(t *testing.T) { + binary := fakeYtdlp(t, `echo "ERROR: Video unavailable" >&2; exit 1`) + d := newTestClient(binary, t.TempDir()) + + handle, err := d.CreateTask(context.Background(), "https://example.com/v", nil) + require.NoError(t, err) + + status := awaitInfo(t, d, handle, downloader.StatusError) + require.NotNil(t, status) + require.Contains(t, status.ErrorMessage, "Video unavailable") +} + +func TestCancel(t *testing.T) { + binary := fakeYtdlp(t, `sleep 30`) + d := newTestClient(binary, t.TempDir()) + + handle, err := d.CreateTask(context.Background(), "https://example.com/v", nil) + require.NoError(t, err) + require.NoError(t, d.Cancel(context.Background(), handle)) + + _, err = d.Info(context.Background(), handle) + require.ErrorIs(t, err, downloader.ErrTaskNotFount) +} + +func TestCustomOutputTemplate(t *testing.T) { + binary := fakeYtdlp(t, ` +dir="" +tmpl="" +prev="" +for a in "$@"; do + if [ "$prev" = "-P" ]; then dir="$a"; fi + if [ "$prev" = "-o" ]; then tmpl="$a"; fi + prev="$a" +done +echo data > "$dir/$tmpl" +`) + d := newTestClient(binary, t.TempDir()) + handle, err := d.CreateTask(context.Background(), "https://example.com/v", map[string]interface{}{ + "output": "custom-name.mp4", + }) + require.NoError(t, err) + status := awaitInfo(t, d, handle, downloader.StatusCompleted) + require.NotNil(t, status) + require.Len(t, status.Files, 1) + require.Equal(t, "custom-name.mp4", status.Files[0].Name) +} diff --git a/pkg/filemanager/workflows/remote_download.go b/pkg/filemanager/workflows/remote_download.go index 2f49d381..facc75b4 100644 --- a/pkg/filemanager/workflows/remote_download.go +++ b/pkg/filemanager/workflows/remote_download.go @@ -299,6 +299,19 @@ func (m *RemoteDownloadTask) buildDownloadOptions(ctx context.Context, base map[ } } } + case types.DownloaderProviderYtDlp: + if isHttpSrc { + if m.state.FileName != "" { + options["output"] = m.state.FileName + } + if m.state.HTTPUsername != "" { + options["username"] = m.state.HTTPUsername + options["password"] = m.state.HTTPPassword + } + if len(m.state.HTTPHeaders) > 0 { + options["add_headers"] = m.state.HTTPHeaders + } + } default: if isHttpSrc { if m.state.FileName != "" {