Fix cloud playback and organize ingest pipeline

This commit is contained in:
ShukeBta
2026-06-13 23:21:56 +08:00
parent 4b2839d753
commit 52b772ea45
18 changed files with 729 additions and 22 deletions
+1
View File
@@ -330,6 +330,7 @@ func serveCloudResolvedLink(svc *service.Container, c *gin.Context, typ, ref str
}
if !link.Proxy {
// Pure offload: send the client straight to the cloud CDN.
setRedirectNoStoreHeaders(c)
logCloudPlayback(svc, "cloud playback redirect",
append(cloudPlaybackLogFields(typ, ref, link, resolveDur),
zap.String("mode", "redirect"),
+46
View File
@@ -1059,6 +1059,15 @@ func embyVideoStreamHandler(svc *service.Container, cloudMode string) gin.Handle
c.Status(http.StatusNotFound)
return
}
if embyShouldRedirectVideoStreamToSTRM(c, svc, c.Param("id"), cloudMode) {
target := "/api/stream/" + url.PathEscape(strings.TrimSpace(c.Param("id")))
if token := embyPlaybackRedirectToken(c, svc); token != "" {
target = embyAppendAPIKey(target, token)
}
setRedirectNoStoreHeaders(c)
c.Redirect(http.StatusFound, absoluteRequestURL(c, target))
return
}
// 直接调用 Stream service 写入 response。
// 此前这里把所有错误一律吞成 404:云盘 Cookie 过期、直链解析失败、
// STRM 播放被关闭……在第三方播放器上全部表现为「404 不存在」,
@@ -1080,6 +1089,43 @@ func embyVideoStreamHandler(svc *service.Container, cloudMode string) gin.Handle
}
}
func embyPlaybackRedirectToken(c *gin.Context, svc *service.Container) string {
if token := embyRequestToken(c); token != "" {
return token
}
if c == nil || svc == nil || svc.Auth == nil || svc.Repo == nil || svc.Repo.User == nil {
return ""
}
uid := embyUserID(c)
if uid == "" {
return ""
}
u, err := svc.Repo.User.FindByID(c.Request.Context(), uid)
if err != nil || u == nil {
return ""
}
token, err := svc.Auth.IssueEmbyToken(u)
if err != nil {
return ""
}
return token
}
func embyShouldRedirectVideoStreamToSTRM(c *gin.Context, svc *service.Container, mediaID, cloudMode string) bool {
if c == nil || svc == nil || svc.Repo == nil || svc.Repo.Media == nil || cloudMode != service.CloudPlaybackModeRedirectProxy {
return false
}
settings := service.CloudPlaybackSettings(c.Request.Context(), svc.Repo)
if settings.PreferredMode != service.CloudPlaybackModeSTRM || !settings.STRMEnabled {
return false
}
m, err := svc.Repo.Media.FindByID(c.Request.Context(), mediaID)
if err != nil || m == nil {
return false
}
return strings.TrimSpace(m.STRMURL) != ""
}
func embyVideoHLSPlaylistHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
uid := embyUserID(c)
+146
View File
@@ -1063,6 +1063,152 @@ func TestEmbyItemsTokenizesEmbeddedCloudMediaSources(t *testing.T) {
}
}
func TestEmbyVideoStreamUsesSTRMWhenRedirectProxyDisabled(t *testing.T) {
gin.SetMode(gin.TestMode)
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatalf("open db: %v", err)
}
if err := db.AutoMigrate(model.AllModels()...); err != nil {
t.Fatalf("migrate: %v", err)
}
repos := repository.New(db)
if err := repos.Setting.Set(t.Context(), service.CloudPlaybackModeSettingKey, service.CloudPlaybackModeSTRM); err != nil {
t.Fatalf("set cloud playback mode: %v", err)
}
if err := repos.Setting.Set(t.Context(), service.CloudPlaybackSTRMEnabledSettingKey, "true"); err != nil {
t.Fatalf("enable strm playback: %v", err)
}
if err := repos.Setting.Set(t.Context(), service.CloudPlaybackRedirectEnabledSettingKey, "false"); err != nil {
t.Fatalf("disable redirect playback: %v", err)
}
if err := repos.User.Create(t.Context(), &model.User{
Base: model.Base{ID: "user-1"},
Username: "tester",
PasswordHash: "x",
Role: "admin",
Tier: "plus",
IsActive: true,
}); err != nil {
t.Fatalf("create user: %v", err)
}
lib := model.Library{Name: "OpenList", Path: "cloud://openlist/Movies", Type: "movie", Enabled: true}
if err := repos.Library.Create(t.Context(), &lib); err != nil {
t.Fatalf("create library: %v", err)
}
if err := db.Create(&model.Media{
Base: model.Base{ID: "cloud-1"},
LibraryID: lib.ID,
Title: "Cloud Movie",
Path: "cloud://openlist/Movies/Movie.mkv",
STRMURL: "/api/cloud/play/openlist?ref=%2FMovies%2FMovie.mkv",
Container: "mkv",
}).Error; err != nil {
t.Fatalf("create media: %v", err)
}
const secret = "test-secret"
router := gin.New()
cfg := &config.Config{Secrets: config.SecretsConfig{JWTSecret: secret}}
registerEmbyRoutes(router, secret, &service.Container{
Repo: repos,
Emby: service.NewEmbyService(cfg, zap.NewNop(), repos),
Stream: service.NewStreamService(cfg, zap.NewNop(), repos, nil),
})
token := signedTestToken(t, secret)
req := httptest.NewRequest(http.MethodGet, "/videos/cloud-1/stream?api_key="+token, nil)
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if w.Code != http.StatusFound {
t.Fatalf("unexpected status: %d body=%s", w.Code, w.Body.String())
}
loc := w.Header().Get("Location")
if !strings.Contains(loc, "/api/stream/cloud-1") || !strings.Contains(loc, "api_key=") {
t.Fatalf("STRM mode should redirect /Videos fallback to tokenized /api/stream, got %q", loc)
}
if got := w.Header().Get("Cache-Control"); !strings.Contains(got, "no-store") {
t.Fatalf("STRM fallback redirect Cache-Control = %q, want no-store", got)
}
if strings.Contains(loc, "/api/cloud/play/") {
t.Fatalf("STRM mode should not expose cloud play directly from /Videos fallback: %q", loc)
}
}
func TestEmbyVideoStreamIssuesTokenForSessionFallbackSTRMRedirect(t *testing.T) {
gin.SetMode(gin.TestMode)
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatalf("open db: %v", err)
}
if err := db.AutoMigrate(model.AllModels()...); err != nil {
t.Fatalf("migrate: %v", err)
}
repos := repository.New(db)
if err := repos.Setting.Set(t.Context(), service.CloudPlaybackModeSettingKey, service.CloudPlaybackModeSTRM); err != nil {
t.Fatalf("set cloud playback mode: %v", err)
}
if err := repos.Setting.Set(t.Context(), service.CloudPlaybackSTRMEnabledSettingKey, "true"); err != nil {
t.Fatalf("enable strm playback: %v", err)
}
if err := repos.Setting.Set(t.Context(), service.CloudPlaybackRedirectEnabledSettingKey, "false"); err != nil {
t.Fatalf("disable redirect playback: %v", err)
}
user := model.User{
Base: model.Base{ID: "user-1"},
Username: "tester",
PasswordHash: "x",
Role: "admin",
Tier: "plus",
IsActive: true,
}
if err := repos.User.Create(t.Context(), &user); err != nil {
t.Fatalf("create user: %v", err)
}
lib := model.Library{Name: "OpenList", Path: "cloud://openlist/Movies", Type: "movie", Enabled: true}
if err := repos.Library.Create(t.Context(), &lib); err != nil {
t.Fatalf("create library: %v", err)
}
if err := db.Create(&model.Media{
Base: model.Base{ID: "cloud-1"},
LibraryID: lib.ID,
Title: "Cloud Movie",
Path: "cloud://openlist/Movies/Movie.mkv",
STRMURL: "/api/cloud/play/openlist?ref=%2FMovies%2FMovie.mkv",
Container: "mkv",
}).Error; err != nil {
t.Fatalf("create media: %v", err)
}
const secret = "test-secret"
cfg := &config.Config{Secrets: config.SecretsConfig{JWTSecret: secret}}
svc := &service.Container{
Repo: repos,
Auth: service.NewAuthService(cfg, zap.NewNop(), repos, nil, nil),
Emby: service.NewEmbyService(cfg, zap.NewNop(), repos),
Stream: service.NewStreamService(cfg, zap.NewNop(), repos, nil),
}
router := gin.New()
router.GET("/videos/:id/stream", func(c *gin.Context) {
c.Set(middleware.CtxUserID, user.ID)
c.Set(middleware.CtxUserRole, user.Role)
embyVideoStreamHandler(svc, service.CloudPlaybackModeRedirectProxy)(c)
})
req := httptest.NewRequest(http.MethodGet, "/videos/cloud-1/stream", nil)
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if w.Code != http.StatusFound {
t.Fatalf("unexpected status: %d body=%s", w.Code, w.Body.String())
}
loc := w.Header().Get("Location")
if !strings.Contains(loc, "/api/stream/cloud-1") || !strings.Contains(loc, "api_key=") {
t.Fatalf("session fallback redirect should include api_key for /api/stream, got %q", loc)
}
}
func TestEmbyLowercaseVideoStreamRouteServesMedia(t *testing.T) {
gin.SetMode(gin.TestMode)
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
+9
View File
@@ -159,6 +159,15 @@ func absoluteRequestURL(c *gin.Context, path string) string {
return scheme + "://" + host + path
}
func setRedirectNoStoreHeaders(c *gin.Context) {
if c == nil {
return
}
c.Header("Cache-Control", "no-store, no-cache, must-revalidate, max-age=0")
c.Header("Pragma", "no-cache")
c.Header("Expires", "0")
}
// transcodeStatusHandler reports the live status of one transcode job.
// We surface the active jobs the transcoder knows about.
func transcodeStatusHandler(svc *service.Container) gin.HandlerFunc {
+5 -1
View File
@@ -118,10 +118,14 @@ func generateSTRMHandler(svc *service.Container) gin.HandlerFunc {
if strmSvc == nil {
strmSvc = service.NewSTRMService(svc.Log, svc.Repo, svc.Cfg)
}
baseURL := strings.TrimRight(strings.TrimSpace(req.BaseURL), "/")
if baseURL == "" {
baseURL = strings.TrimRight(absoluteRequestURL(c, "/"), "/")
}
res, err := strmSvc.GenerateForLibrary(c.Request.Context(), service.GenerateSTRMOptions{
LibraryID: req.LibraryID,
OutputDir: req.OutputDir,
BaseURL: req.BaseURL,
BaseURL: baseURL,
Enabled: req.Enabled,
Overwrite: req.Overwrite,
IncludeLocal: true,
+200 -10
View File
@@ -28,6 +28,13 @@ import (
"github.com/ShukeBta/MediaStationGo/internal/model"
)
const (
organizeSkipAlreadyOrganized = "already organized"
organizeSkipDuplicateLibrary = "duplicate in library"
organizeSkipTargetExists = "target file exists"
organizeSkipSampleClip = "sample/trailer clip"
)
// OrganizeSourceCandidate is a selectable organize source directory surfaced to
// the UI so operators can organize an arbitrary directory (such as the download
// directory) and not only registered libraries.
@@ -127,6 +134,20 @@ func (o *OrganizerService) OrganizeDirectory(ctx context.Context, opts OrganizeO
if _, ok := videoExtensions[ext]; !ok {
return nil, fmt.Errorf("source is not a supported video file: %s", source)
}
if skipped, reason := shouldSkipOrganizeSourceVideo(source, filepath.Dir(source)); skipped {
res.Skipped++
res.Items = append(res.Items, OrganizePreviewItem{Source: source, Action: "skip", Reason: reason})
o.log.Info("organize file finished",
zap.String("source", source),
zap.String("dest", dest),
zap.String("mode", string(mode)),
zap.Int("organized", res.Organized),
zap.Int("replaced", res.Replaced),
zap.Int("skipped", res.Skipped),
zap.Any("skip_reasons", OrganizeSkipReasonCounts(res)),
)
return res, nil
}
if err := o.organizeSourceFile(ctx, source, filepath.Dir(source), dest, mode, opts.MediaType, opts.MediaCategory, opts.DryRun, opts.AllowReplaceExisting, res); err != nil {
res.Errors = append(res.Errors, fmt.Sprintf("%s: %s", filepath.Base(source), err.Error()))
res.Items = append(res.Items, OrganizePreviewItem{Source: source, Action: "error", Reason: err.Error()})
@@ -138,6 +159,7 @@ func (o *OrganizerService) OrganizeDirectory(ctx context.Context, opts OrganizeO
zap.Int("organized", res.Organized),
zap.Int("replaced", res.Replaced),
zap.Int("skipped", res.Skipped),
zap.Any("skip_reasons", OrganizeSkipReasonCounts(res)),
)
return res, nil
}
@@ -149,6 +171,11 @@ func (o *OrganizerService) OrganizeDirectory(ctx context.Context, opts OrganizeO
if _, ok := videoExtensions[ext]; !ok {
return nil
}
if skipped, reason := shouldSkipOrganizeSourceVideo(path, source); skipped {
res.Skipped++
res.Items = append(res.Items, OrganizePreviewItem{Source: path, Action: "skip", Reason: reason})
return nil
}
if err := o.organizeSourceFile(ctx, path, source, dest, mode, opts.MediaType, opts.MediaCategory, opts.DryRun, opts.AllowReplaceExisting, res); err != nil {
res.Errors = append(res.Errors, fmt.Sprintf("%s: %s", filepath.Base(path), err.Error()))
res.Items = append(res.Items, OrganizePreviewItem{Source: path, Action: "error", Reason: err.Error()})
@@ -165,6 +192,7 @@ func (o *OrganizerService) OrganizeDirectory(ctx context.Context, opts OrganizeO
zap.Int("organized", res.Organized),
zap.Int("replaced", res.Replaced),
zap.Int("skipped", res.Skipped),
zap.Any("skip_reasons", OrganizeSkipReasonCounts(res)),
)
return res, nil
}
@@ -212,6 +240,9 @@ func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRo
if !matchedLibrary && layout.Category != "" {
layoutRoot = categoryRoot(layoutRoot, sanitizeFilename(layout.Category))
}
if !matchedLibrary && !dryRun {
o.ensureOrganizeLibraryForRoot(ctx, layoutRoot, layout.MediaType, layout.Category)
}
var destDir, dst, episodeTag string
isSeries := season > 0 || episode > 0
@@ -237,7 +268,7 @@ func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRo
if filepath.Clean(src) == filepath.Clean(dst) {
res.Skipped++
res.Items = append(res.Items, OrganizePreviewItem{
Source: src, Target: dst, Action: "skip", Reason: "already organized",
Source: src, Target: dst, Action: "skip", Reason: organizeSkipAlreadyOrganized,
MediaType: layout.MediaType, Category: layout.Category, Title: title,
})
return nil
@@ -245,7 +276,9 @@ func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRo
// 去重候选:合并「目的地媒体库已扫描入库的同一媒体(按标题/年份/季集匹配,
// 不受目录大小写或布局影响)」与「目标文件夹内已存在的同名视频文件」。
existing := o.existingVersionPaths(ctx, destRoot, destDir, parsedTitle, episodeTag, year, season, episode)
identityExisting := o.existingByIdentity(ctx, destRoot, parsedTitle, year, season, episode)
folderExisting := o.existingByFolder(destDir, episodeTag)
existing := mergeExistingVersionPaths(identityExisting, folderExisting)
if len(existing) > 0 {
srcArea := o.resolutionArea(ctx, src)
bestArea := 0
@@ -278,11 +311,15 @@ func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRo
return nil
}
// 去重:目的地已存在同一媒体且不低于来源分辨率,跳过不再整理过去。
reason := organizeSkipTargetExists
if len(identityExisting) > 0 || o.allExistingPathsInDB(ctx, existing) {
reason = organizeSkipDuplicateLibrary
}
o.log.Debug("organize skip duplicate",
zap.String("src", src), zap.String("dest_dir", destDir))
zap.String("src", src), zap.String("dest_dir", destDir), zap.String("reason", reason))
res.Skipped++
res.Items = append(res.Items, OrganizePreviewItem{
Source: src, Target: dst, Action: "skip", Reason: "duplicate exists",
Source: src, Target: dst, Action: "skip", Reason: reason,
MediaType: layout.MediaType, Category: layout.Category, Title: title,
})
return nil
@@ -303,7 +340,7 @@ func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRo
res.Skipped++
if len(res.Items) > 0 {
res.Items[len(res.Items)-1].Action = "skip"
res.Items[len(res.Items)-1].Reason = "target exists"
res.Items[len(res.Items)-1].Reason = organizeSkipTargetExists
}
return nil
}
@@ -318,6 +355,37 @@ func (o *OrganizerService) organizeSourceFile(ctx context.Context, src, sourceRo
return nil
}
func shouldSkipOrganizeSourceVideo(path, sourceRoot string) (bool, string) {
cleanPath := filepath.Clean(path)
cleanRoot := filepath.Clean(sourceRoot)
if rel, err := filepath.Rel(cleanRoot, cleanPath); err == nil && rel != "." && !strings.HasPrefix(rel, "..") {
dir := filepath.Dir(rel)
if dir != "." {
for _, part := range strings.Split(dir, string(os.PathSeparator)) {
switch normalizeOrganizeCategoryKey(part) {
case "sample", "samples", "trailer", "trailers", "preview", "previews", "teaser", "teasers":
return true, organizeSkipSampleClip
}
}
}
}
base := strings.ToLower(strings.TrimSuffix(filepath.Base(cleanPath), filepath.Ext(cleanPath)))
normalized := strings.NewReplacer("_", " ", "-", " ", ".", " ").Replace(base)
fields := strings.Fields(normalized)
if len(fields) == 0 {
return false, ""
}
if len(fields) == 1 && strings.HasPrefix(fields[0], "sample") {
return true, organizeSkipSampleClip
}
last := fields[len(fields)-1]
switch last {
case "sample", "trailer", "preview", "teaser":
return true, organizeSkipSampleClip
}
return false, ""
}
func normalizeOrganizeMediaType(mediaType string) string {
switch strings.ToLower(strings.TrimSpace(mediaType)) {
case "movie", "film":
@@ -507,6 +575,92 @@ func (o *OrganizerService) organizeLibraryRootForLayout(ctx context.Context, des
return filepath.Clean(bestPath), true
}
func (o *OrganizerService) ensureOrganizeLibraryForRoot(ctx context.Context, root, mediaType, category string) {
if o == nil || o.repo == nil || o.repo.Library == nil {
return
}
root = filepath.Clean(strings.TrimSpace(root))
if root == "" || root == "." {
return
}
if _, ok := ParseCloudLibraryMount(root); ok {
return
}
libraries, err := o.repo.Library.List(ctx)
if err != nil {
if o.log != nil {
o.log.Debug("list libraries before organize auto-create failed", zap.Error(err))
}
return
}
for _, lib := range libraries {
if !lib.Enabled || strings.TrimSpace(lib.Path) == "" {
continue
}
if _, ok := ParseCloudLibraryMount(lib.Path); ok {
continue
}
if pathWithin(root, lib.Path) {
return
}
}
name := strings.TrimSpace(category)
if name == "" {
name = filepath.Base(root)
}
if name == "" || name == "." || name == string(os.PathSeparator) {
name = organizeLibraryTypeName(mediaType)
}
lib := model.Library{
Name: name,
Path: root,
Type: organizeLibraryModelType(mediaType),
Enabled: true,
}
if err := o.repo.Library.Create(ctx, &lib); err != nil {
if o.log != nil {
o.log.Warn("organize auto-create library failed",
zap.String("path", root),
zap.String("type", lib.Type),
zap.String("name", lib.Name),
zap.Error(err))
}
return
}
if o.log != nil {
o.log.Info("organize auto-created missing library",
zap.String("path", root),
zap.String("type", lib.Type),
zap.String("name", lib.Name))
}
}
func organizeLibraryModelType(mediaType string) string {
switch normalizeOrganizeMediaType(mediaType) {
case "tv", "anime", "variety":
return "tv"
case "adult", "movie":
return "movie"
default:
return "movie"
}
}
func organizeLibraryTypeName(mediaType string) string {
switch normalizeOrganizeMediaType(mediaType) {
case "tv":
return "电视剧"
case "anime":
return "动漫"
case "variety":
return "综艺"
case "adult":
return "成人"
default:
return "电影"
}
}
func (o *OrganizerService) organizeCategoryAliases(mediaType, category string) map[string]struct{} {
aliases := map[string]struct{}{}
add := func(values ...string) {
@@ -603,6 +757,13 @@ func pathDepth(path string) int {
// 2. Filesystem: video files inside the computed destination folder (matching
// the SxxExx tag for episodes). Covers destinations that were not scanned.
func (o *OrganizerService) existingVersionPaths(ctx context.Context, destRoot, destDir, title, episodeTag string, year, season, episode int) []string {
return mergeExistingVersionPaths(
o.existingByIdentity(ctx, destRoot, title, year, season, episode),
o.existingByFolder(destDir, episodeTag),
)
}
func mergeExistingVersionPaths(groups ...[]string) []string {
seen := map[string]struct{}{}
var out []string
add := func(p string) {
@@ -619,15 +780,44 @@ func (o *OrganizerService) existingVersionPaths(ctx context.Context, destRoot, d
seen[c] = struct{}{}
out = append(out, c)
}
for _, p := range o.existingByIdentity(ctx, destRoot, title, year, season, episode) {
add(p)
}
for _, p := range o.existingByFolder(destDir, episodeTag) {
add(p)
for _, group := range groups {
for _, p := range group {
add(p)
}
}
return out
}
func (o *OrganizerService) allExistingPathsInDB(ctx context.Context, paths []string) bool {
if o == nil || o.repo == nil || o.repo.DB == nil || len(paths) == 0 {
return false
}
cleaned := make([]string, 0, len(paths))
seen := map[string]struct{}{}
for _, path := range paths {
path = filepath.Clean(strings.TrimSpace(path))
if path == "" || path == "." {
continue
}
if _, ok := seen[path]; ok {
continue
}
seen[path] = struct{}{}
cleaned = append(cleaned, path)
}
if len(cleaned) == 0 {
return false
}
var count int64
if err := o.repo.DB.WithContext(ctx).
Model(&model.Media{}).
Where("path IN ?", cleaned).
Count(&count).Error; err != nil {
return false
}
return count == int64(len(cleaned))
}
// existingByIdentity finds scanned destination media matching the parsed
// identity (case-insensitive title + year for movies; title + season/episode
// for episodes), located under destRoot.
@@ -313,6 +313,73 @@ func TestOrganizeDirectoryDedup(t *testing.T) {
}
}
func TestOrganizeDirectoryTreatsScannedTargetPathAsLibraryDuplicate(t *testing.T) {
root := t.TempDir()
src := filepath.Join(root, "downloads")
dest := filepath.Join(root, "media")
writeOrgFile(t, filepath.Join(src, "Weird.Release.2024.1080p.mkv"), "source")
existing := filepath.Join(dest, "电影", "Weird Release (2024)", "Weird Release (2024).mkv")
writeOrgFile(t, existing, "existing")
repos := newOrganizerTestRepo(t)
row := model.Media{Title: "刮削后的正式片名", Path: existing, Year: 2024, Container: "mkv", Width: 1920, Height: 1080}
if err := repos.Media.Upsert(t.Context(), &row); err != nil {
t.Fatal(err)
}
org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos)
res, err := org.OrganizeDirectory(t.Context(), OrganizeOptions{
SourcePath: src,
DestPath: dest,
TransferMode: TransferCopy,
})
if err != nil {
t.Fatalf("organize directory: %v", err)
}
if res.Organized != 0 || res.Skipped != 1 {
t.Fatalf("result = %+v, want organized=0 skipped=1", res)
}
if len(res.Items) != 1 || res.Items[0].Reason != organizeSkipDuplicateLibrary {
t.Fatalf("items = %+v, want duplicate in library", res.Items)
}
if OrganizeResultNeedsVisibilitySync(res) {
t.Fatal("path already present in DB should not trigger another visibility scan")
}
}
func TestOrganizeDirectorySkipsSampleClips(t *testing.T) {
root := t.TempDir()
src := filepath.Join(root, "downloads")
dest := filepath.Join(root, "media")
writeOrgFile(t, filepath.Join(src, "Demo.Movie.2024.1080p.mkv"), "main")
writeOrgFile(t, filepath.Join(src, "Samples", "Sample1.mkv"), "sample-dir")
writeOrgFile(t, filepath.Join(src, "Demo.Movie.2024.sample.mkv"), "sample-file")
org := NewOrganizerService(&config.Config{}, zap.NewNop(), newOrganizerTestRepo(t))
res, err := org.OrganizeDirectory(t.Context(), OrganizeOptions{
SourcePath: src,
DestPath: dest,
TransferMode: TransferCopy,
})
if err != nil {
t.Fatalf("organize directory: %v", err)
}
if res.Organized != 1 || res.Skipped != 2 {
t.Fatalf("result = %+v, want organized=1 skipped=2", res)
}
if got := OrganizeSkipReasonCounts(res)[organizeSkipSampleClip]; got != 2 {
t.Fatalf("sample skip count = %d, want 2; items=%+v", got, res.Items)
}
if _, err := os.Stat(filepath.Join(dest, "电影", "Sample1", "Sample1.mkv")); !os.IsNotExist(err) {
t.Fatalf("sample directory clip should not be organized, stat err=%v", err)
}
if _, err := os.Stat(filepath.Join(dest, "电影", "Demo Movie Sample (2024)", "Demo Movie Sample (2024).mkv")); !os.IsNotExist(err) {
t.Fatalf("sample suffix clip should not be organized, stat err=%v", err)
}
}
// TestOrganizeDirectorySkipsHigherResolutionWhenReplacementDisabled verifies
// that dedup wins by default: even a higher-resolution source must not replace
// an existing library item unless the caller explicitly allows washing.
@@ -549,6 +616,54 @@ func TestOrganizeDirectoryUsesExplicitCategoryLibraryRoot(t *testing.T) {
}
}
func TestOrganizeDirectoryCreatesMissingCategoryLibraryForVisibility(t *testing.T) {
root := t.TempDir()
srcRoot := filepath.Join(root, "downloads")
dest := filepath.Join(root, "media")
source := filepath.Join(srcRoot, "Gourd.Brothers.S01E01.2026.1080p.mkv")
target := filepath.Join(dest, "电视剧", "未分类", "Gourd Brothers", "Season 01", "Gourd Brothers - S01E01.mkv")
writeOrgFile(t, source, "source")
writeOrgFile(t, target, "already-there")
repos := newOrganizerTestRepo(t)
org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos)
res, err := org.OrganizeDirectory(t.Context(), OrganizeOptions{
SourcePath: srcRoot,
DestPath: dest,
MediaType: "tv",
MediaCategory: "未分类",
TransferMode: TransferCopy,
AllowReplaceExisting: false,
})
if err != nil {
t.Fatalf("organize missing category: %v", err)
}
if res.Organized != 0 || res.Skipped != 1 || len(res.Items) != 1 || res.Items[0].Reason != organizeSkipTargetExists {
t.Fatalf("result = %+v, want skipped target exists", res)
}
var lib model.Library
if err := repos.DB.Where("path = ?", filepath.Join(dest, "电视剧", "未分类")).First(&lib).Error; err != nil {
t.Fatalf("missing auto-created category library: %v", err)
}
if lib.Name != "未分类" || lib.Type != "tv" || !lib.Enabled {
t.Fatalf("auto-created library = %+v, want enabled tv 未分类", lib)
}
scanner := NewScannerService(&config.Config{}, zap.NewNop(), repos, NewHub(zap.NewNop()), nil, nil)
scans := scanner.ScanLibrariesForPath(t.Context(), res.DestPath, "")
added := 0
for _, scan := range scans {
if scan.Error != "" {
t.Fatalf("scan failed: %#v", scan)
}
added += scan.Added
}
if added != 1 {
t.Fatalf("scan added = %d, want 1; scans=%#v", added, scans)
}
}
func TestOrganizeDirectorySmartClassifiesUncategorizedSources(t *testing.T) {
root := t.TempDir()
src := filepath.Join(root, "downloads")
+24 -5
View File
@@ -57,11 +57,27 @@ func OrganizeResultHasChanges(res *OrganizeResult) bool {
}
// 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.
// should make sure destination files exist in the DB. New/replaced files always
// need it. Skips only need it when the target file exists on disk but was not
// proven to be an already scanned library duplicate; otherwise an automatic
// organize loop with only seeding duplicates would rescan libraries forever.
func OrganizeResultNeedsVisibilitySync(res *OrganizeResult) bool {
return res != nil && (res.Organized > 0 || res.Replaced > 0 || res.Skipped > 0)
if OrganizeResultHasChanges(res) {
return true
}
if res == nil {
return false
}
for _, item := range res.Items {
if item.Action != "skip" {
continue
}
switch item.Reason {
case organizeSkipAlreadyOrganized, organizeSkipTargetExists, "duplicate exists", "target exists":
return true
}
}
return false
}
// ScanLibrariesForPath recursively scans libraries affected by an organize
@@ -144,7 +160,10 @@ func (s *ScannerService) scrapeOrganizeTargets(ctx context.Context, targets []mo
Name: lib.Name,
Path: lib.Path,
}
matched, err := s.scraper.EnrichLibrary(ctx, lib.ID)
// Organize is an explicit ingest workflow: after rename/classification,
// previously failed no_match rows should be retried so the operator does
// not need to run a separate manual scrape.
matched, err := s.scraper.EnrichLibrary(ctx, lib.ID, true)
if err != nil {
summary.Error = err.Error()
} else {
+69
View File
@@ -69,6 +69,75 @@ func TestOrganizeDirectoryScanAndScrapeAfter(t *testing.T) {
}
}
func TestOrganizeScanAndScrapeRetriesNoMatchRows(t *testing.T) {
scraper, repos, closeServer := newTestScraper(t)
defer closeServer()
root := t.TempDir()
libRoot := filepath.Join(root, "media", "电视剧")
mediaPath := filepath.Join(libRoot, "间谍过家家", "Season 02", "间谍过家家 - S02E02.mkv")
writeOrgFile(t, mediaPath, "episode")
lib := model.Library{
Name: "剧集",
Path: libRoot,
Type: "tv",
Enabled: true,
}
if err := repos.Library.Create(t.Context(), &lib); err != nil {
t.Fatal(err)
}
media := model.Media{
LibraryID: lib.ID,
Title: "间谍过家家",
Path: mediaPath,
SeasonNum: 2,
EpisodeNum: 2,
ScrapeStatus: "no_match",
}
if err := repos.Media.Upsert(t.Context(), &media); err != nil {
t.Fatal(err)
}
scanner := NewScannerService(&config.Config{}, zap.NewNop(), repos, NewHub(zap.NewNop()), nil, scraper)
_, scrapes := scanner.ScanAndScrapeLibrariesForPath(t.Context(), filepath.Join(root, "media"), "", true)
if len(scrapes) != 1 || scrapes[0].Matched != 1 || scrapes[0].Error != "" || scrapes[0].Skipped {
t.Fatalf("scrapes = %#v, want one retried no_match row", scrapes)
}
var got model.Media
if err := repos.DB.First(&got, "path = ?", mediaPath).Error; err != nil {
t.Fatal(err)
}
if got.ScrapeStatus != "matched" || got.TMDbID != 12345 {
t.Fatalf("media scrape status=%q tmdb=%d, want matched/12345", got.ScrapeStatus, got.TMDbID)
}
}
func TestOrganizeResultNeedsVisibilitySyncIgnoresScannedDuplicates(t *testing.T) {
if OrganizeResultNeedsVisibilitySync(&OrganizeResult{
Skipped: 1,
Items: []OrganizePreviewItem{{Action: "skip", Reason: organizeSkipDuplicateLibrary}},
}) {
t.Fatal("already-scanned duplicate should not trigger another visibility scan")
}
if OrganizeResultNeedsVisibilitySync(&OrganizeResult{
Skipped: 1,
Items: []OrganizePreviewItem{{Action: "skip", Reason: organizeSkipSampleClip}},
}) {
t.Fatal("sample clip skip should not trigger visibility scan")
}
if !OrganizeResultNeedsVisibilitySync(&OrganizeResult{
Skipped: 1,
Items: []OrganizePreviewItem{{Action: "skip", Reason: organizeSkipTargetExists}},
}) {
t.Fatal("unscanned target file should trigger visibility scan")
}
if !OrganizeResultNeedsVisibilitySync(&OrganizeResult{Organized: 1}) {
t.Fatal("organized files must trigger visibility scan")
}
}
func TestOrganizeScrapeAfterEnabledDefaultsOn(t *testing.T) {
if !OrganizeScrapeAfterEnabled(t.Context(), nil) {
t.Fatalf("organize scrape-after should default on without a repo")
+1 -3
View File
@@ -494,10 +494,8 @@ func cloudResolveHotRefreshWindow(ttl time.Duration) time.Duration {
func cloudResolveCacheTTL(typ string) time.Duration {
switch typ {
case cloud.TypeQuark, cloud.Type115:
case cloud.TypeQuark, cloud.Type115, cloud.TypeCloudDrive2, cloud.TypeOpenList:
return 2 * time.Minute
case cloud.TypeCloudDrive2, cloud.TypeOpenList:
return 15 * time.Minute
default:
return 5 * time.Minute
}
@@ -79,3 +79,11 @@ func TestCloudResolveHotCacheRefreshesInBackground(t *testing.T) {
t.Fatalf("refreshed link = %s, want second URL", link.URL)
}
}
func TestCloudResolveCacheTTLUsesShortTTLForCloudPlaybackLinks(t *testing.T) {
for _, typ := range []string{"quark", "cloud115", "clouddrive2", "openlist"} {
if got := cloudResolveCacheTTL(typ); got != 2*time.Minute {
t.Fatalf("%s cloud resolve cache ttl = %v, want 2m", typ, got)
}
}
}
+10
View File
@@ -380,6 +380,7 @@ func (s *StreamService) ServeFileWithCloudMode(w http.ResponseWriter, r *http.Re
// 云盘播放 URL 先规范化为相对路径,免疫扫描时固化的旧 host。
target := normalizeCloudPlayTarget(strmURL)
target = withAuthTokenForInternalRedirect(target, r, PublicServerURL(r.Context(), s.repo, s.cfg))
setCloudRedirectNoStore(w)
http.Redirect(w, r, absoluteInternalRedirect(target, r), http.StatusFound)
return nil
}
@@ -405,6 +406,15 @@ func (s *StreamService) ServeFileWithCloudMode(w http.ResponseWriter, r *http.Re
return nil
}
func setCloudRedirectNoStore(w http.ResponseWriter) {
if w == nil {
return
}
w.Header().Set("Cache-Control", "no-store, no-cache, must-revalidate, max-age=0")
w.Header().Set("Pragma", "no-cache")
w.Header().Set("Expires", "0")
}
func isCloudPlaybackTarget(raw string) bool {
_, _, ok := parseCloudMediaPlaybackURL(raw)
return ok
+3
View File
@@ -84,6 +84,9 @@ func TestServeFileRedirectsInternalSTRMAsAbsoluteURLWithToken(t *testing.T) {
!strings.Contains(loc, "token=jwt123") {
t.Fatalf("redirect Location should be absolute and tokenized, got %q", loc)
}
if got := w.Header().Get("Cache-Control"); !strings.Contains(got, "no-store") {
t.Fatalf("cloud redirect Cache-Control = %q, want no-store", got)
}
}
func TestServeFileRedirectUsesForwardedTunnelHost(t *testing.T) {
+3 -1
View File
@@ -98,7 +98,9 @@ func (s *STRMService) GenerateForLibrary(ctx context.Context, opts GenerateSTRMO
return nil, errors.New("output_dir required")
}
if strings.TrimSpace(opts.BaseURL) != "" && s.repo.Setting != nil {
_ = s.repo.Setting.Set(ctx, "app.server_url", strings.TrimRight(strings.TrimSpace(opts.BaseURL), "/"))
baseURL := strings.TrimRight(strings.TrimSpace(opts.BaseURL), "/")
_ = s.repo.Setting.Set(ctx, "app.server_url", baseURL)
_ = s.repo.Setting.Set(ctx, "strm.base_url", baseURL)
}
if s.repo.Setting != nil {
_ = s.repo.Setting.Set(ctx, "strm.auto_generate_enabled", strconv.FormatBool(opts.Enabled))
+6
View File
@@ -59,6 +59,12 @@ func TestGenerateSTRMForLibraryWritesFilesAndRecords(t *testing.T) {
localSTRM := filepath.Join(outDir, "本地电影 (2025)", "本地电影 (2025).strm")
assertFileContains(t, cloudSTRM, "http://nas.example:18080/api/stream/cloud-media?token=strm-token")
assertFileContains(t, localSTRM, "http://nas.example:18080/api/stream/local-media?token=strm-token")
if got, err := repos.Setting.Get(t.Context(), "app.server_url"); err != nil || got != "http://nas.example:18080" {
t.Fatalf("app.server_url = %q, %v; want generated base url", got, err)
}
if got, err := repos.Setting.Get(t.Context(), "strm.base_url"); err != nil || got != "http://nas.example:18080" {
t.Fatalf("strm.base_url = %q, %v; want generated base url", got, err)
}
var count int64
if err := repos.DB.Model(&model.STRMRecord{}).Count(&count).Error; err != nil {
+45
View File
@@ -1,6 +1,7 @@
package service
import (
"strings"
"sync"
"time"
@@ -257,6 +258,9 @@ func OrganizeTaskMetrics(res *OrganizeResult) map[string]int64 {
metrics["scan_updated"] = scanUpdated
metrics["scan_removed"] = scanRemoved
}
for reason, count := range OrganizeSkipReasonCounts(res) {
metrics["skip_"+organizeMetricKey(reason)] = int64(count)
}
var scrapeMatched int64
for _, scrape := range res.Scrapes {
scrapeMatched += int64(scrape.Matched)
@@ -273,3 +277,44 @@ func OrganizeTaskMetrics(res *OrganizeResult) map[string]int64 {
}
return metrics
}
func OrganizeSkipReasonCounts(res *OrganizeResult) map[string]int {
if res == nil || len(res.Items) == 0 {
return nil
}
counts := map[string]int{}
for _, item := range res.Items {
if item.Action != "skip" {
continue
}
reason := strings.TrimSpace(item.Reason)
if reason == "" {
reason = "unknown"
}
counts[reason]++
}
if len(counts) == 0 {
return nil
}
return counts
}
func organizeMetricKey(value string) string {
value = strings.ToLower(strings.TrimSpace(value))
if value == "" {
return "unknown"
}
replacer := strings.NewReplacer(" ", "_", "-", "_", "/", "_", "\\", "_")
value = replacer.Replace(value)
var b strings.Builder
for _, r := range value {
if (r >= 'a' && r <= 'z') || (r >= '0' && r <= '9') || r == '_' {
b.WriteRune(r)
}
}
out := strings.Trim(b.String(), "_")
if out == "" {
return "unknown"
}
return out
}
+32 -2
View File
@@ -10,6 +10,29 @@ import type { Library, Media } from '../types'
type CloudPlaybackMode = 'strm' | 'redirect_proxy'
function currentOrigin() {
if (typeof window === 'undefined') return ''
return window.location.origin.replace(/\/+$/, '')
}
function isLocalPlaybackBase(raw: string) {
try {
const u = new URL(raw)
const host = u.hostname.toLowerCase()
return host === 'localhost' || host === '127.0.0.1' || host === '::1'
} catch {
return false
}
}
function preferredSTRMBaseURL(saved: string) {
const current = currentOrigin()
const trimmed = saved.trim().replace(/\/+$/, '')
if (!trimmed) return current
if (current && isLocalPlaybackBase(trimmed) && trimmed !== current) return current
return trimmed
}
// StrmPage exposes the URL-as-file admin tooling backed by the Go server:
// - import a brand-new media row directly from a (library, title, url)
// tuple — useful for streaming-only entries with no on-disk file.
@@ -50,7 +73,7 @@ export function StrmPage() {
.listSettings()
.then((rows) => {
const settings = Object.fromEntries(rows.map((row) => [row.key, row.value]))
setBaseURL(settings['app.server_url'] || settings['strm.base_url'] || '')
setBaseURL(preferredSTRMBaseURL(settings['strm.base_url'] || settings['app.server_url'] || ''))
setOutputDir(settings['strm.output_dir'] || '')
const mode = settings['cloud.playback_mode']
const nextMode =
@@ -336,6 +359,13 @@ export function StrmPage() {
value={baseURL}
onChange={(e) => setBaseURL(e.target.value)}
/>
<button
type="button"
className="rounded-2xl border border-primary-400/40 px-3 py-2 text-sm text-brand-500 transition hover:bg-primary-400/10"
onClick={() => setBaseURL(currentOrigin())}
>
使用当前访问地址
</button>
<label className="flex items-center gap-2 rounded-2xl border border-gray-200 bg-white/70 px-3 py-2 text-sm text-ink-50">
<input
type="checkbox"
@@ -350,7 +380,7 @@ export function StrmPage() {
value={outputDir}
onChange={(e) => setOutputDir(e.target.value)}
/>
<button type="submit" disabled={generating || !generateLibraryID || !baseURL.trim()} className="neon-button">
<button type="submit" disabled={generating || !generateLibraryID || !baseURL.trim()} className="neon-button md:col-span-4">
{generating ? <Loader2 size={16} className="animate-spin" /> : <Wand2 size={16} />}
{generating ? '生成中…' : '批量生成 STRM'}
</button>
+6
View File
@@ -30,6 +30,12 @@ const metricLabels: Record<string, string> = {
scrape_matched: '匹配',
scrape_skipped: '刮削跳过',
scrape_errors: '刮削错误',
skip_already_organized: '已在目标',
skip_duplicate_in_library: '已入库去重',
skip_target_file_exists: '目标已存在',
skip_sample_trailer_clip: '样片过滤',
skip_duplicate_exists: '重复跳过',
skip_target_exists: '目标已存在',
visited: '访问',
added: '入库',
updated: '更新',