From a666245c44e229ad059d051bbb86cd2fdd980405 Mon Sep 17 00:00:00 2001 From: ShukeBta Date: Thu, 28 May 2026 20:04:51 +0800 Subject: [PATCH] Fix subscription search to queue full series --- internal/service/subscription.go | 188 +++++++++++++++++++------- internal/service/subscription_test.go | 58 ++++++++ scripts/smoke-test.sh | 7 +- web/src/api/stats.ts | 3 +- web/src/pages/SearchPage.tsx | 2 +- web/src/pages/StatsPage.tsx | 150 ++++++++++++++++---- 6 files changed, 332 insertions(+), 76 deletions(-) create mode 100644 internal/service/subscription_test.go diff --git a/internal/service/subscription.go b/internal/service/subscription.go index b730f86..8839d67 100644 --- a/internal/service/subscription.go +++ b/internal/service/subscription.go @@ -25,6 +25,20 @@ import ( "github.com/ShukeBta/MediaStationGo/internal/repository" ) +var ( + seriesPackRE = regexp.MustCompile(`(?i)(complete|batch|合集|全集|全\s*\d+\s*[集话話期]|整季|全季|s\d{1,2}\s*(?:complete|batch|pack)|season\s*\d{1,2}\s*(?:complete|batch|pack)|s\d{1,2}e\d{1,3}\s*[-~–—]\s*(?:e)?\d{1,3}|第\s*\d+\s*[-~–—]\s*\d+\s*[集话話期])`) + seasonOnlyRE = regexp.MustCompile(`(?i)(?:^|[\s._-])(?:s|season)\s*\d{1,2}(?:[\s._-]|$)|第\s*\d+\s*季`) +) + +type siteSearchCandidate struct { + Item SearchResult + Download string + GUID string + Season int + Episode int + Pack bool +} + // SubscriptionService runs the polling loop. type SubscriptionService struct { cfg *config.Config @@ -248,39 +262,19 @@ func (s *SubscriptionService) runSiteSearch(ctx context.Context, sub *model.Subs seenSet[g] = struct{}{} } + candidates := selectSiteSearchCandidates(results, sub, seenSet) var lastEnqueueErr error - for _, item := range results { - download := strings.TrimSpace(item.DownloadURL) - if download == "" { - download = strings.TrimSpace(item.TorrentURL) - } - if download == "" { - continue - } - guid := download - if _, ok := seenSet[guid]; ok { - continue - } - if s.downloads != nil && s.downloads.TorrentExistsByName(ctx, item.Title) { - seen = append(seen, guid) - if len(seen) > 200 { - seen = seen[len(seen)-200:] - } - _ = s.repo.Setting.Set(ctx, guidKey, strings.Join(seen, "\n")) - now := time.Now() - _ = s.repo.DB.Model(sub).Updates(map[string]any{"last_run_at": &now}).Error - s.hub.Publish("subscription", map[string]any{ - "id": sub.ID, - "name": sub.Name, - "queued": 0, - "keyword": keyword, - "resource": item.Title, - "existing": true, - }) - return 0, nil - } - realURL := s.site.ResolveDownloadURL(ctx, download) + queued := 0 + var resources []string + for _, candidate := range candidates { + item := candidate.Item mediaType, mediaCategory := s.classifySubscriptionItem(ctx, sub, item.Title, item.Category) + if s.shouldSkipExistingTorrent(ctx, mediaType, candidate) { + seen = append(seen, candidate.GUID) + seenSet[candidate.GUID] = struct{}{} + continue + } + realURL := s.site.ResolveDownloadURL(ctx, candidate.Download) savePath := s.resolveSubscriptionSavePath(ctx, sub, mediaType, mediaCategory) if _, err := s.downloads.AddDownload(ctx, sub.UserID, realURL, savePath); err != nil { lastEnqueueErr = err @@ -294,31 +288,131 @@ func (s *SubscriptionService) runSiteSearch(ctx context.Context, sub *model.Subs zap.Error(err)) continue } - seen = append(seen, guid) - if len(seen) > 200 { - seen = seen[len(seen)-200:] - } - _ = s.repo.Setting.Set(ctx, guidKey, strings.Join(seen, "\n")) - now := time.Now() - _ = s.repo.DB.Model(sub).Updates(map[string]any{"last_run_at": &now}).Error - s.hub.Publish("subscription", map[string]any{ - "id": sub.ID, - "name": sub.Name, - "queued": 1, - "keyword": keyword, - "resource": item.Title, - }) - return 1, nil + queued++ + resources = append(resources, item.Title) + seen = append(seen, candidate.GUID) + seenSet[candidate.GUID] = struct{}{} } - + if len(seen) > 200 { + seen = seen[len(seen)-200:] + } + _ = s.repo.Setting.Set(ctx, guidKey, strings.Join(seen, "\n")) now := time.Now() _ = s.repo.DB.Model(sub).Updates(map[string]any{"last_run_at": &now}).Error + if queued > 0 { + s.hub.Publish("subscription", map[string]any{ + "id": sub.ID, + "name": sub.Name, + "queued": queued, + "keyword": keyword, + "resources": resources, + }) + return queued, nil + } if lastEnqueueErr != nil { return 0, fmt.Errorf("找到 PT 资源但加入下载器失败: %w", lastEnqueueErr) } return 0, nil } +func selectSiteSearchCandidates(results []SearchResult, sub *model.Subscription, seenSet map[string]struct{}) []siteSearchCandidate { + candidates := make([]siteSearchCandidate, 0, len(results)) + for _, item := range results { + download := strings.TrimSpace(item.DownloadURL) + if download == "" { + download = strings.TrimSpace(item.TorrentURL) + } + if download == "" { + continue + } + guid := download + if _, ok := seenSet[guid]; ok { + continue + } + season, episode := ParseEpisode(item.Title) + candidates = append(candidates, siteSearchCandidate{ + Item: item, + Download: download, + GUID: guid, + Season: season, + Episode: episode, + Pack: isSeriesPackTitle(item.Title), + }) + } + if len(candidates) <= 1 { + return candidates + } + + mediaType := normalizeMediaType(sub.MediaType, sub.Name+" "+sub.Filter, "") + if !isSubscriptionSeriesType(mediaType) { + return candidates[:1] + } + + for _, candidate := range candidates { + if candidate.Pack { + return []siteSearchCandidate{candidate} + } + } + + byEpisode := make(map[string]siteSearchCandidate) + order := make([]string, 0, len(candidates)) + for _, candidate := range candidates { + if candidate.Episode <= 0 { + continue + } + season := candidate.Season + if season <= 0 { + season = 1 + } + key := fmt.Sprintf("%02dE%03d", season, candidate.Episode) + if _, ok := byEpisode[key]; ok { + continue + } + byEpisode[key] = candidate + order = append(order, key) + } + if len(order) == 0 { + return candidates[:1] + } + + selected := make([]siteSearchCandidate, 0, len(order)) + for _, key := range order { + selected = append(selected, byEpisode[key]) + } + return selected +} + +func isSubscriptionSeriesType(mediaType string) bool { + switch normalizeMediaType(mediaType, "", "") { + case "tv", "anime", "variety": + return true + default: + return false + } +} + +func isSeriesPackTitle(title string) bool { + title = strings.TrimSpace(title) + if title == "" { + return false + } + if seriesPackRE.MatchString(title) { + return true + } + _, episode := ParseEpisode(title) + return episode == 0 && seasonOnlyRE.MatchString(title) +} + +func (s *SubscriptionService) shouldSkipExistingTorrent(ctx context.Context, mediaType string, candidate siteSearchCandidate) bool { + if s == nil || s.downloads == nil { + return false + } + if isSubscriptionSeriesType(mediaType) && !candidate.Pack && candidate.Episode > 0 { + return false + } + return s.downloads.TorrentExistsByName(ctx, candidate.Item.Title) +} + func siteSearchKeyword(sub *model.Subscription) string { if sub == nil { return "" diff --git a/internal/service/subscription_test.go b/internal/service/subscription_test.go new file mode 100644 index 0000000..86a4df2 --- /dev/null +++ b/internal/service/subscription_test.go @@ -0,0 +1,58 @@ +package service + +import ( + "testing" + + "github.com/ShukeBta/MediaStationGo/internal/model" +) + +func TestSelectSiteSearchCandidatesPrefersSeriesPack(t *testing.T) { + sub := &model.Subscription{Name: "间谍过家家 自动订阅", Filter: "间谍过家家 2022", MediaType: "tv"} + results := []SearchResult{ + {Title: "间谍过家家 S01E01 1080p", DownloadURL: "https://pt/download/1", Seeders: 80}, + {Title: "间谍过家家 S01 Complete 1080p", DownloadURL: "https://pt/download/pack", Seeders: 50}, + {Title: "间谍过家家 S01E02 1080p", DownloadURL: "https://pt/download/2", Seeders: 70}, + } + + got := selectSiteSearchCandidates(results, sub, map[string]struct{}{}) + if len(got) != 1 { + t.Fatalf("selected %d candidates, want 1", len(got)) + } + if got[0].Download != "https://pt/download/pack" || !got[0].Pack { + t.Fatalf("selected %#v, want complete pack", got[0]) + } +} + +func TestSelectSiteSearchCandidatesQueuesDistinctEpisodesWhenNoPack(t *testing.T) { + sub := &model.Subscription{Name: "葬送的芙莉莲 自动订阅", Filter: "葬送的芙莉莲", MediaType: "anime"} + results := []SearchResult{ + {Title: "葬送的芙莉莲 S01E01 1080p", DownloadURL: "https://pt/download/1a", Seeders: 90}, + {Title: "葬送的芙莉莲 S01E01 2160p", DownloadURL: "https://pt/download/1b", Seeders: 80}, + {Title: "葬送的芙莉莲 S01E02 1080p", DownloadURL: "https://pt/download/2", Seeders: 70}, + {Title: "葬送的芙莉莲 S01E03 1080p", DownloadURL: "https://pt/download/3", Seeders: 60}, + } + + got := selectSiteSearchCandidates(results, sub, map[string]struct{}{}) + if len(got) != 3 { + t.Fatalf("selected %d candidates, want 3", len(got)) + } + if got[0].Episode != 1 || got[1].Episode != 2 || got[2].Episode != 3 { + t.Fatalf("episodes = %d,%d,%d; want 1,2,3", got[0].Episode, got[1].Episode, got[2].Episode) + } + if got[0].Download != "https://pt/download/1a" { + t.Fatalf("duplicate episode should keep first/best result, got %q", got[0].Download) + } +} + +func TestSelectSiteSearchCandidatesKeepsMovieSingleBest(t *testing.T) { + sub := &model.Subscription{Name: "Inception 自动订阅", Filter: "Inception 2010", MediaType: "movie"} + results := []SearchResult{ + {Title: "Inception 2010 1080p", DownloadURL: "https://pt/download/1080", Seeders: 90}, + {Title: "Inception 2010 2160p", DownloadURL: "https://pt/download/2160", Seeders: 80}, + } + + got := selectSiteSearchCandidates(results, sub, map[string]struct{}{}) + if len(got) != 1 || got[0].Download != "https://pt/download/1080" { + t.Fatalf("selected %#v, want movie best only", got) + } +} diff --git a/scripts/smoke-test.sh b/scripts/smoke-test.sh index 7c04098..17b23fc 100755 --- a/scripts/smoke-test.sh +++ b/scripts/smoke-test.sh @@ -42,6 +42,9 @@ trap cleanup EXIT ok() { printf " \033[32m✓\033[0m %s\n" "$1"; PASS=$((PASS+1)); } fail() { printf " \033[31m✗\033[0m %s\n" "$1"; FAIL=$((FAIL+1)); } hdr() { printf "\n\033[1;36m==> %s\033[0m\n" "$1"; } +json_token() { + python3 -c 'import json,sys; d=json.load(sys.stdin); print(d.get("token") or d.get("access_token") or (d.get("tokens") or {}).get("access_token") or "")' +} require() { command -v "$1" >/dev/null || { echo "missing dependency: $1"; exit 2; } @@ -111,7 +114,7 @@ curl -s "http://127.0.0.1:$PORT/api/health" | grep -q '"ok"' && ok "/api/health" hdr "Auth" TOKEN=$(curl -s -X POST -H 'Content-Type: application/json' \ -d '{"username":"admin","password":"smoketest12345"}' \ - "http://127.0.0.1:$PORT/api/auth/login" | python3 -c 'import json,sys;print(json.load(sys.stdin)["token"])') + "http://127.0.0.1:$PORT/api/auth/login" | json_token) [ -n "$TOKEN" ] && ok "login as admin" || fail "login as admin" H="Authorization: Bearer $TOKEN" curl -s -o /dev/null -w "%{http_code}" -H "$H" "http://127.0.0.1:$PORT/api/me" | grep -q 200 \ @@ -207,7 +210,7 @@ curl -s -X POST -H 'Content-Type: application/json' \ ATOK=$(curl -s -X POST -H 'Content-Type: application/json' \ -d '{"username":"alice","password":"alice12345"}' \ "http://127.0.0.1:$PORT/api/auth/login" \ - | python3 -c 'import json,sys;print(json.load(sys.stdin)["token"])') + | json_token) curl -s -o /dev/null -w "%{http_code}" -H "Authorization: Bearer $ATOK" \ -X POST -H 'Content-Type: application/json' \ -d "{\"name\":\"x\",\"path\":\"$MEDIA\",\"type\":\"movie\"}" \ diff --git a/web/src/api/stats.ts b/web/src/api/stats.ts index fbafc98..d9a568d 100644 --- a/web/src/api/stats.ts +++ b/web/src/api/stats.ts @@ -1,6 +1,7 @@ import { api } from './client' -import type { StatsSnapshot } from '../types' +import type { Hardware, StatsSnapshot } from '../types' export const statsAPI = { snapshot: () => api.get('/stats').then((r) => r.data), + monitor: () => api.get('/stats/monitor').then((r) => r.data), } diff --git a/web/src/pages/SearchPage.tsx b/web/src/pages/SearchPage.tsx index 1a95581..e1f51c5 100644 --- a/web/src/pages/SearchPage.tsx +++ b/web/src/pages/SearchPage.tsx @@ -223,7 +223,7 @@ function ExternalResults({

外部数据源

- 来自 TMDb / 豆瓣 / Bangumi。订阅后会定期搜索已配置 PT 站点,并只入队最佳新资源。 + 来自 TMDb / 豆瓣 / Bangumi。电影入队最佳资源;剧集/动漫优先整季或全集包,否则按集批量入队。

diff --git a/web/src/pages/StatsPage.tsx b/web/src/pages/StatsPage.tsx index c3cdcc4..684c6ce 100644 --- a/web/src/pages/StatsPage.tsx +++ b/web/src/pages/StatsPage.tsx @@ -1,9 +1,9 @@ import { useEffect, useState } from 'react' -import { Activity, Cpu, Database, Film, HardDrive, Users } from 'lucide-react' +import { Activity, Cpu, Database, Film, HardDrive, Radio, Users } from 'lucide-react' import { statsAPI } from '../api/stats' import { MediaCard } from '../components/MediaCard' -import type { StatsSnapshot } from '../types' +import type { Hardware, StatsSnapshot } from '../types' import { groupSeries } from '../utils/groupSeries' // fmtBytes is a tiny helper shared by the dashboard cards. @@ -25,19 +25,52 @@ function fmtHours(seconds: number): string { return `${h.toLocaleString()} h` } -// StatsPage renders the operator dashboard. Refreshes every 10 s. +// StatsPage renders the operator dashboard. Aggregate stats refresh less often, +// while hardware metrics are polled every 2 s for a real-time monitoring feel. export function StatsPage() { const [snap, setSnap] = useState(null) + const [hardware, setHardware] = useState(null) + const [lastMonitorAt, setLastMonitorAt] = useState('') + const [monitorError, setMonitorError] = useState('') const [loading, setLoading] = useState(true) useEffect(() => { let cancelled = false const tick = () => statsAPI.snapshot().then((s) => { - if (!cancelled) setSnap(s) + if (!cancelled) { + setSnap(s) + setHardware((current) => current ?? s.hardware) + setLastMonitorAt((current) => current || s.generated_at) + } }) tick().finally(() => setLoading(false)) - const id = window.setInterval(tick, 10_000) + const id = window.setInterval(tick, 30_000) + return () => { + cancelled = true + window.clearInterval(id) + } + }, []) + + useEffect(() => { + let cancelled = false + const tick = () => + statsAPI + .monitor() + .then((m) => { + if (!cancelled) { + setHardware(m) + setLastMonitorAt(new Date().toISOString()) + setMonitorError('') + } + }) + .catch((err: unknown) => { + if (!cancelled) { + setMonitorError((err as { message?: string })?.message ?? '实时监控暂不可用') + } + }) + tick() + const id = window.setInterval(tick, 2_000) return () => { cancelled = true window.clearInterval(id) @@ -47,23 +80,40 @@ export function StatsPage() { if (loading) return

加载中…

if (!snap) return

无法获取统计数据

+ const live = hardware ?? snap.hardware const memPct = - snap.hardware.memory_total > 0 - ? (snap.hardware.memory_used / snap.hardware.memory_total) * 100 + live.memory_total > 0 + ? (live.memory_used / live.memory_total) * 100 : 0 const diskPct = - snap.hardware.disk_total > 0 - ? (snap.hardware.disk_used / snap.hardware.disk_total) * 100 + live.disk_total > 0 + ? (live.disk_used / live.disk_total) * 100 : 0 const recentlyAddedCards = groupSeries(snap.recently_added) return (
-
-

运行状态

-

- 快照时间:{new Date(snap.generated_at).toLocaleString()} -

+
+
+

运行状态

+

+ 聚合快照:{new Date(snap.generated_at).toLocaleString()} · 实时监控每 2 秒刷新 +

+
+
+ + + + +
+

+ {monitorError ? '监控重试中' : '实时监控中'} +

+

+ 最近刷新:{lastMonitorAt ? new Date(lastMonitorAt).toLocaleTimeString() : '—'} +

+
+
@@ -72,24 +122,49 @@ export function StatsPage() { } label="用户" value={snap.users_count.toLocaleString()} /> } label="入库容量" value={fmtBytes(snap.total_size_bytes)} /> } label="累计时长" value={fmtHours(snap.total_seconds)} /> - } label="CPU 占用" value={`${snap.hardware.cpu_percent.toFixed(1)}%`} /> + } label="CPU 占用" value={`${live.cpu_percent.toFixed(1)}%`} meter={live.cpu_percent} /> } label="内存占用" value={`${memPct.toFixed(1)}%`} /> } label="数据盘占用" value={`${diskPct.toFixed(1)}%`} />
-

系统

+
+

系统实时监控

+ + + Live + +
+
+
+ + + + +
+
+ + + + {monitorError &&

实时监控错误:{monitorError}

} +
+
+
+ +
+

聚合统计

- - - +
@@ -116,24 +191,49 @@ function Tile({ icon, label, value, + meter, }: { icon: React.ReactNode label: string value: string + meter?: number }) { return (
{icon}
-
+

{label}

{value}

+ {typeof meter === 'number' && }
) } +function Meter({ label, value }: { label: string; value: number }) { + return ( +
+
+ {label} + {value.toFixed(1)}% +
+ +
+ ) +} + +function Bar({ value }: { value: number }) { + const pct = Math.max(0, Math.min(100, value || 0)) + const color = pct > 85 ? 'bg-red-500' : pct > 65 ? 'bg-amber-500' : 'bg-emerald-500' + return ( +
+
+
+ ) +} + function Row({ label, value }: { label: string; value: string }) { return (