mirror of
https://github.com/truewhile/MeBox.git
synced 2026-09-29 19:36:36 +08:00
fix organize trigger task flow
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 返回调度器运行状态。
|
||||
|
||||
@@ -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": "任务已在后台触发"})
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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, "动画"):
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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" {
|
||||
|
||||
@@ -5,6 +5,8 @@ export interface JobStatus {
|
||||
interval: string
|
||||
last_run?: string
|
||||
last_err?: string
|
||||
running?: boolean
|
||||
started_at?: string
|
||||
}
|
||||
|
||||
export const schedulerAPI = {
|
||||
|
||||
Reference in New Issue
Block a user