diff --git a/internal/handler/routes_admin.go b/internal/handler/routes_admin.go index 4801052..12b5eae 100644 --- a/internal/handler/routes_admin.go +++ b/internal/handler/routes_admin.go @@ -20,6 +20,24 @@ func registerAdminRoutes(api *gin.RouterGroup, cfg *config.Config, svc *service. registerAdminAPIConfigRoutes(admin, svc) registerAdminRecognitionWordRoutes(admin, svc) registerAdminStrmRoutes(admin, svc) + registerAdminScraperRoutes(admin, svc) +} + +func registerAdminScraperRoutes(admin *gin.RouterGroup, svc *service.Container) { + admin.GET("/scraper/queue", listScrapeQueueHandler(svc)) + admin.POST("/scraper/queue/:id/cancel", cancelScrapeTaskHandler(svc)) + admin.POST("/scraper/queue/:id/retry", retryScrapeTaskHandler(svc)) + admin.DELETE("/scraper/queue/:id", deleteScrapeTaskHandler(svc)) + admin.POST("/scraper/queue/batch", batchActionScrapeTasksHandler(svc)) + admin.POST("/scraper/queue/clear-done", clearDoneScrapeTasksHandler(svc)) + admin.POST("/scraper/queue/clear-finished", clearFinishedScrapeTasksHandler(svc)) + admin.POST("/scraper/queue/clear-canceled", clearCanceledScrapeTasksHandler(svc)) + admin.POST("/scraper/queue/retry-failed", retryAllFailedScrapeTasksHandler(svc)) + admin.POST("/scraper/queue/cancel-pending", cancelPendingScrapeTasksHandler(svc)) + admin.POST("/scraper/queue/enqueue-library/:id", enqueueLibraryScrapeHandler(svc)) + admin.POST("/scraper/queue/enqueue-all", enqueueAllScrapeHandler(svc)) + admin.POST("/media/repair-rescrape", enqueueAllScrapeHandler(svc)) + admin.POST("/libraries/:id/repair-rescrape", enqueueLibraryScrapeHandler(svc)) } func registerAdminStrmRoutes(admin *gin.RouterGroup, svc *service.Container) { diff --git a/internal/handler/scraper_queue.go b/internal/handler/scraper_queue.go new file mode 100644 index 0000000..5b6a21e --- /dev/null +++ b/internal/handler/scraper_queue.go @@ -0,0 +1,167 @@ +package handler + +import ( + "net/http" + "strconv" + + "github.com/gin-gonic/gin" + + "github.com/ShukeBta/MMTL/internal/service" +) + +func listScrapeQueueHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + page, _ := strconv.Atoi(c.DefaultQuery("page", "1")) + pageSize, _ := strconv.Atoi(c.DefaultQuery("page_size", "50")) + snap, err := svc.Scraper.ScrapeQueueSnapshot(c.Request.Context(), c.Query("status"), page, pageSize) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, snap) + } +} + +func cancelScrapeTaskHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + if err := svc.Scraper.CancelScrapeTask(c.Request.Context(), c.Param("id")); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) + } +} + +func retryScrapeTaskHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + if err := svc.Scraper.RetryScrapeTask(c.Request.Context(), c.Param("id")); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) + } +} + +func deleteScrapeTaskHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + if err := svc.Scraper.DeleteScrapeTask(c.Request.Context(), c.Param("id")); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) + } +} + +func batchActionScrapeTasksHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + var req queueBatchReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + n, err := svc.Scraper.BatchActionScrapeTasks(c.Request.Context(), req.Action, req.IDs) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"affected": n, "action": req.Action}) + } +} + +func clearDoneScrapeTasksHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + n, err := svc.Scraper.ClearDoneScrapeTasks(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"deleted": n}) + } +} + +func clearFinishedScrapeTasksHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + n, err := svc.Scraper.ClearFinishedScrapeTasks(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"deleted": n}) + } +} + +func clearCanceledScrapeTasksHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + n, err := svc.Scraper.ClearCanceledScrapeTasks(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"deleted": n}) + } +} + +func retryAllFailedScrapeTasksHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + n, err := svc.Scraper.RetryAllFailedScrapeTasks(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"retried": n}) + } +} + +func cancelPendingScrapeTasksHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + n, err := svc.Scraper.CancelPendingScrapeTasks(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"canceled": n}) + } +} + +type enqueueScrapeReq struct { + EpisodeImages bool `json:"episode_images"` + EpisodeArtwork bool `json:"episode_artwork"` + RefreshMatched bool `json:"refresh_matched"` + IncludeMatched bool `json:"include_matched"` +} + +func (r enqueueScrapeReq) toOptions() service.ScrapeOptions { + epArtwork := r.EpisodeImages || r.EpisodeArtwork + return service.ScrapeOptions{ + EpisodeArtwork: &epArtwork, + IncludeMatched: r.IncludeMatched || r.RefreshMatched, + RetryNoMatch: true, + } +} + +func enqueueLibraryScrapeHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + var req enqueueScrapeReq + _ = c.ShouldBindJSON(&req) + libID := c.Param("id") + n, err := svc.Scraper.EnqueueLibrary(c.Request.Context(), libID, req.toOptions()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"enqueued": n}) + } +} + +func enqueueAllScrapeHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + var req enqueueScrapeReq + _ = c.ShouldBindJSON(&req) + n, err := svc.Scraper.EnqueueAll(c.Request.Context(), req.toOptions()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"enqueued": n}) + } +} diff --git a/internal/handler/streaming.go b/internal/handler/streaming.go index 170134b..23ba58c 100644 --- a/internal/handler/streaming.go +++ b/internal/handler/streaming.go @@ -2,7 +2,6 @@ package handler import ( - "context" "errors" "io" "net/http" @@ -135,28 +134,12 @@ func scrapeOneHandler(svc *service.Container) gin.HandlerFunc { return } options.IncludeMatched = true - m, err := svc.Repo.Media.FindByID(c.Request.Context(), c.Param("id")) - if err != nil || m == nil { - c.JSON(http.StatusNotFound, gin.H{"error": "not found"}) - return - } - task := startScrapeHTTPTask(svc, "手动刮削媒体", m.Title, m.Path) - if err := svc.Scraper.EnrichOneWithOptions(c.Request.Context(), m, options); err != nil { - finishHTTPTask(task, err, "scrape", "手动刮削媒体失败", nil, nil) + task, err := svc.Scraper.EnqueueMedia(c.Request.Context(), c.Param("id"), options) + if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } - reclassified := reclassifyMediaAfterScrape(c.Request.Context(), svc, m.ID) - refreshed, _ := svc.Repo.Media.FindByID(c.Request.Context(), m.ID) - metrics := map[string]int64{"processed": 1} - if refreshed != nil && refreshed.ScrapeStatus == "matched" { - metrics["matched"] = 1 - } - if reclassified > 0 { - metrics["reclassified"] = int64(reclassified) - } - finishHTTPTask(task, nil, "completed", "手动刮削媒体结束", metrics, nil) - c.JSON(http.StatusOK, refreshed) + c.JSON(http.StatusOK, task) } } @@ -170,40 +153,12 @@ func scrapeLibraryHandler(svc *service.Container) gin.HandlerFunc { return } options.IncludeMatched = true - var task *service.TaskHandle - if lib, err := svc.Repo.Library.FindByID(c.Request.Context(), libID); err == nil && lib != nil { - task = startScrapeHTTPTask(svc, "手动刮削媒体库", lib.Name, lib.Path) - } else { - task = startScrapeHTTPTask(svc, "手动刮削媒体库", libID, "") + n, err := svc.Scraper.EnqueueLibrary(c.Request.Context(), libID, options) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return } - // Run in the background so HTTP returns instantly; the WS hub - // pushes per-item progress on the "scrape" topic. - go func(libID string, task *service.TaskHandle, options service.ScrapeOptions) { - result, err := svc.Scraper.EnrichLibraryDetailedWithOptions(context.Background(), libID, options) - reclassified := 0 - if result.Processed > 0 { - reclassified = reclassifyLibraryAfterScrape(context.Background(), svc, libID) - } - metrics := map[string]int64{ - "matched": int64(result.Matched), - "processed": int64(result.Processed), - "candidates": int64(result.Candidates), - } - if reclassified > 0 { - metrics["reclassified"] = int64(reclassified) - } - if result.Failed > 0 { - metrics["errors"] = int64(result.Failed) - } - stage := "completed" - message := "手动刮削媒体库结束" - if err != nil { - stage = "scrape" - message = "手动刮削媒体库失败" - } - finishHTTPTask(task, err, stage, message, metrics, nil) - }(libID, task, options) - c.JSON(http.StatusAccepted, gin.H{"status": "scraping"}) + c.JSON(http.StatusOK, gin.H{"status": "queued", "enqueued": n}) } } diff --git a/internal/model/library_media.go b/internal/model/library_media.go index 0061aed..677325d 100644 --- a/internal/model/library_media.go +++ b/internal/model/library_media.go @@ -9,7 +9,7 @@ type Library struct { CoverURL string `gorm:"size:1024" json:"cover_url,omitempty"` Enabled bool `gorm:"default:true" json:"enabled"` SortOrder int `gorm:"index;default:0" json:"sort_order"` // 手动拖拽排序用,越小越靠前 - CarouselEnabled bool `gorm:"default:true" json:"carousel_enabled"` // 是否参与首页海报轮播 + CarouselEnabled bool `gorm:"default:false" json:"carousel_enabled"` // 是否参与首页海报轮播(默认不参与) Roots []LibraryRoot `gorm:"foreignKey:LibraryID" json:"roots,omitempty"` } diff --git a/internal/model/model.go b/internal/model/model.go index 6deee22..6d7d355 100644 --- a/internal/model/model.go +++ b/internal/model/model.go @@ -57,5 +57,6 @@ func AllModels() []interface{} { &StrmDownloadTask{}, &StrmUploadTask{}, &StrmDirCache{}, + &ScrapeTask{}, } } diff --git a/internal/model/scrape_task.go b/internal/model/scrape_task.go new file mode 100644 index 0000000..9df33e8 --- /dev/null +++ b/internal/model/scrape_task.go @@ -0,0 +1,34 @@ +package model + +import "time" + +const ( + ScrapeTaskPending = "pending" + ScrapeTaskRunning = "running" + ScrapeTaskDone = "done" + ScrapeTaskFailed = "failed" + ScrapeTaskCanceled = "canceled" +) + +// ScrapeTask 表示一条持久化的媒体刮削任务。 +type ScrapeTask struct { + Base + MediaID string `gorm:"index;size:36" json:"media_id"` + LibraryID string `gorm:"index;size:36" json:"library_id"` + LibraryName string `gorm:"size:128" json:"library_name"` + MediaTitle string `gorm:"size:255;not null" json:"media_title"` + MediaPath string `gorm:"size:1024;not null" json:"media_path"` + MediaType string `gorm:"size:16" json:"media_type"` // movie / tv / anime / adult + Provider string `gorm:"size:32" json:"provider"` // tmdb / douban / bangumi / thetvdb / metatube + MatchedTitle string `gorm:"size:255" json:"matched_title"` + MatchedYear int `json:"matched_year"` + PosterURL string `gorm:"size:1024" json:"poster_url"` + BackdropURL string `gorm:"size:1024" json:"backdrop_url"` + Status string `gorm:"index;size:16;default:pending" json:"status"` // pending / running / done / failed / canceled + Error string `gorm:"type:text" json:"error"` + RetryCount int `gorm:"default:0" json:"retry_count"` + EpisodeImages bool `gorm:"default:true" json:"episode_images"` + RefreshMatched bool `gorm:"default:false" json:"refresh_matched"` + StartedAt *time.Time `json:"started_at,omitempty"` + FinishedAt *time.Time `json:"finished_at,omitempty"` +} diff --git a/internal/repository/library_repository.go b/internal/repository/library_repository.go index 4c0c714..44a2311 100644 --- a/internal/repository/library_repository.go +++ b/internal/repository/library_repository.go @@ -15,6 +15,11 @@ type LibraryRepository struct{ db *gorm.DB } // Create persists a new library row. func (r *LibraryRepository) Create(ctx context.Context, l *model.Library) error { + if l != nil && l.SortOrder == 0 { + var maxSort int + _ = r.db.WithContext(ctx).Model(&model.Library{}).Select("COALESCE(MAX(sort_order), -1)").Scan(&maxSort) + l.SortOrder = maxSort + 1 + } return r.db.WithContext(ctx).Create(l).Error } @@ -23,6 +28,11 @@ func (r *LibraryRepository) CreateWithRoots(ctx context.Context, l *model.Librar return r.Create(ctx, l) } return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if l != nil && l.SortOrder == 0 { + var maxSort int + _ = tx.Model(&model.Library{}).Select("COALESCE(MAX(sort_order), -1)").Scan(&maxSort) + l.SortOrder = maxSort + 1 + } if err := tx.Create(l).Error; err != nil { return err } diff --git a/internal/repository/repository.go b/internal/repository/repository.go index 332acfa..2ebe721 100644 --- a/internal/repository/repository.go +++ b/internal/repository/repository.go @@ -30,9 +30,10 @@ type Container struct { StrmSyncPath *StrmSyncPathRepository StrmSyncRecord *StrmSyncRecordRepository StrmDownload *StrmDownloadTaskRepository - StrmUpload *StrmUploadTaskRepository - StrmDirCache *StrmDirCacheRepository -} + StrmUpload *StrmUploadTaskRepository + StrmDirCache *StrmDirCacheRepository + ScrapeTask *ScrapeTaskRepository + } // New 将每个 repository 连接到单个 *gorm.DB。 func New(db *gorm.DB) *Container { @@ -60,5 +61,6 @@ func New(db *gorm.DB) *Container { StrmDownload: &StrmDownloadTaskRepository{db: db}, StrmUpload: &StrmUploadTaskRepository{db: db}, StrmDirCache: &StrmDirCacheRepository{db: db}, + ScrapeTask: &ScrapeTaskRepository{db: db}, } } diff --git a/internal/repository/scrape_task_repository.go b/internal/repository/scrape_task_repository.go new file mode 100644 index 0000000..bf6a793 --- /dev/null +++ b/internal/repository/scrape_task_repository.go @@ -0,0 +1,274 @@ +package repository + +import ( + "context" + "errors" + "strings" + "sync" + "time" + + "gorm.io/gorm" + + "github.com/ShukeBta/MMTL/internal/model" +) + +var scrapeClaimMu sync.Mutex + +// ScrapeTaskRepository persists model.ScrapeTask. +type ScrapeTaskRepository struct{ db *gorm.DB } + +func (r *ScrapeTaskRepository) Create(ctx context.Context, t *model.ScrapeTask) error { + return withSQLiteBusyRetry(ctx, func() error { + return r.db.WithContext(ctx).Create(t).Error + }) +} + +func (r *ScrapeTaskRepository) CreateBatch(ctx context.Context, tasks []model.ScrapeTask) error { + if len(tasks) == 0 { + return nil + } + return withSQLiteBusyRetry(ctx, func() error { + return r.db.WithContext(ctx).CreateInBatches(tasks, 100).Error + }) +} + +func (r *ScrapeTaskRepository) FindByID(ctx context.Context, id string) (*model.ScrapeTask, error) { + var t model.ScrapeTask + err := r.db.WithContext(ctx).Where("id = ?", id).First(&t).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + return &t, err +} + +func (r *ScrapeTaskRepository) FindActiveByMediaID(ctx context.Context, mediaID string) (*model.ScrapeTask, error) { + var t model.ScrapeTask + err := r.db.WithContext(ctx). + Where("media_id = ? AND status IN ?", mediaID, []string{model.ScrapeTaskPending, model.ScrapeTaskRunning}). + First(&t).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + return &t, err +} + +func (r *ScrapeTaskRepository) List(ctx context.Context, status string, page, pageSize int) ([]model.ScrapeTask, int64, error) { + if page < 1 { + page = 1 + } + if pageSize < 1 || pageSize > 200 { + pageSize = 50 + } + q := r.db.WithContext(ctx).Model(&model.ScrapeTask{}) + if strings.TrimSpace(status) != "" && status != "all" { + q = q.Where("status = ?", strings.TrimSpace(status)) + } + var total int64 + if err := q.Count(&total).Error; err != nil { + return nil, 0, err + } + var rows []model.ScrapeTask + err := q.Order("created_at desc"). + Offset((page - 1) * pageSize). + Limit(pageSize). + Find(&rows).Error + return rows, total, err +} + +func (r *ScrapeTaskRepository) CountByStatus(ctx context.Context) (map[string]int64, error) { + var rows []struct { + Status string + Count int64 + } + err := r.db.WithContext(ctx).Model(&model.ScrapeTask{}). + Select("status, count(*) as count"). + Group("status").Scan(&rows).Error + if err != nil { + return nil, err + } + out := map[string]int64{} + for _, row := range rows { + out[row.Status] = row.Count + } + return out, nil +} + +// ClaimPending picks pending scrape tasks and marks them running. +func (r *ScrapeTaskRepository) ClaimPending(ctx context.Context, limit int) ([]model.ScrapeTask, error) { + scrapeClaimMu.Lock() + defer scrapeClaimMu.Unlock() + + var rows []model.ScrapeTask + err := withSQLiteBusyRetry(ctx, func() error { + return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if err := tx.Where("status = ?", model.ScrapeTaskPending). + Order("created_at asc").Limit(limit).Find(&rows).Error; err != nil { + return err + } + if len(rows) == 0 { + return nil + } + ids := make([]string, 0, len(rows)) + now := time.Now() + for i := range rows { + ids = append(ids, rows[i].ID) + rows[i].Status = model.ScrapeTaskRunning + rows[i].StartedAt = &now + } + return tx.Model(&model.ScrapeTask{}).Where("id IN ?", ids). + Updates(map[string]any{"status": model.ScrapeTaskRunning, "started_at": now}).Error + }) + }) + if err != nil { + return nil, err + } + return rows, nil +} + +func (r *ScrapeTaskRepository) Update(ctx context.Context, t *model.ScrapeTask) error { + return withSQLiteBusyRetry(ctx, func() error { + return r.db.WithContext(ctx).Model(&model.ScrapeTask{}).Where("id = ?", t.ID).Updates(map[string]any{ + "status": t.Status, + "error": t.Error, + "provider": t.Provider, + "matched_title": t.MatchedTitle, + "matched_year": t.MatchedYear, + "poster_url": t.PosterURL, + "backdrop_url": t.BackdropURL, + "retry_count": t.RetryCount, + "started_at": t.StartedAt, + "finished_at": t.FinishedAt, + "updated_at": time.Now(), + }).Error + }) +} + +func (r *ScrapeTaskRepository) Delete(ctx context.Context, id string) error { + return withSQLiteBusyRetry(ctx, func() error { + return r.db.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.ScrapeTask{}).Error + }) +} + +func (r *ScrapeTaskRepository) DeleteBatch(ctx context.Context, ids []string) (int64, error) { + if len(ids) == 0 { + return 0, nil + } + var count int64 + err := withSQLiteBusyRetry(ctx, func() error { + res := r.db.WithContext(ctx).Unscoped().Where("id IN ?", ids).Delete(&model.ScrapeTask{}) + count = res.RowsAffected + return res.Error + }) + return count, err +} + +func (r *ScrapeTaskRepository) RetryBatch(ctx context.Context, ids []string) (int64, error) { + if len(ids) == 0 { + return 0, nil + } + var count int64 + err := withSQLiteBusyRetry(ctx, func() error { + res := r.db.WithContext(ctx).Model(&model.ScrapeTask{}). + Where("id IN ? AND status IN ?", ids, []string{model.ScrapeTaskFailed, model.ScrapeTaskCanceled}). + Updates(map[string]any{ + "status": model.ScrapeTaskPending, + "error": "", + "retry_count": 0, + "started_at": nil, + "finished_at": nil, + "updated_at": time.Now(), + }) + count = res.RowsAffected + return res.Error + }) + return count, err +} + +func (r *ScrapeTaskRepository) CancelBatch(ctx context.Context, ids []string) (int64, error) { + if len(ids) == 0 { + return 0, nil + } + now := time.Now() + var count int64 + err := withSQLiteBusyRetry(ctx, func() error { + res := r.db.WithContext(ctx).Model(&model.ScrapeTask{}). + Where("id IN ? AND status IN ?", ids, []string{model.ScrapeTaskPending, model.ScrapeTaskRunning}). + Updates(map[string]any{ + "status": model.ScrapeTaskCanceled, + "error": "已批量取消", + "finished_at": now, + "updated_at": now, + }) + count = res.RowsAffected + return res.Error + }) + return count, err +} + +func (r *ScrapeTaskRepository) ClearDone(ctx context.Context) (int64, error) { + var count int64 + err := withSQLiteBusyRetry(ctx, func() error { + res := r.db.WithContext(ctx).Unscoped().Where("status = ?", model.ScrapeTaskDone).Delete(&model.ScrapeTask{}) + count = res.RowsAffected + return res.Error + }) + return count, err +} + +func (r *ScrapeTaskRepository) ClearFinished(ctx context.Context) (int64, error) { + var count int64 + err := withSQLiteBusyRetry(ctx, func() error { + res := r.db.WithContext(ctx).Unscoped().Where("status IN ?", []string{model.ScrapeTaskDone, model.ScrapeTaskFailed, model.ScrapeTaskCanceled}). + Delete(&model.ScrapeTask{}) + count = res.RowsAffected + return res.Error + }) + return count, err +} + +func (r *ScrapeTaskRepository) ClearCanceled(ctx context.Context) (int64, error) { + var count int64 + err := withSQLiteBusyRetry(ctx, func() error { + res := r.db.WithContext(ctx).Unscoped().Where("status = ?", model.ScrapeTaskCanceled).Delete(&model.ScrapeTask{}) + count = res.RowsAffected + return res.Error + }) + return count, err +} + +func (r *ScrapeTaskRepository) RetryAllFailed(ctx context.Context) (int64, error) { + var count int64 + err := withSQLiteBusyRetry(ctx, func() error { + res := r.db.WithContext(ctx).Model(&model.ScrapeTask{}). + Where("status = ?", model.ScrapeTaskFailed). + Updates(map[string]any{ + "status": model.ScrapeTaskPending, + "error": "", + "retry_count": 0, + "started_at": nil, + "finished_at": nil, + "updated_at": time.Now(), + }) + count = res.RowsAffected + return res.Error + }) + return count, err +} + +func (r *ScrapeTaskRepository) CancelPending(ctx context.Context) (int64, error) { + now := time.Now() + var count int64 + err := withSQLiteBusyRetry(ctx, func() error { + res := r.db.WithContext(ctx).Model(&model.ScrapeTask{}). + Where("status IN ?", []string{model.ScrapeTaskPending, model.ScrapeTaskRunning}). + Updates(map[string]any{ + "status": model.ScrapeTaskCanceled, + "error": "已批量取消", + "finished_at": now, + "updated_at": now, + }) + count = res.RowsAffected + return res.Error + }) + return count, err +} diff --git a/internal/service/manual_scrape_providers.go b/internal/service/manual_scrape_providers.go index ec66883..e15294c 100644 --- a/internal/service/manual_scrape_providers.go +++ b/internal/service/manual_scrape_providers.go @@ -32,18 +32,21 @@ func (s *ScraperService) manualTMDbCandidates(ctx context.Context, query string, for _, typ := range manualTMDbSearchTypes(mediaType) { switch typ { case "movie": - if matches, err := s.tmdb.SearchMovieCandidates(ctx, query, year); err == nil { + if matches, err := s.tmdb.SearchMovieCandidates(ctx, query, year); err == nil && len(matches) > 0 { for _, match := range matches { out = append(out, manualTMDbCandidate{MediaType: "movie", Match: match}) } } case "tv": - if matches, err := s.tmdb.SearchTVCandidates(ctx, query, year); err == nil { + if matches, err := s.tmdb.SearchTVCandidates(ctx, query, year); err == nil && len(matches) > 0 { for _, match := range matches { out = append(out, manualTMDbCandidate{MediaType: "tv", Match: match}) } } } + if len(out) > 0 { + break + } } return out } @@ -70,15 +73,15 @@ func manualTMDbIDSearchTypes(mediaType string) []string { func manualTMDbSearchTypes(mediaType string) []string { if strings.TrimSpace(mediaType) == "" { - return []string{"movie", "tv"} + return []string{"tv", "movie"} } switch normalizeMediaType(mediaType, "", "") { case "tv", "anime", "variety": return []string{"tv", "movie"} case "movie", "adult": - return []string{"movie"} - default: return []string{"movie", "tv"} + default: + return []string{"tv", "movie"} } } diff --git a/internal/service/manual_scrape_test.go b/internal/service/manual_scrape_test.go index 8dae327..ca3326d 100644 --- a/internal/service/manual_scrape_test.go +++ b/internal/service/manual_scrape_test.go @@ -179,9 +179,9 @@ func TestManualSearchFallsBackToMovieFolderForGenericQuery(t *testing.T) { if len(results) != 1 || results[0].TMDbID != 27205 { t.Fatalf("manual search results=%#v, want folder fallback candidate; queries=%v", results, queries) } - if len(queries) < 2 || queries[0] != "00000" || queries[1] != "inception" { - t.Fatalf("manual search queries=%v, want explicit query then folder fallback", queries) - } + if len(queries) < 2 || queries[0] != "00000" || queries[len(queries)-1] != "inception" { + t.Fatalf("manual search queries=%v, want explicit query then folder fallback", queries) + } } func TestManualSearchReturnsMovieFallbackForTVTypedTMDbSearch(t *testing.T) { diff --git a/internal/service/media_library_roots.go b/internal/service/media_library_roots.go index 7493456..33b6d90 100644 --- a/internal/service/media_library_roots.go +++ b/internal/service/media_library_roots.go @@ -52,7 +52,14 @@ func (s *MediaService) CreateLibraryWithRootsAndCover(ctx context.Context, name, s.invalidateMediaCache(ctx) return lib, nil } - lib := &model.Library{Name: strings.TrimSpace(name), Path: roots[0].Path, Type: kind, CoverURL: strings.TrimSpace(coverURL), Enabled: true} + lib := &model.Library{ + Name: strings.TrimSpace(name), + Path: roots[0].Path, + Type: kind, + CoverURL: strings.TrimSpace(coverURL), + Enabled: true, + CarouselEnabled: false, + } if err := s.repo.Library.CreateWithRoots(ctx, lib, roots); err != nil { return nil, err } diff --git a/internal/service/scraper_queue.go b/internal/service/scraper_queue.go new file mode 100644 index 0000000..75a5604 --- /dev/null +++ b/internal/service/scraper_queue.go @@ -0,0 +1,354 @@ +package service + +import ( + "context" + "errors" + "fmt" + "strings" + "sync" + "time" + + "go.uber.org/zap" + + "github.com/ShukeBta/MMTL/internal/model" +) + +type ScrapeQueueCounts struct { + Pending int64 `json:"pending"` + Running int64 `json:"running"` + Done int64 `json:"done"` + Failed int64 `json:"failed"` + Canceled int64 `json:"canceled"` +} + +type ScrapeQueueSnapshot struct { + Counts ScrapeQueueCounts `json:"counts"` + Tasks []model.ScrapeTask `json:"tasks"` + Total int64 `json:"total"` + Page int `json:"page"` + PageSize int `json:"page_size"` +} + +// Start 启动刮削任务队列的后台消费者。 +func (s *ScraperService) Start(ctx context.Context) { + if s == nil { + return + } + go s.queueWorker(ctx) +} + +func (s *ScraperService) queueWorker(ctx context.Context) { + const claimBatch = 4 + sem := make(chan struct{}, 2) // 最大并发刮削数:2 + + for { + select { + case <-ctx.Done(): + return + default: + } + + tasks, err := s.repo.ScrapeTask.ClaimPending(ctx, claimBatch) + if err != nil { + if s.log != nil { + s.log.Warn("claim pending scrape task failed", zap.Error(err)) + } + sleepContext(ctx, 3*time.Second) + continue + } + if len(tasks) == 0 { + sleepContext(ctx, 2*time.Second) + continue + } + + var wg sync.WaitGroup + for i := range tasks { + wg.Add(1) + go func(t *model.ScrapeTask) { + defer wg.Done() + select { + case <-ctx.Done(): + return + case sem <- struct{}{}: + } + defer func() { <-sem }() + + s.processScrapeTask(ctx, t) + }(&tasks[i]) + } + wg.Wait() + } +} + +func (s *ScraperService) processScrapeTask(ctx context.Context, task *model.ScrapeTask) { + media, err := s.repo.Media.FindByID(ctx, task.MediaID) + if err != nil || media == nil { + now := time.Now() + task.Status = model.ScrapeTaskFailed + task.Error = "媒体项已不存在或被删除" + task.FinishedAt = &now + _ = s.repo.ScrapeTask.Update(ctx, task) + return + } + + epArtwork := task.EpisodeImages + options := ScrapeOptions{ + EpisodeArtwork: &epArtwork, + IncludeMatched: task.RefreshMatched, + RetryNoMatch: true, + } + + enrichErr := s.EnrichOneWithOptions(ctx, media, options) + now := time.Now() + task.FinishedAt = &now + + refreshed, _ := s.repo.Media.FindByID(ctx, media.ID) + if refreshed != nil && refreshed.ScrapeStatus == "matched" { + task.Status = model.ScrapeTaskDone + task.Error = "" + task.MatchedTitle = refreshed.Title + task.MatchedYear = refreshed.Year + task.PosterURL = refreshed.PosterURL + task.BackdropURL = refreshed.BackdropURL + if refreshed.TMDbID > 0 { + task.Provider = "tmdb" + } else if strings.TrimSpace(refreshed.DoubanID) != "" { + task.Provider = "douban" + } else if refreshed.BangumiID > 0 { + task.Provider = "bangumi" + } else if strings.TrimSpace(refreshed.TheTVDBID) != "" { + task.Provider = "thetvdb" + } else { + task.Provider = "metatube" + } + } else { + task.Status = model.ScrapeTaskFailed + if enrichErr != nil { + task.Error = enrichErr.Error() + } else if refreshed != nil && refreshed.ScrapeStatus == "no_match" { + task.Error = "未搜索到匹配的元数据" + } else { + task.Error = "刮削未完成匹配" + } + } + + _ = s.repo.ScrapeTask.Update(ctx, task) + + if s.hub != nil { + s.hub.Publish("scraper_queue", map[string]any{ + "task_id": task.ID, + "status": task.Status, + "title": task.MediaTitle, + }) + } +} + +func (s *ScraperService) mediaKind(m *model.Media, lib *model.Library) string { + if m == nil { + return "" + } + if lib != nil && lib.Type != "" { + return lib.Type + } + if mediaIsEpisodic(m, lib) { + return "tv" + } + return "movie" +} + +// EnqueueMedia 把单个媒体项放入刮削队列。 +func (s *ScraperService) EnqueueMedia(ctx context.Context, mediaID string, options ScrapeOptions) (*model.ScrapeTask, error) { + if s == nil || s.repo == nil { + return nil, errors.New("scraper service not initialized") + } + media, err := s.repo.Media.FindByID(ctx, mediaID) + if err != nil || media == nil { + return nil, errors.New("media not found") + } + + if active, _ := s.repo.ScrapeTask.FindActiveByMediaID(ctx, mediaID); active != nil { + return active, nil + } + + libName := "" + var lib *model.Library + if strings.TrimSpace(media.LibraryID) != "" { + lib, _ = s.repo.Library.FindByID(ctx, media.LibraryID) + if lib != nil { + libName = lib.Name + } + } + + task := &model.ScrapeTask{ + MediaID: media.ID, + LibraryID: media.LibraryID, + LibraryName: libName, + MediaTitle: media.Title, + MediaPath: media.Path, + MediaType: s.mediaKind(media, lib), + Status: model.ScrapeTaskPending, + EpisodeImages: options.episodeArtworkEnabled(), + RefreshMatched: options.IncludeMatched || options.RefreshWeakMatched, + } + + if err := s.repo.ScrapeTask.Create(ctx, task); err != nil { + return nil, err + } + return task, nil +} + +// EnqueueLibrary 把指定媒体库内的所有候选媒体批量推入刮削队列。 +func (s *ScraperService) EnqueueLibrary(ctx context.Context, libraryID string, options ScrapeOptions) (int, error) { + if s == nil || s.repo == nil { + return 0, errors.New("scraper service not initialized") + } + + lib, err := s.repo.Library.FindByID(ctx, libraryID) + if err != nil || lib == nil { + return 0, errors.New("library not found") + } + + rows, err := s.scrapeCandidateRows(ctx, libraryID, options) + if err != nil { + return 0, err + } + if len(rows) == 0 { + return 0, nil + } + + tasks := make([]model.ScrapeTask, 0, len(rows)) + for _, m := range rows { + tasks = append(tasks, model.ScrapeTask{ + MediaID: m.ID, + LibraryID: lib.ID, + LibraryName: lib.Name, + MediaTitle: m.Title, + MediaPath: m.Path, + MediaType: s.mediaKind(&m, lib), + Status: model.ScrapeTaskPending, + EpisodeImages: options.episodeArtworkEnabled(), + RefreshMatched: options.IncludeMatched || options.RefreshWeakMatched, + }) + } + + if err := s.repo.ScrapeTask.CreateBatch(ctx, tasks); err != nil { + return 0, err + } + return len(tasks), nil +} + +// EnqueueAll 把所有已启用媒体库的媒体推入刮削队列。 +func (s *ScraperService) EnqueueAll(ctx context.Context, options ScrapeOptions) (int, error) { + libs, err := s.repo.Library.List(ctx) + if err != nil { + return 0, err + } + total := 0 + for _, lib := range libs { + if !lib.Enabled { + continue + } + n, err := s.EnqueueLibrary(ctx, lib.ID, options) + if err != nil { + if s.log != nil { + s.log.Warn("enqueue library for scrape failed", zap.String("library", lib.ID), zap.Error(err)) + } + continue + } + total += n + } + return total, nil +} + +func (s *ScraperService) ScrapeQueueSnapshot(ctx context.Context, status string, page, pageSize int) (*ScrapeQueueSnapshot, error) { + tasks, total, err := s.repo.ScrapeTask.List(ctx, status, page, pageSize) + if err != nil { + return nil, err + } + countsMap, err := s.repo.ScrapeTask.CountByStatus(ctx) + if err != nil { + return nil, err + } + snap := &ScrapeQueueSnapshot{ + Counts: ScrapeQueueCounts{ + Pending: countsMap[model.ScrapeTaskPending], + Running: countsMap[model.ScrapeTaskRunning], + Done: countsMap[model.ScrapeTaskDone], + Failed: countsMap[model.ScrapeTaskFailed], + Canceled: countsMap[model.ScrapeTaskCanceled], + }, + Tasks: tasks, + Total: total, + Page: page, + PageSize: pageSize, + } + return snap, nil +} + +func (s *ScraperService) CancelScrapeTask(ctx context.Context, id string) error { + task, err := s.repo.ScrapeTask.FindByID(ctx, id) + if err != nil || task == nil { + return errors.New("刮削任务不存在") + } + if task.Status != model.ScrapeTaskPending && task.Status != model.ScrapeTaskRunning { + return errors.New("任务已完成或已终止,无法取消") + } + now := time.Now() + task.Status = model.ScrapeTaskCanceled + task.Error = "已取消" + task.FinishedAt = &now + return s.repo.ScrapeTask.Update(ctx, task) +} + +func (s *ScraperService) RetryScrapeTask(ctx context.Context, id string) error { + task, err := s.repo.ScrapeTask.FindByID(ctx, id) + if err != nil || task == nil { + return errors.New("刮削任务不存在") + } + if task.Status != model.ScrapeTaskFailed && task.Status != model.ScrapeTaskCanceled { + return errors.New("只有失败或已取消的任务可以重试") + } + task.Status = model.ScrapeTaskPending + task.Error = "" + task.RetryCount = 0 + task.StartedAt = nil + task.FinishedAt = nil + return s.repo.ScrapeTask.Update(ctx, task) +} + +func (s *ScraperService) DeleteScrapeTask(ctx context.Context, id string) error { + return s.repo.ScrapeTask.Delete(ctx, id) +} + +func (s *ScraperService) BatchActionScrapeTasks(ctx context.Context, action string, ids []string) (int64, error) { + switch action { + case "delete": + return s.repo.ScrapeTask.DeleteBatch(ctx, ids) + case "retry": + return s.repo.ScrapeTask.RetryBatch(ctx, ids) + case "cancel": + return s.repo.ScrapeTask.CancelBatch(ctx, ids) + default: + return 0, fmt.Errorf("不支持的操作: %s", action) + } +} + +func (s *ScraperService) ClearDoneScrapeTasks(ctx context.Context) (int64, error) { + return s.repo.ScrapeTask.ClearDone(ctx) +} + +func (s *ScraperService) ClearFinishedScrapeTasks(ctx context.Context) (int64, error) { + return s.repo.ScrapeTask.ClearFinished(ctx) +} + +func (s *ScraperService) ClearCanceledScrapeTasks(ctx context.Context) (int64, error) { + return s.repo.ScrapeTask.ClearCanceled(ctx) +} + +func (s *ScraperService) RetryAllFailedScrapeTasks(ctx context.Context) (int64, error) { + return s.repo.ScrapeTask.RetryAllFailed(ctx) +} + +func (s *ScraperService) CancelPendingScrapeTasks(ctx context.Context) (int64, error) { + return s.repo.ScrapeTask.CancelPending(ctx) +} diff --git a/internal/service/service.go b/internal/service/service.go index e528545..caf3381 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -102,7 +102,12 @@ func (c *Container) Boot() { c.Strm.Start(c.stopCtx) } - // Mgo 保号规则巡检:默认关闭,由管理员通过 Telegram Bot 命令开启。 + // 启动刮削队列后台消费者 + if c.Scraper != nil { + c.Scraper.Start(c.stopCtx) + } + + // Mgo 保号规则巡检:默认关闭,由管理员通过 Telegram Bot 命令开启。 // 每天触发一次评估;规则里的窗口可随机,不固定。 if c.Device != nil { go c.runInactivitySweeper(c.stopCtx) diff --git a/internal/service/tmdb_search.go b/internal/service/tmdb_search.go index 44e19dd..42de1fa 100644 --- a/internal/service/tmdb_search.go +++ b/internal/service/tmdb_search.go @@ -57,7 +57,7 @@ func (t *TMDbProvider) searchMovieCandidates(ctx context.Context, query string, apiKey := t.resolveAPIKey(ctx) if apiKey == "" { - return nil, nil + return nil, errors.New("TMDb API Key 未配置,请先在「系统设置 → API配置」中填写") } base := t.resolveBaseURL(ctx) @@ -65,7 +65,7 @@ func (t *TMDbProvider) searchMovieCandidates(ctx context.Context, query string, q.Set("api_key", apiKey) q.Set("query", query) q.Set("language", language) - q.Set("include_adult", "false") + q.Set("include_adult", "true") if year > 0 { q.Set("year", fmt.Sprintf("%d", year)) } @@ -79,6 +79,12 @@ func (t *TMDbProvider) searchMovieCandidates(ctx context.Context, query string, if err := t.getJSON(ctx, u, &p); err != nil { return nil, err } + if len(p.Results) == 0 && year > 0 { + // 移除年份限制重试一次,避免年份微小差异(如 2019 vs 2020)导致无搜索结果 + q.Del("year") + u = base + "/search/movie?" + q.Encode() + _ = t.getJSON(ctx, u, &p) + } if len(p.Results) == 0 { return nil, nil } @@ -137,7 +143,7 @@ func (t *TMDbProvider) searchTVCandidates(ctx context.Context, query string, yea apiKey := t.resolveAPIKey(ctx) if apiKey == "" { - return nil, nil + return nil, errors.New("TMDb API Key 未配置,请先在「系统设置 → API配置」中填写") } base := t.resolveBaseURL(ctx) @@ -145,7 +151,7 @@ func (t *TMDbProvider) searchTVCandidates(ctx context.Context, query string, yea q.Set("api_key", apiKey) q.Set("query", query) q.Set("language", language) - q.Set("include_adult", "false") + q.Set("include_adult", "true") if year > 0 { q.Set("first_air_date_year", fmt.Sprintf("%d", year)) } @@ -159,6 +165,12 @@ func (t *TMDbProvider) searchTVCandidates(ctx context.Context, query string, yea if err := t.getJSON(ctx, u, &p); err != nil { return nil, err } + if len(p.Results) == 0 && year > 0 { + // 移除年份限制重试一次,避免年份微小差异导致无搜索结果 + q.Del("first_air_date_year") + u = base + "/search/tv?" + q.Encode() + _ = t.getJSON(ctx, u, &p) + } if len(p.Results) == 0 { return nil, nil } diff --git a/web/src/api/scraper.ts b/web/src/api/scraper.ts new file mode 100644 index 0000000..5fd4108 --- /dev/null +++ b/web/src/api/scraper.ts @@ -0,0 +1,57 @@ +import { api } from './client' +import type { ScrapeQueueSnapshot } from '../types/scraper' + +export interface EnqueueScrapeOptions { + episode_images?: boolean + episode_artwork?: boolean + refresh_matched?: boolean + include_matched?: boolean +} + +export const scraperAPI = { + queue: (status?: string, page = 1, pageSize = 50) => + api + .get('/admin/scraper/queue', { + params: { status, page, page_size: pageSize }, + }) + .then((r) => r.data), + + cancelTask: (id: string) => + api.post(`/admin/scraper/queue/${id}/cancel`).then((r) => r.data), + + retryTask: (id: string) => + api.post(`/admin/scraper/queue/${id}/retry`).then((r) => r.data), + + deleteTask: (id: string) => + api.delete(`/admin/scraper/queue/${id}`).then((r) => r.data), + + batchAction: (action: 'delete' | 'retry' | 'cancel', ids: string[]) => + api + .post<{ affected: number; action: string }>('/admin/scraper/queue/batch', { action, ids }) + .then((r) => r.data), + + clearDone: () => + api.post<{ deleted: number }>('/admin/scraper/queue/clear-done').then((r) => r.data), + + clearFinished: () => + api.post<{ deleted: number }>('/admin/scraper/queue/clear-finished').then((r) => r.data), + + clearCanceled: () => + api.post<{ deleted: number }>('/admin/scraper/queue/clear-canceled').then((r) => r.data), + + retryFailed: () => + api.post<{ retried: number }>('/admin/scraper/queue/retry-failed').then((r) => r.data), + + cancelPending: () => + api.post<{ canceled: number }>('/admin/scraper/queue/cancel-pending').then((r) => r.data), + + enqueueLibrary: (libraryId: string, options?: EnqueueScrapeOptions) => + api + .post<{ enqueued: number }>(`/admin/scraper/queue/enqueue-library/${libraryId}`, options ?? {}) + .then((r) => r.data), + + enqueueAll: (options?: EnqueueScrapeOptions) => + api + .post<{ enqueued: number }>('/admin/scraper/queue/enqueue-all', options ?? {}) + .then((r) => r.data), +} diff --git a/web/src/appRoutes.tsx b/web/src/appRoutes.tsx index 67a14c0..ff1fe8d 100644 --- a/web/src/appRoutes.tsx +++ b/web/src/appRoutes.tsx @@ -33,6 +33,9 @@ const StrmDownloadQueuePage = lazy(() => const StrmUploadQueuePage = lazy(() => import('./pages/StrmQueuePage').then((m) => ({ default: m.StrmUploadQueuePage })), ) +const ScraperQueuePage = lazy(() => + import('./pages/ScraperQueuePage').then((m) => ({ default: m.ScraperQueuePage })), +) export type AppRoute = { path?: string @@ -62,5 +65,6 @@ export const appRoutes: AppRoute[] = [ { path: 'strm', element: , adminOnly: true }, { path: 'strm/downloads', element: , adminOnly: true }, { path: 'strm/uploads', element: , adminOnly: true }, + { path: 'scraper/queue', element: , adminOnly: true }, { path: 'admin', element: , adminOnly: true }, ] diff --git a/web/src/components/layoutNavigation.ts b/web/src/components/layoutNavigation.ts index bd787f0..bc3c711 100644 --- a/web/src/components/layoutNavigation.ts +++ b/web/src/components/layoutNavigation.ts @@ -5,6 +5,7 @@ import { FolderOpen, Library, Settings, + Sparkles, Upload, User, Users, @@ -22,6 +23,7 @@ export type LayoutNavItem = { export const LAYOUT_NAV_ITEMS: LayoutNavItem[] = [ { to: '/profile', label: '个人资料', icon: User }, { to: '/libraries?from=admin', label: '媒体库', icon: Library }, + { to: '/scraper/queue', label: '刮削队列', icon: Sparkles, adminOnly: true }, { to: '/admin', label: '用户管理', icon: Users, adminOnly: true }, { to: '/files', label: '文件管理', icon: FolderOpen, adminOnly: true }, { to: '/settings', label: '系统设置', icon: Settings, adminOnly: true }, diff --git a/web/src/pages/AdminLibraryTable.tsx b/web/src/pages/AdminLibraryTable.tsx index 8e2e91b..4a44262 100644 --- a/web/src/pages/AdminLibraryTable.tsx +++ b/web/src/pages/AdminLibraryTable.tsx @@ -188,7 +188,7 @@ function LibraryTableRow({ library, dragging, dragOver, onDragStart, onDragOver, } function CarouselToggle({ library, onToggleCarousel }: { library: Library; onToggleCarousel: (library: Library) => void }) { - const on = library.carousel_enabled ?? true + const on = Boolean(library.carousel_enabled) return ( + + + 刮削队列 + diff --git a/web/src/pages/LibraryPage.tsx b/web/src/pages/LibraryPage.tsx index 917e5e9..beca73a 100644 --- a/web/src/pages/LibraryPage.tsx +++ b/web/src/pages/LibraryPage.tsx @@ -113,21 +113,23 @@ export function LibraryPage() { return (
- setScrapeDialogOpen(true)} - onRepairRescrape={handleRepairRescrape} - /> + {!selectedSeries && ( + setScrapeDialogOpen(true)} + onRepairRescrape={handleRepairRescrape} + /> + )} -
)} diff --git a/web/src/pages/ScraperQueuePage.tsx b/web/src/pages/ScraperQueuePage.tsx new file mode 100644 index 0000000..f6902a0 --- /dev/null +++ b/web/src/pages/ScraperQueuePage.tsx @@ -0,0 +1,972 @@ +import { useCallback, useEffect, useMemo, useState, type ReactNode } from 'react' +import { Link } from 'react-router-dom' +import toast from 'react-hot-toast' +import { + AlertCircle, + Ban, + CheckCircle2, + Clock, + Copy, + ExternalLink, + Eye, + Film, + Image as ImageIcon, + Layers, + Loader2, + PlayCircle, + RefreshCw, + Search, + Sparkles, + Trash2, + Tv, + X, +} from 'lucide-react' + +import { imageURL } from '../api/client' +import { scraperAPI } from '../api/scraper' +import type { ScrapeQueueSnapshot, ScrapeTask, ScrapeTaskStatus } from '../types/scraper' +import { apiErrorMessage, formatTime, taskStatusMeta } from './StrmManagePage' + +const FILTERS: { key: 'all' | ScrapeTaskStatus; label: string; icon: typeof Clock; color: string }[] = [ + { key: 'all', label: '全部', icon: Sparkles, color: 'text-ink-600' }, + { key: 'pending', label: '排队中', icon: Clock, color: 'text-gray-500' }, + { key: 'running', label: '刮削中', icon: PlayCircle, color: 'text-brand-500' }, + { key: 'done', label: '已匹配', icon: CheckCircle2, color: 'text-emerald-500' }, + { key: 'failed', label: '未匹配/失败', icon: AlertCircle, color: 'text-rose-500' }, + { key: 'canceled', label: '已取消', icon: Ban, color: 'text-amber-500' }, +] + +const PROVIDER_LABELS: Record = { + tmdb: 'TheMovieDB', + douban: '豆瓣 Douban', + bangumi: 'Bangumi 番组计划', + thetvdb: 'TheTVDB', + metatube: 'MetaTube', +} + +const TYPE_ICONS: Record = { + movie: , + tv: , + anime: , + adult: , +} + +const TYPE_LABELS: Record = { + movie: '电影', + tv: '剧集', + anime: '动漫', + adult: 'Adult', +} + +const PAGE_SIZE = 50 + +export function ScraperQueuePage() { + const [snapshot, setSnapshot] = useState(null) + const [filter, setFilter] = useState<'all' | ScrapeTaskStatus>('all') + const [search, setSearch] = useState('') + const [page, setPage] = useState(1) + const [totalPages, setTotalPages] = useState(1) + const [loading, setLoading] = useState(true) + const [isRefreshing, setIsRefreshing] = useState(false) + const [autoRefresh, setAutoRefresh] = useState(true) + const [batchBusy, setBatchBusy] = useState(false) + const [selectedIds, setSelectedIds] = useState>(new Set()) + const [detailTask, setDetailTask] = useState(null) + + const refresh = useCallback( + async (showLoading = false) => { + if (showLoading) setIsRefreshing(true) + try { + const status = filter === 'all' ? undefined : filter + const data = await scraperAPI.queue(status, page, PAGE_SIZE) + const tp = Math.max(1, Math.ceil((data.total ?? data.tasks.length) / PAGE_SIZE)) + if (page > tp) { + setPage(tp) + return + } + setTotalPages(tp) + setSnapshot(data) + } catch { + /* keep existing */ + } finally { + setLoading(false) + if (showLoading) setIsRefreshing(false) + } + }, + [filter, page], + ) + + useEffect(() => { + refresh().catch(() => undefined) + }, [refresh]) + + useEffect(() => { + if (!autoRefresh) return + const timer = setInterval(() => { + refresh().catch(() => undefined) + }, 3000) + return () => clearInterval(timer) + }, [autoRefresh, refresh]) + + useEffect(() => { + setSelectedIds(new Set()) + }, [filter, page]) + + const copyText = (text: string, label: string) => { + navigator.clipboard.writeText(text) + toast.success(`已复制${label}`) + } + + // Row actions + const cancelTask = async (task: ScrapeTask) => { + try { + await scraperAPI.cancelTask(task.id) + toast.success('已取消刮削任务') + await refresh() + } catch (err) { + toast.error(apiErrorMessage(err)) + } + } + + const retryTask = async (task: ScrapeTask) => { + try { + await scraperAPI.retryTask(task.id) + toast.success('已重新推入刮削队列') + await refresh() + } catch (err) { + toast.error(apiErrorMessage(err)) + } + } + + const deleteTask = async (task: ScrapeTask) => { + try { + await scraperAPI.deleteTask(task.id) + toast.success('已删除刮削记录') + if (detailTask?.id === task.id) setDetailTask(null) + await refresh() + } catch (err) { + toast.error(apiErrorMessage(err)) + } + } + + // Batch actions + const runSelectedBatch = async (action: 'retry' | 'cancel' | 'delete') => { + const ids = Array.from(selectedIds) + if (ids.length === 0) return + + const actionText = action === 'retry' ? '重新刮削' : action === 'cancel' ? '取消' : '删除' + if (action === 'delete' && !window.confirm(`确定删除选中的 ${ids.length} 条刮削记录?`)) return + if (action === 'cancel' && !window.confirm(`确定取消选中的 ${ids.length} 个刮削任务?`)) return + + setBatchBusy(true) + try { + const res = await scraperAPI.batchAction(action, ids) + toast.success(`已成功${actionText} ${res.affected} 项`) + setSelectedIds(new Set()) + await refresh() + } catch (err) { + toast.error(apiErrorMessage(err)) + } finally { + setBatchBusy(false) + } + } + + const runGlobalBatch = async ( + action: () => Promise<{ deleted?: number; retried?: number; canceled?: number }>, + confirmMsg?: string, + ) => { + if (confirmMsg && !window.confirm(confirmMsg)) return + setBatchBusy(true) + try { + const res = await action() + if (res.deleted !== undefined) toast.success(`已清空 ${res.deleted} 条记录`) + else if (res.retried !== undefined) toast.success(`已重新入队 ${res.retried} 个任务`) + else if (res.canceled !== undefined) toast.success(`已取消 ${res.canceled} 个任务`) + await refresh() + } catch (err) { + toast.error(apiErrorMessage(err)) + } finally { + setBatchBusy(false) + } + } + + // Enqueue all libraries + const handleEnqueueAll = async () => { + if (!window.confirm('确定将全库所有未匹配或需要更新的媒体重新推入刮削队列?')) return + setBatchBusy(true) + try { + const res = await scraperAPI.enqueueAll({ include_matched: false, refresh_matched: false, episode_images: true }) + toast.success(`已将 ${res.enqueued} 个媒体项推入刮削队列`) + await refresh() + } catch (err) { + toast.error(apiErrorMessage(err)) + } finally { + setBatchBusy(false) + } + } + + // Filter and search + const tasks = snapshot?.tasks ?? [] + const filteredTasks = useMemo(() => { + let list = tasks + if (filter !== 'all') { + list = list.filter((t) => t.status === filter) + } + if (search.trim()) { + const q = search.trim().toLowerCase() + list = list.filter( + (t) => + t.media_title.toLowerCase().includes(q) || + t.matched_title.toLowerCase().includes(q) || + t.library_name.toLowerCase().includes(q) || + t.media_path.toLowerCase().includes(q) || + (t.error && t.error.toLowerCase().includes(q)), + ) + } + return list + }, [tasks, filter, search]) + + const counts = snapshot?.counts + const activeTaskCount = (counts?.pending ?? 0) + (counts?.running ?? 0) + const failedCount = counts?.failed ?? 0 + const allCurrentChecked = + filteredTasks.length > 0 && filteredTasks.every((t) => selectedIds.has(t.id)) + + const toggleSelectAll = () => { + if (allCurrentChecked) { + setSelectedIds(new Set()) + } else { + setSelectedIds(new Set(filteredTasks.map((t) => t.id))) + } + } + + const toggleSelectRow = (id: string) => { + setSelectedIds((prev) => { + const next = new Set(prev) + if (next.has(id)) next.delete(id) + else next.add(id) + return next + }) + } + + return ( +
+ {/* 1. Header */} +
+
+
+ +
+
+
+

刮削队列

+ {autoRefresh && ( + + + 实时同步 + + )} +
+

+ 媒体元数据在线识别与海报/剧照下载进度(TMDb / 豆瓣 / Bangumi / TheTVDB) +

+
+
+ +
+ + + + + + +
+ + + 批量操作 + +
+ {failedCount > 0 && ( + + )} + {activeTaskCount > 0 && ( + + )} +
+ + + +
+
+
+
+ + {/* 2. Status Cards */} +
+ {FILTERS.map((item) => { + const count = + item.key === 'all' + ? (counts?.pending ?? 0) + + (counts?.running ?? 0) + + (counts?.done ?? 0) + + (counts?.failed ?? 0) + + (counts?.canceled ?? 0) + : counts?.[item.key] ?? 0 + const isActive = filter === item.key + const ItemIcon = item.icon + + return ( + + ) + })} +
+ + {/* 3. Search & Batch Actions */} +
+
+ + setSearch(e.target.value)} + placeholder="搜索媒体标题、匹配结果、媒体库或错误信息…" + className="h-9 w-full rounded-xl border border-gray-200 bg-white pl-9 pr-8 text-xs text-ink-600 placeholder:text-gray-400 outline-none transition focus:border-brand-500 focus:ring-2 focus:ring-brand-100/60" + /> + {search && ( + + )} +
+ + {selectedIds.size > 0 && ( +
+ 已选中 {selectedIds.size} 项 +
+ + + + +
+ )} +
+ + {/* 4. Table */} +
+ {loading ? ( +
+ +
+ ) : filteredTasks.length === 0 ? ( +
+ {search + ? '没有找到符合搜索条件的刮削任务' + : filter === 'all' + ? '刮削队列为空,暂无进行或排队中的任务' + : `「${FILTERS.find((f) => f.key === filter)?.label}」状态下暂无任务`} +
+ ) : ( +
+ + + + + + + + + + + + + + + {filteredTasks.map((task) => { + const status = taskStatusMeta(task.status) + const isSelected = selectedIds.has(task.id) + + return ( + + + + {/* Media title & path */} + + + {/* Library */} + + + {/* Scraped matched result */} + + + {/* Provider */} + + + {/* Status & Error */} + + + {/* Time */} + + + {/* Actions */} + + + ) + })} + +
+ + 媒体文件所属媒体库刮削匹配结果识别源状态时间操作
+ toggleSelectRow(task.id)} + className="h-3.5 w-3.5 rounded border-gray-300 text-brand-500 focus:ring-brand-400 cursor-pointer" + /> + +
+ {TYPE_ICONS[task.media_type] || } +
+ setDetailTask(task)} + className="cursor-pointer truncate font-medium text-ink-600 hover:text-brand-500 hover:underline block" + title={task.media_title} + > + {task.media_title} + + + {task.media_path} + +
+
+
+ + {task.library_name || '媒体库'} + + + {task.matched_title ? ( +
+ {task.poster_url ? ( + { + e.currentTarget.style.display = 'none' + }} + /> + ) : ( +
+ +
+ )} +
+ + {task.matched_title} + + {task.matched_year > 0 && ( + + {task.matched_year} 年 + + )} +
+
+ ) : ( + + {task.status === 'pending' || task.status === 'running' + ? '等待识别…' + : '未匹配到结果'} + + )} +
+ {task.provider ? ( + + {PROVIDER_LABELS[task.provider] ?? task.provider} + + ) : ( + — + )} + +
+ + {task.status === 'running' && ( + + )} + {task.status === 'done' + ? '已匹配' + : task.status === 'failed' + ? '未匹配' + : status.label} + + {task.error && ( + setDetailTask(task)} + className="cursor-pointer truncate max-w-[180px] text-[10px] text-rose-500 hover:underline" + title={task.error} + > + {task.error} + + )} +
+
+ {formatTime(task.created_at)} + +
+ + + {(task.status === 'pending' || task.status === 'running') && ( + + )} + + {(task.status === 'failed' || + task.status === 'canceled' || + task.status === 'done') && ( + + )} + + {(task.status === 'done' || + task.status === 'failed' || + task.status === 'canceled') && ( + + )} +
+
+
+ )} + + {/* 5. Pagination */} + {(snapshot?.total ?? 0) > 0 && ( +
+ + 共 {snapshot?.total ?? 0} 条 · 第 {page} / {totalPages} 页 + +
+ + +
+
+ )} +
+ + {/* 6. Task Detail Modal */} + {detailTask && ( + setDetailTask(null)} + onRetry={retryTask} + onCancel={cancelTask} + onDelete={deleteTask} + onCopy={copyText} + /> + )} +
+ ) +} + +function ScrapeDetailModal({ + task, + onClose, + onRetry, + onCancel, + onDelete, + onCopy, +}: { + task: ScrapeTask + onClose: () => void + onRetry: (t: ScrapeTask) => void + onCancel: (t: ScrapeTask) => void + onDelete: (t: ScrapeTask) => void + onCopy: (text: string, label: string) => void +}) { + const status = taskStatusMeta(task.status) + + return ( +
+
e.stopPropagation()} + > +
+
+ +

刮削任务详情

+
+ +
+ +
+ {/* Matched Poster / Info Banner */} + {task.matched_title ? ( +
+ {task.poster_url && ( + + )} +
+
+ + 已匹配 + + {task.provider && ( + + {PROVIDER_LABELS[task.provider] ?? task.provider} + + )} +
+

+ {task.matched_title} +

+
+ {task.matched_year > 0 && 年份:{task.matched_year}} + 类型:{TYPE_LABELS[task.media_type] ?? task.media_type} +
+ {task.media_id && ( + + 在媒体详情中查看 + + + )} +
+
+ ) : null} + + {/* Media Info Box */} +
+
+ 原始媒体标题 + {task.media_title} +
+
+ 所属媒体库 + {task.library_name} +
+
+ 媒体库类型 + + {TYPE_LABELS[task.media_type] ?? task.media_type} + +
+
+ 当前状态 + + {task.status === 'done' ? '已匹配' : task.status === 'failed' ? '未匹配' : status.label} + +
+
+ 剧照/海报刮削 + + {task.episode_images ? '开启' : '关闭'} + +
+
+ + {/* File path */} +
+
+ 磁盘文件路径 + +
+
+ {task.media_path} +
+
+ + {/* Error Message Box */} + {task.error && ( +
+
+ + 刮削未匹配 / 异常详情 + + +
+
+ {task.error} +
+
+ )} + + {/* Timeline */} +
+
入队时间:{formatTime(task.created_at)}
+ {task.started_at &&
开始刮削:{formatTime(task.started_at)}
} + {task.finished_at &&
完成时间:{formatTime(task.finished_at)}
} +
+
+ + {/* Footer Actions */} +
+
+ {(task.status === 'done' || + task.status === 'failed' || + task.status === 'canceled') && ( + + )} +
+ +
+ {(task.status === 'pending' || task.status === 'running') && ( + + )} + + + + +
+
+
+
+ ) +} diff --git a/web/src/pages/useAdminLibraryPanel.ts b/web/src/pages/useAdminLibraryPanel.ts index 0276f08..94c696b 100644 --- a/web/src/pages/useAdminLibraryPanel.ts +++ b/web/src/pages/useAdminLibraryPanel.ts @@ -165,7 +165,7 @@ function useLibraryActions(refresh: () => Promise) { } const toggleCarouselLibrary = async (library: Library) => { - const next = !(library.carousel_enabled ?? true) + const next = !Boolean(library.carousel_enabled) await libraryAPI.update(library.id, { carousel_enabled: next }) toast.success(next ? `「${library.name}」已加入首页轮播` : `「${library.name}」已移出首页轮播`) await refresh() diff --git a/web/src/types/scraper.ts b/web/src/types/scraper.ts new file mode 100644 index 0000000..04a5e0c --- /dev/null +++ b/web/src/types/scraper.ts @@ -0,0 +1,40 @@ +export type ScrapeTaskStatus = 'pending' | 'running' | 'done' | 'failed' | 'canceled' + +export interface ScrapeTask { + id: string + media_id: string + library_id: string + library_name: string + media_title: string + media_path: string + media_type: string + provider: string + matched_title: string + matched_year: number + poster_url: string + backdrop_url: string + status: ScrapeTaskStatus + error: string + retry_count: number + episode_images: boolean + refresh_matched: boolean + created_at: string + started_at?: string | null + finished_at?: string | null +} + +export interface ScrapeQueueCounts { + pending: number + running: number + done: number + failed: number + canceled: number +} + +export interface ScrapeQueueSnapshot { + counts: ScrapeQueueCounts + tasks: ScrapeTask[] + total: number + page: number + page_size: number +}