缓存弹幕抓取结果并合并并发请求

This commit is contained in:
truewhile
2026-09-12 11:02:24 +08:00
parent 6797394c5b
commit 8f64f961ac
4 changed files with 277 additions and 6 deletions
+1 -1
View File
@@ -20,6 +20,7 @@ require (
go.uber.org/zap v1.27.0
golang.org/x/crypto v0.49.0
golang.org/x/image v0.37.0
golang.org/x/sync v0.20.0
golang.org/x/sys v0.42.0
golang.org/x/time v0.15.0
gopkg.in/yaml.v3 v3.0.1
@@ -88,7 +89,6 @@ require (
golang.org/x/arch v0.25.0 // indirect
golang.org/x/exp v0.0.0-20251023183803-a4bb9ffd2546 // indirect
golang.org/x/net v0.52.0 // indirect
golang.org/x/sync v0.20.0 // indirect
golang.org/x/text v0.35.0 // indirect
google.golang.org/protobuf v1.36.11 // indirect
gopkg.in/ini.v1 v1.67.0 // indirect
+35
View File
@@ -824,3 +824,38 @@ func TestDanmakuFetchEmbyRemoteStreamFailedFallsBackToSearch(t *testing.T) {
require.Equal(t, "降级搜索番剧", res.AnimeTitle)
require.Contains(t, res.Raw, "降级搜索弹幕")
}
// 即使刮削元数据完整,也应先走准确率最高的 hash 层。
func TestDanmakuFetchPrefersHashOverCompleteMetadata(t *testing.T) {
videoPath, wantHash := writeDanmakuTestVideo(t, "测试动画.第01话.mkv")
var seen string
official := danmakuOfficialServer(t,
`{"success":true,"isMatched":true,"matches":[{"episodeId":25484,"animeId":1001,"animeTitle":"官方测试动画","episodeTitle":"第1话"}]}`,
`<?xml version="1.0"?><i><d p="0.5,1,16777215,user1">hash命中</d></i>`,
&seen)
overrideDanmakuOfficialBase(t, official.URL)
source := newDanmakuSourceServerWithSearch(t,
`{"hasMore":false,"animes":[{"animeId":2002,"animeTitle":"官方测试动画","episodes":[{"episodeId":25484,"episodeTitle":"第1话"}]}]}`)
svc := newDanmakuTestService(t)
ctx := context.Background()
require.NoError(t, svc.repo.Setting.Set(ctx, DanmakuSourceKey, source.URL()))
m := model.Media{
Title: "测试动画",
Path: videoPath,
SizeBytes: 123,
EpisodeNum: 1,
EpisodeTitle: "起始与终结的序章",
Year: 2011,
}
m.ID = "hash-first"
require.NoError(t, svc.repo.DB.Create(&m).Error)
res, err := svc.Fetch(ctx, "hash-first", "", "")
require.NoError(t, err)
require.Equal(t, "hash", res.MatchMode)
require.NotEmpty(t, res.Raw)
require.Contains(t, seen, wantHash)
}
+163 -5
View File
@@ -20,6 +20,8 @@ import (
"time"
"unicode"
"golang.org/x/sync/singleflight"
"go.uber.org/zap"
"github.com/truewhile/MeBox/internal/model"
@@ -136,6 +138,12 @@ type DanmakuService struct {
hashCacheMu sync.Mutex
hashCache map[string]string // stamp → 16MB-prefix MD5
// resultCache 缓存整条弹幕抓取结果,避免同一集重复播放时重走
// 「16MB 哈希 + 上游搜索 + 评论拉取」这条高延迟链路。
resultCacheMu sync.Mutex
resultCache map[string]danmakuResultCacheEntry
fetchGroup singleflight.Group
}
func danmakuHTTPClient() *http.Client {
@@ -147,10 +155,11 @@ func NewDanmakuService(log *zap.Logger, repo *repository.Container) *DanmakuServ
log = zap.NewNop()
}
return &DanmakuService{
log: log,
repo: repo,
client: danmakuHTTPClient(),
hashCache: make(map[string]string),
log: log,
repo: repo,
client: danmakuHTTPClient(),
hashCache: make(map[string]string),
resultCache: make(map[string]danmakuResultCacheEntry),
}
}
@@ -170,6 +179,121 @@ func (s *DanmakuService) SetRemoteMediaResolver(resolve DanmakuRemoteMediaResolv
}
}
// 弹幕抓取结果缓存:同一集在 TTL 内重复播放时直接返回,避免重复执行
// 「16MB 前置哈希 → 上游搜索 → 评论拉取」这条高延迟链路。弹幕库内容变化
// 很慢,而用户通常会在短时间内反复切集/回看,因此 TTL 取 6 小时。
const (
danmakuResultCacheTTL = 6 * time.Hour
danmakuResultCacheMaxEntries = 64
// 候选列表由上游搜索决定,相对稳定但可能随源更新,缓存 10 分钟。
danmakuResultCacheCandidateTTL = 10 * time.Minute
// 空结果可能只是上游临时抖动,只短暂缓存,避免长时间看不到弹幕。
danmakuResultCacheEmptyTTL = time.Minute
)
type danmakuResultCacheEntry struct {
result *DanmakuFetchResult
expiresAt time.Time
storedAt time.Time
}
func danmakuResultCacheKey(source, mediaID, keyword, episodeID string, merge bool) string {
return strings.Join([]string{
strings.TrimSpace(source),
mediaID,
strings.TrimSpace(keyword),
strings.TrimSpace(episodeID),
strconv.FormatBool(merge),
}, "\x00")
}
func (s *DanmakuService) resultCacheGet(key string) (*DanmakuFetchResult, bool) {
if s == nil || key == "" {
return nil, false
}
now := time.Now()
s.resultCacheMu.Lock()
defer s.resultCacheMu.Unlock()
entry, ok := s.resultCache[key]
if !ok {
return nil, false
}
if now.After(entry.expiresAt) {
delete(s.resultCache, key)
return nil, false
}
return cloneDanmakuFetchResult(entry.result), true
}
func (s *DanmakuService) resultCachePut(key string, result *DanmakuFetchResult) {
if s == nil || key == "" || result == nil {
return
}
now := time.Now()
s.resultCacheMu.Lock()
defer s.resultCacheMu.Unlock()
if s.resultCache == nil {
s.resultCache = make(map[string]danmakuResultCacheEntry)
}
ttl := danmakuResultCacheTTLFor(result)
if ttl <= 0 {
return
}
if _, exists := s.resultCache[key]; !exists && len(s.resultCache) >= danmakuResultCacheMaxEntries {
oldestKey := ""
var oldest time.Time
for k, entry := range s.resultCache {
if oldestKey == "" || entry.storedAt.Before(oldest) {
oldestKey, oldest = k, entry.storedAt
}
}
delete(s.resultCache, oldestKey)
}
s.resultCache[key] = danmakuResultCacheEntry{
result: cloneDanmakuFetchResult(result),
expiresAt: now.Add(ttl),
storedAt: now,
}
}
// danmakuResultCacheTTLFor 按结果完整性选择缓存时长:拿到弹幕正文才值得
// 长缓存;只有候选列表时短缓存;空结果只缓存一分钟。
func danmakuResultCacheTTLFor(result *DanmakuFetchResult) time.Duration {
if result == nil {
return 0
}
if strings.TrimSpace(result.Raw) != "" {
return danmakuResultCacheTTL
}
if len(result.Candidates) > 0 {
return danmakuResultCacheCandidateTTL
}
return danmakuResultCacheEmptyTTL
}
// cloneDanmakuFetchResult 深拷贝切片字段,避免缓存命中后调用方修改共享数据。
func cloneDanmakuFetchResult(in *DanmakuFetchResult) *DanmakuFetchResult {
if in == nil {
return nil
}
out := *in
out.Candidates = cloneDanmakuAnimeList(in.Candidates)
out.Alternatives = cloneDanmakuAnimeList(in.Alternatives)
return &out
}
func cloneDanmakuAnimeList(in []DanmakuAnime) []DanmakuAnime {
if in == nil {
return nil
}
out := make([]DanmakuAnime, len(in))
for i, anime := range in {
out[i] = anime
out[i].Episodes = append([]DanmakuEpisode(nil), anime.Episodes...)
}
return out
}
// Config reads danmaku settings from the runtime settings table.
func (s *DanmakuService) Config(ctx context.Context) DanmakuRenderConfig {
cfg := DanmakuRenderConfig{
@@ -257,8 +381,42 @@ func (s *DanmakuService) Fetch(ctx context.Context, mediaID, keyword, episodeID
return s.FetchWithOptions(ctx, mediaID, keyword, episodeID, DanmakuFetchOptions{})
}
// FetchWithOptions 是 Fetch 的带偏好版本。
// FetchWithOptions 是 Fetch 的带偏好版本。结果按「配置源 + 媒体 + 关键词 +
// 指定集 + 合并开关」缓存,并用 singleflight 合并并发请求,避免同一集被重复抓取。
func (s *DanmakuService) FetchWithOptions(ctx context.Context, mediaID, keyword, episodeID string, opts DanmakuFetchOptions) (*DanmakuFetchResult, error) {
if s == nil {
return nil, errors.New("danmaku service unavailable")
}
cfg := s.Config(ctx)
if !cfg.Enabled {
return &DanmakuFetchResult{DanmakuRenderConfig: cfg, SourceType: "auto"}, nil
}
key := danmakuResultCacheKey(cfg.Source, mediaID, keyword, episodeID, opts.MergeSources)
if cached, ok := s.resultCacheGet(key); ok {
return cached, nil
}
value, err, _ := s.fetchGroup.Do(key, func() (any, error) {
// 等待期间可能已有同一 key 的请求写入缓存。
if cached, ok := s.resultCacheGet(key); ok {
return cached, nil
}
res, err := s.fetchWithOptionsUncached(ctx, mediaID, keyword, episodeID, opts)
if err != nil {
// 与原实现一致:失败时仍把已填充的渲染配置/匹配信息交给调用方。
return res, err
}
s.resultCachePut(key, res)
return res, nil
})
res, _ := value.(*DanmakuFetchResult)
if err != nil {
return cloneDanmakuFetchResult(res), err
}
return cloneDanmakuFetchResult(res), nil
}
// fetchWithOptionsUncached 是未命中缓存时执行的原始抓取流程。
func (s *DanmakuService) fetchWithOptionsUncached(ctx context.Context, mediaID, keyword, episodeID string, opts DanmakuFetchOptions) (*DanmakuFetchResult, error) {
res := &DanmakuFetchResult{DanmakuRenderConfig: s.Config(ctx), SourceType: "auto"}
if !res.Enabled {
return res, nil
+78
View File
@@ -6,6 +6,8 @@ import (
"net/http"
"net/http/httptest"
"path/filepath"
"sync"
"sync/atomic"
"testing"
"github.com/glebarez/sqlite"
@@ -326,3 +328,79 @@ func TestDanmakuFetchDetectsJSONSource(t *testing.T) {
require.Equal(t, "xml", res2.SourceType)
require.Contains(t, res2.Raw, "弹幕A")
}
// 同一集在缓存 TTL 内重复请求时不应再访问上游(搜索和评论都只发一次)。
func TestDanmakuFetchCachesRepeatedRequests(t *testing.T) {
var searchCalls, commentCalls int32
mux := http.NewServeMux()
mux.HandleFunc("/api/v2/search/episodes", func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(&searchCalls, 1)
w.Header().Set("Content-Type", "application/json")
fmt.Fprint(w, `{"hasMore":false,"animes":[{"animeId":1001,"animeTitle":"测试动画","episodes":[{"episodeId":25484,"episodeTitle":"第1话"}]}]}`)
})
mux.HandleFunc("/api/v2/comment/25484", func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(&commentCalls, 1)
w.Header().Set("Content-Type", "application/xml")
fmt.Fprint(w, `<?xml version="1.0"?><i><d p="0.5,1,16777215,user1">缓存测试</d></i>`)
})
srv := httptest.NewServer(mux)
defer srv.Close()
svc := newDanmakuTestService(t)
ctx := context.Background()
require.NoError(t, svc.repo.Setting.Set(ctx, DanmakuSourceKey, srv.URL))
seedDanmakuMedia(t, svc, "cache-media", "测试动画", "", 1)
first, err := svc.Fetch(ctx, "cache-media", "", "")
require.NoError(t, err)
require.NotEmpty(t, first.Raw)
second, err := svc.Fetch(ctx, "cache-media", "", "")
require.NoError(t, err)
require.Equal(t, first.Raw, second.Raw)
require.EqualValues(t, 1, atomic.LoadInt32(&searchCalls))
require.EqualValues(t, 1, atomic.LoadInt32(&commentCalls))
}
// 并发的同一集请求应被 singleflight 合并,上游只被访问一次。
func TestDanmakuFetchCoalescesConcurrentRequests(t *testing.T) {
var searchCalls int32
mux := http.NewServeMux()
mux.HandleFunc("/api/v2/search/episodes", func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(&searchCalls, 1)
w.Header().Set("Content-Type", "application/json")
fmt.Fprint(w, `{"hasMore":false,"animes":[{"animeId":1001,"animeTitle":"测试动画","episodes":[{"episodeId":25484,"episodeTitle":"第1话"}]}]}`)
})
mux.HandleFunc("/api/v2/comment/25484", 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.5,1,16777215,user1">并发测试</d></i>`)
})
srv := httptest.NewServer(mux)
defer srv.Close()
svc := newDanmakuTestService(t)
ctx := context.Background()
require.NoError(t, svc.repo.Setting.Set(ctx, DanmakuSourceKey, srv.URL))
seedDanmakuMedia(t, svc, "concurrent-media", "测试动画", "", 1)
const workers = 8
start := make(chan struct{})
errs := make(chan error, workers)
var wg sync.WaitGroup
for i := 0; i < workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
<-start
_, err := svc.Fetch(ctx, "concurrent-media", "", "")
errs <- err
}()
}
close(start)
wg.Wait()
close(errs)
for err := range errs {
require.NoError(t, err)
}
require.EqualValues(t, 1, atomic.LoadInt32(&searchCalls))
}