package service import ( "bufio" "database/sql" "encoding/json" "fmt" "log" "os" "os/exec" "path/filepath" "sort" "strings" "sync" "syscall" "time" "github.com/BurntSushi/toml" "github.com/gabriel-vasile/mimetype" "vidarchive/internal/config" "vidarchive/internal/models" "vidarchive/internal/repository" ) type DownloadService struct { repo *repository.DownloadRepository librarySvc *LibraryService presetSvc *PresetService settingsSvc *SettingsService cfg *config.Config cache *ProgressCache processMu sync.Mutex processes map[int64]*os.Process } func NewDownloadService(repo *repository.DownloadRepository, librarySvc *LibraryService, presetSvc *PresetService, settingsSvc *SettingsService, cfg *config.Config) *DownloadService { return &DownloadService{ repo: repo, librarySvc: librarySvc, presetSvc: presetSvc, settingsSvc: settingsSvc, cfg: cfg, cache: NewProgressCache(), processes: make(map[int64]*os.Process), } } func (s *DownloadService) Create(url string, presetID *int64, formatOverride, customFlags, outputDir string) (*models.Download, error) { d := &models.Download{ URL: url, Status: "queued", FormatOverride: formatOverride, CustomFlags: customFlags, OutputDir: sql.NullString{String: outputDir, Valid: outputDir != ""}, } if presetID != nil { d.PresetID = sqlNullInt64(*presetID) } if err := s.repo.Create(d); err != nil { return nil, err } return d, nil } func (s *DownloadService) GetByID(id int64) (*models.Download, error) { if live, ok := s.cache.Get(id); ok { d, err := s.repo.GetByID(id) if err != nil { return nil, err } logs := live.Logs.String() if logs != "" { d.Logs = sql.NullString{String: logs, Valid: true} } return d, nil } return s.repo.GetByID(id) } func (s *DownloadService) GetAll(status, sortBy string) ([]*models.Download, error) { downloads, err := s.repo.GetAll(status, sortBy) if err != nil { return nil, err } for _, d := range downloads { if live, ok := s.cache.Get(d.ID); ok { logs := live.Logs.String() if logs != "" { d.Logs = sql.NullString{String: logs, Valid: true} } } } return downloads, nil } func (s *DownloadService) GetQueued(limit int) ([]*models.Download, error) { return s.repo.GetQueued(limit) } func (s *DownloadService) Delete(id int64) error { s.killProcess(id) s.cache.Delete(id) return s.repo.Delete(id) } func (s *DownloadService) killProcess(id int64) { s.processMu.Lock() proc, ok := s.processes[id] delete(s.processes, id) s.processMu.Unlock() if !ok || proc == nil { return } _ = syscall.Kill(-proc.Pid, syscall.SIGKILL) } func (s *DownloadService) DeleteAll() error { return s.repo.DeleteAll() } func (s *DownloadService) ResetStalledDownloads() error { return s.repo.UpdateStatusWhere("downloading", "queued") } func (s *DownloadService) ListFormats(url string) ([]*models.FormatInfo, error) { cmd := exec.Command(s.cfg.YTDLPPath, "-F", "--no-warnings", url) output, err := cmd.CombinedOutput() if err != nil { return nil, fmt.Errorf("yt-dlp -F failed: %w\nOutput: %s", err, string(output)) } return parseFormatList(string(output)), nil } func (s *DownloadService) ExecuteDownload(d *models.Download) error { if err := s.repo.MarkStarted(d.ID); err != nil { return err } s.cache.Set(d.ID, &LiveDownload{LastUpdate: time.Now()}) defer s.cache.Delete(d.ID) var preset *models.Preset var err error if d.PresetID.Valid { preset, err = s.presetSvc.GetByID(d.PresetID.Int64) if err != nil { preset, _ = s.presetSvc.GetDefault() } } else { preset, _ = s.presetSvc.GetDefault() } if preset == nil { preset = &models.Preset{} } tempDownloadDir := filepath.Join(s.cfg.TempDir, fmt.Sprintf("%d", d.ID)) if err := os.MkdirAll(tempDownloadDir, 0755); err != nil { return fmt.Errorf("create temp download dir: %w", err) } args := s.presetSvc.BuildArgs(preset, d.FormatOverride, d.CustomFlags) cookies, err := s.settingsSvc.GetCookies() if err == nil && strings.TrimSpace(cookies) != "" { tmpFile, err := os.CreateTemp("", "cookies-*.txt") if err == nil { tmpFile.WriteString(cookies) tmpFile.Close() args = append(args, "--cookies", tmpFile.Name()) defer os.Remove(tmpFile.Name()) } } args = append(args, "-P", tempDownloadDir) args = append(args, "-o", "item-%(autonumber)05d/%(title)s.%(ext)s") args = append(args, d.URL) cmd := exec.Command(s.cfg.YTDLPPath, args...) cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} stdout, err := cmd.StdoutPipe() if err != nil { s.finalizeError(d.ID, err) return err } cmd.Stderr = cmd.Stdout if err := cmd.Start(); err != nil { s.finalizeError(d.ID, err) return err } s.processMu.Lock() s.processes[d.ID] = cmd.Process s.processMu.Unlock() defer func() { s.processMu.Lock() delete(s.processes, d.ID) s.processMu.Unlock() }() scanner := bufio.NewScanner(stdout) ticker := time.NewTicker(10 * time.Second) defer ticker.Stop() done := make(chan struct{}) go func() { for { select { case <-ticker.C: logs := s.cache.FlushLogs(d.ID) if logs != "" { s.repo.UpdateLogs(d.ID, logs) } case <-done: return } } }() for scanner.Scan() { line := scanner.Text() s.cache.AppendLog(d.ID, line) } close(done) logs := s.cache.FlushLogs(d.ID) if logs != "" { s.repo.UpdateLogs(d.ID, logs) } if err := cmd.Wait(); err != nil { s.finalizeError(d.ID, err) return err } if err := s.repo.MarkCompleted(d.ID, "completed"); err != nil { return err } if err := s.importDownloadedItems(d, tempDownloadDir); err != nil { log.Printf("Download %d completed but import failed: %v", d.ID, err) } return nil } func (s *DownloadService) finalizeError(id int64, err error) { logs := s.cache.FlushLogs(id) if logs != "" { s.repo.UpdateLogs(id, logs) } s.repo.MarkError(id, err.Error()) } func (s *DownloadService) importDownloadedItems(d *models.Download, tempDownloadDir string) error { entries, err := os.ReadDir(tempDownloadDir) if err != nil { return err } baseLibraryDir := s.cfg.LibraryDir if d.OutputDir.Valid && d.OutputDir.String != "" { cleanDir := filepath.Clean(d.OutputDir.String) fullPath := filepath.Join(baseLibraryDir, cleanDir) resolvedPath, err := filepath.Abs(fullPath) if err != nil { return fmt.Errorf("invalid output directory: %w", err) } resolvedLibraryDir, _ := filepath.Abs(baseLibraryDir) if !strings.HasPrefix(resolvedPath, resolvedLibraryDir+string(filepath.Separator)) && resolvedPath != resolvedLibraryDir { return fmt.Errorf("invalid output directory: path traversal attempt detected") } baseLibraryDir = fullPath } if err := os.MkdirAll(baseLibraryDir, 0755); err != nil { return err } var itemDirs []string for _, entry := range entries { if !entry.IsDir() { continue } name := entry.Name() if strings.HasPrefix(name, "item-") { itemDirs = append(itemDirs, filepath.Join(tempDownloadDir, name)) } } sort.Strings(itemDirs) for _, itemDir := range itemDirs { if err := s.importItemDir(d.URL, itemDir, baseLibraryDir); err != nil { log.Printf("warning: failed to import item %s: %v", itemDir, err) } } os.Remove(tempDownloadDir) return nil } func (s *DownloadService) importItemDir(url, itemDir, baseLibraryDir string) error { entries, err := os.ReadDir(itemDir) if err != nil { return err } var mediaFiles []os.DirEntry var infoJSONPath string var subtitleFiles []string for _, entry := range entries { if entry.IsDir() { continue } name := entry.Name() path := filepath.Join(itemDir, name) ext := strings.ToLower(filepath.Ext(name)) if name == "info.json" || strings.HasSuffix(name, ".info.json") { infoJSONPath = path continue } if ext == ".vtt" || ext == ".srt" || ext == ".ass" || ext == ".ssa" { subtitleFiles = append(subtitleFiles, path) continue } mtype, err := mimetype.DetectFile(path) if err == nil && mtype != nil && (strings.HasPrefix(mtype.String(), "audio/") || strings.HasPrefix(mtype.String(), "video/")) { mediaFiles = append(mediaFiles, entry) } } if len(mediaFiles) == 0 { return fmt.Errorf("no media files found in %s", itemDir) } name := s.deriveItemName(itemDir, infoJSONPath, mediaFiles) targetDir := s.uniqueDir(baseLibraryDir, name) if err := os.MkdirAll(targetDir, 0755); err != nil { return err } if infoJSONPath != "" { if err := os.Rename(infoJSONPath, filepath.Join(targetDir, "info.json")); err != nil { return err } } for _, entry := range mediaFiles { if err := os.Rename(filepath.Join(itemDir, entry.Name()), filepath.Join(targetDir, entry.Name())); err != nil { return err } } if len(subtitleFiles) > 0 { subtitlesDir := filepath.Join(targetDir, subtitlesDirName) if err := os.MkdirAll(subtitlesDir, 0755); err != nil { return err } for _, sf := range subtitleFiles { if err := os.Rename(sf, filepath.Join(subtitlesDir, filepath.Base(sf))); err != nil { return err } } } metadata := models.ItemMetadata{ Name: name, SourceURL: url, Duration: -1, } markerPath := filepath.Join(targetDir, itemMarkerName) f, err := os.Create(markerPath) if err != nil { return err } defer f.Close() if err := toml.NewEncoder(f).Encode(metadata); err != nil { return err } return nil } func (s *DownloadService) deriveItemName(itemDir, infoJSONPath string, mediaFiles []os.DirEntry) string { if infoJSONPath != "" { data, err := os.ReadFile(infoJSONPath) if err == nil { var info struct { Title string `json:"title"` } if err := json.Unmarshal(data, &info); err == nil && info.Title != "" { return sanitizeDirName(info.Title) } } } sort.Slice(mediaFiles, func(i, j int) bool { ii, _ := os.Stat(filepath.Join(itemDir, mediaFiles[i].Name())) jj, _ := os.Stat(filepath.Join(itemDir, mediaFiles[j].Name())) if ii == nil || jj == nil { return false } return ii.Size() > jj.Size() }) base := strings.TrimSuffix(mediaFiles[0].Name(), filepath.Ext(mediaFiles[0].Name())) return sanitizeDirName(base) } func (s *DownloadService) uniqueDir(base, name string) string { dir := filepath.Join(base, name) if _, err := os.Stat(dir); os.IsNotExist(err) { return dir } for i := 1; ; i++ { candidate := fmt.Sprintf("%s-%d", dir, i) if _, err := os.Stat(candidate); os.IsNotExist(err) { return candidate } } } func sanitizeDirName(name string) string { name = strings.TrimSpace(name) replacer := strings.NewReplacer( "/", "-", "\\", "-", ":", "-", "*", "-", "?", "-", "\"", "-", "<", "-", ">", "-", "|", "-", ) name = replacer.Replace(name) name = strings.TrimSpace(name) if name == "" { name = "untitled" } return name } func parseFormatList(output string) []*models.FormatInfo { lines := strings.Split(output, "\n") var formats []*models.FormatInfo inFormats := false for _, line := range lines { line = strings.TrimSpace(line) if line == "" { continue } if strings.Contains(line, "ID") && strings.Contains(line, "EXT") { inFormats = true continue } if !inFormats { continue } parts := strings.Fields(line) if len(parts) >= 4 { format := &models.FormatInfo{ ID: parts[0], Ext: parts[1], } for i, part := range parts { if strings.Contains(part, "x") && !strings.Contains(part, "http") { format.Resolution = part if i+1 < len(parts) { format.FPS = parts[i+1] } break } } if len(parts) > 3 { format.Note = strings.Join(parts[3:], " ") } formats = append(formats, format) } } return formats } func sqlNullInt64(v int64) sql.NullInt64 { return sql.NullInt64{Int64: v, Valid: true} }