From 5f42836bbf6983a65b4a80dc8d7d0a2800fff509 Mon Sep 17 00:00:00 2001 From: ShukeBta <272197458+ShukeBta@users.noreply.github.com> Date: Sat, 13 Jun 2026 17:08:56 +0800 Subject: [PATCH] fix organize trigger task flow --- internal/config/config.go | 2 +- internal/config/config_test.go | 3 + internal/handler/organizer.go | 12 ++- internal/handler/scheduler.go | 21 +++- internal/handler/scheduler_handler.go | 9 +- internal/handler/system_extra.go | 5 +- internal/service/downloads.go | 4 +- internal/service/downloads_test.go | 6 +- internal/service/media_classifier.go | 26 ++++- internal/service/media_classifier_test.go | 32 ++++++ internal/service/organizer_scan.go | 16 ++- internal/service/organizer_scrape_test.go | 16 +++ internal/service/scheduler.go | 81 +++++++++++++-- internal/service/scheduler_test.go | 116 ++++++++++++++++++++++ web/src/api/scheduler.ts | 2 + 15 files changed, 312 insertions(+), 39 deletions(-) diff --git a/internal/config/config.go b/internal/config/config.go index 778e579..fd09429 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -260,7 +260,7 @@ func setDefaults(v *viper.Viper) { v.SetDefault("flaresolverr.timeout", 60) v.SetDefault("downloads.smart_classify", true) - v.SetDefault("organizer.smart_classify", false) + v.SetDefault("organizer.smart_classify", true) v.SetDefault("organizer.auto_after_download", false) v.SetDefault("organize.scrape_after", true) v.SetDefault("scrape.delay_min_ms", 250) diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 202dac6..3298acb 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -40,6 +40,9 @@ func TestLoadDefaults(t *testing.T) { if cfg.Secrets.JWTSecret == "" { t.Fatalf("expected auto-generated JWT secret") } + if !cfg.Organizer.SmartClassify { + t.Fatalf("expected organizer smart classify enabled by default") + } // Re-loading must reuse the persisted secret on disk. cfg2, err := Load() if err != nil { diff --git a/internal/handler/organizer.go b/internal/handler/organizer.go index 62c21dc..33c894d 100644 --- a/internal/handler/organizer.go +++ b/internal/handler/organizer.go @@ -21,7 +21,7 @@ type organizeReq struct { TransferMode string `json:"transfer_mode"` MediaType string `json:"media_type"` MediaCategory string `json:"media_category"` - ScanAfter bool `json:"scan_after"` + ScanAfter *bool `json:"scan_after"` ScrapeAfter *bool `json:"scrape_after"` LibraryID string `json:"library_id"` DryRun bool `json:"dry_run"` @@ -66,7 +66,7 @@ func organizeMediaHandler(svc *service.Container) gin.HandlerFunc { return } payload := gin.H{"path": dst} - if req.ScanAfter && !req.DryRun && svc.Scan != nil { + if organizeScanAfter(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 @@ -89,7 +89,7 @@ func organizeLibraryHandler(svc *service.Container) gin.HandlerFunc { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } - if req.ScanAfter && !req.DryRun && svc.Scan != nil { + if organizeScanAfter(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) } @@ -121,7 +121,7 @@ func organizeDirectoryHandler(svc *service.Container) gin.HandlerFunc { return } updateHTTPTask(task, "organize", "手动整理完成,准备扫描入库", service.OrganizeTaskMetrics(res)) - if req.ScanAfter && !req.DryRun && svc.Scan != nil { + if organizeScanAfter(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) } @@ -167,3 +167,7 @@ func scanAndScrapeAfterOrganize(c *gin.Context, svc *service.Container, destRoot } return svc.Scan.ScanAndScrapeLibrariesForPath(c.Request.Context(), destRoot, preferredLibraryID, scrapeAfter) } + +func organizeScanAfter(value *bool) bool { + return value == nil || *value +} diff --git a/internal/handler/scheduler.go b/internal/handler/scheduler.go index 6e0e3de..77ca3d8 100644 --- a/internal/handler/scheduler.go +++ b/internal/handler/scheduler.go @@ -2,6 +2,7 @@ package handler import ( + "errors" "net/http" "github.com/gin-gonic/gin" @@ -18,10 +19,24 @@ func schedulerStatusHandler(svc *service.Container) gin.HandlerFunc { func schedulerRunHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { name := c.Param("name") - if err := svc.Scheduler.RunNow(c.Request.Context(), name); err != nil { - c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + if !triggerSchedulerJob(c, svc, name) { return } - c.Status(http.StatusNoContent) + c.JSON(http.StatusAccepted, gin.H{"ok": true, "message": "任务已在后台触发"}) } } + +func triggerSchedulerJob(c *gin.Context, svc *service.Container, name string) bool { + if err := svc.Scheduler.RunNowAsync(c.Request.Context(), name); err != nil { + switch { + case errors.Is(err, service.ErrSchedulerJobAlreadyRunning): + c.JSON(http.StatusConflict, gin.H{"error": "任务正在运行,请稍后到实时任务查看进度"}) + case errors.Is(err, service.ErrSchedulerJobNotFound): + c.JSON(http.StatusNotFound, gin.H{"error": "任务不存在"}) + default: + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + } + return false + } + return true +} diff --git a/internal/handler/scheduler_handler.go b/internal/handler/scheduler_handler.go index dccaaa5..75ec673 100644 --- a/internal/handler/scheduler_handler.go +++ b/internal/handler/scheduler_handler.go @@ -2,8 +2,6 @@ package handler import ( - "net/http" - "github.com/gin-gonic/gin" "go.uber.org/zap" @@ -30,14 +28,11 @@ func (h *SchedulerHandler) ListTasks(c *gin.Context) { // RunTask 手动触发指定任务。 func (h *SchedulerHandler) RunTask(c *gin.Context) { name := c.Param("id") - ctx := c.Request.Context() - - if err := h.svc.Scheduler.RunNow(ctx, name); err != nil { - Error(c, http.StatusInternalServerError, ErrInternal, "任务执行失败: "+err.Error()) + if !triggerSchedulerJob(c, h.svc, name) { return } - SuccessWithMessage(c, "任务已触发执行", nil) + SuccessWithMessage(c, "任务已触发后台执行", nil) } // GetStatus 返回调度器运行状态。 diff --git a/internal/handler/system_extra.go b/internal/handler/system_extra.go index dd3418d..3d9755a 100644 --- a/internal/handler/system_extra.go +++ b/internal/handler/system_extra.go @@ -159,11 +159,10 @@ func schemaHandler(_ *service.Container) gin.HandlerFunc { // schedulerTriggerHandler is the alternate path for /admin/scheduler/:name/run. func schedulerTriggerHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { - if err := svc.Scheduler.RunNow(c.Request.Context(), c.Param("name")); err != nil { - c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + if !triggerSchedulerJob(c, svc, c.Param("name")) { return } - c.JSON(http.StatusOK, gin.H{"ok": true}) + c.JSON(http.StatusAccepted, gin.H{"ok": true, "message": "任务已在后台触发"}) } } diff --git a/internal/service/downloads.go b/internal/service/downloads.go index 430ce3f..c508eda 100644 --- a/internal/service/downloads.go +++ b/internal/service/downloads.go @@ -1138,7 +1138,7 @@ func (d *DownloadService) onTorrentComplete(ctx context.Context, torrent QBitTor Metrics: OrganizeTaskMetrics(res), }) } - if d.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" && OrganizeResultHasChanges(res) { + if d.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" && OrganizeResultNeedsVisibilitySync(res) { if taskHandle != nil { taskHandle.Update(TaskUpdate{ Stage: "scan_scrape", @@ -1147,7 +1147,7 @@ func (d *DownloadService) onTorrentComplete(ctx context.Context, torrent QBitTor }) } res.Scans, res.Scrapes = d.scanner.ScanAndScrapeLibrariesForPath(ctx, res.DestPath, "", OrganizeScrapeAfterEnabled(ctx, d.repo)) - } else if d.log != nil && res != nil && !OrganizeResultHasChanges(res) { + } else if d.log != nil && res != nil && !OrganizeResultNeedsVisibilitySync(res) { d.log.Info("auto organize completed torrent skipped scan; no destination changes", zap.String("hash", torrent.Hash), zap.String("source", source), diff --git a/internal/service/downloads_test.go b/internal/service/downloads_test.go index 0434b4a..7ed04ef 100644 --- a/internal/service/downloads_test.go +++ b/internal/service/downloads_test.go @@ -369,7 +369,7 @@ func TestDownloadPollSkipsRecordedCompletedTorrentCatchup(t *testing.T) { } } -func TestAutoOrganizeSkipsScanWhenNoFilesChanged(t *testing.T) { +func TestAutoOrganizeSyncsVisibilityWhenTargetAlreadyExists(t *testing.T) { root := t.TempDir() src := filepath.Join(root, "downloads", "国产剧", "狂飙.S01E01.2023.1080p.mkv") dest := filepath.Join(root, "media") @@ -413,8 +413,8 @@ func TestAutoOrganizeSkipsScanWhenNoFilesChanged(t *testing.T) { if err := repos.DB.Model(&model.Media{}).Count(&count).Error; err != nil { t.Fatal(err) } - if count != 0 { - t.Fatalf("no-op auto organize triggered scan and created %d media rows, want 0", count) + if count != 1 { + t.Fatalf("target already exists should still be scanned into DB, count=%d want 1", count) } } diff --git a/internal/service/media_classifier.go b/internal/service/media_classifier.go index b6de33b..f949718 100644 --- a/internal/service/media_classifier.go +++ b/internal/service/media_classifier.go @@ -13,6 +13,7 @@ import ( var ( classifierEpisodeRE = regexp.MustCompile(`(?i)\bS\d{1,2}E\d{1,3}\b|第\s*\d+\s*[集期]|(?:^|[\s._-])E\d{1,3}(?:[\s._-]|$)`) classifierSeasonRE = regexp.MustCompile(`(?i)\bS\d{1,2}\b|第\s*\d+\s*季`) + classifierJAVCodeRE = regexp.MustCompile(`(?:^|[\s._\-/\[\]()])[A-Z]{2,6}[-_]?\d{3,5}(?:[\s._\-/\[\]()]|$)`) ) const DownloadSmartClassifySettingKey = "downloads.smart_classify" @@ -34,7 +35,11 @@ func classifyMediaCategory(input mediaClassifyInput, categories map[string]strin rawText := input.Title + " " + input.Category + " " + strings.Join(input.Genres, " ") text := strings.ToLower(rawText) - isChinese := hasAny(languages, "ZH", "ZH-CN", "ZH-TW", "CN") || hasAny(countries, "CN", "TW", "HK", "MO") || containsHan(rawText) || containsAnyText(text, "华语", "国产", "国剧", "国漫") + isChinesePlatform := containsAnyText(text, + "iqiyi", "qiyi", "youku", "tencent", "wetv", "mgtv", "mango", "hunantv", "cctv", + "bilibili", "bili", "芒果tv", "腾讯视频", "优酷", "爱奇艺", + ) + isChinese := hasAny(languages, "ZH", "ZH-CN", "ZH-TW", "CN") || hasAny(countries, "CN", "TW", "HK", "MO") || containsHan(rawText) || containsAnyText(text, "华语", "国产", "国剧", "国漫") || isChinesePlatform isJapanese := hasAny(languages, "JA", "JP") || hasAny(countries, "JP") || containsJapaneseKana(rawText) || strings.Contains(text, "日番") isKorean := hasAny(languages, "KO", "KR") || hasAny(countries, "KR", "KP") || containsKoreanHangul(rawText) isEastAsian := isJapanese || isKorean || hasAny(countries, "TH", "IN", "SG") @@ -44,7 +49,10 @@ func classifyMediaCategory(input mediaClassifyInput, categories map[string]strin ) isLatinFallback := containsLatin(rawText) && !containsHan(rawText) && !containsJapaneseKana(rawText) && !containsKoreanHangul(rawText) isWestern := isWesternByMetadata || (mediaType == "tv" && isLatinFallback) - hasAnimeText := containsAnyText(text, "动画", "动漫", "番剧", "年番", "国漫", "日番", "bangumi", "anime") + hasAnimeText := containsAnyText(text, "动画", "动漫", "番剧", "年番", "国漫", "日番", "bangumi", "anime", "b-global", "ani-one", "crunchyroll") + hasVarietyText := containsAnyText(text, "综艺", "真人秀", "脱口秀", "晚会", "春晚", "gala", "festival gala", "reality", "talk show") + hasDocumentaryText := containsAnyText(text, "纪录", "纪录片", "documentary", "docu", "national geographic", "natgeo") + isAdultText := containsAnyText(text, "adult", "nsfw", "成人", "番号", "jav", "9kg", "uncensored", "无码", "有码") || classifierJAVCodeRE.MatchString(strings.ToUpper(rawText)) hasGenre := func(values ...string) bool { for _, value := range values { @@ -63,6 +71,9 @@ func classifyMediaCategory(input mediaClassifyInput, categories map[string]strin switch mediaType { case "movie": + if isAdultText { + return categoryName(categories, "adult", "成人") + } if hasGenre("16", "ANIMATION", "动画", "动漫") { return categoryName(categories, "animation_movie", "动画电影") } @@ -84,10 +95,13 @@ func classifyMediaCategory(input mediaClassifyInput, categories map[string]strin case "variety": return categoryName(categories, "variety", "综艺") case "tv": - if hasGenre("10764", "10767", "REALITY", "TALK", "综艺", "真人秀", "脱口秀") { + if isAdultText { + return categoryName(categories, "adult", "成人") + } + if hasGenre("10764", "10767", "REALITY", "TALK", "综艺", "真人秀", "脱口秀") || hasVarietyText { return categoryName(categories, "variety", "综艺") } - if hasGenre("99", "DOCUMENTARY", "纪录", "纪录片") { + if hasGenre("99", "DOCUMENTARY", "纪录", "纪录片") || hasDocumentaryText { return categoryName(categories, "documentary", "纪录片") } if hasGenre("10762", "KIDS", "儿童") { @@ -131,8 +145,10 @@ func normalizeMediaType(mediaType, title, category string) string { } text := strings.ToLower(title + " " + category) switch { - case strings.Contains(text, "adult") || strings.Contains(text, "nsfw") || strings.Contains(text, "成人") || strings.Contains(text, "番号") || strings.Contains(text, "jav") || strings.Contains(text, "9kg"): + case strings.Contains(text, "adult") || strings.Contains(text, "nsfw") || strings.Contains(text, "成人") || strings.Contains(text, "番号") || strings.Contains(text, "jav") || strings.Contains(text, "9kg") || classifierJAVCodeRE.MatchString(strings.ToUpper(title+" "+category)): return "adult" + case containsAnyText(text, "综艺", "真人秀", "脱口秀", "晚会", "春晚", "gala", "festival gala", "reality", "talk show"): + return "variety" case strings.Contains(text, "movie") || strings.Contains(text, "电影"): return "movie" case strings.Contains(text, "anime") || strings.Contains(text, "bangumi") || strings.Contains(text, "动漫") || strings.Contains(text, "动画"): diff --git a/internal/service/media_classifier_test.go b/internal/service/media_classifier_test.go index aac4d21..fb74716 100644 --- a/internal/service/media_classifier_test.go +++ b/internal/service/media_classifier_test.go @@ -89,6 +89,38 @@ func TestClassifyMediaCategoryMatchesMoviePilotStyleRules(t *testing.T) { }, want: "欧美剧", }, + { + name: "iqiyi romanized chinese drama", + input: mediaClassifyInput{ + MediaType: "tv", + Title: "Motherhood.of.Taihang.S01E01.2026.1080p.iQIYI.WEB-DL", + }, + want: "国产剧", + }, + { + name: "youku romanized chinese drama", + input: mediaClassifyInput{ + MediaType: "tv", + Title: "Ashes.to.Crown.S01E15.2160p.YOUKU.WEB-DL", + }, + want: "国产剧", + }, + { + name: "gala is variety", + input: mediaClassifyInput{ + MediaType: "tv", + Title: "HNTV Spring Festival Gala 2026 2160p WEB-DL", + }, + want: "综艺", + }, + { + name: "jav code is adult", + input: mediaClassifyInput{ + MediaType: "movie", + Title: "IPZZ-293-UC 1080p", + }, + want: "成人", + }, } for _, tt := range tests { diff --git a/internal/service/organizer_scan.go b/internal/service/organizer_scan.go index c547b09..93d3b8a 100644 --- a/internal/service/organizer_scan.go +++ b/internal/service/organizer_scan.go @@ -37,15 +37,15 @@ type OrganizeScrapeSummary struct { // metadata scraping after scanning organized files into libraries. func OrganizeScrapeAfterEnabled(ctx context.Context, repo *repository.Container) bool { if repo == nil || repo.Setting == nil { - return false + return true } if value, err := repo.Setting.Get(ctx, "organize.scrape_after"); err == nil && strings.TrimSpace(value) != "" { - return parseBoolSetting(value, false) + return parseBoolSetting(value, true) } if value, err := repo.Setting.Get(ctx, "scrape.auto_on_scan"); err == nil && strings.TrimSpace(value) != "" { - return parseBoolSetting(value, false) + return parseBoolSetting(value, true) } - return false + return true } // OrganizeResultHasChanges reports whether an organize run actually changed @@ -56,6 +56,14 @@ func OrganizeResultHasChanges(res *OrganizeResult) bool { return res != nil && (res.Organized > 0 || res.Replaced > 0) } +// OrganizeResultNeedsVisibilitySync reports whether a just-finished organize +// should make sure the destination media exists in the DB. Skipped target files +// can still be invisible after a restart or an older organize run, so download +// completion should run one visibility sync even when no bytes were moved. +func OrganizeResultNeedsVisibilitySync(res *OrganizeResult) bool { + return res != nil && (res.Organized > 0 || res.Replaced > 0 || res.Skipped > 0) +} + // ScanLibrariesForPath recursively scans libraries affected by an organize // destination. If preferredLibraryID is set, only that library is scanned. // Otherwise every enabled library whose path intersects destRoot is scanned; diff --git a/internal/service/organizer_scrape_test.go b/internal/service/organizer_scrape_test.go index 07a852c..ed50656 100644 --- a/internal/service/organizer_scrape_test.go +++ b/internal/service/organizer_scrape_test.go @@ -68,3 +68,19 @@ func TestOrganizeDirectoryScanAndScrapeAfter(t *testing.T) { t.Fatalf("organized file missing at %q: %v", media.Path, err) } } + +func TestOrganizeScrapeAfterEnabledDefaultsOn(t *testing.T) { + if !OrganizeScrapeAfterEnabled(t.Context(), nil) { + t.Fatalf("organize scrape-after should default on without a repo") + } + repos := newOrganizerTestRepo(t) + if !OrganizeScrapeAfterEnabled(t.Context(), repos) { + t.Fatalf("organize scrape-after should default on when setting is absent") + } + if err := repos.Setting.Set(t.Context(), "organize.scrape_after", "false"); err != nil { + t.Fatal(err) + } + if OrganizeScrapeAfterEnabled(t.Context(), repos) { + t.Fatalf("explicit organize.scrape_after=false should be respected") + } +} diff --git a/internal/service/scheduler.go b/internal/service/scheduler.go index 0a03215..b2be194 100644 --- a/internal/service/scheduler.go +++ b/internal/service/scheduler.go @@ -25,6 +25,7 @@ package service import ( "context" + "errors" "strconv" "strings" "sync" @@ -55,6 +56,11 @@ type SchedulerService struct { jobs []*scheduledJob } +var ( + ErrSchedulerJobNotFound = errors.New("scheduled job not found") + ErrSchedulerJobAlreadyRunning = errors.New("scheduled job already running") +) + func (s *SchedulerService) SetTaskTracker(tasks *TaskTrackerService) { s.tasks = tasks } @@ -66,6 +72,8 @@ type scheduledJob struct { run func(ctx context.Context) error lastRun time.Time lastErr string + running bool + started time.Time } type schedulerManualRunKey struct{} @@ -168,6 +176,8 @@ type JobStatus struct { Interval string `json:"interval"` LastRun time.Time `json:"last_run,omitempty"` LastErr string `json:"last_err,omitempty"` + Running bool `json:"running,omitempty"` + Started time.Time `json:"started_at,omitempty"` } // Status returns the current state of every registered job. @@ -181,6 +191,8 @@ func (s *SchedulerService) Status() []JobStatus { Interval: j.interval.String(), LastRun: j.lastRun, LastErr: j.lastErr, + Running: j.running, + Started: j.started, }) } return out @@ -188,11 +200,34 @@ func (s *SchedulerService) Status() []JobStatus { // RunNow triggers a single run of the named job synchronously. func (s *SchedulerService) RunNow(ctx context.Context, name string) error { - for _, j := range s.jobs { - if j.name == name { - return s.runOnce(context.WithValue(ctx, schedulerManualRunKey{}, true), j) - } + j := s.jobByName(name) + if j == nil { + return ErrSchedulerJobNotFound } + return s.runOnce(context.WithValue(ctx, schedulerManualRunKey{}, true), j) +} + +// RunNowAsync triggers a named job in the background and returns immediately. +// The job is detached from the HTTP request cancellation so a browser timeout, +// route change, or reverse-proxy disconnect cannot kill long organize/scan work. +func (s *SchedulerService) RunNowAsync(ctx context.Context, name string) error { + j := s.jobByName(name) + if j == nil { + return ErrSchedulerJobNotFound + } + runCtx := context.Background() + if ctx != nil { + runCtx = context.WithoutCancel(ctx) + } + runCtx = context.WithValue(runCtx, schedulerManualRunKey{}, true) + if err := s.beginRun(j); err != nil { + return err + } + go func() { + if err := s.runReserved(runCtx, j); err != nil && s.log != nil { + s.log.Warn("manual scheduled job failed", zap.String("name", name), zap.Error(err)) + } + }() return nil } @@ -235,6 +270,35 @@ func (s *SchedulerService) loopWithInitialDelay(ctx context.Context, j *schedule } func (s *SchedulerService) runOnce(ctx context.Context, j *scheduledJob) error { + if err := s.beginRun(j); err != nil { + return err + } + return s.runReserved(ctx, j) +} + +func (s *SchedulerService) jobByName(name string) *scheduledJob { + s.mu.Lock() + defer s.mu.Unlock() + for _, j := range s.jobs { + if j.name == name { + return j + } + } + return nil +} + +func (s *SchedulerService) beginRun(j *scheduledJob) error { + s.mu.Lock() + defer s.mu.Unlock() + if j.running { + return ErrSchedulerJobAlreadyRunning + } + j.running = true + j.started = s.currentTime() + return nil +} + +func (s *SchedulerService) runReserved(ctx context.Context, j *scheduledJob) error { err := j.run(ctx) s.mu.Lock() j.lastRun = s.currentTime() @@ -243,12 +307,15 @@ func (s *SchedulerService) runOnce(ctx context.Context, j *scheduledJob) error { } else { j.lastErr = "" } + j.running = false + j.started = time.Time{} + lastErr := j.lastErr s.mu.Unlock() if s.hub != nil { s.hub.Publish("scheduler", map[string]any{ "name": j.name, "ok": err == nil, - "error": j.lastErr, + "error": lastErr, }) } return err @@ -504,7 +571,7 @@ func (s *SchedulerService) jobOrganizeSource(ctx context.Context) error { Metrics: OrganizeTaskMetrics(res), }) } - if s.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" && OrganizeResultHasChanges(res) { + if s.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" && OrganizeResultNeedsVisibilitySync(res) { if task != nil { task.Update(TaskUpdate{ Stage: "scan_scrape", @@ -513,7 +580,7 @@ func (s *SchedulerService) jobOrganizeSource(ctx context.Context) error { }) } res.Scans, res.Scrapes = s.scanner.ScanAndScrapeLibrariesForPath(ctx, res.DestPath, "", OrganizeScrapeAfterEnabled(ctx, s.repo)) - } else if s.log != nil && res != nil && !OrganizeResultHasChanges(res) { + } else if s.log != nil && res != nil && !OrganizeResultNeedsVisibilitySync(res) { s.log.Info("scheduled source organize skipped scan; no destination changes", zap.String("source", res.SourcePath), zap.String("dest", res.DestPath), diff --git a/internal/service/scheduler_test.go b/internal/service/scheduler_test.go index fe27579..7649aed 100644 --- a/internal/service/scheduler_test.go +++ b/internal/service/scheduler_test.go @@ -2,6 +2,7 @@ package service import ( "context" + "errors" "net/http" "net/http/httptest" "os" @@ -114,6 +115,121 @@ func TestSchedulerRunNowOrganizeSourceBypassesDisabledSwitch(t *testing.T) { } } +func TestSchedulerRunNowAsyncSurvivesCallerCancellation(t *testing.T) { + scheduler := NewSchedulerService(zap.NewNop(), nil, nil, nil, nil, nil, nil, "") + started := make(chan struct{}) + release := make(chan struct{}) + finished := make(chan struct{}) + scheduler.jobs = []*scheduledJob{{ + name: "organize_source", + interval: time.Minute, + run: func(ctx context.Context) error { + close(started) + defer close(finished) + select { + case <-ctx.Done(): + return ctx.Err() + case <-release: + return nil + } + }, + }} + + ctx, cancel := context.WithCancel(t.Context()) + if err := scheduler.RunNowAsync(ctx, "organize_source"); err != nil { + t.Fatalf("run now async: %v", err) + } + <-started + cancel() + select { + case <-finished: + t.Fatal("manual scheduled job was canceled with the HTTP caller context") + case <-time.After(50 * time.Millisecond): + } + close(release) + select { + case <-finished: + case <-time.After(time.Second): + t.Fatal("manual scheduled job did not finish after release") + } + status := scheduler.Status() + if len(status) != 1 || status[0].Running || status[0].LastErr != "" { + t.Fatalf("unexpected status after async run: %+v", status) + } +} + +func TestSchedulerRunNowAsyncRejectsDuplicateRun(t *testing.T) { + scheduler := NewSchedulerService(zap.NewNop(), nil, nil, nil, nil, nil, nil, "") + started := make(chan struct{}) + release := make(chan struct{}) + scheduler.jobs = []*scheduledJob{{ + name: "organize_source", + interval: time.Minute, + run: func(ctx context.Context) error { + close(started) + <-release + return nil + }, + }} + + if err := scheduler.RunNowAsync(t.Context(), "organize_source"); err != nil { + t.Fatalf("first run now async: %v", err) + } + <-started + if err := scheduler.RunNowAsync(t.Context(), "organize_source"); !errors.Is(err, ErrSchedulerJobAlreadyRunning) { + t.Fatalf("duplicate run error = %v, want %v", err, ErrSchedulerJobAlreadyRunning) + } + close(release) +} + +func TestSchedulerOrganizeSourceSyncsVisibilityWhenTargetAlreadyExists(t *testing.T) { + root := t.TempDir() + src := filepath.Join(root, "downloads") + dest := filepath.Join(root, "media") + writeOrgFile(t, filepath.Join(src, "国产剧", "狂飙.S01E01.2023.1080p.mkv"), "episode") + + repos := newOrganizerTestRepo(t) + for key, value := range map[string]string{ + "organize.source_dir": src, + "organize.target_dir": dest, + "organize.transfer_mode": "copy", + } { + if err := repos.Setting.Set(t.Context(), key, value); err != nil { + t.Fatal(err) + } + } + lib := model.Library{Name: "国产剧", Path: filepath.Join(dest, "电视剧", "国产剧"), Type: "tv", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + organizer := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + if _, err := organizer.OrganizeDirectory(t.Context(), OrganizeOptions{ + SourcePath: src, + DestPath: dest, + TransferMode: TransferCopy, + }); err != nil { + t.Fatalf("seed organize destination: %v", err) + } + scanner := NewScannerService(&config.Config{}, zap.NewNop(), repos, NewHub(zap.NewNop()), nil, nil) + scheduler := NewSchedulerService(zap.NewNop(), repos, scanner, nil, organizer, nil, NewHub(zap.NewNop()), "") + scheduler.jobs = []*scheduledJob{{ + name: "organize_source", + interval: time.Minute, + run: scheduler.jobOrganizeSource, + }} + + if err := scheduler.RunNow(t.Context(), "organize_source"); err != nil { + t.Fatalf("run now organize source: %v", err) + } + var count int64 + if err := repos.DB.Model(&model.Media{}).Count(&count).Error; err != nil { + t.Fatal(err) + } + if count != 1 { + t.Fatalf("target already exists should still be scanned into DB, count=%d want 1", count) + } +} + func TestSchedulerCloudSyncImportsMountedCloudLibrary(t *testing.T) { upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if r.URL.Path != "/file/sort" || r.URL.Query().Get("pdir_fid") != "0" { diff --git a/web/src/api/scheduler.ts b/web/src/api/scheduler.ts index 14ff7c3..527d748 100644 --- a/web/src/api/scheduler.ts +++ b/web/src/api/scheduler.ts @@ -5,6 +5,8 @@ export interface JobStatus { interval: string last_run?: string last_err?: string + running?: boolean + started_at?: string } export const schedulerAPI = {