This commit is contained in:
truewhile
2026-09-01 18:55:05 +08:00
parent 5c9e7fcaa6
commit db64a6c093
49 changed files with 1140 additions and 162 deletions
+1
View File
@@ -43,6 +43,7 @@ func setDefaults(v *viper.Viper) {
v.SetDefault("logging.max_backups", 10)
v.SetDefault("cache.cache_dir", "./cache")
v.SetDefault("cache.images_max_size_mb", 500)
v.SetDefault("cache.cleanup_interval_min", 60)
v.SetDefault("cache.redis_url", "")
v.SetDefault("cache.redis_prefix", "mmtl")
+3
View File
@@ -44,6 +44,9 @@ func (c *Config) normalize() error {
if c.Cache.CacheDir == "" {
c.Cache.CacheDir = filepath.Join(c.App.DataDir, "cache")
}
if c.Cache.ImagesMaxSizeMB < 0 {
c.Cache.ImagesMaxSizeMB = 0
}
if c.Cache.RedisPrefix == "" {
c.Cache.RedisPrefix = "mmtl"
}
+1
View File
@@ -116,6 +116,7 @@ type LoggingConfig struct {
// CacheConfig 控制磁盘转码/刮削缓存。
type CacheConfig struct {
CacheDir string `mapstructure:"cache_dir"`
ImagesMaxSizeMB int `mapstructure:"images_max_size_mb"`
MaxDiskUsageMB int `mapstructure:"max_disk_usage_mb"`
TTLHours int `mapstructure:"ttl_hours"`
AutoCleanup bool `mapstructure:"auto_cleanup"`
+3
View File
@@ -64,6 +64,9 @@ func updateSettingHandler(svc *service.Container) gin.HandlerFunc {
if req.Key == "transcode.hw_enabled" || req.Key == "transcode.hw_accel" || req.Key == "transcoder.hardware_accel" || req.Key == "transcoder.encoder" {
svc.Transcoder.StopAll()
}
if req.Key == "cache.images_max_size_mb" && svc.Scheduler != nil {
_ = svc.Scheduler.RunNowAsync(c.Request.Context(), "image_cache_cleanup")
}
c.Status(http.StatusNoContent)
}
}
+24
View File
@@ -195,4 +195,28 @@ func deleteEmbyMountHandler(svc *service.Container) gin.HandlerFunc {
}
c.JSON(http.StatusOK, gin.H{"ok": true})
}
}
type reorderEmbyMountsReq struct {
IDs []string `json:"ids" binding:"required"`
}
// reorderEmbyMountsHandler 批量重排挂载媒体库顺序。
func reorderEmbyMountsHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
var req reorderEmbyMountsReq
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
if svc.EmbyRemote == nil {
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "emby remote service not available"})
return
}
if err := svc.EmbyRemote.ReorderMounts(c.Request.Context(), req.IDs); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"ok": true})
}
}
@@ -0,0 +1,65 @@
package handler
import (
"bytes"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"github.com/gin-gonic/gin"
"github.com/glebarez/sqlite"
"go.uber.org/zap"
"gorm.io/gorm"
"github.com/ShukeBta/MMTL/internal/database"
"github.com/ShukeBta/MMTL/internal/model"
"github.com/ShukeBta/MMTL/internal/repository"
"github.com/ShukeBta/MMTL/internal/service"
)
func TestReorderEmbyMountsHandler(t *testing.T) {
gin.SetMode(gin.TestMode)
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err := database.AutoMigrate(db); err != nil {
t.Fatal(err)
}
repos := repository.New(db)
ctx := t.Context()
m1 := &model.EmbyMount{AccountID: "acct-1", RemoteViewID: "view-1", Name: "Mount 1"}
m2 := &model.EmbyMount{AccountID: "acct-1", RemoteViewID: "view-2", Name: "Mount 2"}
_ = repos.EmbyMount.Create(ctx, m1)
_ = repos.EmbyMount.Create(ctx, m2)
svc := &service.Container{
Repo: repos,
EmbyRemote: service.NewEmbyRemoteService(nil, zap.NewNop(), repos, nil),
}
router := gin.New()
router.PUT("/admin/emby/mounts/reorder", reorderEmbyMountsHandler(svc))
body, _ := json.Marshal(map[string]any{
"ids": []string{m2.ID, m1.ID},
})
req := httptest.NewRequest(http.MethodPut, "/admin/emby/mounts/reorder", bytes.NewReader(body))
req.Header.Set("Content-Type", "application/json")
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if w.Code != http.StatusOK {
t.Fatalf("expected status 200, got %d: %s", w.Code, w.Body.String())
}
list, err := repos.EmbyMount.List(ctx)
if err != nil {
t.Fatal(err)
}
if len(list) != 2 || list[0].ID != m2.ID || list[1].ID != m1.ID {
t.Fatalf("expected order [m2, m1], got [m%s, m%s]", list[0].ID, list[1].ID)
}
}
+5 -4
View File
@@ -49,10 +49,11 @@ func registerAdminStrmRoutes(admin *gin.RouterGroup, svc *service.Container) {
// Emby 挂载管理:远程 Emby 媒体库挂载(账号复用 strm/accounts)
admin.GET("/emby/accounts/:id/views", embyAccountViewsHandler(svc))
admin.POST("/emby/accounts/:id/full-mount", fullMountEmbyAccountHandler(svc))
admin.GET("/emby/mounts", listEmbyMountsHandler(svc))
admin.POST("/emby/mounts", createEmbyMountsHandler(svc))
admin.PUT("/emby/mounts/:id", updateEmbyMountHandler(svc))
admin.DELETE("/emby/mounts/:id", deleteEmbyMountHandler(svc))
admin.GET("/emby/mounts", listEmbyMountsHandler(svc))
admin.POST("/emby/mounts", createEmbyMountsHandler(svc))
admin.PUT("/emby/mounts/reorder", reorderEmbyMountsHandler(svc))
admin.PUT("/emby/mounts/:id", updateEmbyMountHandler(svc))
admin.DELETE("/emby/mounts/:id", deleteEmbyMountHandler(svc))
admin.GET("/strm/accounts", listStrmAccountsHandler(svc))
admin.POST("/strm/accounts", createStrmAccountHandler(svc))
+1
View File
@@ -14,6 +14,7 @@ type EmbyMount struct {
RemoteViewName string `gorm:"size:255" json:"remote_view_name"` // 远程媒体库原名(展示冗余)
CollectionType string `gorm:"size:32" json:"collection_type"` // movies / tvshows / music ...
Name string `gorm:"size:255" json:"name,omitempty"` // 覆盖显示名(可选,默认「账号 · 库名」)
SortOrder int `gorm:"default:0;index" json:"sort_order"` // 手动排序用,越小越靠前
ProxyPlay bool `gorm:"default:false" json:"proxy_play"` // 该挂载播放流量是否经 MMTL 反向代理
Enabled bool `gorm:"default:true" json:"enabled"` // 是否在媒体库中展示
}
+38 -4
View File
@@ -15,7 +15,14 @@ type EmbyMountRepository struct{ db *gorm.DB }
func (r *EmbyMountRepository) Create(ctx context.Context, m *model.EmbyMount) error {
return withSQLiteBusyRetry(ctx, func() error {
return r.db.WithContext(ctx).Create(m).Error
return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if m != nil && m.SortOrder == 0 {
var maxSort int
_ = tx.Model(&model.EmbyMount{}).Select("COALESCE(MAX(sort_order), -1)").Scan(&maxSort)
m.SortOrder = maxSort + 1
}
return tx.Create(m).Error
})
})
}
@@ -27,7 +34,17 @@ func (r *EmbyMountRepository) CreateInBatches(ctx context.Context, mounts []*mod
batchSize = 50
}
return withSQLiteBusyRetry(ctx, func() error {
return r.db.WithContext(ctx).CreateInBatches(mounts, batchSize).Error
return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var maxSort int
_ = tx.Model(&model.EmbyMount{}).Select("COALESCE(MAX(sort_order), -1)").Scan(&maxSort)
for _, m := range mounts {
if m != nil && m.SortOrder == 0 {
maxSort++
m.SortOrder = maxSort
}
}
return tx.CreateInBatches(mounts, batchSize).Error
})
})
}
@@ -45,16 +62,33 @@ func (r *EmbyMountRepository) FindByID(ctx context.Context, id string) (*model.E
func (r *EmbyMountRepository) List(ctx context.Context) ([]model.EmbyMount, error) {
var rows []model.EmbyMount
err := r.db.WithContext(ctx).Order("created_at desc").Find(&rows).Error
err := r.db.WithContext(ctx).Order("sort_order asc, created_at asc").Find(&rows).Error
return rows, err
}
func (r *EmbyMountRepository) ListByAccountID(ctx context.Context, accountID string) ([]model.EmbyMount, error) {
var rows []model.EmbyMount
err := r.db.WithContext(ctx).Where("account_id = ?", accountID).Order("created_at asc").Find(&rows).Error
err := r.db.WithContext(ctx).Where("account_id = ?", accountID).Order("sort_order asc, created_at asc").Find(&rows).Error
return rows, err
}
func (r *EmbyMountRepository) SetSortOrder(ctx context.Context, ids []string) error {
if len(ids) == 0 {
return nil
}
return withSQLiteBusyRetry(ctx, func() error {
return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
for i, id := range ids {
if err := tx.Model(&model.EmbyMount{}).Where("id = ?", id).
Update("sort_order", i).Error; err != nil {
return err
}
}
return nil
})
})
}
func (r *EmbyMountRepository) CountByAccountID(ctx context.Context, accountID string) (int64, error) {
var count int64
err := r.db.WithContext(ctx).Model(&model.EmbyMount{}).Where("account_id = ?", accountID).Count(&count).Error
+74
View File
@@ -0,0 +1,74 @@
package repository
import (
"testing"
"github.com/glebarez/sqlite"
"gorm.io/gorm"
"github.com/ShukeBta/MMTL/internal/database"
"github.com/ShukeBta/MMTL/internal/model"
)
func TestEmbyMountSortOrderAndReorder(t *testing.T) {
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err := database.AutoMigrate(db); err != nil {
t.Fatalf("migrate: %v", err)
}
repos := New(db)
ctx := t.Context()
// 1. Create mounts and verify auto-assigned sort_order
m1 := &model.EmbyMount{AccountID: "acct-1", RemoteViewID: "view-1", Name: "Mount 1"}
m2 := &model.EmbyMount{AccountID: "acct-1", RemoteViewID: "view-2", Name: "Mount 2"}
m3 := &model.EmbyMount{AccountID: "acct-1", RemoteViewID: "view-3", Name: "Mount 3"}
if err := repos.EmbyMount.Create(ctx, m1); err != nil {
t.Fatalf("create m1: %v", err)
}
if err := repos.EmbyMount.Create(ctx, m2); err != nil {
t.Fatalf("create m2: %v", err)
}
if err := repos.EmbyMount.Create(ctx, m3); err != nil {
t.Fatalf("create m3: %v", err)
}
if m1.SortOrder >= m2.SortOrder || m2.SortOrder >= m3.SortOrder {
t.Fatalf("expected ascending sort order on create: m1=%d, m2=%d, m3=%d",
m1.SortOrder, m2.SortOrder, m3.SortOrder)
}
// 2. Query list and verify initial order
list, err := repos.EmbyMount.List(ctx)
if err != nil {
t.Fatalf("list mounts: %v", err)
}
if len(list) != 3 || list[0].ID != m1.ID || list[1].ID != m2.ID || list[2].ID != m3.ID {
t.Fatalf("unexpected list order: %+v", list)
}
// 3. Reorder: m3, m1, m2
if err := repos.EmbyMount.SetSortOrder(ctx, []string{m3.ID, m1.ID, m2.ID}); err != nil {
t.Fatalf("SetSortOrder failed: %v", err)
}
// 4. Query list again and verify updated order
reordered, err := repos.EmbyMount.List(ctx)
if err != nil {
t.Fatalf("list mounts after reorder: %v", err)
}
if len(reordered) != 3 {
t.Fatalf("expected 3 mounts, got %d", len(reordered))
}
if reordered[0].ID != m3.ID || reordered[1].ID != m1.ID || reordered[2].ID != m2.ID {
t.Fatalf("expected order [m3, m1, m2], got: %s, %s, %s",
reordered[0].ID, reordered[1].ID, reordered[2].ID)
}
if reordered[0].SortOrder != 0 || reordered[1].SortOrder != 1 || reordered[2].SortOrder != 2 {
t.Fatalf("unexpected sort orders: %d, %d, %d",
reordered[0].SortOrder, reordered[1].SortOrder, reordered[2].SortOrder)
}
}
+97
View File
@@ -4,8 +4,11 @@
package service
import (
"errors"
"os"
"path/filepath"
"sort"
"strings"
"time"
)
@@ -42,3 +45,97 @@ func walkAndPrune(root string, cutoff time.Time) error {
}
return nil
}
// PruneImageCacheResult holds stats from an image cache prune operation.
type PruneImageCacheResult struct {
TotalFilesBefore int
TotalBytesBefore int64
DeletedFiles int
FreedBytes int64
RemainingBytes int64
}
type imageCacheFileEntry struct {
path string
size int64
modTime time.Time
}
// PruneImageCache scans imagesDir for cached image files. If the total disk usage
// exceeds maxSizeBytes, it removes files starting from the oldest (by ModTime)
// until disk usage falls to or below targetSizeBytes (80% of maxSizeBytes).
//
// In-flight temporary files (*.tmp) are skipped to avoid corrupting concurrent writes.
// Empty subdirectories left behind are best-effort removed.
func PruneImageCache(imagesDir string, maxSizeBytes int64) (PruneImageCacheResult, error) {
var result PruneImageCacheResult
if imagesDir == "" || maxSizeBytes <= 0 {
return result, nil
}
if _, err := os.Stat(imagesDir); err != nil {
return result, nil
}
var (
dirs []string
entries []imageCacheFileEntry
)
_ = filepath.Walk(imagesDir, func(path string, info os.FileInfo, err error) error {
if err != nil {
return nil
}
if info.IsDir() {
if path != imagesDir {
dirs = append(dirs, path)
}
return nil
}
// Skip temporary files created during image download.
name := info.Name()
if strings.HasSuffix(name, ".tmp") || strings.HasPrefix(name, "img-") && strings.Contains(name, ".tmp") {
return nil
}
size := info.Size()
result.TotalFilesBefore++
result.TotalBytesBefore += size
entries = append(entries, imageCacheFileEntry{
path: path,
size: size,
modTime: info.ModTime(),
})
return nil
})
result.RemainingBytes = result.TotalBytesBefore
if result.TotalBytesBefore <= maxSizeBytes {
return result, nil
}
// High/Low watermark: prune down to 80% of max size to leave headroom
// and prevent disk thrashing on consecutive writes.
targetSizeBytes := maxSizeBytes * 80 / 100
sort.Slice(entries, func(i, j int) bool {
return entries[i].modTime.Before(entries[j].modTime)
})
for _, entry := range entries {
if result.RemainingBytes <= targetSizeBytes {
break
}
if err := os.Remove(entry.path); err == nil || errors.Is(err, os.ErrNotExist) {
result.DeletedFiles++
result.FreedBytes += entry.size
result.RemainingBytes -= entry.size
}
}
// Clean up emptied subdirectories from deepest to shallowest.
for i := len(dirs) - 1; i >= 0; i-- {
_ = os.Remove(dirs[i])
}
return result, nil
}
+162
View File
@@ -0,0 +1,162 @@
package service
import (
"context"
"os"
"path/filepath"
"testing"
"time"
"go.uber.org/zap"
)
func TestPruneImageCache_UnderLimit(t *testing.T) {
dir := t.TempDir()
file1 := filepath.Join(dir, "img1")
file2 := filepath.Join(dir, "img2")
if err := os.WriteFile(file1, make([]byte, 100), 0o600); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(file2, make([]byte, 200), 0o600); err != nil {
t.Fatal(err)
}
// Max limit is 500 bytes, total is 300 bytes -> no prune
res, err := PruneImageCache(dir, 500)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if res.DeletedFiles != 0 {
t.Fatalf("expected 0 deleted files, got %d", res.DeletedFiles)
}
if res.TotalFilesBefore != 2 || res.TotalBytesBefore != 300 || res.RemainingBytes != 300 {
t.Fatalf("unexpected stats: %+v", res)
}
}
func TestPruneImageCache_OverLimitLRU(t *testing.T) {
dir := t.TempDir()
now := time.Now()
// Create 4 files of 100 bytes each, with distinct mtime
fOldest := filepath.Join(dir, "oldest")
fMidOld := filepath.Join(dir, "mid_old")
fMidNew := filepath.Join(dir, "mid_new")
fNewest := filepath.Join(dir, "newest")
for _, f := range []string{fOldest, fMidOld, fMidNew, fNewest} {
if err := os.WriteFile(f, make([]byte, 100), 0o600); err != nil {
t.Fatal(err)
}
}
_ = os.Chtimes(fOldest, now.Add(-4*time.Hour), now.Add(-4*time.Hour))
_ = os.Chtimes(fMidOld, now.Add(-3*time.Hour), now.Add(-3*time.Hour))
_ = os.Chtimes(fMidNew, now.Add(-2*time.Hour), now.Add(-2*time.Hour))
_ = os.Chtimes(fNewest, now.Add(-1*time.Hour), now.Add(-1*time.Hour))
// Total = 400 bytes. Max limit = 300 bytes.
// Target = 300 * 80 / 100 = 240 bytes.
// Deleting oldest (100) brings total to 300 (> 240).
// Deleting mid_old (100) brings total to 200 (<= 240).
// Total deleted = 2 files (200 bytes), remaining = 200 bytes.
res, err := PruneImageCache(dir, 300)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if res.DeletedFiles != 2 {
t.Fatalf("expected 2 deleted files, got %d", res.DeletedFiles)
}
if res.FreedBytes != 200 {
t.Fatalf("expected 200 freed bytes, got %d", res.FreedBytes)
}
if res.RemainingBytes != 200 {
t.Fatalf("expected 200 remaining bytes, got %d", res.RemainingBytes)
}
// Verify oldest and mid_old were deleted, mid_new and newest still exist
if _, err := os.Stat(fOldest); !os.IsNotExist(err) {
t.Fatalf("expected oldest file to be deleted, got err=%v", err)
}
if _, err := os.Stat(fMidOld); !os.IsNotExist(err) {
t.Fatalf("expected mid_old file to be deleted, got err=%v", err)
}
if _, err := os.Stat(fMidNew); err != nil {
t.Fatalf("expected mid_new file to exist, got err=%v", err)
}
if _, err := os.Stat(fNewest); err != nil {
t.Fatalf("expected newest file to exist, got err=%v", err)
}
}
func TestPruneImageCache_SkipsTmpFiles(t *testing.T) {
dir := t.TempDir()
fTmp := filepath.Join(dir, "img-123.tmp")
fImg := filepath.Join(dir, "cached_img")
if err := os.WriteFile(fTmp, make([]byte, 500), 0o600); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(fImg, make([]byte, 100), 0o600); err != nil {
t.Fatal(err)
}
// Limit is 200 bytes. fTmp (500) is ignored, only fImg (100) is counted <= 200.
res, err := PruneImageCache(dir, 200)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if res.DeletedFiles != 0 {
t.Fatalf("expected 0 deleted files, got %d", res.DeletedFiles)
}
if _, err := os.Stat(fTmp); err != nil {
t.Fatalf("expected tmp file to remain untouched, got %v", err)
}
}
func TestPruneImageCache_ZeroOrNegativeLimit(t *testing.T) {
dir := t.TempDir()
f := filepath.Join(dir, "img")
if err := os.WriteFile(f, make([]byte, 100), 0o600); err != nil {
t.Fatal(err)
}
res, err := PruneImageCache(dir, 0)
if err != nil || res.DeletedFiles != 0 {
t.Fatalf("expected no-op for 0 limit, got %+v, err=%v", res, err)
}
res, err = PruneImageCache(dir, -10)
if err != nil || res.DeletedFiles != 0 {
t.Fatalf("expected no-op for negative limit, got %+v, err=%v", res, err)
}
}
func TestSchedulerJobCleanImageCache(t *testing.T) {
cacheRoot := t.TempDir()
imagesDir := filepath.Join(cacheRoot, "images")
if err := os.MkdirAll(imagesDir, 0o750); err != nil {
t.Fatal(err)
}
f := filepath.Join(imagesDir, "old_poster")
if err := os.WriteFile(f, make([]byte, 2*1024*1024), 0o600); err != nil {
t.Fatal(err)
}
scheduler := NewSchedulerService(zap.NewNop(), nil, nil, nil, nil, nil, cacheRoot)
// Set limit to 1MB; our file is 2MB -> should be pruned
scheduler.SetImagesMaxSizeMBProvider(func() int {
return 1
})
if err := scheduler.jobCleanImageCache(context.Background()); err != nil {
t.Fatalf("jobCleanImageCache failed: %v", err)
}
if _, err := os.Stat(f); !os.IsNotExist(err) {
t.Fatalf("expected file to be pruned, got err=%v", err)
}
}
+103 -1
View File
@@ -343,4 +343,106 @@ func TestDanmakuSameBase(t *testing.T) {
require.Equal(t, int64(99999), resManual.EpisodeID)
require.Equal(t, "手动搜索动画B", resManual.AnimeTitle)
require.Contains(t, resManual.Raw, "手动搜索弹幕")
}
}
// Emby 远程挂载条目:通过伪装 ID 解析出流直链,通过 Range 提取 16MB 前缀计算 hash 并匹配弹幕。
func TestDanmakuFetchEmbyRemoteHashViaDirectLink(t *testing.T) {
content := bytes.Repeat([]byte("emby-remote-video-bytes-9876543210"), 300)
sum := md5.Sum(content)
wantHash := hex.EncodeToString(sum[:])
var gotRange string
var rangeHits int
rangeSrv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
rangeHits++
gotRange = r.Header.Get("Range")
w.Header().Set("Content-Type", "application/octet-stream")
_, _ = w.Write(content)
}))
t.Cleanup(rangeSrv.Close)
var seen string
official := danmakuOfficialServer(t,
`{"success":true,"isMatched":true,"matches":[{"episodeId":25484,"animeId":2001,"animeTitle":"芙莉莲","episodeTitle":"第1话"}]}`,
`<?xml version="1.0"?><i><d p="1.2,1,16777215,user1">Emby远程弹幕命中</d></i>`,
&seen)
overrideDanmakuOfficialBase(t, official.URL)
remoteMediaID := EncodeEmbyRemoteID("mount-123", "remote-item-456")
svc := newDanmakuTestService(t)
svc.SetRemoteMediaResolver(func(_ context.Context, encodedID string) (*model.Media, string, error) {
require.Equal(t, remoteMediaID, encodedID)
return &model.Media{
Base: model.Base{ID: remoteMediaID},
Title: "葬送的芙莉莲",
EpisodeTitle: "第1话",
EpisodeNum: 1,
Path: "/mnt/emby/anime/Frieren/S01E01.mkv",
SizeBytes: int64(len(content)),
DurationSec: 1400,
}, rangeSrv.URL, nil
})
ctx := context.Background()
res, err := svc.Fetch(ctx, remoteMediaID, "", "")
require.NoError(t, err)
require.True(t, res.Enabled)
require.Equal(t, "hash", res.MatchMode)
require.Equal(t, int64(25484), res.EpisodeID)
require.Equal(t, "芙莉莲", res.AnimeTitle)
require.Contains(t, res.Raw, "Emby远程弹幕命中")
require.Contains(t, gotRange, "bytes=0-")
require.Contains(t, seen, `"fileHash":"`+wantHash+`"`)
require.Contains(t, seen, `"fileName":"`+url.QueryEscape("S01E01")+`"`)
require.Equal(t, 1, rangeHits)
// 第二次拉取验证 hashCache 命中,不重复请求 rangeSrv
res2, err := svc.Fetch(ctx, remoteMediaID, "", "")
require.NoError(t, err)
require.Equal(t, "hash", res2.MatchMode)
require.Equal(t, 1, rangeHits)
}
// Emby 远程直链拉取失败时(如网络异常),能平滑降级走番剧原名/标题关键词搜索。
func TestDanmakuFetchEmbyRemoteStreamFailedFallsBackToSearch(t *testing.T) {
mux := http.NewServeMux()
mux.HandleFunc("/api/v2/search/episodes", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
// 文件名搜索 ep01 时无结果,模拟文件名未匹配
if r.URL.Query().Get("anime") == "ep01" {
fmt.Fprint(w, `{"hasMore":false,"animes":[]}`)
return
}
// 降级到番剧名搜索命中
fmt.Fprint(w, `{"hasMore":false,"animes":[{"animeId":3001,"animeTitle":"降级搜索番剧","episodes":[{"episodeId":7799,"episodeTitle":"第1话"}]}]}`)
})
mux.HandleFunc("/api/v2/comment/7799", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/xml")
fmt.Fprint(w, `<?xml version="1.0"?><i><d p="0.8,1,16777215,user1">降级搜索弹幕</d></i>`)
})
official := httptest.NewServer(mux)
t.Cleanup(official.Close)
overrideDanmakuOfficialBase(t, official.URL)
remoteMediaID := EncodeEmbyRemoteID("mount-123", "remote-item-789")
svc := newDanmakuTestService(t)
// 返回一个不存在的流服务地址模拟 Range 拉取失败
svc.SetRemoteMediaResolver(func(_ context.Context, encodedID string) (*model.Media, string, error) {
return &model.Media{
Base: model.Base{ID: remoteMediaID},
Title: "降级搜索番剧",
EpisodeNum: 1,
Path: "/mnt/emby/anime/fallback/ep01.mkv",
DurationSec: 1200,
}, "http://127.0.0.1:1/invalid-stream", nil
})
ctx := context.Background()
res, err := svc.Fetch(ctx, remoteMediaID, "", "")
require.NoError(t, err)
require.True(t, res.Enabled)
require.Equal(t, "search", res.MatchMode)
require.Equal(t, int64(7799), res.EpisodeID)
require.Equal(t, "降级搜索番剧", res.AnimeTitle)
require.Contains(t, res.Raw, "降级搜索弹幕")
}
+101 -8
View File
@@ -94,6 +94,10 @@ type DanmakuEpisode struct {
EpisodeTitle string `json:"episodeTitle"`
}
// DanmakuRemoteMediaResolver resolves an Emby remote pseudo-ID (e.g. embyremote~mount~id)
// into a memory model.Media and a direct stream URL.
type DanmakuRemoteMediaResolver func(ctx context.Context, encodedID string) (*model.Media, string, error)
// DanmakuService fetches danmaku for a media item through the dandanplay
// protocol: match by 16MB-prefix hash, then search for an episode id by the
// video's name, then fetch the comment library XML. The React player parses
@@ -108,6 +112,10 @@ type DanmakuService struct {
// StrmService.ResolvePlay; nil means strm sources are skipped.
strmResolve func(ctx context.Context, provider string, q url.Values) (*StrmPlayResult, error)
// remoteResolve resolves an Emby remote pseudo-ID into *model.Media and
// direct stream URL for range hashing.
remoteResolve DanmakuRemoteMediaResolver
hashCacheMu sync.Mutex
hashCache map[string]string // stamp → 16MB-prefix MD5
}
@@ -136,6 +144,14 @@ func (s *DanmakuService) SetStrmResolver(resolve func(ctx context.Context, provi
}
}
// SetRemoteMediaResolver wires the resolver used to fetch metadata and direct
// stream URLs for Emby remote mounted media.
func (s *DanmakuService) SetRemoteMediaResolver(resolve DanmakuRemoteMediaResolver) {
if s != nil {
s.remoteResolve = resolve
}
}
// Config reads danmaku settings from the runtime settings table.
func (s *DanmakuService) Config(ctx context.Context) DanmakuRenderConfig {
cfg := DanmakuRenderConfig{
@@ -220,13 +236,17 @@ func (s *DanmakuService) Fetch(ctx context.Context, mediaID, keyword, episodeID
target := ""
// 1) hash 识别:始终走官方 /api/v2/match(keyword 手动覆盖时跳过,直接走第 3 层)。
if target == "" && !manualKeyword && media != nil && media.Path != "" {
if target == "" && !manualKeyword && media != nil && (media.Path != "" || IsEmbyRemoteID(media.ID)) {
if hash, ok := s.mediaHash(ctx, media); ok {
fileSize := media.SizeBytes
if strings.EqualFold(filepath.Ext(media.Path), ".strm") {
if media.Path != "" && strings.EqualFold(filepath.Ext(media.Path), ".strm") {
fileSize = 0 // strm 行的 SizeBytes 是文本大小,不是视频大小
}
matches, err := s.matchOfficial(ctx, danmakuMatchFileName(media.Path), hash, fileSize, media.DurationSec)
matchName := danmakuMatchFileName(media.Path)
if matchName == "" {
matchName = term.name
}
matches, err := s.matchOfficial(ctx, matchName, hash, fileSize, media.DurationSec)
if err != nil {
s.log.Warn("danmaku hash match failed", zap.String("media_id", mediaID), zap.Error(err))
} else if len(matches) > 0 {
@@ -311,6 +331,29 @@ type danmakuSearchTerms struct {
// (movies / unknown) is left empty so the search does not filter by episode.
func (s *DanmakuService) searchTerms(ctx context.Context, mediaID string) (danmakuSearchTerms, *model.Media, error) {
var term danmakuSearchTerms
if IsEmbyRemoteID(mediaID) {
if s == nil || s.remoteResolve == nil {
return term, nil, errors.New("remote emby resolver unavailable")
}
m, _, err := s.remoteResolve(ctx, mediaID)
if err != nil || m == nil {
if err != nil {
return term, nil, err
}
return term, nil, errors.New("media not found")
}
if name := strings.TrimSpace(m.OriginalName); name != "" {
term.name = name
} else if name := strings.TrimSpace(m.Title); name != "" {
term.name = name
} else {
term.name = danmakuMatchFileName(m.Path)
}
if m.EpisodeNum > 0 {
term.episode = strconv.Itoa(m.EpisodeNum)
}
return term, m, nil
}
if s == nil || s.repo == nil || s.repo.Media == nil {
return term, nil, errors.New("media repository unavailable")
}
@@ -484,7 +527,14 @@ func (s *DanmakuService) hashCachePut(stamp, hash string) {
// ("xxx.mkv.strm") — so a second strip removes a real video extension only
// (filepath.Ext would misread names like "xxx.第01话" as having an extension).
func danmakuMatchFileName(path string) string {
base := filepath.Base(path)
if path == "" {
return ""
}
clean := strings.ReplaceAll(path, "\\", "/")
if idx := strings.LastIndex(clean, "/"); idx >= 0 {
clean = clean[idx+1:]
}
base := filepath.Base(clean)
if ext := filepath.Ext(base); ext != "" {
base = strings.TrimSuffix(base, ext)
}
@@ -493,14 +543,20 @@ func danmakuMatchFileName(path string) string {
base = strings.TrimSuffix(base, filepath.Ext(base))
}
}
return base
return strings.TrimSpace(base)
}
// mediaHash returns the dandanplay match hash (MD5 of the first 16MB of the
// video). Local videos are hashed straight from disk; .strm indirections are
// resolved (local path / direct link) and only the 16MB prefix is downloaded.
// video). Local videos are hashed straight from disk; .strm indirections and
// remote Emby streams are range-fetched and only the 16MB prefix is downloaded.
func (s *DanmakuService) mediaHash(ctx context.Context, media *model.Media) (string, bool) {
if media == nil || media.Path == "" {
if media == nil {
return "", false
}
if IsEmbyRemoteID(media.ID) {
return s.hashEmbyRemote(ctx, media)
}
if media.Path == "" {
return "", false
}
if strings.EqualFold(filepath.Ext(media.Path), ".strm") {
@@ -517,6 +573,43 @@ func (s *DanmakuService) mediaHash(ctx context.Context, media *model.Media) (str
return s.hashLocalFile(media.Path)
}
// hashEmbyRemote computes the 16MB-prefix MD5 of a remote Emby stream via HTTP Range.
func (s *DanmakuService) hashEmbyRemote(ctx context.Context, media *model.Media) (string, bool) {
if media == nil || media.ID == "" {
return "", false
}
if h, ok := s.hashCacheGet("e|" + media.ID); ok {
return h, true
}
if s.remoteResolve == nil {
return "", false
}
_, streamURL, err := s.remoteResolve(ctx, media.ID)
if err != nil || strings.TrimSpace(streamURL) == "" {
if err != nil {
s.log.Warn("danmaku emby stream url resolve failed, hash layer skipped",
zap.String("media_id", media.ID), zap.Error(err))
}
return "", false
}
body, err := s.openRangeBody(ctx, streamURL, nil)
if err != nil || body == nil {
if err != nil {
s.log.Warn("danmaku emby range fetch failed, hash layer skipped",
zap.String("media_id", media.ID), zap.Error(err))
}
return "", false
}
defer body.Close()
h := md5.New()
if _, err := io.Copy(h, io.LimitReader(body, danmakuHashPrefixBytes)); err != nil {
return "", false
}
hash := hex.EncodeToString(h.Sum(nil))
s.hashCachePut("e|"+media.ID, hash)
return hash, true
}
// hashLocalFile computes the MD5 of the first 16MB of a local video, cached
// by path+size+mtime so repeated danmaku loads skip the disk read.
func (s *DanmakuService) hashLocalFile(path string) (string, bool) {
+12
View File
@@ -231,6 +231,18 @@ func (r *EmbyRemoteService) DeleteMount(ctx context.Context, id string) error {
return err
}
// ReorderMounts 批量重排挂载媒体库顺序。
func (r *EmbyRemoteService) ReorderMounts(ctx context.Context, ids []string) error {
if len(ids) == 0 {
return nil
}
if err := r.repo.EmbyMount.SetSortOrder(ctx, ids); err != nil {
return err
}
r.invalidateRemoteMediaCache(ctx)
return nil
}
// FullMountAccount 把账号的全部远程媒体库(View)挂载进来(幂等,已存在跳过)。
func (r *EmbyRemoteService) FullMountAccount(ctx context.Context, acct *model.StrmAccount, proxyPlayDefault bool) (int, error) {
views, err := r.RemoteViews(ctx, acct)
+53 -32
View File
@@ -36,23 +36,28 @@ func (r *EmbyRemoteService) RemoteLibraries(ctx context.Context) ([]RemoteLibrar
if err != nil || len(mounts) == 0 {
return nil, err
}
// 按账号分组,每账号拉一次 Views 做匹配。
byAccount := map[string][]*model.EmbyMount{}
type accountData struct {
acct *model.StrmAccount
cfg *EmbyRemoteConfig
viewByName map[string]map[string]any
}
acctData := map[string]*accountData{}
for i := range mounts {
m := mounts[i]
m := &mounts[i]
if !m.Enabled {
continue
}
byAccount[m.AccountID] = append(byAccount[m.AccountID], &mounts[i])
}
out := make([]RemoteLibraryView, 0, len(mounts))
for accountID, accountMounts := range byAccount {
acct := r.AccountByID(ctx, accountID)
if _, ok := acctData[m.AccountID]; ok {
continue
}
acct := r.AccountByID(ctx, m.AccountID)
if acct == nil {
acctData[m.AccountID] = nil
continue
}
cfg, cfgErr := r.configOf(acct)
if cfgErr != nil {
acctData[m.AccountID] = nil
continue
}
views, viewErr := r.RemoteViews(ctx, acct)
@@ -61,31 +66,47 @@ func (r *EmbyRemoteService) RemoteLibraries(ctx context.Context) ([]RemoteLibrar
r.log.Warn("web remote emby views failed",
zap.String("account", acct.Name), zap.Error(viewErr))
}
acctData[m.AccountID] = nil
continue
}
viewByName := map[string]map[string]any{}
for _, v := range views {
viewByName[remoteItemString(v, "Id")] = v
}
for _, mount := range accountMounts {
v, ok := viewByName[mount.RemoteViewID]
if !ok {
continue // 远程已删除该媒体库
}
lib := r.mapRemoteMountToLibrary(mount, acct, cfg, v)
if lib == nil {
continue
}
out = append(out, RemoteLibraryView{
Library: *lib,
MountID: mount.ID,
AccountID: acct.ID,
RemoteID: mount.RemoteViewID,
CollectionType: mount.CollectionType,
AccountName: acct.Name,
})
acctData[m.AccountID] = &accountData{
acct: acct,
cfg: cfg,
viewByName: viewByName,
}
}
out := make([]RemoteLibraryView, 0, len(mounts))
for i := range mounts {
m := &mounts[i]
if !m.Enabled {
continue
}
data := acctData[m.AccountID]
if data == nil || data.viewByName == nil {
continue
}
v, ok := data.viewByName[m.RemoteViewID]
if !ok {
continue // 远程已删除该媒体库
}
lib := r.mapRemoteMountToLibrary(m, data.acct, data.cfg, v)
if lib == nil {
continue
}
out = append(out, RemoteLibraryView{
Library: *lib,
MountID: m.ID,
AccountID: data.acct.ID,
RemoteID: m.RemoteViewID,
CollectionType: m.CollectionType,
AccountName: data.acct.Name,
})
}
return out, nil
}
@@ -125,13 +146,13 @@ func (r *EmbyRemoteService) mapRemoteMountToLibrary(mount *model.EmbyMount, acct
case "music":
libType = "music"
}
lib := &model.Library{
Base: model.Base{ID: EncodeEmbyRemoteID(mount.ID, mount.RemoteViewID)},
Name: name,
Type: libType,
Enabled: true,
SortOrder: 1000, // 远程库排在本地库之后
}
lib := &model.Library{
Base: model.Base{ID: EncodeEmbyRemoteID(mount.ID, mount.RemoteViewID)},
Name: name,
Type: libType,
Enabled: true,
SortOrder: 1000 + mount.SortOrder, // 远程库排在本地库之后,且保持挂载库排序
}
// 远程媒体库封面只有真实存在图片标签才下发。
if remoteItemHasImageTag(item, "Primary") {
lib.CoverURL = r.remoteItemImageURL(cfg, mount.RemoteViewID, "Primary")
+10
View File
@@ -84,3 +84,13 @@ func (p *ImageProxy) libraryRoots() []string {
p.libRootsAt = time.Now()
return p.libRootsCache
}
// Prune removes oldest cached images until disk usage is within the configured limit.
func (p *ImageProxy) Prune() (PruneImageCacheResult, error) {
if p.cfg == nil || p.cfg.Cache.ImagesMaxSizeMB <= 0 {
return PruneImageCacheResult{}, nil
}
maxBytes := int64(p.cfg.Cache.ImagesMaxSizeMB) * 1024 * 1024
return PruneImageCache(p.cacheDir, maxBytes)
}
+12 -5
View File
@@ -106,11 +106,18 @@ func ApplyRuntimeSetting(cfg *config.Config, key, value string) {
cfg.App.SSLCert = value
case "https.key":
cfg.App.SSLKey = value
case "https.cert_path":
cfg.App.SSLCertPath = strings.TrimSpace(value)
case "https.key_path":
cfg.App.SSLKeyPath = strings.TrimSpace(value)
}
case "https.cert_path":
cfg.App.SSLCertPath = strings.TrimSpace(value)
case "https.key_path":
cfg.App.SSLKeyPath = strings.TrimSpace(value)
case "cache.images_max_size_mb":
if n, err := strconv.Atoi(value); err == nil {
if n < 0 {
n = 0
}
cfg.Cache.ImagesMaxSizeMB = n
}
}
}
// ParseBoolSetting is the exported variant of parseBoolSetting for handlers
+18
View File
@@ -41,6 +41,8 @@ type SchedulerService struct {
cacheDir string
now func() time.Time
imagesMaxSizeMBProvider func() int
mu sync.Mutex
stopCh chan struct{}
jobs []*scheduledJob
@@ -59,6 +61,17 @@ func (s *SchedulerService) SetOrganizePipeline(pipeline *OrganizePipelineService
s.organizePipeline = pipeline
}
func (s *SchedulerService) SetImagesMaxSizeMBProvider(fn func() int) {
s.imagesMaxSizeMBProvider = fn
}
func (s *SchedulerService) imagesMaxSizeMB() int {
if s.imagesMaxSizeMBProvider != nil {
return s.imagesMaxSizeMBProvider()
}
return 0
}
// scheduledJob is one recurring task.
type scheduledJob struct {
name string
@@ -120,6 +133,11 @@ func (s *SchedulerService) Start(ctx context.Context) {
interval: 24 * time.Hour,
run: s.jobPurgeRecycleBin,
},
{
name: "image_cache_cleanup",
interval: 1 * time.Hour,
run: s.jobCleanImageCache,
},
}
for _, j := range s.jobs {
initialDelay := 15 * time.Second
+30
View File
@@ -2,6 +2,7 @@ package service
import (
"context"
"path/filepath"
"strconv"
"strings"
"time"
@@ -202,3 +203,32 @@ func isMissingTableErr(err error) bool {
}
return err == gorm.ErrInvalidDB
}
// jobCleanImageCache prunes image proxy cache files when disk usage exceeds the configured limit.
func (s *SchedulerService) jobCleanImageCache(ctx context.Context) error {
if s.cacheDir == "" {
return nil
}
maxMB := s.imagesMaxSizeMB()
if maxMB <= 0 {
return nil
}
imagesDir := filepath.Join(s.cacheDir, "images")
maxSizeBytes := int64(maxMB) * 1024 * 1024
res, err := PruneImageCache(imagesDir, maxSizeBytes)
if err != nil {
if s.log != nil {
s.log.Warn("scheduled image cache cleanup failed", zap.Error(err))
}
return err
}
if res.DeletedFiles > 0 && s.log != nil {
s.log.Info("scheduled image cache cleanup completed",
zap.Int("deleted_files", res.DeletedFiles),
zap.Int64("freed_bytes", res.FreedBytes),
zap.Int64("remaining_bytes", res.RemainingBytes),
)
}
return nil
}
+30
View File
@@ -2,11 +2,13 @@ package service
import (
"context"
"errors"
"strings"
"go.uber.org/zap"
"github.com/ShukeBta/MMTL/internal/config"
"github.com/ShukeBta/MMTL/internal/model"
"github.com/ShukeBta/MMTL/internal/repository"
)
@@ -117,6 +119,28 @@ func (b *serviceContainerBuilder) initContentServices() {
b.c.FFTools = NewFFmpegToolsService(b.cfg, b.log, b.repos)
// 弹幕 hash 识别需要把 strm 指向解析成可拉取的直链/本地路径。
b.c.Danmaku.SetStrmResolver(b.c.Strm.ResolvePlay)
// 弹幕识别需要把远程 Emby 条目解析为 Media 元数据及可拉取前 16MB 的直链 URL。
if b.c.EmbyRemote != nil {
b.c.Danmaku.SetRemoteMediaResolver(func(ctx context.Context, encodedID string) (*model.Media, string, error) {
mountID, remoteID, ok := DecodeEmbyRemoteID(encodedID)
if !ok {
return nil, "", errors.New("invalid emby remote id")
}
mount, acct, err := b.c.EmbyRemote.ResolveMount(ctx, mountID)
if err != nil || mount == nil || acct == nil {
if err != nil {
return nil, "", err
}
return nil, "", errors.New("emby mount or account not found")
}
m, err := b.c.EmbyRemote.RemoteMediaDetail(ctx, mount, acct, remoteID)
if err != nil || m == nil {
return nil, "", err
}
streamURL, _ := b.c.EmbyRemote.WebStreamURL(ctx, acct, remoteID)
return m, streamURL, nil
})
}
}
func (b *serviceContainerBuilder) initAccessAndStorageServices() {
@@ -131,6 +155,12 @@ func (b *serviceContainerBuilder) initAccessAndStorageServices() {
)
b.c.Scheduler.SetTaskTracker(b.c.Tasks)
b.c.Scheduler.SetOrganizePipeline(b.c.OrganizePipeline)
b.c.Scheduler.SetImagesMaxSizeMBProvider(func() int {
if b.cfg == nil {
return 0
}
return b.cfg.Cache.ImagesMaxSizeMB
})
}
func (b *serviceContainerBuilder) initIdentityServices() {