This commit is contained in:
truewhile
2026-09-03 16:53:12 +08:00
parent 4150d03852
commit c702a68cfb
24 changed files with 691 additions and 168 deletions
+2 -2
View File
@@ -43,8 +43,8 @@ func TestLoadDefaults(t *testing.T) {
if cfg.Cache.RedisPrefix != "mebox" {
t.Fatalf("expected default redis prefix, got %q", cfg.Cache.RedisPrefix)
}
if cfg.Cache.MediaTTLSeconds != 15 {
t.Fatalf("expected default media cache ttl 15, got %d", cfg.Cache.MediaTTLSeconds)
if cfg.Cache.MediaTTLSeconds != 90 {
t.Fatalf("expected default media cache ttl 90, got %d", cfg.Cache.MediaTTLSeconds)
}
if cfg.Search.Index != "mebox_media" {
t.Fatalf("expected default search index, got %q", cfg.Search.Index)
+1 -1
View File
@@ -47,7 +47,7 @@ func setDefaults(v *viper.Viper) {
v.SetDefault("cache.cleanup_interval_min", 60)
v.SetDefault("cache.redis_url", "")
v.SetDefault("cache.redis_prefix", "mebox")
v.SetDefault("cache.media_ttl_seconds", 15)
v.SetDefault("cache.media_ttl_seconds", 90)
v.SetDefault("search.backend", "")
v.SetDefault("search.opensearch_url", "")
+1 -1
View File
@@ -51,7 +51,7 @@ func (c *Config) normalize() error {
c.Cache.RedisPrefix = "mebox"
}
if c.Cache.MediaTTLSeconds < 1 {
c.Cache.MediaTTLSeconds = 15
c.Cache.MediaTTLSeconds = 90
}
c.Search.Backend = strings.ToLower(strings.TrimSpace(c.Search.Backend))
if c.Search.Index == "" {
+9 -1
View File
@@ -60,8 +60,16 @@ func ensurePostgresColumnCompatibility(db *gorm.DB) error {
func ensurePerformanceIndexes(db *gorm.DB) error {
statements := []string{
// 完整多级排序索引:媒体库分页与首页预览的 ORDER BY
// (release_date, year, updated_at, created_at, id) DESC 与索引列完全一致,
// LIMIT 分页沿索引顺序直取,免去对整库行做临时 B-tree 排序。
`CREATE INDEX IF NOT EXISTS idx_media_library_recent_active ON media(library_id, release_date DESC, year DESC, updated_at DESC, created_at DESC, id DESC) WHERE deleted_at IS NULL`,
// 计数覆盖索引:首页 CountByLibraries 的 GROUP BY library_id + nsfw 谓词
// 全部落在索引键/部分索引条件上,纯索引扫描即可完成,不回表。
`CREATE INDEX IF NOT EXISTS idx_media_library_nsfw_active ON media(library_id, nsfw) WHERE deleted_at IS NULL`,
// 旧的两键前缀索引被上面的完整排序索引完全覆盖,删除以降低写放大。
`DROP INDEX IF EXISTS idx_media_library_release_active`,
`CREATE INDEX IF NOT EXISTS idx_media_library_created_active ON media(library_id, created_at DESC) WHERE deleted_at IS NULL`,
`CREATE INDEX IF NOT EXISTS idx_media_library_release_active ON media(library_id, release_date DESC, year DESC) WHERE deleted_at IS NULL`,
`CREATE INDEX IF NOT EXISTS idx_media_library_episode_active ON media(library_id, season_num, episode_num, created_at DESC) WHERE deleted_at IS NULL`,
`CREATE INDEX IF NOT EXISTS idx_media_library_root_active ON media(library_id, library_root_id) WHERE deleted_at IS NULL`,
`CREATE INDEX IF NOT EXISTS idx_media_series_active ON media(series_id, season_num, episode_num) WHERE deleted_at IS NULL`,
+2
View File
@@ -9,12 +9,14 @@ import (
"go.uber.org/zap"
"github.com/truewhile/MeBox/internal/config"
"github.com/truewhile/MeBox/internal/middleware"
"github.com/truewhile/MeBox/internal/service"
)
// Register attaches every API route to the engine.
func Register(r *gin.Engine, cfg *config.Config, log *zap.Logger, svc *service.Container) {
api := r.Group("/api")
api.Use(middleware.GzipAPI())
{
api.GET("/health", healthCheck)
api.GET("/version", versionInfo)
+47
View File
@@ -0,0 +1,47 @@
package middleware
import (
"github.com/gin-contrib/gzip"
"github.com/gin-gonic/gin"
)
// gzipExcludedExtensions 已经是压缩格式(图片/字体/媒体/归档)的响应体,
// 再 gzip 只浪费 CPU 不省流量。
var gzipExcludedExtensions = []string{
".png", ".gif", ".jpeg", ".jpg", ".webp", ".avif", ".ico", ".svg",
".woff", ".woff2", ".ttf", ".otf",
".mp4", ".mkv", ".webm", ".ts", ".m4s", ".m3u8",
".mp3", ".flac", ".aac", ".ogg",
".zip", ".gz", ".xz", ".7z", ".rar",
}
// apiGzipExcludedPrefixes 大文件流式传输(Range 语义)与 WS/SSE 长连接
// 不参与 gzip:压缩会破坏 Range / 逐块推送语义。
var apiGzipExcludedPrefixes = []string{
"/api/stream/",
"/api/hls/",
"/api/img",
"/api/subtitles/",
"/api/strm/play/",
"/api/ws",
"/api/events",
}
// GzipAPI 压缩 /api 下的 JSON 响应(媒体列表动辄数 MB,压缩率 85%+)。
// gin-contrib/gzip 的路径排除是前缀匹配,且会自动校验 Accept-Encoding
// 与 Connection: Upgrade(WebSocket 安全)。
func GzipAPI() gin.HandlerFunc {
return gzip.Gzip(
gzip.DefaultCompression,
gzip.WithExcludedExtensions(gzipExcludedExtensions),
gzip.WithExcludedPaths(apiGzipExcludedPrefixes),
)
}
// GzipStatic 压缩 SPA 静态资源(JS/CSS/HTML 是构建产物的大头)。
func GzipStatic() gin.HandlerFunc {
return gzip.Gzip(
gzip.DefaultCompression,
gzip.WithExcludedExtensions(gzipExcludedExtensions),
)
}
+84
View File
@@ -0,0 +1,84 @@
package middleware
import (
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/gin-gonic/gin"
)
func newGzipTestRouter() *gin.Engine {
gin.SetMode(gin.TestMode)
router := gin.New()
api := router.Group("/api")
api.Use(GzipAPI())
api.GET("/libraries", func(c *gin.Context) {
c.JSON(http.StatusOK, gin.H{"items": strings.Repeat("mebox", 200)})
})
api.GET("/stream/:id", func(c *gin.Context) {
c.String(http.StatusOK, strings.Repeat("video-bytes", 200))
})
router.GET("/assets/app.js", GzipStatic(), func(c *gin.Context) {
c.Data(http.StatusOK, "text/javascript", []byte(strings.Repeat("console.log(1);", 200)))
})
return router
}
func TestGzipAPICompressesJSONWhenAccepted(t *testing.T) {
router := newGzipTestRouter()
req := httptest.NewRequest(http.MethodGet, "/api/libraries", nil)
req.Header.Set("Accept-Encoding", "gzip")
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if w.Code != http.StatusOK {
t.Fatalf("status = %d, want 200", w.Code)
}
if got := w.Header().Get("Content-Encoding"); got != "gzip" {
t.Fatalf("Content-Encoding = %q, want gzip", got)
}
if raw := strings.Repeat("mebox", 200); w.Body.Len() >= len(raw) {
t.Fatalf("body not compressed: len = %d, raw = %d", w.Body.Len(), len(raw))
}
}
func TestGzipAPISkipsWhenNotAccepted(t *testing.T) {
router := newGzipTestRouter()
req := httptest.NewRequest(http.MethodGet, "/api/libraries", nil)
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if got := w.Header().Get("Content-Encoding"); got == "gzip" {
t.Fatal("Content-Encoding should not be gzip without Accept-Encoding")
}
}
func TestGzipAPIExcludesStreamPath(t *testing.T) {
router := newGzipTestRouter()
req := httptest.NewRequest(http.MethodGet, "/api/stream/abc", nil)
req.Header.Set("Accept-Encoding", "gzip")
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if got := w.Header().Get("Content-Encoding"); got == "gzip" {
t.Fatal("stream path must not be gzipped (Range semantics)")
}
}
func TestGzipStaticCompressesAssets(t *testing.T) {
router := newGzipTestRouter()
req := httptest.NewRequest(http.MethodGet, "/assets/app.js", nil)
req.Header.Set("Accept-Encoding", "gzip")
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if got := w.Header().Get("Content-Encoding"); got != "gzip" {
t.Fatalf("Content-Encoding = %q, want gzip for static assets", got)
}
}
+32 -10
View File
@@ -23,12 +23,34 @@ import (
// 永远捞不到数据。
func (r *MediaRepository) Upsert(ctx context.Context, m *model.Media) error {
return withSQLiteBusyRetry(ctx, func() error {
return r.upsert(ctx, m)
return r.upsertWithDB(ctx, r.db, m)
})
}
func (r *MediaRepository) upsert(ctx context.Context, m *model.Media) error {
existing, created, err := r.findOrCreateMediaByPath(ctx, m)
// UpsertBatch 在单个事务里逐条执行 Upsert:扫描一批只提交(fsync)一次,
// 而不是每条一个隐式事务。任一条目落库失败不影响批内已成功的条目——
// 事务回滚后由调用方退回逐条 Upsert 兜底。
func (r *MediaRepository) UpsertBatch(ctx context.Context, items []*model.Media) error {
if len(items) == 0 {
return nil
}
return withSQLiteBusyRetry(ctx, func() error {
return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
for _, m := range items {
if m == nil {
continue
}
if err := r.upsertWithDB(ctx, tx, m); err != nil {
return err
}
}
return nil
})
})
}
func (r *MediaRepository) upsertWithDB(ctx context.Context, db *gorm.DB, m *model.Media) error {
existing, created, err := r.findOrCreateMediaByPath(ctx, db, m)
if err != nil {
return err
}
@@ -38,20 +60,20 @@ func (r *MediaRepository) upsert(ctx context.Context, m *model.Media) error {
}
updates := mediaUpsertUpdates(existing, *m)
return r.applyMediaUpsertUpdates(ctx, m, existing, updates)
return r.applyMediaUpsertUpdates(ctx, db, m, existing, updates)
}
func (r *MediaRepository) findOrCreateMediaByPath(ctx context.Context, m *model.Media) (model.Media, bool, error) {
func (r *MediaRepository) findOrCreateMediaByPath(ctx context.Context, db *gorm.DB, m *model.Media) (model.Media, bool, error) {
var existing model.Media
err := r.db.WithContext(ctx).Unscoped().Where("path = ?", m.Path).First(&existing).Error
err := db.WithContext(ctx).Unscoped().Where("path = ?", m.Path).First(&existing).Error
if errors.Is(err, gorm.ErrRecordNotFound) {
// 新行:保证 scrape_status 走 GORM default:pending(即留空让数据库填)。
if m.ScrapeStatus == "" {
m.ScrapeStatus = "pending"
}
if createErr := r.db.WithContext(ctx).Create(m).Error; createErr == nil {
if createErr := db.WithContext(ctx).Create(m).Error; createErr == nil {
return *m, true, nil
} else if retryErr := r.db.WithContext(ctx).Unscoped().Where("path = ?", m.Path).First(&existing).Error; retryErr != nil {
} else if retryErr := db.WithContext(ctx).Unscoped().Where("path = ?", m.Path).First(&existing).Error; retryErr != nil {
return model.Media{}, false, createErr
}
}
@@ -237,12 +259,12 @@ func setNonEmptyMediaString(updates map[string]any, key, current, next string) {
}
}
func (r *MediaRepository) applyMediaUpsertUpdates(ctx context.Context, m *model.Media, existing model.Media, updates map[string]any) error {
func (r *MediaRepository) applyMediaUpsertUpdates(ctx context.Context, db *gorm.DB, m *model.Media, existing model.Media, updates map[string]any) error {
if len(updates) == 0 {
*m = existing
return nil
}
if err := r.db.WithContext(ctx).Unscoped().Model(&model.Media{}).
if err := db.WithContext(ctx).Unscoped().Model(&model.Media{}).
Where("id = ?", existing.ID).Updates(updates).Error; err != nil {
return err
}
+1 -1
View File
@@ -49,7 +49,7 @@ func (e *EmbyService) embyLatestCacheKey(userID, parentID string, limit int) str
func (e *EmbyService) mediaCacheTTLSeconds() int {
if e == nil || e.cfg == nil || e.cfg.Cache.MediaTTLSeconds < 1 {
return 15
return 90
}
return e.cfg.Cache.MediaTTLSeconds
}
+51 -1
View File
@@ -7,6 +7,7 @@ import (
"fmt"
"sort"
"strings"
"time"
"github.com/truewhile/MeBox/internal/model"
"github.com/truewhile/MeBox/internal/repository"
@@ -70,11 +71,60 @@ func (s *MediaService) seriesCardsCacheKey(libraryID string, visibility MediaVis
func (s *MediaService) mediaCacheTTLSeconds() int {
if s == nil || s.cfg == nil || s.cfg.Cache.MediaTTLSeconds < 1 {
return 15
return 90
}
return s.cfg.Cache.MediaTTLSeconds
}
// mediaObjectTTL 对象缓存与字节缓存同 TTL,失效同走 invalidateMediaCache
// 的 DeletePrefix("media:")(对象存储一并清除)。
func (s *MediaService) mediaObjectTTL() time.Duration {
return time.Duration(s.mediaCacheTTLSeconds()) * time.Second
}
func hashObjectCacheKey(parts []string) string {
sum := sha1.Sum([]byte(strings.Join(parts, "|")))
return hex.EncodeToString(sum[:])
}
// mediaGroupedRowsCacheKey 版本分组源行(已挂库元数据)的对象缓存。
func (s *MediaService) mediaGroupedRowsCacheKey(libraryID string, libraryIDs []string, filter repository.MediaQueryFilter) string {
libs := append([]string(nil), libraryIDs...)
sort.Strings(libs)
allowed := append([]string(nil), filter.AllowedLibraryIDs...)
hidden := append([]string(nil), filter.HiddenLibraryIDs...)
sort.Strings(allowed)
sort.Strings(hidden)
return "media:obj:grouped-rows:" + hashObjectCacheKey([]string{
libraryID,
strings.Join(libs, ","),
fmt.Sprintf("%t", filter.IncludeNSFW),
strings.Join(allowed, ","),
strings.Join(hidden, ","),
})
}
// libraryRowsCacheKey 整库可见行 + 预计算剧集索引的对象缓存(剧集卡片/剧集
// 列表共用一次 SQL 加载与一次 key 解析)。
func (s *MediaService) libraryRowsCacheKey(libraryID string, visibility MediaVisibility) string {
allowed := append([]string(nil), visibility.AllowedLibraryIDs...)
hidden := append([]string(nil), visibility.HiddenLibraryIDs...)
sort.Strings(allowed)
sort.Strings(hidden)
return "media:obj:library-rows:" + hashObjectCacheKey([]string{
libraryID,
fmt.Sprintf("%t", visibility.IncludeNSFW),
strings.Join(allowed, ","),
strings.Join(hidden, ","),
})
}
// libraryCardsObjectKey 系列卡片结果的对象缓存(与 seriesCardsCacheKey 同维度,
// 换独立前缀避免与字节缓存 key 冲突)。
func (s *MediaService) libraryCardsObjectKey(libraryID string, visibility MediaVisibility) string {
return s.seriesCardsCacheKey(libraryID, visibility) + ":obj"
}
func (s *MediaService) invalidateMediaCache(ctx context.Context) {
if s != nil && s.cache != nil {
s.cache.DeletePrefix(ctx, "media:")
+10 -6
View File
@@ -73,11 +73,15 @@ func (s *MediaService) listMediaVisibleForGrouping(ctx context.Context, libraryI
AllowedLibraryIDs: visibility.AllowedLibraryIDs,
HiddenLibraryIDs: visibility.HiddenLibraryIDs,
}
cacheKey := s.mediaListCacheKey(libraryID, libraryIDs, 0, maxMediaSearchLimit, filter) + ":group-source"
var cached mediaListCacheValue
if s.cache != nil && s.cache.GetJSON(ctx, cacheKey, &cached) {
s.attachLibraryMetadata(ctx, cached.Items)
return cached.Items, nil
cacheKey := s.mediaGroupedRowsCacheKey(libraryID, libraryIDs, filter)
if s.cache != nil {
if cachedObj, ok := s.cache.GetObject(cacheKey); ok {
if cached, ok := cachedObj.([]model.Media); ok {
// 对象缓存中的切片视为不可变;attachLibraryMetadata 会在填充时
// 执行过,命中路径直接返回副本即可(调用方只读)。
return cached, nil
}
}
}
items, total, err := s.repo.Media.ListByLibrariesFiltered(ctx, libraryIDs, 0, maxMediaSearchLimit, filter)
if err != nil {
@@ -91,7 +95,7 @@ func (s *MediaService) listMediaVisibleForGrouping(ctx context.Context, libraryI
}
s.attachLibraryMetadata(ctx, items)
if s.cache != nil {
s.cache.SetJSON(ctx, cacheKey, mediaListCacheValue{Items: items, Total: total}, time.Duration(s.mediaCacheTTLSeconds())*time.Second)
s.cache.SetObject(cacheKey, items, s.mediaObjectTTL())
}
return items, nil
}
+82 -22
View File
@@ -17,6 +17,15 @@ type seriesCardsCacheValue struct {
Total int64 `json:"total"`
}
// libraryRowsCacheValue 整库可见行 + 预计算索引的对象缓存。填充时一次性完成
// 整库 SQL 加载与 series key 解析,之后剧集卡片、剧集列表、媒体详情剧集
// 共享这份结果:命中路径零 SQL、零逐行 key 解析。Rows/Episodes 视为不可变。
type libraryRowsCacheValue struct {
Rows []model.Media
Resolver mediaSeriesKeyResolver
Episodes map[string][]model.Media
}
type SeriesCard struct {
Key string `json:"key"`
Rep model.Media `json:"rep"`
@@ -29,24 +38,59 @@ type seriesCardGroup struct {
latest time.Time
}
func (s *MediaService) ListLibrarySeriesCards(ctx context.Context, libraryID string, visibility MediaVisibility) ([]SeriesCard, int64, error) {
visibility = ExpandMediaVisibilityForMergedCloudLibraries(ctx, s.repo, visibility)
cacheKey := s.seriesCardsCacheKey(libraryID, visibility)
var cached seriesCardsCacheValue
if s.cache != nil && s.cache.GetJSON(ctx, cacheKey, &cached) {
return cached.Cards, cached.Total, nil
// libraryRowsWithIndex 返回整库行与预计算剧集索引(带对象缓存)。
func (s *MediaService) libraryRowsWithIndex(ctx context.Context, libraryID string, visibility MediaVisibility) (*libraryRowsCacheValue, error) {
cacheKey := s.libraryRowsCacheKey(libraryID, visibility)
if s.cache != nil {
if obj, ok := s.cache.GetObject(cacheKey); ok {
if cached, ok := obj.(*libraryRowsCacheValue); ok {
return cached, nil
}
}
}
rows, _, err := s.listAllMediaVisible(ctx, libraryID, visibility)
if err != nil {
return nil, err
}
// listAllMediaVisible 走 ListMediaVisible,行已带库元数据(resolver 的
// key 计算依赖 DisplayLibraryPath/ID)。
resolver := newMediaSeriesKeyResolver(rows)
episodes := make(map[string][]model.Media, len(rows)/4+1)
for _, row := range rows {
k := resolver.key(row)
if k == "" {
continue
}
episodes[k] = append(episodes[k], row)
}
value := &libraryRowsCacheValue{Rows: rows, Resolver: resolver, Episodes: episodes}
if s.cache != nil {
s.cache.SetObject(cacheKey, value, s.mediaObjectTTL())
}
return value, nil
}
func (s *MediaService) ListLibrarySeriesCards(ctx context.Context, libraryID string, visibility MediaVisibility) ([]SeriesCard, int64, error) {
visibility = ExpandMediaVisibilityForMergedCloudLibraries(ctx, s.repo, visibility)
cacheKey := s.libraryCardsObjectKey(libraryID, visibility)
if s.cache != nil {
if obj, ok := s.cache.GetObject(cacheKey); ok {
if cached, ok := obj.(*seriesCardsCacheValue); ok {
return cached.Cards, cached.Total, nil
}
}
}
rows, err := s.libraryRowsWithIndex(ctx, libraryID, visibility)
if err != nil {
return nil, 0, err
}
cards := groupMediaSeriesCards(rows)
cards := groupMediaSeriesCards(rows.Rows)
if cards == nil {
cards = []SeriesCard{}
}
total := int64(len(cards))
if s.cache != nil {
s.cache.SetJSON(ctx, cacheKey, seriesCardsCacheValue{Cards: cards, Total: total}, time.Duration(s.mediaCacheTTLSeconds())*time.Second)
s.cache.SetObject(cacheKey, &seriesCardsCacheValue{Cards: cards, Total: total}, s.mediaObjectTTL())
}
return cards, total, nil
}
@@ -72,17 +116,35 @@ func (s *MediaService) ListRecentSeriesCards(ctx context.Context, limit int, vis
}
func (s *MediaService) ListLibrarySeriesEpisodes(ctx context.Context, libraryID, key string, visibility MediaVisibility) ([]model.Media, error) {
rows, _, err := s.listAllMediaVisible(ctx, libraryID, visibility)
rows, err := s.libraryRowsWithIndex(ctx, libraryID, visibility)
if err != nil {
return nil, err
}
// 预计算索引命中:O(1) 查找,拷贝后排序避免改动共享缓存。
if episodes, ok := rows.Episodes[key]; ok && len(episodes) > 0 {
out := make([]model.Media, len(episodes))
copy(out, episodes)
sortEpisodesForDisplay(out)
return out, nil
}
// 兜底:索引未命中(如卡片 key 与行缓存短暂跨代),保持原线性匹配逻辑。
all := rows.Rows
out := make([]model.Media, 0)
resolver := newMediaSeriesKeyResolver(rows)
for _, row := range rows {
if resolver.key(row) == key {
for _, row := range all {
if rows.Resolver.key(row) == key {
out = append(out, row)
}
}
if len(out) == 0 {
return []model.Media{}, nil
}
sortEpisodesForDisplay(out)
return out, nil
}
// sortEpisodesForDisplay 与历史行为一致:季/集号升序,再按入库时间兜底。
func sortEpisodesForDisplay(out []model.Media) {
sort.SliceStable(out, func(i, j int) bool {
if out[i].SeasonNum != out[j].SeasonNum {
return out[i].SeasonNum < out[j].SeasonNum
@@ -92,7 +154,6 @@ func (s *MediaService) ListLibrarySeriesEpisodes(ctx context.Context, libraryID,
}
return out[i].CreatedAt.Before(out[j].CreatedAt)
})
return out, nil
}
func (s *MediaService) ListMediaEpisodes(ctx context.Context, mediaID string, visibility MediaVisibility) ([]model.Media, error) {
@@ -109,23 +170,22 @@ func (s *MediaService) ListMediaEpisodes(ctx context.Context, mediaID string, vi
if target.LibraryID == "" {
return []model.Media{*target}, nil
}
rows, _, err := s.listAllMediaVisible(ctx, target.LibraryID, visibility)
cache, err := s.libraryRowsWithIndex(ctx, target.LibraryID, visibility)
if err != nil {
return nil, err
}
rows := cache.Rows
if len(rows) == 0 {
return []model.Media{*target}, nil
}
resolver := newMediaSeriesKeyResolver(rows)
targetKey := resolver.key(*target)
out := make([]model.Media, 0)
var out []model.Media
targetKey := cache.Resolver.key(*target)
if targetKey != "" {
for _, row := range rows {
if resolver.key(row) == targetKey {
out = append(out, row)
}
// 预计算索引命中时拷贝,避免改动共享缓存。
if episodes, ok := cache.Episodes[targetKey]; ok {
out = make([]model.Media, len(episodes))
copy(out, episodes)
}
}
+61 -1
View File
@@ -20,6 +20,7 @@ type RuntimeCacheService struct {
mu sync.RWMutex
memory map[string]runtimeCacheItem
obj map[string]runtimeObjectItem
limit int
}
@@ -28,8 +29,18 @@ type runtimeCacheItem struct {
expiresAt time.Time
}
// runtimeObjectItem 直存 Go 对象,跳过 JSON 编解码。热点路径(整库行、
// 分组结果)每次请求都要完整反序列化,字节缓存避免了 SQL 却没避免解码;
// 对象缓存命中时零解码零分配。存储的值视为不可变:读取方如需修改必须
// 自行浅拷贝。仅进程内生效(Redis 只支持字节),多实例部署退化为各实例
// 独立缓存,与现有内存 L1 语义一致。
type runtimeObjectItem struct {
value any
expiresAt time.Time
}
func NewRuntimeCacheService(cfg *config.Config, log *zap.Logger) *RuntimeCacheService {
c := &RuntimeCacheService{log: log, memory: map[string]runtimeCacheItem{}, limit: 2048}
c := &RuntimeCacheService{log: log, memory: map[string]runtimeCacheItem{}, obj: map[string]runtimeObjectItem{}, limit: 2048}
if cfg == nil {
return c
}
@@ -112,6 +123,50 @@ func (c *RuntimeCacheService) SetJSON(ctx context.Context, key string, value any
}
}
// GetObject 返回缓存中的对象。返回值不可变:调用方需要修改时必须先自行拷贝。
func (c *RuntimeCacheService) GetObject(key string) (any, bool) {
if !c.Enabled() || strings.TrimSpace(key) == "" {
return nil, false
}
fullKey := c.key(key)
now := time.Now()
c.mu.RLock()
item, ok := c.obj[fullKey]
c.mu.RUnlock()
if !ok {
return nil, false
}
if now.After(item.expiresAt) {
c.mu.Lock()
delete(c.obj, fullKey)
c.mu.Unlock()
return nil, false
}
return item.value, true
}
// SetObject 存入一个此后视为不可变的对象。
func (c *RuntimeCacheService) SetObject(key string, value any, ttl time.Duration) {
if !c.Enabled() || strings.TrimSpace(key) == "" || value == nil || ttl <= 0 {
return
}
fullKey := c.key(key)
now := time.Now()
c.mu.Lock()
defer c.mu.Unlock()
if len(c.obj) >= c.limit {
for k, item := range c.obj {
if now.After(item.expiresAt) || len(c.obj) >= c.limit {
delete(c.obj, k)
}
if len(c.obj) < c.limit {
break
}
}
}
c.obj[fullKey] = runtimeObjectItem{value: value, expiresAt: now.Add(ttl)}
}
func (c *RuntimeCacheService) DeletePrefix(ctx context.Context, prefix string) {
if !c.Enabled() || strings.TrimSpace(prefix) == "" {
return
@@ -190,4 +245,9 @@ func (c *RuntimeCacheService) deleteMemoryPrefix(prefix string) {
delete(c.memory, key)
}
}
for key := range c.obj {
if strings.HasPrefix(key, prefix) {
delete(c.obj, key)
}
}
}
+36 -1
View File
@@ -63,6 +63,8 @@ func (b *localMediaWriteBatch) Flush() {
return
}
existingPaths := b.existingPaths(items)
upsertItems := make([]*model.Media, 0, len(items))
upsertAfter := make([]func(), 0, len(items))
createItems := make([]localMediaWriteItem, 0, len(items))
createMedia := make([]model.Media, 0, len(items))
for _, item := range items {
@@ -70,12 +72,16 @@ func (b *localMediaWriteBatch) Flush() {
continue
}
if existingPaths[filepath.Clean(item.media.Path)] {
b.upsertExistingItem(item)
// 已存在行:攒起来在一个事务里逐条 upsert(一批一次提交)。
after := item.after
upsertItems = append(upsertItems, item.media)
upsertAfter = append(upsertAfter, after)
continue
}
createItems = append(createItems, item)
createMedia = append(createMedia, *item.media)
}
b.flushUpserts(items, upsertItems, upsertAfter)
if len(createMedia) == 0 {
b.publish()
return
@@ -142,6 +148,35 @@ func (b *localMediaWriteBatch) existingPaths(items []localMediaWriteItem) map[st
return out
}
// flushUpserts 把已存在行的 upsert 攒成一个事务(一次提交/一组 fsync)。
// 整批失败(如单条数据触发约束)时退回逐条 Upsert,只丢真正坏的那几条。
func (b *localMediaWriteBatch) flushUpserts(allItems []localMediaWriteItem, upsertItems []*model.Media, upsertAfter []func()) {
if len(upsertItems) == 0 {
return
}
if err := b.scanner.repo.Media.UpsertBatch(b.ctx, upsertItems); err == nil {
b.res.Updated += len(upsertItems)
for _, after := range upsertAfter {
if after != nil {
after()
}
}
return
} else if b.scanner.log != nil {
b.scanner.log.Warn("batch upsert failed; falling back to per-item upsert",
zap.Int("items", len(upsertItems)))
}
// 兜底:按原始顺序找回每个条目的 path/after(两个切片同序但可能含 nil)。
idx := 0
for _, item := range allItems {
if item.media == nil || idx >= len(upsertItems) || upsertItems[idx] != item.media {
continue
}
idx++
b.upsertExistingItem(item)
}
}
func (b *localMediaWriteBatch) upsertExistingItem(item localMediaWriteItem) {
if item.media == nil {
return