From 495292ab3b631bf29f8daf22af3d4f7abf078f31 Mon Sep 17 00:00:00 2001 From: truewhile <62226914+truewhile@users.noreply.github.com> Date: Thu, 20 Aug 2026 15:43:39 +0800 Subject: [PATCH] =?UTF-8?q?bug=E5=A4=84=E7=90=86=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit bug处理优化 --- cmd/server/log_rotate.go | 17 +- cmd/server/logging.go | 25 +- cmd/server/logging_test.go | 7 +- cmd/server/network.go | 12 +- internal/handler/profile.go | 8 +- internal/service/organize_pipeline.go | 2 + internal/service/organize_pipeline_tasks.go | 1 + .../organizer_directory_source_existing.go | 4 + .../service/organizer_directory_versions.go | 105 +++++- internal/service/organizer_media.go | 7 + internal/service/organizer_sidecar.go | 63 ++++ internal/service/organizer_sidecar_test.go | 78 +++++ internal/service/scraper_local_metadata.go | 9 +- internal/service/task_tracker.go | 8 + .../service/task_tracker_organize_items.go | 146 ++++++++ internal/service/task_tracker_test.go | 66 ++++ internal/service/transcoder_ffmpeg_probe.go | 8 +- internal/service/transfer_test.go | 3 + internal/service/transmission_rpc.go | 70 ++-- web/src/api/tasks.ts | 14 + web/src/hooks/useSSE.ts | 7 +- web/src/hooks/useWebSocket.ts | 6 +- web/src/pages/TaskItemRetryDialog.tsx | 171 +++++++++ web/src/pages/TasksPage.tsx | 327 +++++++++++++++--- 24 files changed, 1043 insertions(+), 121 deletions(-) create mode 100644 internal/service/organizer_sidecar_test.go create mode 100644 internal/service/task_tracker_organize_items.go create mode 100644 web/src/pages/TaskItemRetryDialog.tsx diff --git a/cmd/server/log_rotate.go b/cmd/server/log_rotate.go index ee606d8..ace5547 100644 --- a/cmd/server/log_rotate.go +++ b/cmd/server/log_rotate.go @@ -76,14 +76,19 @@ func (w *rotatingFileWriter) Sync() error { if w.file == nil { return nil } - err := w.file.Sync() - closeErr := w.file.Close() + return w.file.Sync() +} + +func (w *rotatingFileWriter) Close() error { + w.mu.Lock() + defer w.mu.Unlock() + if w.file == nil { + return nil + } + err := w.file.Close() w.file = nil w.size = 0 - if err != nil { - return err - } - return closeErr + return err } func (w *rotatingFileWriter) open() error { diff --git a/cmd/server/logging.go b/cmd/server/logging.go index 3d67994..cc61329 100644 --- a/cmd/server/logging.go +++ b/cmd/server/logging.go @@ -13,8 +13,14 @@ import ( // newLogger 根据 cfg.Logging 构建 Zap。 func newLogger(cfg *config.Config) (*zap.Logger, error) { + log, _, err := newLoggerWithCloser(cfg) + return log, err +} + +func newLoggerWithCloser(cfg *config.Config) (*zap.Logger, func(), error) { if cfg.App.Debug { - return zap.NewDevelopment() + log, err := zap.NewDevelopment() + return log, func() {}, err } level := configuredLogLevel(cfg.Logging.Level) encoderCfg := zap.NewProductionEncoderConfig() @@ -28,33 +34,42 @@ func newLogger(cfg *config.Config) (*zap.Logger, error) { cores := []zapcore.Core{ zapcore.NewCore(encoder, zapcore.Lock(os.Stdout), level), } + var closers []func() error appPath, warnPath, errorPath := logFilePaths(cfg) if appPath != "" { appWriter, err := newRotatingFileWriter(appPath, cfg.Logging) if err != nil { - return nil, err + return nil, nil, err } cores = append(cores, zapcore.NewCore(encoder, appWriter, level)) + closers = append(closers, appWriter.Close) } if warnPath != "" { warnWriter, err := newRotatingFileWriter(warnPath, cfg.Logging) if err != nil { - return nil, err + return nil, nil, err } cores = append(cores, zapcore.NewCore(encoder, warnWriter, zap.LevelEnablerFunc(func(lvl zapcore.Level) bool { return lvl == zapcore.WarnLevel && level.Enabled(lvl) }))) + closers = append(closers, warnWriter.Close) } if errorPath != "" { errorWriter, err := newRotatingFileWriter(errorPath, cfg.Logging) if err != nil { - return nil, err + return nil, nil, err } cores = append(cores, zapcore.NewCore(encoder, errorWriter, zap.LevelEnablerFunc(func(lvl zapcore.Level) bool { return lvl >= zapcore.ErrorLevel && level.Enabled(lvl) }))) + closers = append(closers, errorWriter.Close) } - return zap.New(zapcore.NewTee(cores...), zap.AddCaller(), zap.AddStacktrace(zapcore.ErrorLevel), zap.ErrorOutput(zapcore.Lock(os.Stderr))), nil + closeFn := func() { + for _, c := range closers { + _ = c() + } + } + return zap.New(zapcore.NewTee(cores...), zap.AddCaller(), zap.AddStacktrace(zapcore.ErrorLevel), zap.ErrorOutput(zapcore.Lock(os.Stderr))), closeFn, nil } func configuredLogLevel(raw string) zapcore.Level { diff --git a/cmd/server/logging_test.go b/cmd/server/logging_test.go index 9d82d16..5532753 100644 --- a/cmd/server/logging_test.go +++ b/cmd/server/logging_test.go @@ -22,10 +22,11 @@ func TestProductionLoggerWritesConfiguredInfoToAppLogAndSplitsWarnError(t *testi cfg.Logging.MaxSizeMB = 1 cfg.Logging.MaxBackups = 2 - log, err := newLogger(cfg) + log, closeFn, err := newLoggerWithCloser(cfg) if err != nil { t.Fatal(err) } + defer closeFn() log.Info("info should be stored") log.Warn("warning only", zap.String("kind", "warn")) log.Error("error only", zap.String("kind", "error")) @@ -70,10 +71,11 @@ func TestProductionLoggerDefaultsToWarnInAppLog(t *testing.T) { cfg.Logging.OutputPath = filepath.Join(dir, "logs") cfg.Logging.EnableRotation = true - log, err := newLogger(cfg) + log, closeFn, err := newLoggerWithCloser(cfg) if err != nil { t.Fatal(err) } + defer closeFn() log.Info("info should stay quiet by default") log.Warn("warning should be stored") _ = log.Sync() @@ -101,6 +103,7 @@ func TestRotatingFileWriterCapsFileSize(t *testing.T) { if err != nil { t.Fatal(err) } + defer writer.Close() chunk := strings.Repeat("x", 700*1024) if _, err := writer.Write([]byte(chunk)); err != nil { t.Fatal(err) diff --git a/cmd/server/network.go b/cmd/server/network.go index 2b78ebe..a72e960 100644 --- a/cmd/server/network.go +++ b/cmd/server/network.go @@ -1,8 +1,10 @@ package main import ( + "io" "net" "net/http" + "strings" "time" ) @@ -42,10 +44,12 @@ func getPublicIP(timeout time.Duration) string { return "" } defer resp.Body.Close() - buf := make([]byte, 64) - n, err := resp.Body.Read(buf) - if err != nil || n == 0 || resp.StatusCode != http.StatusOK { + if resp.StatusCode != http.StatusOK { return "" } - return string(buf[:n]) + data, err := io.ReadAll(io.LimitReader(resp.Body, 64)) + if err != nil || len(data) == 0 { + return "" + } + return strings.TrimSpace(string(data)) } diff --git a/internal/handler/profile.go b/internal/handler/profile.go index 54c51e6..2b32752 100644 --- a/internal/handler/profile.go +++ b/internal/handler/profile.go @@ -20,8 +20,12 @@ func updateProfileHandler(svc *service.Container) gin.HandlerFunc { c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) return } - uid, _ := c.Get(middleware.CtxUserID) - userID := uid.(string) + uid, exists := c.Get(middleware.CtxUserID) + userID, ok := uid.(string) + if !exists || !ok || userID == "" { + c.JSON(http.StatusUnauthorized, gin.H{"error": "unauthorized"}) + return + } hideAdultChanged, err := profileHideAdultChanged(c.Request.Context(), svc, userID, patch) if err != nil { c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) diff --git a/internal/service/organize_pipeline.go b/internal/service/organize_pipeline.go index 62fe9fa..6b819b7 100644 --- a/internal/service/organize_pipeline.go +++ b/internal/service/organize_pipeline.go @@ -118,6 +118,7 @@ func (p *OrganizePipelineService) Run(ctx context.Context, req OrganizePipelineR Message: "整理/重命名完成,准备扫描入库", Metrics: OrganizeTaskMetrics(res), Details: OrganizeTaskDetails(res, 8), + Items: organizeItemsFromResult(res), }) } @@ -128,6 +129,7 @@ func (p *OrganizePipelineService) Run(ctx context.Context, req OrganizePipelineR Message: "正在扫描入库并按设置刮削", Metrics: OrganizeTaskMetrics(res), Details: OrganizeTaskDetails(res, 8), + Items: combineOrganizeItems(res), }) } scanRoot := organizeScanRoot(res, path) diff --git a/internal/service/organize_pipeline_tasks.go b/internal/service/organize_pipeline_tasks.go index 4cf600b..7f62c54 100644 --- a/internal/service/organize_pipeline_tasks.go +++ b/internal/service/organize_pipeline_tasks.go @@ -34,6 +34,7 @@ func (p *OrganizePipelineService) finishTask(task *TaskHandle, err error, stage, Message: message, Metrics: OrganizeTaskMetrics(res), Details: OrganizeTaskDetails(res, 8), + Items: combineOrganizeItems(res), }) } diff --git a/internal/service/organizer_directory_source_existing.go b/internal/service/organizer_directory_source_existing.go index 5880d7f..230c750 100644 --- a/internal/service/organizer_directory_source_existing.go +++ b/internal/service/organizer_directory_source_existing.go @@ -127,6 +127,10 @@ func (o *OrganizerService) writeOrganizedSourceFile(ctx context.Context, req org o.log.Warn("organize sidecar nfo failed", zap.String("from", req.Source), zap.String("to", plan.Target.Path), zap.Error(err)) } + if err := transferSidecarArtwork(req.Source, plan.Target.Path, req.Mode); err != nil { + o.log.Warn("organize sidecar artwork failed", + zap.String("from", req.Source), zap.String("to", plan.Target.Path), zap.Error(err)) + } o.persistOrganizedSourceMetadata(ctx, plan) req.Result.Organized++ return nil diff --git a/internal/service/organizer_directory_versions.go b/internal/service/organizer_directory_versions.go index a60238a..fcc6b57 100644 --- a/internal/service/organizer_directory_versions.go +++ b/internal/service/organizer_directory_versions.go @@ -2,10 +2,11 @@ package service import ( "context" - "fmt" "os" "path/filepath" + "strconv" "strings" + "time" "go.uber.org/zap" @@ -185,33 +186,109 @@ func (o *OrganizerService) existingByFolder(destDir, episodeTag string) []string return out } -// replaceVersions removes the existing lower-resolution files (and their NFO -// sidecars + DB rows) and transfers src into dst. +// replaceVersions promotes a higher-resolution src into dst, then removes the +// old lower-resolution files (and their NFO sidecars + DB rows). dst may +// coincide with one of the existing versions being superseded, so the new file +// is first transferred to a temporary staging name inside the destination +// directory. Only after the transfer succeeds are the old versions removed and +// the staged file renamed into place. This keeps the lower-res versions intact +// if the transfer fails (hardlink across mounts, disk full, etc.) instead of +// destroying them before the new file exists. func (o *OrganizerService) replaceVersions(ctx context.Context, src string, existing []string, dst string, mode TransferMode) error { + dstDir := filepath.Dir(dst) + if err := os.MkdirAll(dstDir, 0o755); err != nil { // #nosec G301 -- organized media directories must remain readable by NAS/player users. + return err + } + // Unique staging name so dst (which may already exist as the lower-res + // version) is never the transfer target and never clobbered early. + stage := dst + ".replace" + randomSuffix() + cleanup := func() { + _ = os.Remove(stage) + _ = os.Remove(nfoPath(stage)) + removeStagedArtwork(stage) + } + if err := transferFile(src, stage, mode); err != nil { + cleanup() + return err + } + if err := transferSidecarNFO(src, stage, mode); err != nil { + o.log.Warn("organize replace sidecar nfo failed", + zap.String("from", src), zap.String("to", dst), zap.Error(err)) + } + if err := transferSidecarArtwork(src, stage, mode); err != nil { + o.log.Warn("organize replace sidecar artwork failed", + zap.String("from", src), zap.String("to", dst), zap.Error(err)) + } + // New file is safely staged; the transfer succeeded so it is now safe to + // supersede the existing lower-res versions. for _, e := range existing { - if err := os.Remove(e); err != nil && !os.IsNotExist(err) { - return fmt.Errorf("remove existing %s: %w", e, err) - } if nfo := nfoPath(e); nfo != "" { _ = os.Remove(nfo) } + if err := os.Remove(e); err != nil && !os.IsNotExist(err) { + o.log.Warn("organize replace remove existing failed", + zap.String("path", e), zap.Error(err)) + } if o.repo != nil && o.repo.DB != nil { _ = o.repo.DB.WithContext(ctx).Where("path = ?", e).Delete(&model.Media{}).Error } } - if err := os.MkdirAll(filepath.Dir(dst), 0o755); err != nil { // #nosec G301 -- organized media directories must remain readable by NAS/player users. + // Move staged file + sidecars into the final path. + if err := os.Rename(stage, dst); err != nil { + cleanup() return err } - if err := transferFile(src, dst, mode); err != nil { - return err - } - if err := transferSidecarNFO(src, dst, mode); err != nil { - o.log.Warn("organize sidecar nfo failed", - zap.String("from", src), zap.String("to", dst), zap.Error(err)) - } + moveSidecarRename(nfoPath(stage), nfoPath(dst)) + moveStagedArtwork(stage, dst) return nil } +// randomSuffix returns a short random suffix for staging filenames. +func randomSuffix() string { + return strconv.Itoa(int(time.Now().UnixNano()) & 0xffffff) +} + +// moveStagedArtwork renames artwork sidecars staged alongside `stage` into +// their final names next to `dst`. +func moveStagedArtwork(stage, dst string) { + stageBase := strings.TrimSuffix(filepath.Base(stage), filepath.Ext(stage)) + dstBase := strings.TrimSuffix(filepath.Base(dst), filepath.Ext(dst)) + stageDir := filepath.Dir(stage) + for _, suffix := range artworkSidecarSuffixes { + for _, ext := range artworkSidecarExtensions { + srcPath := filepath.Join(stageDir, stageBase+suffix+ext) + if _, err := os.Stat(srcPath); err != nil { + continue + } + _ = os.Rename(srcPath, filepath.Join(stageDir, dstBase+suffix+ext)) + } + } +} + +// moveSidecarRename renames a staged sidecar to its final name if it exists. +func moveSidecarRename(from, to string) { + if from == "" || to == "" { + return + } + if _, err := os.Stat(from); err != nil { + return + } + _ = os.Rename(from, to) +} + +// removeStagedArtwork removes artwork sidecars that were staged alongside +// `stage`, used when a replace fails and its staged outputs must be cleaned up. +func removeStagedArtwork(stage string) { + stageBase := strings.TrimSuffix(filepath.Base(stage), filepath.Ext(stage)) + stageDir := filepath.Dir(stage) + for _, suffix := range artworkSidecarSuffixes { + for _, ext := range artworkSidecarExtensions { + path := filepath.Join(stageDir, stageBase+suffix+ext) + _ = os.Remove(path) + } + } +} + // resolutionArea returns the pixel area (width*height) of a video file for 洗版 // comparison. It prefers ffprobe; when unavailable it falls back to a // resolution token in the filename (2160p/1080p/720p). Returns 0 when the diff --git a/internal/service/organizer_media.go b/internal/service/organizer_media.go index 7449c16..d918d03 100644 --- a/internal/service/organizer_media.go +++ b/internal/service/organizer_media.go @@ -226,6 +226,13 @@ func (o *OrganizerService) applyOrganizeMedia(ctx context.Context, req organizeM zap.String("to", nfoPath(dst.path)), zap.Error(err)) } + if err := transferSidecarArtwork(m.Path, dst.path, req.transferMode); err != nil { + o.log.Warn("organize sidecar artwork failed", + zap.String("media", m.ID), + zap.String("from", m.Path), + zap.String("to", dst.path), + zap.Error(err)) + } o.log.Info("organized", zap.String("media", m.ID), zap.String("from", m.Path), diff --git a/internal/service/organizer_sidecar.go b/internal/service/organizer_sidecar.go index 7b8495b..0d11e53 100644 --- a/internal/service/organizer_sidecar.go +++ b/internal/service/organizer_sidecar.go @@ -3,6 +3,7 @@ package service import ( "os" "path/filepath" + "strings" ) // transferSidecarNFO moves/copies/links the .nfo sidecar alongside its media @@ -27,3 +28,65 @@ func transferSidecarNFO(srcMedia, dstMedia string, mode TransferMode) error { } return transferFile(src, dst, mode) } + +// artworkSidecarSuffixes lists the poster/backdrop sidecar name suffixes that +// scrapers write next to a media file. It is deliberately conservative: we only +// follow the canonical "-poster" / "-backdrop" names and their +// Emby/Jellyfin variants, all scoped to the same base name as the media. +var artworkSidecarSuffixes = []string{ + "-poster", ".poster", + "-backdrop", ".backdrop", + "-fanart", ".fanart", "-landscape", + "-cover", ".cover", "-thumb", ".thumb", +} + +// artworkSidecarExtensions are the image extensions a sidecar may use. Both +// "base-poster.jpg" (bare separator) and "base.poster.jpg" (dotted separator) +// rely on suffix matching, so we probe the common image extensions. +var artworkSidecarExtensions = []string{".jpg", ".jpeg", ".png", ".webp", ".gif", ".bmp", ".tbn"} + +// transferSidecarArtwork moves/copies/links the scraped poster/backdrop +// sidecar files alongside its media using the same transfer mode, mirroring +// transferSidecarNFO. Previously organize moved only the .nfo and left posters +// behind in the old folder; this keeps artwork with the organized file. +func transferSidecarArtwork(srcMedia, dstMedia string, mode TransferMode) error { + srcDir := filepath.Dir(srcMedia) + base := strings.TrimSuffix(filepath.Base(srcMedia), filepath.Ext(srcMedia)) + if base == "" || base == "." { + return nil + } + // Find every existing sidecar by probing suffix + extension combinations. + sources := make([]string, 0, len(artworkSidecarSuffixes)*len(artworkSidecarExtensions)) + seen := map[string]struct{}{} + for _, suffix := range artworkSidecarSuffixes { + for _, ext := range artworkSidecarExtensions { + path := filepath.Join(srcDir, base+suffix+ext) + key := strings.ToLower(filepath.Clean(path)) + if _, ok := seen[key]; ok { + continue + } + seen[key] = struct{}{} + if _, err := os.Stat(path); err == nil { + sources = append(sources, path) + } + } + } + if len(sources) == 0 { + return nil + } + dstDir := filepath.Dir(dstMedia) + if err := os.MkdirAll(dstDir, 0o755); err != nil { // #nosec G301 -- sidecar media directories must remain readable by NAS/player users. + return err + } + var firstErr error + for _, src := range sources { + dst := filepath.Join(dstDir, filepath.Base(src)) + if _, err := os.Stat(dst); err == nil { + continue // never clobber an existing artwork at the destination + } + if err := transferFile(src, dst, mode); err != nil && firstErr == nil { + firstErr = err + } + } + return firstErr +} diff --git a/internal/service/organizer_sidecar_test.go b/internal/service/organizer_sidecar_test.go new file mode 100644 index 0000000..58f2461 --- /dev/null +++ b/internal/service/organizer_sidecar_test.go @@ -0,0 +1,78 @@ +package service + +import ( + "os" + "path/filepath" + "testing" +) + +func TestTransferSidecarArtworkMovesPosterAndBackdrop(t *testing.T) { + srcDir := t.TempDir() + dstDir := t.TempDir() + + srcMedia := filepath.Join(srcDir, "ADN-188.mkv") + // Write a fake media + its scraped sidecar artwork. + if err := os.WriteFile(srcMedia, []byte("media"), 0o644); err != nil { + t.Fatal(err) + } + poster := filepath.Join(srcDir, "ADN-188-poster.jpg") + backdrop := filepath.Join(srcDir, "ADN-188-backdrop.jpg") + if err := os.WriteFile(poster, []byte("poster"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(backdrop, []byte("backdrop"), 0o644); err != nil { + t.Fatal(err) + } + + dstMedia := filepath.Join(dstDir, "出口, of a title (2018).mkv") + if err := transferSidecarArtwork(srcMedia, dstMedia, TransferCopy); err != nil { + t.Fatalf("transferSidecarArtwork: %v", err) + } + + // Poster + backdrop relocated with the same base file name. + dstPoster := filepath.Join(dstDir, "ADN-188-poster.jpg") + dstBackdrop := filepath.Join(dstDir, "ADN-188-backdrop.jpg") + for _, p := range []string{dstPoster, dstBackdrop} { + if _, err := os.Stat(p); err != nil { + t.Fatalf("expected %s to exist after transfer: %v", p, err) + } + } +} + +func TestTransferSidecarArtworkNoopWhenNoArtwork(t *testing.T) { + srcDir := t.TempDir() + dstDir := t.TempDir() + srcMedia := filepath.Join(srcDir, "A.mkv") + if err := os.WriteFile(srcMedia, []byte("x"), 0o644); err != nil { + t.Fatal(err) + } + if err := transferSidecarArtwork(srcMedia, filepath.Join(dstDir, "B.mkv"), TransferCopy); err != nil { + t.Fatalf("expected no error with no artwork, got %v", err) + } +} + +func TestTransferSidecarArtworkNeverOverwritesExisting(t *testing.T) { + srcDir := t.TempDir() + dstDir := t.TempDir() + srcMedia := filepath.Join(srcDir, "A.mkv") + poster := filepath.Join(srcDir, "A-poster.jpg") + if err := os.WriteFile(srcMedia, []byte("x"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(poster, []byte("new"), 0o644); err != nil { + t.Fatal(err) + } + // Pre-existing poster at destination must be preserved. + dstMedia := filepath.Join(dstDir, "A.mkv") + dstPoster := filepath.Join(dstDir, "A-poster.jpg") + if err := os.WriteFile(dstPoster, []byte("keep"), 0o644); err != nil { + t.Fatal(err) + } + if err := transferSidecarArtwork(srcMedia, dstMedia, TransferCopy); err != nil { + t.Fatalf("transferSidecarArtwork: %v", err) + } + data, _ := os.ReadFile(dstPoster) + if string(data) != "keep" { + t.Fatalf("destination poster overwritten, got %q", data) + } +} diff --git a/internal/service/scraper_local_metadata.go b/internal/service/scraper_local_metadata.go index 7584460..3bdd8da 100644 --- a/internal/service/scraper_local_metadata.go +++ b/internal/service/scraper_local_metadata.go @@ -77,10 +77,15 @@ func mergeLocalMetadataIntoMatch(match *Match, local *LocalMetadata) { if local.Overview != "" { match.Overview = local.Overview } - if local.PosterURL != "" { + // Artwork from local metadata (sidecar NFO / folder poster) should only be + // used when the online match did not already produce a usable HTTP artwork. + // Otherwise a stale local path (e.g. a poster left behind in the library root + // from a previous organize) would clobber the freshly-scraped cover and get + // re-applied on every rescrape. + if local.PosterURL != "" && !isHTTPish(match.PosterURL) { match.PosterURL = local.PosterURL } - if local.BackdropURL != "" { + if local.BackdropURL != "" && !isHTTPish(match.BackdropURL) { match.BackdropURL = local.BackdropURL } if local.Rating > 0 { diff --git a/internal/service/task_tracker.go b/internal/service/task_tracker.go index f62002a..7f079fe 100644 --- a/internal/service/task_tracker.go +++ b/internal/service/task_tracker.go @@ -34,6 +34,7 @@ type BackgroundTask struct { Error string `json:"error,omitempty"` Details []string `json:"details,omitempty"` Metrics map[string]int64 `json:"metrics,omitempty"` + Items []TaskItemRecord `json:"items,omitempty"` StartedAt time.Time `json:"started_at"` UpdatedAt time.Time `json:"updated_at"` FinishedAt *time.Time `json:"finished_at,omitempty"` @@ -46,6 +47,7 @@ type TaskUpdate struct { Message string Details []string Metrics map[string]int64 + Items []TaskItemRecord } type TaskSnapshot struct { @@ -214,6 +216,9 @@ func applyTaskUpdate(task *BackgroundTask, update TaskUpdate) { if update.Metrics != nil { task.Metrics = cloneTaskMetrics(update.Metrics) } + if update.Items != nil { + task.Items = append([]TaskItemRecord(nil), update.Items...) + } } func cloneBackgroundTask(task BackgroundTask) BackgroundTask { @@ -221,6 +226,9 @@ func cloneBackgroundTask(task BackgroundTask) BackgroundTask { if task.Details != nil { task.Details = append([]string(nil), task.Details...) } + if task.Items != nil { + task.Items = append([]TaskItemRecord(nil), task.Items...) + } if task.FinishedAt != nil { finishedAt := *task.FinishedAt task.FinishedAt = &finishedAt diff --git a/internal/service/task_tracker_organize_items.go b/internal/service/task_tracker_organize_items.go new file mode 100644 index 0000000..e95b826 --- /dev/null +++ b/internal/service/task_tracker_organize_items.go @@ -0,0 +1,146 @@ +package service + +import "strings" + +// Item-level task status. A single BackgroundTask (one organize → scan → +// scrape run) is broken down into per-file / per-library TaskItemRecord rows so +// the operator can watch exactly which file is being organized, renamed, ingested +// or scraped, and retry only the failed ones. +const ( + ItemStatusPending = "pending" // 待进行 + ItemStatusRunning = "running" // 进行中 + ItemStatusSucceeded = "succeeded" // 成功 + ItemStatusFailed = "failed" // 失败 +) + +// Item kind / phase. 整理与重命名 share the "organize" phase (per file); +// scan(入库)与 scrape(刮削)are reported per library. +const ( + ItemKindOrganize = "organize" + ItemKindScan = "scan" + ItemKindScrape = "scrape" +) + +// TaskItemRecord is a single row on the live tasks board. Each record maps to +// exactly one file (organize) or one library (ingest / scrape). +type TaskItemRecord struct { + ID string `json:"id"` + Kind string `json:"kind"` // organize / scan / scrape + Status string `json:"status"` // pending / running / succeeded / failed + Name string `json:"name"` // display name (file base name or library name) + Source string `json:"source,omitempty"` // source file path (organize) + DestPath string `json:"dest_path,omitempty"` + LibraryID string `json:"library_id,omitempty"` + Error string `json:"error,omitempty"` +} + +// organizeItemsFromResult converts the final organize result items into +// per-file task records. The organize result is only complete after the whole +// directory walk finishes, so this is called once at stage boundaries rather +// than per file. Each item is keyed by its source path for stable identity. +func organizeItemsFromResult(res *OrganizeResult) []TaskItemRecord { + if res == nil || len(res.Items) == 0 { + return nil + } + out := make([]TaskItemRecord, 0, len(res.Items)) + for _, item := range res.Items { + rec := TaskItemRecord{ + ID: "organize:" + item.Source, + Kind: ItemKindOrganize, + Name: itemBaseName(item.Source, item.Title), + Source: item.Source, + DestPath: item.Target, + } + switch item.Action { + case "error": + rec.Status = ItemStatusFailed + rec.Error = strings.TrimSpace(item.Reason) + default: + // organize / replace / reclassify / cleanup / skip all count as done. + rec.Status = ItemStatusSucceeded + } + out = append(out, rec) + } + return out +} + +// scanItemsFromResult converts per-library scan summaries into task records. +func scanItemsFromResult(res *OrganizeResult) []TaskItemRecord { + if res == nil || len(res.Scans) == 0 { + return nil + } + out := make([]TaskItemRecord, 0, len(res.Scans)) + for _, scan := range res.Scans { + rec := TaskItemRecord{ + ID: "scan:" + scan.LibraryID, + Kind: ItemKindScan, + Name: scan.Name, + LibraryID: scan.LibraryID, + Status: ItemStatusSucceeded, + } + if scan.Error != "" { + rec.Status = ItemStatusFailed + rec.Error = scan.Error + } + out = append(out, rec) + } + return out +} + +// scrapeItemsFromResult converts per-library scrape summaries into task records. +func scrapeItemsFromResult(res *OrganizeResult) []TaskItemRecord { + if res == nil || len(res.Scrapes) == 0 { + return nil + } + out := make([]TaskItemRecord, 0, len(res.Scrapes)) + for _, scrape := range res.Scrapes { + rec := TaskItemRecord{ + ID: "scrape:" + scrape.LibraryID, + Kind: ItemKindScrape, + Name: scrape.Name, + LibraryID: scrape.LibraryID, + Status: ItemStatusSucceeded, + } + if scrape.Error != "" { + rec.Status = ItemStatusFailed + rec.Error = scrape.Error + out = append(out, rec) + continue + } + if scrape.Skipped { + rec.Status = ItemStatusSucceeded + } + out = append(out, rec) + } + return out +} + +// combineOrganizeItems merges organize + scan + scrape item rows into one flat +// list ordered by phase, deduping by ID (scan/scrape share the library ID and +// there is at most one of each per run). +func combineOrganizeItems(res *OrganizeResult) []TaskItemRecord { + items := organizeItemsFromResult(res) + items = append(items, scanItemsFromResult(res)...) + items = append(items, scrapeItemsFromResult(res)...) + return items +} + +func itemBaseName(source, title string) string { + if strings.TrimSpace(title) != "" && !strings.Contains(title, ".") { + return strings.TrimSpace(title) + } + if source == "" { + return "" + } + base := source + if idx := strings.LastIndexAny(base, "/\\"); idx >= 0 { + base = base[idx+1:] + } + if idx := strings.LastIndex(base, "."); idx > 0 { + base = base[:idx] + } + if base == "" { + return source + } + return base +} diff --git a/internal/service/task_tracker_test.go b/internal/service/task_tracker_test.go index b9d3079..1baa581 100644 --- a/internal/service/task_tracker_test.go +++ b/internal/service/task_tracker_test.go @@ -27,3 +27,69 @@ func TestOrganizeTaskMetricsIncludesScrapeProcessed(t *testing.T) { t.Fatalf("scrape_skipped = %d, want 1", metrics["scrape_skipped"]) } } + +func TestCombineOrganizeItemsBuildsPerItemRows(t *testing.T) { + res := &OrganizeResult{ + Items: []OrganizePreviewItem{ + {Source: `/downloads/Show.S01E01.mkv`, Target: `/media/Show/Show.S01E01.mkv`, Action: "organize", Title: "Show"}, + {Source: `/downloads/dup.mkv`, Action: "skip", Reason: "duplicate in library"}, + {Source: `/downloads/broken.mp4`, Action: "error", Reason: "unsupported codec"}, + }, + Scans: []OrganizeScanSummary{ + {LibraryID: "lib1", Name: "电影库"}, + {LibraryID: "lib2", Name: "剧集库", Error: "scan failed"}, + }, + Scrapes: []OrganizeScrapeSummary{ + {LibraryID: "lib1", Name: "电影库", Matched: 2}, + {LibraryID: "lib2", Name: "剧集库", Error: "scrape failed"}, + }, + } + + items := combineOrganizeItems(res) + if len(items) != 7 { + t.Fatalf("len(items) = %d, want 7", len(items)) + } + + // organize items mapped by source path. + bySource := map[string]TaskItemRecord{} + for _, item := range items { + if item.Kind == ItemKindOrganize { + bySource[item.Source] = item + } + } + if got := bySource["/downloads/Show.S01E01.mkv"]; got.Status != ItemStatusSucceeded || got.Kind != ItemKindOrganize { + t.Fatalf("organized item = %#v, want succeeded/organize", got) + } + if got := bySource["/downloads/broken.mp4"]; got.Status != ItemStatusFailed || got.Error == "" { + t.Fatalf("error item = %#v, want failed with error", got) + } + if got := bySource["/downloads/dup.mkv"]; got.Status != ItemStatusSucceeded { + t.Fatalf("skip item = %#v, want succeeded", got) + } + + // scan + scrape items carry library IDs. + var scanFailed, scrapeFailed bool + for _, item := range items { + if item.Kind == ItemKindScan && item.LibraryID == "lib2" { + scanFailed = item.Status == ItemStatusFailed + } + if item.Kind == ItemKindScrape && item.LibraryID == "lib2" { + scrapeFailed = item.Status == ItemStatusFailed + } + } + if !scanFailed { + t.Fatal("scan lib2 should be failed") + } + if !scrapeFailed { + t.Fatal("scrape lib2 should be failed") + } +} + +func TestCombineOrganizeItemsHandlesNoItems(t *testing.T) { + if got := combineOrganizeItems(nil); got != nil { + t.Fatalf("combineOrganizeItems(nil) = %#v, want nil", got) + } + if got := combineOrganizeItems(&OrganizeResult{}); got != nil { + t.Fatalf("combineOrganizeItems(empty) = %#v, want nil", got) + } +} diff --git a/internal/service/transcoder_ffmpeg_probe.go b/internal/service/transcoder_ffmpeg_probe.go index 5d4800c..0901cd7 100644 --- a/internal/service/transcoder_ffmpeg_probe.go +++ b/internal/service/transcoder_ffmpeg_probe.go @@ -66,13 +66,19 @@ func (t *TranscoderService) runFFmpeg(ctx context.Context, job *hlsJob, source s } func (t *TranscoderService) resolveFFmpegPath() (string, error) { + t.mu.Lock() + configuredPath := strings.TrimSpace(t.cfg.App.FFmpegPath) + t.mu.Unlock() + var lastErr error - for _, bin := range executableCandidates(strings.TrimSpace(t.cfg.App.FFmpegPath), "ffmpeg") { + for _, bin := range executableCandidates(configuredPath, "ffmpeg") { if err := validateFFmpegForTranscode(context.Background(), bin, t.effectiveEncoder()); err != nil { lastErr = err continue } + t.mu.Lock() t.cfg.App.FFmpegPath = bin + t.mu.Unlock() return bin, nil } if lastErr != nil { diff --git a/internal/service/transfer_test.go b/internal/service/transfer_test.go index 335bbb9..5c8a305 100644 --- a/internal/service/transfer_test.go +++ b/internal/service/transfer_test.go @@ -109,6 +109,9 @@ func TestTransferFileSymlinkKeepsSource(t *testing.T) { src := writeTemp(t, dir, "src.mkv", "payload") dst := filepath.Join(dir, "dst.mkv") if err := transferFile(src, dst, TransferSymlink); err != nil { + if errors.Is(err, os.ErrPermission) || strings.Contains(strings.ToLower(err.Error()), "privilege") { + t.Skipf("skipping symlink test due to permission: %v", err) + } t.Fatalf("symlink: %v", err) } fi, err := os.Lstat(dst) diff --git a/internal/service/transmission_rpc.go b/internal/service/transmission_rpc.go index 6a10784..f17cbb1 100644 --- a/internal/service/transmission_rpc.go +++ b/internal/service/transmission_rpc.go @@ -93,41 +93,49 @@ func (a *TransmissionAdapter) rpcLocked(ctx context.Context, method string, args } for attempt := 0; attempt < 2; attempt++ { - req, err := newDownloadClientHTTPRequest(ctx, http.MethodPost, rpcURL, bytes.NewReader(body)) + res, retry, err := func() (*transmissionRPCResponse, bool, error) { + req, err := newDownloadClientHTTPRequest(ctx, http.MethodPost, rpcURL, bytes.NewReader(body)) + if err != nil { + return nil, false, err + } + req.Header.Set("Content-Type", "application/json") + if a.sessionID != "" { + req.Header.Set("X-Transmission-Session-Id", a.sessionID) + } + if a.cfg.Username != "" { + req.SetBasicAuth(a.cfg.Username, a.cfg.Password) + } + + resp, err := a.client.Do(req) + if err != nil { + return nil, false, err + } + defer resp.Body.Close() + + if resp.StatusCode == 409 { + a.sessionID = resp.Header.Get("X-Transmission-Session-Id") + return nil, true, nil + } + if resp.StatusCode >= 400 { + raw, _ := io.ReadAll(resp.Body) + return nil, false, fmt.Errorf("transmission rpc error: %d: %s", resp.StatusCode, string(raw)) + } + + var result transmissionRPCResponse + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return nil, false, err + } + if result.Result != "success" { + return nil, false, fmt.Errorf("transmission rpc result: %s", result.Result) + } + return &result, false, nil + }() if err != nil { return nil, err } - req.Header.Set("Content-Type", "application/json") - if a.sessionID != "" { - req.Header.Set("X-Transmission-Session-Id", a.sessionID) + if !retry { + return res, nil } - if a.cfg.Username != "" { - req.SetBasicAuth(a.cfg.Username, a.cfg.Password) - } - - resp, err := a.client.Do(req) - if err != nil { - return nil, err - } - defer resp.Body.Close() - - if resp.StatusCode == 409 { - a.sessionID = resp.Header.Get("X-Transmission-Session-Id") - continue - } - if resp.StatusCode >= 400 { - raw, _ := io.ReadAll(resp.Body) - return nil, fmt.Errorf("transmission rpc error: %d: %s", resp.StatusCode, string(raw)) - } - - var result transmissionRPCResponse - if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { - return nil, err - } - if result.Result != "success" { - return nil, fmt.Errorf("transmission rpc result: %s", result.Result) - } - return &result, nil } return nil, fmt.Errorf("transmission: failed after CSRF retry") } diff --git a/web/src/api/tasks.ts b/web/src/api/tasks.ts index 0ad32e5..9a7600f 100644 --- a/web/src/api/tasks.ts +++ b/web/src/api/tasks.ts @@ -8,6 +8,19 @@ export interface ActiveTranscode { playlist_ok: boolean } +export type TaskItemStatus = 'pending' | 'running' | 'succeeded' | 'failed' + +export interface TaskItem { + id: string + kind: string // organize | scan | scrape + status: TaskItemStatus + name: string + source?: string + dest_path?: string + library_id?: string + error?: string +} + export interface BackgroundTask { id: string kind: string @@ -20,6 +33,7 @@ export interface BackgroundTask { error?: string details?: string[] metrics?: Record + items?: TaskItem[] started_at: string updated_at: string finished_at?: string diff --git a/web/src/hooks/useSSE.ts b/web/src/hooks/useSSE.ts index 3e9e46c..16cfb22 100644 --- a/web/src/hooks/useSSE.ts +++ b/web/src/hooks/useSSE.ts @@ -34,6 +34,9 @@ export function useSSE( options: { autoConnect?: boolean } = {} ) { const { autoConnect = true } = options + const onEventRef = useRef(onEvent) + onEventRef.current = onEvent + const eventSourceRef = useRef(null) const reconnectTimeoutRef = useRef | null>(null) const isConnectedRef = useRef(false) @@ -64,7 +67,7 @@ export function useSSE( eventSource.onmessage = (event) => { try { const data = JSON.parse(event.data) as SSEEvent - onEvent(data) + onEventRef.current(data) } catch (err) { console.error('Failed to parse SSE event:', err) } @@ -84,7 +87,7 @@ export function useSSE( } } - }, [onEvent]) + }, []) const disconnect = useCallback(() => { if (reconnectTimeoutRef.current) { diff --git a/web/src/hooks/useWebSocket.ts b/web/src/hooks/useWebSocket.ts index 269b983..5401062 100644 --- a/web/src/hooks/useWebSocket.ts +++ b/web/src/hooks/useWebSocket.ts @@ -11,6 +11,8 @@ const MAX_RECONNECT_ATTEMPTS = 5 export function useWebSocket(onEvent: (topic: string, payload: unknown) => void) { const ref = useRef(null) const token = useAuthStore((s) => s.token) + const onEventRef = useRef(onEvent) + onEventRef.current = onEvent useEffect(() => { if (!token) return @@ -32,7 +34,7 @@ export function useWebSocket(onEvent: (topic: string, payload: unknown) => void) try { const msg = JSON.parse(ev.data) if (msg && typeof msg.topic === 'string') { - onEvent(msg.topic, msg.payload) + onEventRef.current(msg.topic, msg.payload) } } catch { // ignore malformed frames @@ -52,5 +54,5 @@ export function useWebSocket(onEvent: (topic: string, payload: unknown) => void) if (timer) window.clearTimeout(timer) ref.current?.close() } - }, [token, onEvent]) + }, [token]) } diff --git a/web/src/pages/TaskItemRetryDialog.tsx b/web/src/pages/TaskItemRetryDialog.tsx new file mode 100644 index 0000000..42e1423 --- /dev/null +++ b/web/src/pages/TaskItemRetryDialog.tsx @@ -0,0 +1,171 @@ +import { useState } from 'react' +import toast from 'react-hot-toast' +import { X } from 'lucide-react' + +import { toolsAPI } from '../api/tools' +import { libraryAPI } from '../api/library' +import type { BackgroundTask, TaskItem } from '../api/tasks' + +const kindLabels: Record = { + organize: '整理 / 重命名', + scan: '入库扫描', + scrape: '刮削', +} + +// ItemRetryDialog lets the operator manually re-handle a single failed item. +// It pre-fills the values captured from the failed task so retrying is +// reproducible, while still allowing overrides before re-running. +export function ItemRetryDialog({ + task, + item, + onClose, +}: { + task: BackgroundTask + item: TaskItem + onClose: () => void +}) { + const [sourcePath, setSourcePath] = useState(item.source || task.source_path || '') + const [destPath, setDestPath] = useState(item.dest_path || task.dest_path || '') + const [transferMode, setTransferMode] = useState('hardlink') + const [scanAfter, setScanAfter] = useState(true) + const [scrapeAfter, setScrapeAfter] = useState(true) + const [busy, setBusy] = useState(false) + + const run = async () => { + setBusy(true) + try { + if (item.kind === 'organize') { + if (!sourcePath.trim()) { + toast.error('整理需要来源路径') + return + } + await toolsAPI.organizeDirectory({ + source_path: sourcePath.trim(), + dest_path: destPath.trim() || undefined, + transfer_mode: transferMode, + scan_after: scanAfter, + scrape_after: scanAfter && scrapeAfter, + }) + toast.success(`已重新整理:${item.name}`) + } else if (item.kind === 'scan' && item.library_id) { + await libraryAPI.scan(item.library_id) + toast.success(`已重新入库:${item.name}`) + } else if (item.kind === 'scrape' && item.library_id) { + await libraryAPI.scrape(item.library_id, { + episode_artwork: true, + refresh_matched: true, + }) + toast.success(`已重新刮削:${item.name}`) + } else { + toast.error(`无法处理该类型条目(${item.kind})`) + return + } + onClose() + } catch (err: unknown) { + toast.error((err as { response?: { data?: { error?: string } } })?.response?.data?.error ?? '处理失败') + } finally { + setBusy(false) + } + } + + return ( +
+
e.stopPropagation()} + > +
+
+

手动处理失败条目

+

+ {item.name} · {kindLabels[item.kind] ?? item.kind} +

+
+ +
+ + {item.error && ( +
{item.error}
+ )} + +
+ {item.kind === 'organize' && ( + <> + + + +
+ + +
+ + )} + + {item.kind === 'scan' && ( +

+ 将重新扫描媒体库「{item.name}」以入库。完成后会自动触发刮削(如已启用)。 +

+ )} + + {item.kind === 'scrape' && ( +

+ 将重新刮削媒体库「{item.name}」。会重试上次未匹配的条目并刷新元数据。 +

+ )} +
+ +
+ + +
+
+
+ ) +} diff --git a/web/src/pages/TasksPage.tsx b/web/src/pages/TasksPage.tsx index 6476c11..87a8a4c 100644 --- a/web/src/pages/TasksPage.tsx +++ b/web/src/pages/TasksPage.tsx @@ -1,9 +1,14 @@ -import { useEffect, useState } from 'react' +import { useCallback, useEffect, useState } from 'react' import toast from 'react-hot-toast' -import { Activity, Copy } from 'lucide-react' +import { Activity, Copy, RotateCw, SlidersHorizontal } from 'lucide-react' -import { tasksAPI, type BackgroundTask, type TasksSnapshot } from '../api/tasks' +import { toolsAPI } from '../api/tools' +import { libraryAPI } from '../api/library' +import { tasksAPI, type BackgroundTask, type TaskItem, type TasksSnapshot } from '../api/tasks' +import { useWebSocket } from '../hooks/useWebSocket' +import { confirmAction } from '../components/confirmAction' import { TorrentTaskTable, TranscodeTaskTable } from './TaskRuntimeTables' +import { ItemRetryDialog } from './TaskItemRetryDialog' const metricLabels: Record = { organized: '新增', @@ -51,6 +56,19 @@ function hasTaskIssues(task: BackgroundTask): boolean { return Boolean(task.metrics?.errors || task.metrics?.scan_errors || task.metrics?.scrape_errors) } +const itemKindLabels: Record = { + organize: '整理/重命名', + scan: '入库', + scrape: '刮削', +} + +const itemStatusLabels: Record = { + pending: '待进行', + running: '进行中', + succeeded: '成功', + failed: '失败', +} + function statusBadge(task: BackgroundTask) { if (task.status === 'failed') { return failed @@ -64,6 +82,19 @@ function statusBadge(task: BackgroundTask) { return running } +function itemStatusBadge(item: TaskItem) { + switch (item.status) { + case 'failed': + return 失败 + case 'running': + return 进行中 + case 'succeeded': + return 成功 + default: + return 待进行 + } +} + function taskCopyText(task: BackgroundTask): string { const lines = [ `任务: ${task.name}`, @@ -79,6 +110,17 @@ function taskCopyText(task: BackgroundTask): string { lines.push('详情:') lines.push(...task.details) } + if (task.items?.length) { + lines.push('条目:') + lines.push( + ...task.items.map( + (item) => + `[${itemKindLabels[item.kind] ?? item.kind}] ${item.name} → ${itemStatusLabels[item.status] ?? item.status}${ + item.error ? ` (${item.error})` : '' + }`, + ), + ) + } return lines.join('\n') } @@ -91,69 +133,200 @@ async function copyTask(task: BackgroundTask) { } } -function BackgroundTaskTable({ tasks, empty }: { tasks: BackgroundTask[]; empty: string }) { +// Renders one task's per-item rows, each with status + (for failed items) +// retry / manual-handle actions. +function TaskItemRows({ + task, + onRetry, + onManual, +}: { + task: BackgroundTask + onRetry: (item: TaskItem) => void + onManual: (item: TaskItem) => void +}) { + const items = task.items ?? [] + if (items.length === 0) { + // Fall back to the compact aggregate view for legacy tasks without items. + return ( +
+
+
+ {task.name} +
+
+ {task.source_path || task.dest_path || task.message || '-'} +
+
+
+ {statusBadge(task)} + {hasTaskIssues(task) && ( + + )} +
+
+ ) + } + return ( +
+ + + + + + + + + + + + {items.map((item) => ( + + + + + + + + ))} + +
名称类型状态路径操作
+
+ {item.name} +
+ {item.error && ( +
+ {item.error} +
+ )} +
+ + {itemKindLabels[item.kind] ?? item.kind} + + {itemStatusBadge(item)} +
+ {item.source || item.dest_path || '-'} +
+
+ {item.status === 'failed' && ( +
+ + +
+ )} +
+
+ ) +} + +function TaskSection({ + tasks, + empty, + onRetry, + onManual, +}: { + tasks: BackgroundTask[] + empty: string + onRetry: (task: BackgroundTask, item: TaskItem) => void + onManual: (task: BackgroundTask, item: TaskItem) => void +}) { + const withItems = tasks.filter((task) => (task.items ?? []).length > 0) + const withoutItems = tasks.filter((task) => (task.items ?? []).length === 0) 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 || '-'}
+
+ {withItems.length > 0 && ( +
+ {withItems.map((task) => ( +
+
+
+ {task.name} + {statusBadge(task)} +
+ onRetry(task, item)} onManual={(item) => onManual(task, item)} /> {formatMetrics(task.metrics) && ( -
{formatMetrics(task.metrics)}
+
{formatMetrics(task.metrics)}
)} - {task.details && task.details.length > 0 && ( -
- {task.details.map((line, index) => ( -
- {line} -
- ))} -
- )} -
- {new Date(task.finished_at || task.updated_at || task.started_at).toLocaleTimeString()} -
+ + ))} + + )} + {withoutItems.length > 0 && ( +
+ {withoutItems.map((task) => ( +
+ onRetry(task, item)} + onManual={(item) => onManual(task, item)} + /> +
+ ))} +
+ )} + ) } // TasksPage shows everything the backend is doing right now: ffmpeg -// transcodes + qBittorrent downloads. Refreshes every 3 s. +// transcodes + qBittorrent downloads + item-level organize/ingest/scrape +// progress. Refreshes every 3 s and reconciles against live WS "task" events. export function TasksPage() { const [snap, setSnap] = useState(null) + const [manualTarget, setManualTarget] = useState<{ task: BackgroundTask; item: TaskItem } | null>(null) + + const mergeEvent = useCallback((payload: unknown) => { + const eventTask = payload as BackgroundTask + if (!eventTask || typeof eventTask.id !== 'string') return + setSnap((prev) => { + if (!prev) return prev + const merge = (list: BackgroundTask[]) => + list.map((t) => (t.id === eventTask.id ? { ...eventTask } : t)) + return { + ...prev, + background_tasks: { + active: merge(prev.background_tasks?.active ?? []), + recent: merge(prev.background_tasks?.recent ?? []), + }, + } + }) + }, []) + + useWebSocket((topic, payload) => { + if (topic === 'task') mergeEvent(payload) + }) useEffect(() => { let cancelled = false @@ -174,6 +347,42 @@ export function TasksPage() { const torrents = snap.torrents ?? [] const background = snap.background_tasks ?? { active: [], recent: [] } + const handleRetry = async (task: BackgroundTask, item: TaskItem) => { + const ok = await confirmAction({ + title: '重试失败条目', + message: `确定要重新处理「${item.name}」吗?\n类型:${itemKindLabels[item.kind] ?? item.kind}`, + confirmText: '重试', + }) + if (!ok) return + try { + if (item.kind === 'scan' && item.library_id) { + await libraryAPI.scan(item.library_id) + toast.success(`已重新入库:${item.name}`) + } else if (item.kind === 'scrape' && item.library_id) { + await libraryAPI.scrape(item.library_id) + toast.success(`已重新刮削:${item.name}`) + } else if (item.kind === 'organize' && item.source) { + const dest = task.dest_path || item.dest_path || undefined + await toolsAPI.organizeDirectory({ + source_path: item.source, + dest_path: dest, + scan_after: true, + scrape_after: true, + }) + toast.success(`已重新整理:${item.name}`) + } else { + toast.error(`无法重试:缺少必要信息(${item.kind})`) + return + } + } catch (err: unknown) { + toast.error((err as { response?: { data?: { error?: string } } })?.response?.data?.error ?? '重试失败') + } + } + + const handleManual = (task: BackgroundTask, item: TaskItem) => { + setManualTarget({ task, item }) + } + return (
@@ -186,11 +395,21 @@ export function TasksPage() {

运行中

- + void handleRetry(task, item)} + onManual={handleManual} + />

最近完成

- + void handleRetry(task, item)} + onManual={handleManual} + />
@@ -204,6 +423,14 @@ export function TasksPage() {

下载任务

+ + {manualTarget && ( + setManualTarget(null)} + /> + )}
) }