diff --git a/internal/config/config.go b/internal/config/config.go index 9d5a2bc..778e579 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -262,7 +262,7 @@ func setDefaults(v *viper.Viper) { v.SetDefault("downloads.smart_classify", true) v.SetDefault("organizer.smart_classify", false) v.SetDefault("organizer.auto_after_download", false) - v.SetDefault("organize.scrape_after", false) + v.SetDefault("organize.scrape_after", true) v.SetDefault("scrape.delay_min_ms", 250) v.SetDefault("scrape.delay_max_ms", 500) v.SetDefault("organizer.categories.chinese_movie", "华语电影") diff --git a/internal/handler/media.go b/internal/handler/media.go index 8173f47..0aabdce 100644 --- a/internal/handler/media.go +++ b/internal/handler/media.go @@ -97,6 +97,7 @@ func scanLibraryHandler(svc *service.Container) gin.HandlerFunc { return } if _, ok := service.ParseCloudLibraryMount(lib.Path); ok { + task := startScanHTTPTask(svc, "云盘扫描队列", lib.Name, lib.Path) if svc.WSHub != nil { svc.WSHub.Publish("scan", gin.H{ "library_id": id, @@ -108,6 +109,7 @@ func scanLibraryHandler(svc *service.Container) gin.HandlerFunc { }) } _, _, _ = svc.Scan.StartCloudLibraryScan(id, false) + finishHTTPTask(task, nil, "queued", "云盘扫描已加入后台队列", map[string]int64{"queued": 1}) c.JSON(http.StatusAccepted, gin.H{ "library_id": id, "visited": 0, @@ -121,15 +123,47 @@ func scanLibraryHandler(svc *service.Container) gin.HandlerFunc { }) return } + task := startScanHTTPTask(svc, "手动扫描入库", lib.Name, lib.Path) res, err := svc.Scan.ScanLibrary(c.Request.Context(), id) if err != nil { + finishHTTPTask(task, err, "scan", "手动扫描入库失败", nil) c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } + finishHTTPTask(task, nil, "completed", "手动扫描入库结束", scanTaskMetrics(res)) c.JSON(http.StatusOK, res) } } +func startScanHTTPTask(svc *service.Container, name, libraryName, path string) *service.TaskHandle { + if svc == nil || svc.Tasks == nil { + return nil + } + if libraryName != "" { + name += ":" + libraryName + } + return svc.Tasks.Start(service.TaskKindScan, name, service.TaskUpdate{ + Stage: "scan", + SourcePath: path, + Message: "正在扫描并入库", + }) +} + +func scanTaskMetrics(res *service.ScanResult) map[string]int64 { + if res == nil { + return nil + } + return map[string]int64{ + "visited": int64(res.Visited), + "added": int64(res.Added), + "updated": int64(res.Updated), + "skipped": int64(res.Skipped), + "probed": int64(res.Probed), + "local_metadata": int64(res.LocalMetadata), + "removed": res.Removed, + } +} + func listMediaHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { id := c.Param("id") diff --git a/internal/handler/organizer.go b/internal/handler/organizer.go index 38ff239..62c21dc 100644 --- a/internal/handler/organizer.go +++ b/internal/handler/organizer.go @@ -15,15 +15,16 @@ import ( // source_path = 源目录(待整理),dest_path = 目的地目录(整理输出)。 // target_path 为 dest_path 的向后兼容别名。 type organizeReq struct { - SourcePath string `json:"source_path"` - DestPath string `json:"dest_path"` - TargetPath string `json:"target_path"` // deprecated alias for dest_path - TransferMode string `json:"transfer_mode"` - MediaType string `json:"media_type"` - ScanAfter bool `json:"scan_after"` - ScrapeAfter *bool `json:"scrape_after"` - LibraryID string `json:"library_id"` - DryRun bool `json:"dry_run"` + SourcePath string `json:"source_path"` + DestPath string `json:"dest_path"` + TargetPath string `json:"target_path"` // deprecated alias for dest_path + TransferMode string `json:"transfer_mode"` + MediaType string `json:"media_type"` + MediaCategory string `json:"media_category"` + ScanAfter bool `json:"scan_after"` + ScrapeAfter *bool `json:"scrape_after"` + LibraryID string `json:"library_id"` + DryRun bool `json:"dry_run"` } // bindOrganizeOptions parses the optional JSON body into OrganizeOptions. @@ -40,10 +41,11 @@ func organizeOptionsFromReq(req organizeReq) service.OrganizeOptions { dest = strings.TrimSpace(req.TargetPath) } opts := service.OrganizeOptions{ - SourcePath: strings.TrimSpace(req.SourcePath), - DestPath: dest, - MediaType: strings.TrimSpace(req.MediaType), - DryRun: req.DryRun, + SourcePath: strings.TrimSpace(req.SourcePath), + DestPath: dest, + MediaType: strings.TrimSpace(req.MediaType), + MediaCategory: strings.TrimSpace(req.MediaCategory), + DryRun: req.DryRun, } if m := strings.TrimSpace(req.TransferMode); m != "" { opts.TransferMode = service.TransferMode(m) @@ -56,17 +58,21 @@ func organizeMediaHandler(svc *service.Container) gin.HandlerFunc { var req organizeReq _ = c.ShouldBindJSON(&req) opts := organizeOptionsFromReq(req) + task := startOrganizeHTTPTask(svc, "手动整理媒体", opts) dst, err := svc.Organizer.OrganizeMediaWithOptions(c.Request.Context(), c.Param("id"), opts) if err != nil { + finishHTTPTask(task, err, "organize", "手动整理媒体失败", nil) c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } payload := gin.H{"path": dst} if req.ScanAfter && !req.DryRun && svc.Scan != nil { + updateHTTPTask(task, "scan_scrape", "正在扫描入库并按设置刮削", nil) scans, scrapes := scanAndScrapeAfterOrganize(c, svc, dst, strings.TrimSpace(req.LibraryID), req.ScrapeAfter) payload["scans"] = scans payload["scrapes"] = scrapes } + finishHTTPTask(task, nil, "completed", "手动整理媒体结束", nil) c.JSON(http.StatusOK, payload) } } @@ -76,14 +82,18 @@ func organizeLibraryHandler(svc *service.Container) gin.HandlerFunc { var req organizeReq _ = c.ShouldBindJSON(&req) opts := organizeOptionsFromReq(req) + task := startOrganizeHTTPTask(svc, "手动整理媒体库", opts) res, err := svc.Organizer.OrganizeLibraryWithOptions(c.Request.Context(), c.Param("id"), opts) if err != nil { + finishHTTPTask(task, err, "organize", "手动整理媒体库失败", nil) c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } if req.ScanAfter && !req.DryRun && svc.Scan != nil { + updateHTTPTask(task, "scan_scrape", "正在扫描入库并按设置刮削", nil) res.Scans, res.Scrapes = scanAndScrapeAfterOrganize(c, svc, res.DestPath, c.Param("id"), req.ScrapeAfter) } + finishHTTPTask(task, nil, "completed", "手动整理媒体库结束", service.OrganizeTaskMetrics(res)) c.JSON(http.StatusOK, res) } } @@ -103,18 +113,53 @@ func organizeDirectoryHandler(svc *service.Container) gin.HandlerFunc { var req organizeReq _ = c.ShouldBindJSON(&req) opts := organizeOptionsFromReq(req) + task := startOrganizeHTTPTask(svc, "手动整理入库", opts) res, err := svc.Organizer.OrganizeDirectory(c.Request.Context(), opts) if err != nil { + finishHTTPTask(task, err, "organize", "手动整理入库失败", nil) c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) return } + updateHTTPTask(task, "organize", "手动整理完成,准备扫描入库", service.OrganizeTaskMetrics(res)) if req.ScanAfter && !req.DryRun && svc.Scan != nil { + updateHTTPTask(task, "scan_scrape", "正在扫描入库并按设置刮削", service.OrganizeTaskMetrics(res)) res.Scans, res.Scrapes = scanAndScrapeAfterOrganize(c, svc, res.DestPath, strings.TrimSpace(req.LibraryID), req.ScrapeAfter) } + finishHTTPTask(task, nil, "completed", "手动整理入库结束", service.OrganizeTaskMetrics(res)) c.JSON(http.StatusOK, res) } } +func startOrganizeHTTPTask(svc *service.Container, name string, opts service.OrganizeOptions) *service.TaskHandle { + if svc == nil || svc.Tasks == nil { + return nil + } + message := "正在整理/重命名/入库" + if opts.DryRun { + message = "正在预览整理/重命名" + } + return svc.Tasks.Start(service.TaskKindOrganize, name, service.TaskUpdate{ + Stage: "organize", + SourcePath: opts.SourcePath, + DestPath: opts.DestPath, + Message: message, + }) +} + +func updateHTTPTask(task *service.TaskHandle, stage, message string, metrics map[string]int64) { + if task == nil { + return + } + task.Update(service.TaskUpdate{Stage: stage, Message: message, Metrics: metrics}) +} + +func finishHTTPTask(task *service.TaskHandle, err error, stage, message string, metrics map[string]int64) { + if task == nil { + return + } + task.Finish(err, service.TaskUpdate{Stage: stage, Message: message, Metrics: metrics}) +} + func scanAndScrapeAfterOrganize(c *gin.Context, svc *service.Container, destRoot, preferredLibraryID string, scrapeOverride *bool) ([]service.OrganizeScanSummary, []service.OrganizeScrapeSummary) { scrapeAfter := service.OrganizeScrapeAfterEnabled(c.Request.Context(), svc.Repo) if scrapeOverride != nil { diff --git a/internal/handler/streaming.go b/internal/handler/streaming.go index 6e9b04e..cce91d2 100644 --- a/internal/handler/streaming.go +++ b/internal/handler/streaming.go @@ -81,11 +81,18 @@ func scrapeOneHandler(svc *service.Container) gin.HandlerFunc { c.JSON(http.StatusNotFound, gin.H{"error": "not found"}) return } + task := startScrapeHTTPTask(svc, "手动刮削媒体", m.Title, m.Path) if err := svc.Scraper.EnrichOne(c.Request.Context(), m); err != nil { + finishHTTPTask(task, err, "scrape", "手动刮削媒体失败", nil) c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } refreshed, _ := svc.Repo.Media.FindByID(c.Request.Context(), m.ID) + metrics := map[string]int64{"processed": 1} + if refreshed != nil && refreshed.ScrapeStatus == "matched" { + metrics["matched"] = 1 + } + finishHTTPTask(task, nil, "completed", "手动刮削媒体结束", metrics) c.JSON(http.StatusOK, refreshed) } } @@ -93,15 +100,44 @@ func scrapeOneHandler(svc *service.Container) gin.HandlerFunc { // scrapeLibraryHandler retries every pending/no_match media in a library. func scrapeLibraryHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { + libID := c.Param("id") + var task *service.TaskHandle + if lib, err := svc.Repo.Library.FindByID(c.Request.Context(), libID); err == nil && lib != nil { + task = startScrapeHTTPTask(svc, "手动刮削媒体库", lib.Name, lib.Path) + } else { + task = startScrapeHTTPTask(svc, "手动刮削媒体库", libID, "") + } // Run in the background so HTTP returns instantly; the WS hub // pushes per-item progress on the "scrape" topic. - go func(libID string) { - _, _ = svc.Scraper.EnrichLibrary(context.Background(), libID, true) - }(c.Param("id")) + go func(libID string, task *service.TaskHandle) { + matched, err := svc.Scraper.EnrichLibrary(context.Background(), libID, true) + metrics := map[string]int64{"matched": int64(matched)} + stage := "completed" + message := "手动刮削媒体库结束" + if err != nil { + stage = "scrape" + message = "手动刮削媒体库失败" + } + finishHTTPTask(task, err, stage, message, metrics) + }(libID, task) c.JSON(http.StatusAccepted, gin.H{"status": "scraping"}) } } +func startScrapeHTTPTask(svc *service.Container, name, title, path string) *service.TaskHandle { + if svc == nil || svc.Tasks == nil { + return nil + } + if title != "" { + name += ":" + title + } + return svc.Tasks.Start(service.TaskKindScrape, name, service.TaskUpdate{ + Stage: "scrape", + SourcePath: path, + Message: "正在刮削元数据", + }) +} + // reprobeHandler re-runs ffprobe against a single media. Admin-only. func reprobeHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { diff --git a/internal/handler/tasks.go b/internal/handler/tasks.go index 9e7cb56..e685db4 100644 --- a/internal/handler/tasks.go +++ b/internal/handler/tasks.go @@ -18,9 +18,14 @@ func tasksHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { transcodes := svc.Transcoder.Active() _, torrents, _ := svc.Downloads.List(c.Request.Context()) + background := service.TaskSnapshot{} + if svc.Tasks != nil { + background = svc.Tasks.Snapshot() + } c.JSON(http.StatusOK, gin.H{ - "transcodes": transcodes, - "torrents": torrents, + "transcodes": transcodes, + "torrents": torrents, + "background_tasks": background, }) } } diff --git a/internal/model/model.go b/internal/model/model.go index 90edbba..66a1cce 100644 --- a/internal/model/model.go +++ b/internal/model/model.go @@ -189,16 +189,22 @@ type PlaylistItem struct { // DownloadTask 是待处理(或已完成)的 torrent / HTTP 下载。 type DownloadTask struct { Base - UserID string `gorm:"index;size:36" json:"user_id"` - Source string `gorm:"size:32;not null" json:"source"` // qbittorrent / transmission / http - URL string `gorm:"size:2048;not null" json:"-"` - Title string `gorm:"size:512" json:"title,omitempty"` - PosterURL string `gorm:"size:2048" json:"poster_url,omitempty"` - BackdropURL string `gorm:"size:2048" json:"backdrop_url,omitempty"` - Overview string `gorm:"type:text" json:"overview,omitempty"` - SavePath string `gorm:"size:1024" json:"save_path"` - Status string `gorm:"size:32;default:queued" json:"status"` - Progress float32 `json:"progress"` + UserID string `gorm:"index;size:36" json:"user_id"` + Source string `gorm:"size:32;not null" json:"source"` // qbittorrent / transmission / http + URL string `gorm:"size:2048;not null" json:"-"` + Title string `gorm:"size:512" json:"title,omitempty"` + PosterURL string `gorm:"size:2048" json:"poster_url,omitempty"` + BackdropURL string `gorm:"size:2048" json:"backdrop_url,omitempty"` + Overview string `gorm:"type:text" json:"overview,omitempty"` + SavePath string `gorm:"size:1024" json:"save_path"` + MediaType string `gorm:"size:16" json:"media_type,omitempty"` + MediaCategory string `gorm:"size:128" json:"media_category,omitempty"` + Status string `gorm:"size:32;default:queued" json:"status"` + Progress float32 `json:"progress"` + + // AllowExistingLibrary is true for subscription wash/upgrade tasks that are + // allowed to replace an existing library item after download completion. + AllowExistingLibrary bool `gorm:"default:false" json:"allow_existing_library,omitempty"` } // Subscription 是自动化规则,轮询 RSS 源并将匹配种子排队到配置的下载客户端。 diff --git a/internal/service/downloads.go b/internal/service/downloads.go index 84489b9..430ce3f 100644 --- a/internal/service/downloads.go +++ b/internal/service/downloads.go @@ -46,6 +46,7 @@ type DownloadService struct { organizer *OrganizerService scanner *ScannerService site *SiteService + tasks *TaskTrackerService mu sync.Mutex stopCh chan struct{} @@ -61,6 +62,10 @@ func (d *DownloadService) SetScanner(scanner *ScannerService) { d.scanner = scanner } +func (d *DownloadService) SetTaskTracker(tasks *TaskTrackerService) { + d.tasks = tasks +} + var torrentEpisodeToken = regexp.MustCompile(`(?i)e\d{1,3}`) const settingDownloadClientsManaged = "download_clients.managed" @@ -98,42 +103,46 @@ type DownloadTaskMeta struct { } type DownloadTaskView struct { - ID string `json:"id"` - Source string `json:"source"` - Title string `json:"title"` - PosterURL string `json:"poster_url,omitempty"` - BackdropURL string `json:"backdrop_url,omitempty"` - Overview string `json:"overview,omitempty"` - SavePath string `json:"save_path"` - Status string `json:"status"` - Progress float32 `json:"progress"` - State string `json:"state,omitempty"` - DLSpeed int64 `json:"dlspeed,omitempty"` - UpSpeed int64 `json:"upspeed,omitempty"` - Size int64 `json:"size,omitempty"` - Downloaded int64 `json:"downloaded,omitempty"` - NumSeeds int `json:"num_seeds,omitempty"` - NumLeechs int `json:"num_leechs,omitempty"` - CreatedAt time.Time `json:"created_at"` - UpdatedAt time.Time `json:"updated_at"` + ID string `json:"id"` + Source string `json:"source"` + Title string `json:"title"` + PosterURL string `json:"poster_url,omitempty"` + BackdropURL string `json:"backdrop_url,omitempty"` + Overview string `json:"overview,omitempty"` + SavePath string `json:"save_path"` + MediaType string `json:"media_type,omitempty"` + MediaCategory string `json:"media_category,omitempty"` + Status string `json:"status"` + Progress float32 `json:"progress"` + State string `json:"state,omitempty"` + DLSpeed int64 `json:"dlspeed,omitempty"` + UpSpeed int64 `json:"upspeed,omitempty"` + Size int64 `json:"size,omitempty"` + Downloaded int64 `json:"downloaded,omitempty"` + NumSeeds int `json:"num_seeds,omitempty"` + NumLeechs int `json:"num_leechs,omitempty"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` } type DownloadTorrentView struct { - Hash string `json:"hash"` - Name string `json:"name"` - Title string `json:"title"` - PosterURL string `json:"poster_url,omitempty"` - BackdropURL string `json:"backdrop_url,omitempty"` - Overview string `json:"overview,omitempty"` - State string `json:"state"` - Progress float32 `json:"progress"` - DLSpeed int64 `json:"dlspeed"` - UpSpeed int64 `json:"upspeed"` - NumSeeds int `json:"num_seeds"` - NumLeechs int `json:"num_leechs"` - Size int64 `json:"size"` - Downloaded int64 `json:"downloaded"` - SavePath string `json:"save_path"` + Hash string `json:"hash"` + Name string `json:"name"` + Title string `json:"title"` + PosterURL string `json:"poster_url,omitempty"` + BackdropURL string `json:"backdrop_url,omitempty"` + Overview string `json:"overview,omitempty"` + MediaType string `json:"media_type,omitempty"` + MediaCategory string `json:"media_category,omitempty"` + State string `json:"state"` + Progress float32 `json:"progress"` + DLSpeed int64 `json:"dlspeed"` + UpSpeed int64 `json:"upspeed"` + NumSeeds int `json:"num_seeds"` + NumLeechs int `json:"num_leechs"` + Size int64 `json:"size"` + Downloaded int64 `json:"downloaded"` + SavePath string `json:"save_path"` } // NewDownloadService is the constructor. @@ -469,15 +478,18 @@ func (d *DownloadService) createTask(ctx context.Context, userID, urlStr, savePa title = publicDownloadTitle(urlStr) } t := &model.DownloadTask{ - UserID: userID, - Source: "qbittorrent", - URL: urlStr, - Title: title, - PosterURL: meta.PosterURL, - BackdropURL: meta.BackdropURL, - Overview: meta.Overview, - SavePath: savePath, - Status: "queued", + UserID: userID, + Source: "qbittorrent", + URL: urlStr, + Title: title, + PosterURL: meta.PosterURL, + BackdropURL: meta.BackdropURL, + Overview: meta.Overview, + SavePath: savePath, + MediaType: meta.MediaType, + MediaCategory: meta.MediaCategory, + Status: "queued", + AllowExistingLibrary: meta.AllowExistingLibrary, } if err := d.repo.Download.Create(ctx, t); err != nil { return nil, err @@ -608,24 +620,26 @@ func downloadTaskView(row model.DownloadTask, torrent QBitTorrent) DownloadTaskV } size := torrent.Size return DownloadTaskView{ - ID: row.ID, - Source: row.Source, - Title: firstNonEmpty(row.Title, "下载任务"), - PosterURL: row.PosterURL, - BackdropURL: row.BackdropURL, - Overview: row.Overview, - SavePath: row.SavePath, - Status: row.Status, - Progress: progress, - State: state, - DLSpeed: torrent.DLSpeed, - UpSpeed: torrent.UpSpeed, - Size: size, - Downloaded: downloadedBytes(size, progress), - NumSeeds: torrent.NumSeeds, - NumLeechs: torrent.NumLeech, - CreatedAt: row.CreatedAt, - UpdatedAt: row.UpdatedAt, + ID: row.ID, + Source: row.Source, + Title: firstNonEmpty(row.Title, "下载任务"), + PosterURL: row.PosterURL, + BackdropURL: row.BackdropURL, + Overview: row.Overview, + SavePath: row.SavePath, + MediaType: row.MediaType, + MediaCategory: row.MediaCategory, + Status: row.Status, + Progress: progress, + State: state, + DLSpeed: torrent.DLSpeed, + UpSpeed: torrent.UpSpeed, + Size: size, + Downloaded: downloadedBytes(size, progress), + NumSeeds: torrent.NumSeeds, + NumLeechs: torrent.NumLeech, + CreatedAt: row.CreatedAt, + UpdatedAt: row.UpdatedAt, } } @@ -635,21 +649,23 @@ func downloadTorrentView(torrent QBitTorrent, row model.DownloadTask) DownloadTo title = row.Title } return DownloadTorrentView{ - Hash: torrent.Hash, - Name: torrent.Name, - Title: firstNonEmpty(title, "下载任务"), - PosterURL: row.PosterURL, - BackdropURL: row.BackdropURL, - Overview: row.Overview, - State: torrent.State, - Progress: torrent.Progress, - DLSpeed: torrent.DLSpeed, - UpSpeed: torrent.UpSpeed, - NumSeeds: torrent.NumSeeds, - NumLeechs: torrent.NumLeech, - Size: torrent.Size, - Downloaded: downloadedBytes(torrent.Size, torrent.Progress), - SavePath: torrent.SavePath, + Hash: torrent.Hash, + Name: torrent.Name, + Title: firstNonEmpty(title, "下载任务"), + PosterURL: row.PosterURL, + BackdropURL: row.BackdropURL, + Overview: row.Overview, + MediaType: row.MediaType, + MediaCategory: firstNonEmpty(row.MediaCategory, torrent.Category), + State: torrent.State, + Progress: torrent.Progress, + DLSpeed: torrent.DLSpeed, + UpSpeed: torrent.UpSpeed, + NumSeeds: torrent.NumSeeds, + NumLeechs: torrent.NumLeech, + Size: torrent.Size, + Downloaded: downloadedBytes(torrent.Size, torrent.Progress), + SavePath: torrent.SavePath, } } @@ -1086,19 +1102,50 @@ func (d *DownloadService) onTorrentComplete(ctx context.Context, torrent QBitTor zap.String("content_path", torrent.ContentPath)) return } + taskRow, hasTask := d.completedTorrentTask(ctx, torrent) + allowReplace := hasTask && taskRow.AllowExistingLibrary d.log.Info("download completed, triggering directory organize", zap.String("hash", torrent.Hash), zap.String("name", torrent.Name), - zap.String("source", source)) - res, err := d.organizer.OrganizeDirectory(ctx, OrganizeOptions{SourcePath: source}) + zap.String("source", source), + zap.Bool("allow_replace_existing", allowReplace)) + taskHandle := d.startDownloadOrganizeTask(torrent, source, allowReplace) + res, err := d.organizer.OrganizeDirectory(ctx, OrganizeOptions{ + SourcePath: source, + MediaType: downloadTaskMediaType(taskRow), + MediaCategory: firstNonEmpty(downloadTaskMediaCategory(taskRow), torrent.Category), + AllowReplaceExisting: allowReplace, + }) if err != nil { + if taskHandle != nil { + taskHandle.Finish(err, TaskUpdate{ + Stage: "organize", + Message: "下载完成自动整理失败", + }) + } d.log.Error("auto organize completed torrent failed", zap.String("hash", torrent.Hash), zap.String("source", source), zap.Error(err)) return } + if taskHandle != nil && res != nil { + taskHandle.Update(TaskUpdate{ + Stage: "organize", + SourcePath: res.SourcePath, + DestPath: res.DestPath, + Message: "下载完成整理已完成,准备扫描入库", + Metrics: OrganizeTaskMetrics(res), + }) + } if d.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" && OrganizeResultHasChanges(res) { + if taskHandle != nil { + taskHandle.Update(TaskUpdate{ + Stage: "scan_scrape", + Message: "正在扫描入库并按设置刮削", + Metrics: OrganizeTaskMetrics(res), + }) + } res.Scans, res.Scrapes = d.scanner.ScanAndScrapeLibrariesForPath(ctx, res.DestPath, "", OrganizeScrapeAfterEnabled(ctx, d.repo)) } else if d.log != nil && res != nil && !OrganizeResultHasChanges(res) { d.log.Info("auto organize completed torrent skipped scan; no destination changes", @@ -1108,6 +1155,13 @@ func (d *DownloadService) onTorrentComplete(ctx context.Context, torrent QBitTor zap.Int("replaced", res.Replaced), zap.Int("skipped", res.Skipped)) } + if taskHandle != nil { + taskHandle.Finish(nil, TaskUpdate{ + Stage: "completed", + Message: "下载完成自动整理入库结束", + Metrics: OrganizeTaskMetrics(res), + }) + } d.markCompletedTorrentCatchupRecorded(context.Background(), torrent) d.log.Info("auto organize completed torrent finished", zap.String("hash", torrent.Hash), @@ -1120,6 +1174,59 @@ func (d *DownloadService) onTorrentComplete(ctx context.Context, torrent QBitTor zap.Int("errors", len(res.Errors))) } +func (d *DownloadService) startDownloadOrganizeTask(torrent QBitTorrent, source string, allowReplace bool) *TaskHandle { + if d == nil || d.tasks == nil { + return nil + } + message := "下载完成后自动整理/重命名/入库" + if allowReplace { + message = "下载完成后自动整理/重命名/入库(允许洗版替换)" + } + name := strings.TrimSpace(torrent.Name) + if name == "" { + name = "下载完成自动整理" + } + return d.tasks.Start(TaskKindOrganize, name, TaskUpdate{ + Stage: "organize", + SourcePath: source, + Message: message, + }) +} + +func (d *DownloadService) completedTorrentTask(ctx context.Context, torrent QBitTorrent) (*model.DownloadTask, bool) { + if d == nil || d.repo == nil || d.repo.Download == nil { + return nil, false + } + rows, err := d.repo.Download.List(ctx) + if err != nil || len(rows) == 0 { + return nil, false + } + taskByKey := tasksByIdentity(rows) + if task, ok := findMatchingTaskByIdentity(torrent.Name, taskByKey); ok { + return &task, true + } + if strings.TrimSpace(torrent.ContentPath) != "" { + if task, ok := findMatchingTaskByIdentity(filepath.Base(torrent.ContentPath), taskByKey); ok { + return &task, true + } + } + return nil, false +} + +func downloadTaskMediaType(task *model.DownloadTask) string { + if task == nil { + return "" + } + return strings.TrimSpace(task.MediaType) +} + +func downloadTaskMediaCategory(task *model.DownloadTask) string { + if task == nil { + return "" + } + return strings.TrimSpace(task.MediaCategory) +} + // DownloadPathMappingsSettingKey 允许用户自定义「下载器路径 → 本程序路径」 // 映射,每行一条,格式 `客户端路径=本地路径`(也接受 `=>` 或单个 `:` 分隔)。 // qBittorrent 与本程序常在不同容器/主机里,对同一份数据看到的路径不同; diff --git a/internal/service/downloads_test.go b/internal/service/downloads_test.go index 1d332a4..0434b4a 100644 --- a/internal/service/downloads_test.go +++ b/internal/service/downloads_test.go @@ -100,6 +100,9 @@ func TestDownloadCompleteAutoOrganizesContentPath(t *testing.T) { } repos := newOrganizerTestRepo(t) + if err := repos.DB.AutoMigrate(&model.DownloadTask{}); err != nil { + t.Fatal(err) + } for key, value := range map[string]string{ "organizer.auto_after_download": "true", "organize.target_dir": dest, @@ -128,6 +131,146 @@ func TestDownloadCompleteAutoOrganizesContentPath(t *testing.T) { } } +func TestDownloadCompleteAutoOrganizeUsesTaskMediaCategory(t *testing.T) { + root := t.TempDir() + src := filepath.Join(root, "downloads", "Motherhood.of.Taihang.S01E01.2026.1080p.mkv") + dest := filepath.Join(root, "media") + writeOrgFile(t, src, "episode") + + repos := newOrganizerTestRepo(t) + if err := repos.DB.AutoMigrate(&model.DownloadTask{}); err != nil { + t.Fatal(err) + } + for key, value := range map[string]string{ + "organizer.auto_after_download": "true", + "organize.target_dir": dest, + "organize.transfer_mode": "copy", + } { + if err := repos.Setting.Set(t.Context(), key, value); err != nil { + t.Fatal(err) + } + } + task := &model.DownloadTask{ + UserID: "u1", + Source: "qbittorrent", + URL: "magnet:?xt=urn:btih:motherhood", + Title: "Motherhood.of.Taihang.S01E01.2026.1080p", + SavePath: filepath.Join(root, "downloads"), + MediaType: "tv", + MediaCategory: "国产剧", + Status: "completed", + Progress: 1, + } + if err := repos.Download.Create(t.Context(), task); err != nil { + t.Fatal(err) + } + + org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + svc := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), org) + svc.onTorrentComplete(t.Context(), QBitTorrent{ + Hash: "done-category", + Name: "Motherhood.of.Taihang.S01E01.2026.1080p", + Progress: 1, + SavePath: filepath.Join(root, "downloads"), + ContentPath: src, + }) + + categoryRoot := filepath.Join(dest, "电视剧", "国产剧") + var organized string + err := filepath.WalkDir(categoryRoot, func(path string, entry os.DirEntry, err error) error { + if err != nil || entry.IsDir() { + return err + } + if filepath.Ext(path) == ".mkv" { + organized = path + } + return nil + }) + if err != nil { + t.Fatalf("walk organized category: %v", err) + } + if organized == "" { + t.Fatalf("expected organized file under %q", categoryRoot) + } + wrongRoot := filepath.Join(dest, "电视剧", "Motherhood Of Taihang") + if _, err := os.Stat(wrongRoot); !os.IsNotExist(err) { + t.Fatalf("unexpected uncategorized organize root %q, err=%v", wrongRoot, err) + } +} + +func TestDownloadCompleteOnlyReplacesExistingWhenTaskAllowsWash(t *testing.T) { + for _, tc := range []struct { + name string + allowWash bool + wantContent string + }{ + {name: "wash disabled keeps existing", allowWash: false, wantContent: "inception-1080p"}, + {name: "wash enabled replaces existing", allowWash: true, wantContent: "inception-2160p"}, + } { + t.Run(tc.name, func(t *testing.T) { + root := t.TempDir() + src := filepath.Join(root, "downloads", "Inception 2010 2160p BluRay.mkv") + dest := filepath.Join(root, "media") + existing := filepath.Join(dest, "电影", "Inception (2010)", "Inception (2010).mkv") + writeOrgFile(t, src, "inception-2160p") + writeOrgFile(t, existing, "inception-1080p") + + repos := newOrganizerTestRepo(t) + if err := repos.DB.AutoMigrate(&model.DownloadTask{}); err != nil { + t.Fatal(err) + } + for key, value := range map[string]string{ + "organizer.auto_after_download": "true", + "organize.target_dir": dest, + "organize.transfer_mode": "copy", + } { + if err := repos.Setting.Set(t.Context(), key, value); err != nil { + t.Fatal(err) + } + } + if err := repos.Media.Upsert(t.Context(), &model.Media{ + Title: "Inception", + Path: existing, + Year: 2010, + Container: "mkv", + Width: 1920, + Height: 1080, + }); err != nil { + t.Fatal(err) + } + if err := repos.Download.Create(t.Context(), &model.DownloadTask{ + Source: "qbittorrent", + URL: "magnet:?xt=urn:btih:inception", + Title: "Inception 2010 2160p BluRay", + SavePath: filepath.Dir(src), + Status: "completed", + Progress: 1, + AllowExistingLibrary: tc.allowWash, + }); err != nil { + t.Fatal(err) + } + + org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + svc := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), org) + svc.onTorrentComplete(t.Context(), QBitTorrent{ + Hash: "inception", + Name: "Inception 2010 2160p BluRay", + Progress: 1, + SavePath: filepath.Dir(src), + ContentPath: src, + }) + + got, err := os.ReadFile(existing) + if err != nil { + t.Fatal(err) + } + if string(got) != tc.wantContent { + t.Fatalf("existing content = %q, want %q", string(got), tc.wantContent) + } + }) + } +} + func TestCompletedTorrentSourceDoesNotFallbackToSavePath(t *testing.T) { root := t.TempDir() savePath := filepath.Join(root, "downloads", "日番") diff --git a/internal/service/media.go b/internal/service/media.go index 1cd181e..a0d48ef 100644 --- a/internal/service/media.go +++ b/internal/service/media.go @@ -69,9 +69,7 @@ func (s *MediaService) CreateLibrary(ctx context.Context, name, path, kind strin if err != nil { return nil, err } - if kind == "" { - kind = "movie" - } + kind = inferLibraryKind(name, abs, kind) lib := &model.Library{Name: name, Path: abs, Type: kind, Enabled: true} if err := s.repo.Library.Create(ctx, lib); err != nil { return nil, err @@ -79,19 +77,41 @@ func (s *MediaService) CreateLibrary(ctx context.Context, name, path, kind strin return lib, nil } +func inferLibraryKind(name, path, requested string) string { + requested = normalizeOrganizeMediaType(requested) + text := strings.ToLower(name + " " + filepath.ToSlash(path)) + switch { + case containsAnyText(text, "成人", "番号", "jav", "9kg", "adult", "nsfw"): + return "adult" + case containsAnyText(text, "综艺", "真人秀", "variety"): + return "variety" + case containsAnyText(text, "国漫", "日漫", "日番", "动漫", "动画", "anime", "bangumi") && !containsAnyText(text, "动画电影"): + return "anime" + case containsAnyText(text, "电视剧", "国产剧", "欧美剧", "日韩剧", "日剧", "韩剧", "剧集", "tv", "series"): + return "tv" + case containsAnyText(text, "电影", "movie", "film"): + return "movie" + } + if requested != "" { + return requested + } + return "movie" +} + func resolveAccessibleLibraryPath(path string) (string, error) { - abs, err := filepath.Abs(strings.TrimSpace(path)) - if err != nil { - return "", fmt.Errorf("invalid path: %w", err) + input := strings.TrimSpace(path) + if input == "" { + return "", errors.New("path required") } - if isAccessibleDir(abs) { - return abs, nil - } - for _, candidate := range dockerVolumePathCandidates(abs) { + for _, candidate := range mappedPathCandidates(input) { if isAccessibleDir(candidate) { return filepath.Clean(candidate), nil } } + abs, err := filepath.Abs(input) + if err != nil { + return "", fmt.Errorf("invalid path: %w", err) + } return "", fmt.Errorf("path is not an accessible directory: %s", abs) } @@ -147,12 +167,15 @@ func mappedPathCandidates(input string) []string { } clean := filepath.Clean(input) add(clean) - if slashClean := cleanPathForVolumeMapping(input); slashClean != "" { - add(slashClean) + for _, candidate := range dockerVolumePathCandidates(input) { + add(candidate) } for _, candidate := range dockerVolumePathCandidates(clean) { add(candidate) } + if slashClean := cleanPathForVolumeMapping(input); slashClean != "" { + add(slashClean) + } if abs, err := filepath.Abs(input); err == nil { add(abs) for _, candidate := range dockerVolumePathCandidates(abs) { @@ -228,6 +251,7 @@ func cleanPathForVolumeMapping(path string) string { return "" } path = strings.ReplaceAll(path, "\\", "/") + path = trimEmbeddedWindowsDrive(path) return filepath.ToSlash(filepath.Clean(filepath.FromSlash(path))) } @@ -238,6 +262,18 @@ func pathAfterWindowsDrivePrefix(path string) string { return path } +func trimEmbeddedWindowsDrive(path string) string { + for i := 0; i+2 < len(path); i++ { + if !isASCIIAlpha(path[i]) || path[i+1] != ':' || path[i+2] != '/' { + continue + } + if i == 0 || path[i-1] == '/' { + return path[i:] + } + } + return path +} + func isASCIIAlpha(ch byte) bool { return (ch >= 'a' && ch <= 'z') || (ch >= 'A' && ch <= 'Z') } diff --git a/internal/service/media_test.go b/internal/service/media_test.go index bebabe2..17b61ee 100644 --- a/internal/service/media_test.go +++ b/internal/service/media_test.go @@ -34,6 +34,53 @@ func TestResolveAccessibleLibraryPathMapsConfiguredHostMediaDir(t *testing.T) { } } +func TestResolveAccessibleLibraryPathMapsWindowsDriveBeforeLinuxAbs(t *testing.T) { + root := t.TempDir() + containerRoot := filepath.Join(root, "container", "media") + containerLibrary := filepath.Join(containerRoot, "电视剧", "国产剧") + if err := os.MkdirAll(containerLibrary, 0o755); err != nil { + t.Fatal(err) + } + t.Setenv("MEDIASTATION_MEDIA_DIR", `Q:\media`) + t.Setenv("MEDIASTATION_MEDIA_CONTAINER_DIR", containerRoot) + + for _, input := range []string{ + `Q:\media\电视剧\国产剧`, + `Q:/media/电视剧/国产剧`, + `/app/Q:\media\电视剧\国产剧`, + `/app/Q:/media/电视剧/国产剧`, + } { + t.Run(input, func(t *testing.T) { + got, err := resolveAccessibleLibraryPath(input) + if err != nil { + t.Fatalf("resolveAccessibleLibraryPath() error = %v", err) + } + if got != filepath.Clean(containerLibrary) { + t.Fatalf("resolveAccessibleLibraryPath() = %q, want %q", got, filepath.Clean(containerLibrary)) + } + }) + } +} + +func TestResolveAccessibleLibraryPathRecoversDockerPollutedWindowsDrive(t *testing.T) { + root := t.TempDir() + containerRoot := filepath.Join(root, "container", "media") + containerLibrary := filepath.Join(containerRoot, "电视剧", "国产剧") + if err := os.MkdirAll(containerLibrary, 0o755); err != nil { + t.Fatal(err) + } + t.Setenv("MEDIASTATION_MEDIA_DIR", `F:\media`) + t.Setenv("MEDIASTATION_MEDIA_CONTAINER_DIR", containerRoot) + + got, err := resolveAccessibleLibraryPath(`/app/F:\media\电视剧\国产剧`) + if err != nil { + t.Fatalf("resolveAccessibleLibraryPath() error = %v", err) + } + if got != filepath.Clean(containerLibrary) { + t.Fatalf("resolveAccessibleLibraryPath() = %q, want %q", got, filepath.Clean(containerLibrary)) + } +} + func TestResolveAccessibleLibraryPathKeepsAccessibleContainerPath(t *testing.T) { containerLibrary := filepath.Join(t.TempDir(), "media", "电影") if err := os.MkdirAll(containerLibrary, 0o755); err != nil { @@ -49,6 +96,26 @@ func TestResolveAccessibleLibraryPathKeepsAccessibleContainerPath(t *testing.T) } } +func TestInferLibraryKindFromCategoryPathOverridesMovieDefault(t *testing.T) { + for _, tc := range []struct { + name string + path string + want string + }{ + {name: "国产剧", path: `/media/电视剧/国产剧`, want: "tv"}, + {name: "日漫", path: `/media/电视剧/日漫`, want: "anime"}, + {name: "综艺", path: `/media/电视剧/综艺`, want: "variety"}, + {name: "成人", path: `/media/成人`, want: "adult"}, + {name: "动画电影", path: `/media/电影/动画电影`, want: "movie"}, + } { + t.Run(tc.name, func(t *testing.T) { + if got := inferLibraryKind(tc.name, tc.path, "movie"); got != tc.want { + t.Fatalf("inferLibraryKind() = %q, want %q", got, tc.want) + } + }) + } +} + func TestMappedPathCandidatesMapWindowsDriveDownloadMarker(t *testing.T) { root := t.TempDir() containerDownloads := filepath.Join(root, "container", "downloads") @@ -64,6 +131,51 @@ func TestMappedPathCandidatesMapWindowsDriveDownloadMarker(t *testing.T) { t.Fatalf("mappedPathCandidates() missing %q", want) } +func TestResolveMappedDestinationPathPrefersConfiguredContainerMapping(t *testing.T) { + root := t.TempDir() + containerMedia := filepath.Join(root, "container", "media") + if err := os.MkdirAll(containerMedia, 0o755); err != nil { + t.Fatal(err) + } + t.Setenv("MEDIASTATION_MEDIA_DIR", `Q:\media`) + t.Setenv("MEDIASTATION_MEDIA_CONTAINER_DIR", containerMedia) + + for _, input := range []string{`Q:\media`, `Q:/media`, `/app/Q:\media`} { + t.Run(input, func(t *testing.T) { + got := resolveMappedDestinationPath(input) + if got != filepath.Clean(containerMedia) { + t.Fatalf("resolveMappedDestinationPath() = %q, want %q", got, filepath.Clean(containerMedia)) + } + }) + } +} + +func TestResolveAccessibleMappedPathMapsWindowsDownloadVariants(t *testing.T) { + root := t.TempDir() + containerDownloads := filepath.Join(root, "container", "downloads") + containerSource := filepath.Join(containerDownloads, "国产剧") + if err := os.MkdirAll(containerSource, 0o755); err != nil { + t.Fatal(err) + } + t.Setenv("MEDIASTATION_DOWNLOAD_DIR", `Q:\downloads`) + t.Setenv("MEDIASTATION_DOWNLOAD_CONTAINER_DIR", containerDownloads) + + for _, input := range []string{`Q:\downloads\国产剧`, `Q:/downloads/国产剧`, `/app/Q:\downloads\国产剧`} { + t.Run(input, func(t *testing.T) { + got, info, err := resolveAccessibleMappedPath(input) + if err != nil { + t.Fatalf("resolveAccessibleMappedPath() error = %v", err) + } + if !info.IsDir() { + t.Fatalf("resolved path is not dir") + } + if got != filepath.Clean(containerSource) { + t.Fatalf("resolveAccessibleMappedPath() = %q, want %q", got, filepath.Clean(containerSource)) + } + }) + } +} + func TestDeleteCloudLibraryPurgesMountWithoutRecycleBin(t *testing.T) { db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) if err != nil { diff --git a/internal/service/organizer.go b/internal/service/organizer.go index 9a43ea0..4a8f4ac 100644 --- a/internal/service/organizer.go +++ b/internal/service/organizer.go @@ -87,8 +87,13 @@ type OrganizeOptions struct { TransferMode TransferMode // MediaType 手动整理时由 UI 指定的媒体类型。空值时按文件名/目录推断。 MediaType string + // MediaCategory 由订阅/下载任务或 UI 指定的分类。空值时按目录/NFO/规则推断。 + MediaCategory string // DryRun 仅生成整理预览,不实际移动/复制/硬链接文件。 DryRun bool + // AllowReplaceExisting 允许用本次来源替换目标库中已存在的同一媒体。 + // 默认 false:只去重不洗版,避免未开启洗版的订阅/手动整理留下或替换出多份版本。 + AllowReplaceExisting bool } // OrganizeMedia moves a single media file into the target library directory. diff --git a/internal/service/organizer_directory.go b/internal/service/organizer_directory.go index 2ba6c5c..80e553a 100644 --- a/internal/service/organizer_directory.go +++ b/internal/service/organizer_directory.go @@ -127,7 +127,7 @@ func (o *OrganizerService) OrganizeDirectory(ctx context.Context, opts OrganizeO if _, ok := videoExtensions[ext]; !ok { return nil, fmt.Errorf("source is not a supported video file: %s", source) } - if err := o.organizeSourceFile(ctx, source, filepath.Dir(source), dest, mode, opts.MediaType, opts.DryRun, res); err != nil { + if err := o.organizeSourceFile(ctx, source, filepath.Dir(source), dest, mode, opts.MediaType, opts.MediaCategory, opts.DryRun, opts.AllowReplaceExisting, res); err != nil { res.Errors = append(res.Errors, fmt.Sprintf("%s: %s", filepath.Base(source), err.Error())) res.Items = append(res.Items, OrganizePreviewItem{Source: source, Action: "error", Reason: err.Error()}) } @@ -149,7 +149,7 @@ func (o *OrganizerService) OrganizeDirectory(ctx context.Context, opts OrganizeO if _, ok := videoExtensions[ext]; !ok { return nil } - if err := o.organizeSourceFile(ctx, path, source, dest, mode, opts.MediaType, opts.DryRun, res); err != nil { + if err := o.organizeSourceFile(ctx, path, source, dest, mode, opts.MediaType, opts.MediaCategory, opts.DryRun, opts.AllowReplaceExisting, res); err != nil { res.Errors = append(res.Errors, fmt.Sprintf("%s: %s", filepath.Base(path), err.Error())) res.Items = append(res.Items, OrganizePreviewItem{Source: path, Action: "error", Reason: err.Error()}) } @@ -176,7 +176,7 @@ type organizeDirectoryLayout struct { // organizeSourceFile organizes a single video file from the source directory // into destRoot, applying dedup + 洗版. -func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRoot, destRoot string, mode TransferMode, mediaTypeOverride string, dryRun bool, res *OrganizeResult) error { +func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRoot, destRoot string, mode TransferMode, mediaTypeOverride, mediaCategoryOverride string, dryRun bool, allowReplaceExisting bool, res *OrganizeResult) error { ext := filepath.Ext(src) title, year := CleanQuery(src) if title == "" { @@ -200,14 +200,16 @@ func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRo if layout.MediaType == "" { layout.MediaType = o.inferMediaTypeForSourceFile(src, title, season, episode) } - if layout.Category == "" { + if category := strings.TrimSpace(mediaCategoryOverride); category != "" { + layout.Category = sanitizeFilename(category) + } else if layout.Category == "" { layout.Category = o.smartClassifySourceFile(ctx, src, sourceRoot, layout.MediaType, title, parsedTitle) } - layoutRoot := destRoot - if layout.MediaType != "" { + layoutRoot, matchedLibrary := o.organizeLibraryRootForLayout(ctx, destRoot, layout.MediaType, layout.Category) + if !matchedLibrary && layout.MediaType != "" { layoutRoot = o.organizeRoot(destRoot, layout.MediaType, layout.Category) } - if layout.Category != "" { + if !matchedLibrary && layout.Category != "" { layoutRoot = categoryRoot(layoutRoot, sanitizeFilename(layout.Category)) } @@ -254,7 +256,7 @@ func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRo } // 洗版:仅当来源与已存在版本的分辨率都可判定、且来源更高时才替换; // 任一方分辨率未知时保守跳过,绝不删除无法判定的已存在文件。 - if srcArea > 0 && bestArea > 0 && srcArea > bestArea { + if allowReplaceExisting && srcArea > 0 && bestArea > 0 && srcArea > bestArea { res.Items = append(res.Items, OrganizePreviewItem{ Source: src, Target: dst, Action: "replace", Reason: "higher resolution", MediaType: layout.MediaType, Category: layout.Category, Title: title, @@ -451,6 +453,147 @@ func (o *OrganizerService) directoryCategoryTypes() map[string]organizeDirectory return out } +func (o *OrganizerService) organizeLibraryRootForLayout(ctx context.Context, destRoot, mediaType, category string) (string, bool) { + if o == nil || o.repo == nil || o.repo.Library == nil { + return "", false + } + libraries, err := o.repo.Library.List(ctx) + if err != nil { + if o.log != nil { + o.log.Debug("list libraries for organize target failed", zap.Error(err)) + } + return "", false + } + destRoot = filepath.Clean(strings.TrimSpace(destRoot)) + mediaType = normalizeOrganizeMediaType(mediaType) + aliases := o.organizeCategoryAliases(mediaType, category) + + bestPath := "" + bestScore := -1 + bestDepth := -1 + for _, lib := range libraries { + if !lib.Enabled || strings.TrimSpace(lib.Path) == "" { + continue + } + if _, ok := ParseCloudLibraryMount(lib.Path); ok { + continue + } + if destRoot != "" && destRoot != "." && !pathWithin(lib.Path, destRoot) && !pathWithin(destRoot, lib.Path) { + continue + } + categoryMatch := len(aliases) > 0 && libraryMatchesOrganizeCategory(lib, aliases) + typeScore := organizeLibraryTypeScore(mediaType, lib.Type) + if len(aliases) > 0 { + if !categoryMatch { + continue + } + } else if typeScore <= 0 { + continue + } + score := typeScore + if categoryMatch { + score += 20 + } + depth := pathDepth(lib.Path) + if score > bestScore || (score == bestScore && depth > bestDepth) { + bestScore = score + bestDepth = depth + bestPath = lib.Path + } + } + if bestPath == "" { + return "", false + } + return filepath.Clean(bestPath), true +} + +func (o *OrganizerService) organizeCategoryAliases(mediaType, category string) map[string]struct{} { + aliases := map[string]struct{}{} + add := func(values ...string) { + for _, value := range values { + key := normalizeOrganizeCategoryKey(value) + if key != "" { + aliases[key] = struct{}{} + } + } + } + categories := o.categoryMap() + add(category) + switch normalizeOrganizeCategoryKey(category) { + case normalizeOrganizeCategoryKey(categoryName(categories, "jp_anime", "日番")), "日番", "日漫", "日本动漫", "日本動畫", "日本动画": + add("日番", "日漫", "日本动漫", "日本动画") + case normalizeOrganizeCategoryKey(categoryName(categories, "cn_anime", "国漫")), "国漫", "国产动漫", "國漫": + add("国漫", "国产动漫") + case normalizeOrganizeCategoryKey(categoryName(categories, "domestic_tv", "国产剧")), "国产剧", "国剧", "大陆剧", "国产电视剧": + add("国产剧", "国剧", "大陆剧", "国产电视剧") + case normalizeOrganizeCategoryKey(categoryName(categories, "euus_tv", "欧美剧")), "欧美剧", "欧美电视剧": + add("欧美剧", "欧美电视剧") + case normalizeOrganizeCategoryKey(categoryName(categories, "jk_tv", "日韩剧")), "日韩剧", "日剧", "韩剧": + add("日韩剧", "日剧", "韩剧") + case normalizeOrganizeCategoryKey(categoryName(categories, "variety", "综艺")), "综艺", "真人秀": + add("综艺", "真人秀") + case normalizeOrganizeCategoryKey(categoryName(categories, "documentary", "纪录片")), "纪录片", "纪录": + add("纪录片", "纪录") + case normalizeOrganizeCategoryKey(categoryName(categories, "children", "儿童")), "儿童", "少儿": + add("儿童", "少儿") + case normalizeOrganizeCategoryKey(categoryName(categories, "chinese_movie", "华语电影")), "华语电影", "国产电影", "大陆电影": + add("华语电影", "国产电影", "大陆电影") + case normalizeOrganizeCategoryKey(categoryName(categories, "foreign_movie", "外语电影")), "外语电影": + add("外语电影") + case normalizeOrganizeCategoryKey(categoryName(categories, "animation_movie", "动画电影")), "动画电影", "动漫电影": + add("动画电影", "动漫电影") + case normalizeOrganizeCategoryKey(categoryName(categories, "adult", "成人")), "成人": + add("成人") + case normalizeOrganizeCategoryKey(categoryName(categories, "adult_9kg", "9KG")), "9kg": + add("9KG") + case normalizeOrganizeCategoryKey(categoryName(categories, "adult_jav", "番号")), "番号", "jav": + add("番号", "JAV") + } + return aliases +} + +func libraryMatchesOrganizeCategory(lib model.Library, aliases map[string]struct{}) bool { + for _, value := range []string{lib.Name, filepath.Base(filepath.Clean(lib.Path))} { + if _, ok := aliases[normalizeOrganizeCategoryKey(value)]; ok { + return true + } + } + return false +} + +func organizeLibraryTypeScore(mediaType, libraryType string) int { + libraryType = normalizeOrganizeMediaType(libraryType) + if mediaType == "" || libraryType == "" { + return 1 + } + if mediaType == libraryType { + return 8 + } + if mediaType == "anime" && libraryType == "tv" { + return 5 + } + if mediaType == "variety" && libraryType == "tv" { + return 5 + } + return 0 +} + +func normalizeOrganizeCategoryKey(value string) string { + value = strings.ToLower(strings.TrimSpace(value)) + value = strings.ReplaceAll(value, " ", "") + value = strings.ReplaceAll(value, "_", "") + value = strings.ReplaceAll(value, "-", "") + return value +} + +func pathDepth(path string) int { + path = filepath.Clean(path) + if path == "." || path == string(os.PathSeparator) { + return 0 + } + return len(strings.Split(path, string(os.PathSeparator))) +} + // existingVersionPaths returns existing destination files that represent the // same media, combining two strategies and de-duplicating by path: // diff --git a/internal/service/organizer_directory_test.go b/internal/service/organizer_directory_test.go index 48ced49..5b7cc4b 100644 --- a/internal/service/organizer_directory_test.go +++ b/internal/service/organizer_directory_test.go @@ -313,9 +313,10 @@ func TestOrganizeDirectoryDedup(t *testing.T) { } } -// TestOrganizeDirectoryReplaceHigherResolution verifies 洗版: a higher-resolution -// source replaces the lower-resolution version already in the destination. -func TestOrganizeDirectoryReplaceHigherResolution(t *testing.T) { +// TestOrganizeDirectorySkipsHigherResolutionWhenReplacementDisabled verifies +// that dedup wins by default: even a higher-resolution source must not replace +// an existing library item unless the caller explicitly allows washing. +func TestOrganizeDirectorySkipsHigherResolutionWhenReplacementDisabled(t *testing.T) { root := t.TempDir() src := filepath.Join(root, "downloads") dest := filepath.Join(root, "media") @@ -341,6 +342,50 @@ func TestOrganizeDirectoryReplaceHigherResolution(t *testing.T) { if err != nil { t.Fatalf("organize directory: %v", err) } + if res.Skipped != 1 || res.Replaced != 0 || res.Organized != 0 { + t.Fatalf("expected skipped=1 replaced=0 organized=0 (wash disabled), got %+v", res) + } + got, err := os.ReadFile(existing) + if err != nil || string(got) != "inception-1080p" { + t.Fatalf("destination must keep existing version when wash disabled, got %q err=%v", string(got), err) + } + var count int64 + if err := repos.DB.Model(&model.Media{}).Where("path = ?", existing).Count(&count).Error; err != nil { + t.Fatal(err) + } + if count != 1 { + t.Fatalf("expected existing media DB row kept, found %d", count) + } +} + +// TestOrganizeDirectoryReplaceHigherResolutionWhenAllowed verifies 洗版: a +// higher-resolution source replaces the lower-resolution version already in the +// destination only when the caller explicitly allows replacement. +func TestOrganizeDirectoryReplaceHigherResolutionWhenAllowed(t *testing.T) { + root := t.TempDir() + src := filepath.Join(root, "downloads") + dest := filepath.Join(root, "media") + + writeOrgFile(t, filepath.Join(src, "Inception 2010 2160p BluRay.mkv"), "inception-uhd") + existing := filepath.Join(dest, "电影", "Inception (2010)", "Inception (2010).mkv") + writeOrgFile(t, existing, "inception-1080p") + + repos := newOrganizerTestRepo(t) + row := model.Media{Title: "Inception", Path: existing, Year: 2010, Container: "mkv", Width: 1920, Height: 1080} + if err := repos.Media.Upsert(t.Context(), &row); err != nil { + t.Fatal(err) + } + + org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + res, err := org.OrganizeDirectory(t.Context(), OrganizeOptions{ + SourcePath: src, + DestPath: dest, + TransferMode: TransferCopy, + AllowReplaceExisting: true, + }) + if err != nil { + t.Fatalf("organize directory: %v", err) + } if res.Replaced != 1 || res.Organized != 0 || res.Skipped != 0 { t.Fatalf("expected replaced=1 organized=0 skipped=0 (洗版), got %+v", res) } @@ -465,6 +510,45 @@ func TestOrganizeDirectoryUsesDownloadCategoryLayout(t *testing.T) { } } +func TestOrganizeDirectoryUsesExplicitCategoryLibraryRoot(t *testing.T) { + root := t.TempDir() + src := filepath.Join(root, "downloads", "Motherhood.of.Taihang.S01E01.2026.1080p.mkv") + dest := filepath.Join(root, "media") + writeOrgFile(t, src, "episode") + + repos := newOrganizerTestRepo(t) + libraryRoot := filepath.Join(dest, "电视剧", "国产剧") + wrongType := model.Library{Name: "国产剧", Path: libraryRoot, Type: "movie", Enabled: true} + rightType := model.Library{Name: "国产剧", Path: libraryRoot, Type: "tv", Enabled: true} + if err := repos.Library.Create(t.Context(), &wrongType); err != nil { + t.Fatal(err) + } + if err := repos.Library.Create(t.Context(), &rightType); err != nil { + t.Fatal(err) + } + + org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + res, err := org.OrganizeDirectory(t.Context(), OrganizeOptions{ + SourcePath: src, + DestPath: dest, + MediaType: "tv", + MediaCategory: "国产剧", + TransferMode: TransferCopy, + }) + if err != nil { + t.Fatalf("organize explicit category: %v", err) + } + if res.Organized != 1 || len(res.Items) != 1 { + t.Fatalf("result = %+v, want one organized item", res) + } + if !pathWithin(res.Items[0].Target, libraryRoot) { + t.Fatalf("target = %q, want under %q", res.Items[0].Target, libraryRoot) + } + if pathWithin(res.Items[0].Target, filepath.Join(dest, "电视剧")) && !pathWithin(res.Items[0].Target, libraryRoot) { + t.Fatalf("target landed outside category library: %q", res.Items[0].Target) + } +} + func TestOrganizeDirectorySmartClassifiesUncategorizedSources(t *testing.T) { root := t.TempDir() src := filepath.Join(root, "downloads") @@ -590,3 +674,20 @@ func TestOrganizeDirectoryScanAfterRecursesNestedDownloadFolders(t *testing.T) { t.Fatalf("media rows = %d, want 2", count) } } + +func TestSelectOrganizeScanTargetsDedupesSamePathByPathType(t *testing.T) { + root := t.TempDir() + path := filepath.Join(root, "media", "电视剧", "国产剧") + libraries := []model.Library{ + {Name: "国产剧", Path: path, Type: "movie", Enabled: true}, + {Name: "国产剧", Path: path, Type: "tv", Enabled: true}, + } + + targets := selectOrganizeScanTargets(libraries, filepath.Join(root, "media"), "") + if len(targets) != 1 { + t.Fatalf("targets = %#v, want one deduped target", targets) + } + if targets[0].Type != "tv" { + t.Fatalf("target type = %q, want tv", targets[0].Type) + } +} diff --git a/internal/service/organizer_scan.go b/internal/service/organizer_scan.go index 17f3bdb..c547b09 100644 --- a/internal/service/organizer_scan.go +++ b/internal/service/organizer_scan.go @@ -2,6 +2,7 @@ package service import ( "context" + "path/filepath" "strings" "github.com/ShukeBta/MediaStationGo/internal/model" @@ -175,7 +176,40 @@ func selectOrganizeScanTargets(libraries []model.Library, destRoot, preferredLib } } if len(matched) > 0 { - return matched + return dedupeOrganizeScanTargets(matched) } - return enabled + return dedupeOrganizeScanTargets(enabled) +} + +func dedupeOrganizeScanTargets(libraries []model.Library) []model.Library { + out := make([]model.Library, 0, len(libraries)) + byPath := map[string]int{} + for _, lib := range libraries { + key := strings.ToLower(filepath.Clean(lib.Path)) + if existingIndex, ok := byPath[key]; ok { + if organizeScanTargetScore(lib) > organizeScanTargetScore(out[existingIndex]) { + out[existingIndex] = lib + } + continue + } + byPath[key] = len(out) + out = append(out, lib) + } + return out +} + +func organizeScanTargetScore(lib model.Library) int { + libraryType := normalizeOrganizeMediaType(lib.Type) + inferred := normalizeMediaType("", lib.Name, lib.Path) + score := 0 + if libraryType != "" && libraryType == inferred { + score += 10 + } + if librarySupportsSeasons(&lib) { + score += 2 + } + if libraryType != "" && libraryType != "movie" { + score++ + } + return score } diff --git a/internal/service/qbittorrent.go b/internal/service/qbittorrent.go index c397316..86df98b 100644 --- a/internal/service/qbittorrent.go +++ b/internal/service/qbittorrent.go @@ -54,6 +54,7 @@ type QBitTorrent struct { NumLeech int `json:"num_leechs"` Size int64 `json:"size"` SavePath string `json:"save_path"` + Category string `json:"category"` // ContentPath is qBittorrent's resolved payload path. For single-file // torrents it points at the file; for multi-file torrents it points at the // root folder. Prefer it for automatic organize so we do not scan the whole diff --git a/internal/service/scanner.go b/internal/service/scanner.go index a574dea..2f7bf10 100644 --- a/internal/service/scanner.go +++ b/internal/service/scanner.go @@ -748,6 +748,9 @@ func (s *ScannerService) scanLibrary(ctx context.Context, libraryID string, auto } return res, err } + if err := s.resolveLocalLibraryPath(ctx, lib); err != nil { + return &ScanResult{LibraryID: lib.ID}, err + } res := &ScanResult{LibraryID: lib.ID} seen := make(map[string]struct{}) seenInodes := make(map[string]string) @@ -805,6 +808,9 @@ func (s *ScannerService) IngestPath(ctx context.Context, libraryID, path string) if err != nil || lib == nil { return false, err } + if err := s.resolveLocalLibraryPath(ctx, lib); err != nil { + return false, err + } fi, err := os.Stat(path) if err != nil || fi.IsDir() { return false, err @@ -818,6 +824,37 @@ func (s *ScannerService) IngestPath(ctx context.Context, libraryID, path string) return res.Added+res.Updated > 0, nil } +func (s *ScannerService) resolveLocalLibraryPath(ctx context.Context, lib *model.Library) error { + if lib == nil || strings.TrimSpace(lib.Path) == "" { + return nil + } + resolved, err := resolveAccessibleLibraryPath(lib.Path) + if err != nil { + return err + } + if sameLibraryPath(resolved, lib.Path) { + lib.Path = filepath.Clean(lib.Path) + return nil + } + if s.repo != nil && s.repo.DB != nil { + if updateErr := s.repo.DB.WithContext(ctx).Model(&model.Library{}).Where("id = ?", lib.ID).Update("path", resolved).Error; updateErr != nil && s.log != nil { + s.log.Warn("update mapped library path failed", + zap.String("library_id", lib.ID), + zap.String("from", lib.Path), + zap.String("to", resolved), + zap.Error(updateErr)) + } + } + if s.log != nil { + s.log.Info("mapped library path for scan", + zap.String("library_id", lib.ID), + zap.String("from", lib.Path), + zap.String("to", resolved)) + } + lib.Path = resolved + return nil +} + func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Library, mount CloudMountInfo, autoScrape bool) (*ScanResult, error) { res := &ScanResult{LibraryID: lib.ID} if s.storage == nil { diff --git a/internal/service/scanner_incremental_test.go b/internal/service/scanner_incremental_test.go index 4de3d65..12d7a53 100644 --- a/internal/service/scanner_incremental_test.go +++ b/internal/service/scanner_incremental_test.go @@ -95,6 +95,41 @@ func TestScanLibraryReadsLocalSTRMTarget(t *testing.T) { } } +func TestScanLibraryMapsPersistedHostLibraryPath(t *testing.T) { + sc, repos := newScannerTestEnv(t) + root := t.TempDir() + containerRoot := filepath.Join(root, "container", "media") + containerLibrary := filepath.Join(containerRoot, "电视剧", "国产剧") + if err := os.MkdirAll(containerLibrary, 0o755); err != nil { + t.Fatal(err) + } + file := filepath.Join(containerLibrary, "狂飙.S01E01.2023.mkv") + if err := os.WriteFile(file, []byte("episode"), 0o644); err != nil { + t.Fatal(err) + } + t.Setenv("MEDIASTATION_MEDIA_DIR", `Q:\media`) + t.Setenv("MEDIASTATION_MEDIA_CONTAINER_DIR", containerRoot) + + lib := model.Library{Name: "国产剧", Path: `Q:\media\电视剧\国产剧`, Type: "tv", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + res, err := sc.ScanLibrary(t.Context(), lib.ID) + if err != nil { + t.Fatalf("scan mapped host path: %v", err) + } + if res.Added != 1 { + t.Fatalf("scan result = %#v, want added=1", res) + } + var stored model.Library + if err := repos.DB.First(&stored, "id = ?", lib.ID).Error; err != nil { + t.Fatal(err) + } + if stored.Path != filepath.Clean(containerLibrary) { + t.Fatalf("stored path = %q, want %q", stored.Path, filepath.Clean(containerLibrary)) + } +} + func TestRemovePathDeletesVanishedMedia(t *testing.T) { sc, repos := newScannerTestEnv(t) root := t.TempDir() diff --git a/internal/service/scheduler.go b/internal/service/scheduler.go index 4a90913..0a03215 100644 --- a/internal/service/scheduler.go +++ b/internal/service/scheduler.go @@ -46,6 +46,7 @@ type SchedulerService struct { organizer *OrganizerService storageCfg *StorageConfigService hub *Hub + tasks *TaskTrackerService cacheDir string now func() time.Time @@ -54,6 +55,10 @@ type SchedulerService struct { jobs []*scheduledJob } +func (s *SchedulerService) SetTaskTracker(tasks *TaskTrackerService) { + s.tasks = tasks +} + // scheduledJob is one recurring task. type scheduledJob struct { name string @@ -479,11 +484,34 @@ func (s *SchedulerService) jobOrganizeSource(ctx context.Context) error { if s.organizer == nil || (!manual && !s.autoOrganizeSourceEnabled(ctx)) { return nil } + task := s.startScheduledOrganizeTask(ctx, manual) res, err := s.organizer.OrganizeDirectory(ctx, OrganizeOptions{}) if err != nil { + if task != nil { + task.Finish(err, TaskUpdate{ + Stage: "organize", + Message: "自动整理入库失败", + }) + } return err } + if task != nil && res != nil { + task.Update(TaskUpdate{ + Stage: "organize", + SourcePath: res.SourcePath, + DestPath: res.DestPath, + Message: "自动整理完成,准备扫描入库", + Metrics: OrganizeTaskMetrics(res), + }) + } if s.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" && OrganizeResultHasChanges(res) { + if task != nil { + task.Update(TaskUpdate{ + Stage: "scan_scrape", + Message: "正在扫描入库并按设置刮削", + Metrics: OrganizeTaskMetrics(res), + }) + } res.Scans, res.Scrapes = s.scanner.ScanAndScrapeLibrariesForPath(ctx, res.DestPath, "", OrganizeScrapeAfterEnabled(ctx, s.repo)) } else if s.log != nil && res != nil && !OrganizeResultHasChanges(res) { s.log.Info("scheduled source organize skipped scan; no destination changes", @@ -494,6 +522,13 @@ func (s *SchedulerService) jobOrganizeSource(ctx context.Context) error { zap.Int("skipped", res.Skipped), ) } + if task != nil { + task.Finish(nil, TaskUpdate{ + Stage: "completed", + Message: "自动整理入库结束", + Metrics: OrganizeTaskMetrics(res), + }) + } if s.log != nil && res != nil { s.log.Info("scheduled source organize finished", zap.String("source", res.SourcePath), @@ -508,6 +543,24 @@ func (s *SchedulerService) jobOrganizeSource(ctx context.Context) error { return nil } +func (s *SchedulerService) startScheduledOrganizeTask(ctx context.Context, manual bool) *TaskHandle { + if s == nil || s.tasks == nil { + return nil + } + name := "自动整理重命名入库" + message := "正在执行计划自动整理/重命名/入库" + if manual { + name = "手动触发自动整理重命名入库" + message = "正在执行手动触发的自动整理/重命名/入库" + } + return s.tasks.Start(TaskKindOrganize, name, TaskUpdate{ + Stage: "organize", + SourcePath: s.organizer.defaultSourceRoot(ctx, ""), + DestPath: s.organizer.defaultDestRoot(ctx, ""), + Message: message, + }) +} + func (s *SchedulerService) autoOrganizeSourceEnabled(ctx context.Context) bool { if s.repo == nil || s.repo.Setting == nil { return false diff --git a/internal/service/service.go b/internal/service/service.go index 916967d..5d78259 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -23,6 +23,7 @@ type Container struct { Repo *repository.Container WSHub *Hub SSEHub *SSEHub + Tasks *TaskTrackerService Auth *AuthService Media *MediaService Scan *ScannerService @@ -83,6 +84,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont hub := NewHub(log) go hub.Run() + tasks := NewTaskTrackerService(log, hub) // 初始化 SSE Hub sseHub := NewSSEHub(log) @@ -123,6 +125,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont assistant := NewAssistantService(log, repos, ai) douban := NewDoubanProvider(cfg, log) scheduler := NewSchedulerService(log, repos, scanner, transcoder, organizer, storageCfg, hub, cfg.Cache.CacheDir) + scheduler.SetTaskTracker(tasks) // 初始化认证相关服务 tokenSvc := NewTokenService(cfg, log, repos) @@ -145,6 +148,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont siteSvc := NewSiteService(log, repos, flareSolverrURL) downloads := NewDownloadService(log, repos, hub, organizer, siteSvc) downloads.SetScanner(scanner) + downloads.SetTaskTracker(tasks) subscription := NewSubscriptionService(cfg, log, repos, downloads, siteSvc, hub) // 让图片代理把媒体库根目录视为可读的本地图片位置:海报/封面等 @@ -174,6 +178,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont Repo: repos, WSHub: hub, SSEHub: sseHub, + Tasks: tasks, Auth: authSvc, Media: NewMediaService(cfg, log, repos), Scan: scanner, diff --git a/internal/service/task_tracker.go b/internal/service/task_tracker.go new file mode 100644 index 0000000..d77ba79 --- /dev/null +++ b/internal/service/task_tracker.go @@ -0,0 +1,275 @@ +package service + +import ( + "sync" + "time" + + "github.com/google/uuid" + "go.uber.org/zap" +) + +const ( + TaskStatusRunning = "running" + TaskStatusCompleted = "completed" + TaskStatusFailed = "failed" + + TaskKindOrganize = "organize" + TaskKindScan = "scan" + TaskKindScrape = "scrape" +) + +// BackgroundTask is the compact, operator-facing shape shown on the live tasks +// page. It tracks long-running work that is not represented by a download or +// transcode job, such as organize → scan → scrape ingest flows. +type BackgroundTask struct { + ID string `json:"id"` + Kind string `json:"kind"` + Name string `json:"name"` + Status string `json:"status"` + Stage string `json:"stage,omitempty"` + SourcePath string `json:"source_path,omitempty"` + DestPath string `json:"dest_path,omitempty"` + Message string `json:"message,omitempty"` + Error string `json:"error,omitempty"` + Metrics map[string]int64 `json:"metrics,omitempty"` + StartedAt time.Time `json:"started_at"` + UpdatedAt time.Time `json:"updated_at"` + FinishedAt *time.Time `json:"finished_at,omitempty"` +} + +type TaskUpdate struct { + Stage string + SourcePath string + DestPath string + Message string + Metrics map[string]int64 +} + +type TaskSnapshot struct { + Active []BackgroundTask `json:"active"` + Recent []BackgroundTask `json:"recent"` +} + +type TaskTrackerService struct { + log *zap.Logger + hub *Hub + + mu sync.Mutex + active map[string]*BackgroundTask + recent []BackgroundTask + maxRecent int + now func() time.Time +} + +type TaskHandle struct { + tracker *TaskTrackerService + id string +} + +func NewTaskTrackerService(log *zap.Logger, hub *Hub) *TaskTrackerService { + return &TaskTrackerService{ + log: log, + hub: hub, + active: make(map[string]*BackgroundTask), + maxRecent: 30, + now: time.Now, + } +} + +func (t *TaskTrackerService) Start(kind, name string, update TaskUpdate) *TaskHandle { + if t == nil { + return nil + } + now := t.currentTime() + task := &BackgroundTask{ + ID: uuid.NewString(), + Kind: kind, + Name: name, + Status: TaskStatusRunning, + Stage: update.Stage, + SourcePath: update.SourcePath, + DestPath: update.DestPath, + Message: update.Message, + Metrics: cloneTaskMetrics(update.Metrics), + StartedAt: now, + UpdatedAt: now, + } + t.mu.Lock() + t.active[task.ID] = task + snapshot := cloneBackgroundTask(*task) + t.mu.Unlock() + t.publish(snapshot) + return &TaskHandle{tracker: t, id: task.ID} +} + +func (h *TaskHandle) Update(update TaskUpdate) { + if h == nil || h.tracker == nil { + return + } + h.tracker.update(h.id, update) +} + +func (h *TaskHandle) Finish(err error, update TaskUpdate) { + if h == nil || h.tracker == nil { + return + } + h.tracker.finish(h.id, err, update) +} + +func (t *TaskTrackerService) Snapshot() TaskSnapshot { + if t == nil { + return TaskSnapshot{} + } + t.mu.Lock() + defer t.mu.Unlock() + active := make([]BackgroundTask, 0, len(t.active)) + for _, task := range t.active { + active = append(active, cloneBackgroundTask(*task)) + } + recent := make([]BackgroundTask, 0, len(t.recent)) + for _, task := range t.recent { + recent = append(recent, cloneBackgroundTask(task)) + } + return TaskSnapshot{Active: active, Recent: recent} +} + +func (t *TaskTrackerService) update(id string, update TaskUpdate) { + now := t.currentTime() + t.mu.Lock() + task, ok := t.active[id] + if !ok { + t.mu.Unlock() + return + } + applyTaskUpdate(task, update) + task.UpdatedAt = now + snapshot := cloneBackgroundTask(*task) + t.mu.Unlock() + t.publish(snapshot) +} + +func (t *TaskTrackerService) finish(id string, err error, update TaskUpdate) { + now := t.currentTime() + t.mu.Lock() + task, ok := t.active[id] + if !ok { + t.mu.Unlock() + return + } + applyTaskUpdate(task, update) + task.UpdatedAt = now + task.FinishedAt = &now + if err != nil { + task.Status = TaskStatusFailed + task.Error = err.Error() + } else { + task.Status = TaskStatusCompleted + } + delete(t.active, id) + snapshot := cloneBackgroundTask(*task) + t.recent = append([]BackgroundTask{snapshot}, t.recent...) + if t.maxRecent <= 0 { + t.maxRecent = 30 + } + if len(t.recent) > t.maxRecent { + t.recent = t.recent[:t.maxRecent] + } + t.mu.Unlock() + t.publish(snapshot) +} + +func (t *TaskTrackerService) currentTime() time.Time { + if t != nil && t.now != nil { + return t.now() + } + return time.Now() +} + +func (t *TaskTrackerService) publish(task BackgroundTask) { + if t == nil || t.hub == nil { + return + } + t.hub.Publish("task", task) +} + +func applyTaskUpdate(task *BackgroundTask, update TaskUpdate) { + if update.Stage != "" { + task.Stage = update.Stage + } + if update.SourcePath != "" { + task.SourcePath = update.SourcePath + } + if update.DestPath != "" { + task.DestPath = update.DestPath + } + if update.Message != "" { + task.Message = update.Message + } + if update.Metrics != nil { + task.Metrics = cloneTaskMetrics(update.Metrics) + } +} + +func cloneBackgroundTask(task BackgroundTask) BackgroundTask { + task.Metrics = cloneTaskMetrics(task.Metrics) + if task.FinishedAt != nil { + finishedAt := *task.FinishedAt + task.FinishedAt = &finishedAt + } + return task +} + +func cloneTaskMetrics(metrics map[string]int64) map[string]int64 { + if len(metrics) == 0 { + return nil + } + out := make(map[string]int64, len(metrics)) + for key, value := range metrics { + out[key] = value + } + return out +} + +func OrganizeTaskMetrics(res *OrganizeResult) map[string]int64 { + if res == nil { + return nil + } + metrics := map[string]int64{ + "organized": int64(res.Organized), + "replaced": int64(res.Replaced), + "skipped": int64(res.Skipped), + "errors": int64(len(res.Errors)), + } + var scanVisited, scanAdded, scanUpdated, scanRemoved int64 + for _, scan := range res.Scans { + scanVisited += int64(scan.Visited) + scanAdded += int64(scan.Added) + scanUpdated += int64(scan.Updated) + scanRemoved += scan.Removed + if scan.Error != "" { + metrics["scan_errors"]++ + } + } + if len(res.Scans) > 0 { + metrics["scans"] = int64(len(res.Scans)) + metrics["scan_visited"] = scanVisited + metrics["scan_added"] = scanAdded + metrics["scan_updated"] = scanUpdated + metrics["scan_removed"] = scanRemoved + } + var scrapeMatched int64 + for _, scrape := range res.Scrapes { + scrapeMatched += int64(scrape.Matched) + if scrape.Error != "" { + metrics["scrape_errors"]++ + } + if scrape.Skipped { + metrics["scrape_skipped"]++ + } + } + if len(res.Scrapes) > 0 { + metrics["scrapes"] = int64(len(res.Scrapes)) + metrics["scrape_matched"] = scrapeMatched + } + return metrics +} diff --git a/web/src/api/tasks.ts b/web/src/api/tasks.ts index 1083d44..a2cd4fd 100644 --- a/web/src/api/tasks.ts +++ b/web/src/api/tasks.ts @@ -8,9 +8,31 @@ export interface ActiveTranscode { playlist_ok: boolean } +export interface BackgroundTask { + id: string + kind: string + name: string + status: 'running' | 'completed' | 'failed' + stage?: string + source_path?: string + dest_path?: string + message?: string + error?: string + metrics?: Record + started_at: string + updated_at: string + finished_at?: string +} + +export interface BackgroundTaskSnapshot { + active: BackgroundTask[] + recent: BackgroundTask[] +} + export interface TasksSnapshot { transcodes: ActiveTranscode[] torrents: QBitTorrent[] | null + background_tasks?: BackgroundTaskSnapshot } export const tasksAPI = { diff --git a/web/src/pages/FileManagerPage.tsx b/web/src/pages/FileManagerPage.tsx index 58a7604..64b673f 100644 --- a/web/src/pages/FileManagerPage.tsx +++ b/web/src/pages/FileManagerPage.tsx @@ -67,7 +67,7 @@ type AutoOrganizeConfig = { const AUTO_ORGANIZE_DEFAULTS: AutoOrganizeConfig = { enabled: 'false', afterDownload: 'false', - scrapeAfter: 'false', + scrapeAfter: 'true', sourceDir: '', targetDir: '', transferMode: 'hardlink', @@ -127,7 +127,7 @@ export function FileManagerPage() { const [organizeTransferMode, setOrganizeTransferMode] = useState('hardlink') const [organizeMediaType, setOrganizeMediaType] = useState('auto') const [scanAfter, setScanAfter] = useState(true) - const [scrapeAfter, setScrapeAfter] = useState(false) + const [scrapeAfter, setScrapeAfter] = useState(true) const [organizeBusy, setOrganizeBusy] = useState('') const [previewItems, setPreviewItems] = useState = { + organized: '新增', + replaced: '替换', + skipped: '跳过', + errors: '错误', + scans: '扫描库', + scan_visited: '访问', + scan_added: '入库', + scan_updated: '更新', + scan_removed: '移除', + scan_errors: '扫描错误', + scrapes: '刮削库', + scrape_matched: '匹配', + scrape_skipped: '刮削跳过', + scrape_errors: '刮削错误', + visited: '访问', + added: '入库', + updated: '更新', + probed: '探测', + local_metadata: '本地元数据', + removed: '移除', + matched: '匹配', + processed: '处理', + queued: '排队', +} + +function formatMetrics(metrics?: Record): string { + if (!metrics) return '' + return Object.entries(metrics) + .filter(([, value]) => Number.isFinite(value) && value !== 0) + .map(([key, value]) => `${metricLabels[key] ?? key} ${value}`) + .join(' · ') +} + +function statusBadge(task: BackgroundTask) { + if (task.status === 'failed') { + return failed + } + if (task.status === 'completed') { + return done + } + return running +} + +function BackgroundTaskTable({ tasks, empty }: { tasks: BackgroundTask[]; empty: string }) { + if (tasks.length === 0) return

{empty}

+ return ( + + + + + + + + + + + + {tasks.map((task) => ( + + + + + + + + ))} + +
任务阶段状态结果时间
+
{task.name}
+
+ {task.source_path || task.dest_path || task.message || '-'} +
+
{task.stage || '-'}{statusBadge(task)} +
{task.error || task.message || '-'}
+ {formatMetrics(task.metrics) && ( +
{formatMetrics(task.metrics)}
+ )} +
+ {new Date(task.finished_at || task.updated_at || task.started_at).toLocaleTimeString()} +
+ ) +} + // TasksPage shows everything the backend is doing right now: ffmpeg // transcodes + qBittorrent downloads. Refreshes every 3 s. export function TasksPage() { @@ -37,6 +121,7 @@ export function TasksPage() { if (!snap) return

加载中…

const torrents = snap.torrents ?? [] + const background = snap.background_tasks ?? { active: [], recent: [] } return (
@@ -45,6 +130,20 @@ export function TasksPage() {

实时任务

+
+

整理 / 重命名 / 入库 / 刮削任务

+
+
+

运行中

+ +
+
+

最近完成

+ +
+
+
+

转码任务

{snap.transcodes.length === 0 &&

暂无运行中转码。

}