优化,日志脱敏

This commit is contained in:
truewhile
2026-09-14 15:12:42 +08:00
parent 5c85478883
commit 3aa8662c45
31 changed files with 1220 additions and 1099 deletions
+6 -1
View File
@@ -127,5 +127,10 @@ func (z zapStdLogger) Printf(format string, args ...interface{}) {
if z.log == nil {
return
}
z.log.Sugar().Infof(format, args...)
message := fmt.Sprintf(format, args...)
if strings.Contains(strings.ToLower(message), "context canceled") {
z.log.Debug(message)
return
}
z.log.Info(message)
}
+1
View File
@@ -19,6 +19,7 @@ func Register(r *gin.Engine, cfg *config.Config, log *zap.Logger, svc *service.C
api.Use(middleware.GzipAPI())
{
api.GET("/health", healthCheck)
api.HEAD("/health", healthCheck)
api.GET("/version", versionInfo)
api.GET("/public/ui-config", publicUIConfigHandler(svc))
+8 -4
View File
@@ -1423,7 +1423,9 @@ func (s *DanmakuService) fetchCommentWithFallback(ctx context.Context, primary,
}
lastErr = err
if i < len(bases)-1 {
s.log.Warn("danmaku comment fetch failed on configured source, falling back to official", zap.String("source", base), zap.Error(err))
s.log.Warn("danmaku comment fetch failed on configured source, falling back to official",
zap.String("source", redactSensitiveURL(base)),
zap.Error(redactSensitiveError(err)))
}
}
return "", "auto", lastErr
@@ -1439,7 +1441,9 @@ func (s *DanmakuService) searchCandidatesWithSource(ctx context.Context, configu
if err == nil {
return candidates, configured, nil
}
s.log.Warn("danmaku search failed on configured source, falling back to official", zap.String("source", configured), zap.Error(err))
s.log.Warn("danmaku search failed on configured source, falling back to official",
zap.String("source", redactSensitiveURL(configured)),
zap.Error(redactSensitiveError(err)))
}
candidates, err := s.searchCandidates(ctx, official, name, episode)
return candidates, official, err
@@ -1573,9 +1577,9 @@ func (s *DanmakuService) lookupConfiguredEpisodes(ctx context.Context, configure
candidates, err := s.searchCandidates(ctx, configured, m.AnimeTitle, episodeNum)
if err != nil {
s.log.Debug("danmaku configured lookup failed",
zap.String("source", configured),
zap.String("source", redactSensitiveURL(configured)),
zap.String("anime", m.AnimeTitle),
zap.Error(err))
zap.Error(redactSensitiveError(err)))
return nil, false
}
matched := matchDanmakuEpisodes(candidates, episodeNum, m.EpisodeTitle)
+2 -1
View File
@@ -69,7 +69,8 @@ type EmbyService struct {
// latestFlight collapses the homepage stampede: clients request Latest
// for every library at once, and a shared expiry used to rebuild each
// library in parallel.
latestFlight singleflight.Group
latestFlight singleflight.Group
latestRefresh sync.Map
tmdb *TMDbProvider
adult *AdultProvider
+49 -12
View File
@@ -130,11 +130,27 @@ func (e *EmbyService) LatestItems(ctx context.Context, userID, parentID string,
limit = 20
}
cacheKey := e.embyLatestCacheKey(userID, parentID, limit)
if items, ok := e.cachedLatestItems(ctx, cacheKey); ok {
if items, stale, ok := e.cachedLatestItemsWithStale(ctx, cacheKey); ok {
if stale {
e.refreshLatestItemsAsync(cacheKey, userID, parentID, limit)
}
return items, nil
}
value, err := e.loadLatestItemsCached(ctx, userID, parentID, limit)
if err != nil {
return nil, err
}
if value.Items == nil {
return []map[string]any{}, nil
}
return value.Items, nil
}
func (e *EmbyService) loadLatestItemsCached(ctx context.Context, userID, parentID string, limit int) (embyLatestCacheValue, error) {
cacheKey := e.embyLatestCacheKey(userID, parentID, limit)
// 一个客户端断开不应取消正在为其他客户端填充的共享重建。
loadCtx := context.WithoutCancel(ctx)
loadCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 45*time.Second)
defer cancel()
v, err, _ := e.latestFlight.Do(cacheKey, func() (any, error) {
if items, ok := e.cachedLatestItems(loadCtx, cacheKey); ok {
return embyLatestCacheValue{Items: items}, nil
@@ -144,34 +160,55 @@ func (e *EmbyService) LatestItems(ctx context.Context, userID, parentID string,
return nil, err
}
if e.cache != nil {
e.cache.SetJSON(loadCtx, cacheKey, value, time.Duration(e.embyLatestCacheTTLSeconds())*time.Second)
freshTTL := time.Duration(e.embyLatestCacheTTLSeconds()) * time.Second
e.cache.SetJSONWithStale(loadCtx, cacheKey, value, freshTTL, freshTTL+30*time.Minute)
}
e.rememberArtworkRefs(value.Artwork)
return value, nil
})
if err != nil {
return nil, err
return embyLatestCacheValue{}, err
}
cached, _ := v.(embyLatestCacheValue)
if cached.Items == nil {
return []map[string]any{}, nil
return cached, nil
}
func (e *EmbyService) refreshLatestItemsAsync(cacheKey, userID, parentID string, limit int) {
if e == nil || e.cache == nil {
return
}
return cached.Items, nil
if _, loaded := e.latestRefresh.LoadOrStore(cacheKey, struct{}{}); loaded {
return
}
go func() {
defer e.latestRefresh.Delete(cacheKey)
if _, err := e.loadLatestItemsCached(context.Background(), userID, parentID, limit); err != nil && e.log != nil {
e.log.Debug("background refresh of latest items failed",
zap.String("parent_id", parentID),
zap.Error(redactSensitiveError(err)))
}
}()
}
func (e *EmbyService) cachedLatestItems(ctx context.Context, cacheKey string) ([]map[string]any, bool) {
items, stale, ok := e.cachedLatestItemsWithStale(ctx, cacheKey)
return items, ok && !stale
}
func (e *EmbyService) cachedLatestItemsWithStale(ctx context.Context, cacheKey string) ([]map[string]any, bool, bool) {
if e == nil || e.cache == nil {
return nil, false
return nil, false, false
}
var cached embyLatestCacheValue
if !e.cache.GetJSON(ctx, cacheKey, &cached) {
return nil, false
found, stale := e.cache.GetJSONStale(ctx, cacheKey, &cached)
if !found {
return nil, false, false
}
e.rememberArtworkRefs(cached.Artwork)
if cached.Items == nil {
return []map[string]any{}, true
return []map[string]any{}, stale, true
}
return cached.Items, true
return cached.Items, stale, true
}
func (e *EmbyService) loadLatestItems(ctx context.Context, userID, parentID string, limit int) (embyLatestCacheValue, error) {
+17 -9
View File
@@ -510,12 +510,12 @@ func (r *EmbyRemoteService) ensureTokenOnLine(ctx context.Context, acct *model.S
req.Header.Set("X-Emby-Authorization", `MediaBrowser Client="MeBox", Device="MeBox-Federated", DeviceId="mebox-federated", Version="1.0"`)
resp, err := r.http.Do(req)
if err != nil {
return fmt.Errorf("连接远程 Emby 失败: %w", err)
return redactSensitiveError(fmt.Errorf("连接远程 Emby 失败: %w", err))
}
defer resp.Body.Close()
if resp.StatusCode >= 300 {
data, _ := io.ReadAll(io.LimitReader(resp.Body, 512))
return fmt.Errorf("远程 Emby 登录失败(%d): %s", resp.StatusCode, strings.TrimSpace(string(data)))
return redactSensitiveError(fmt.Errorf("远程 Emby 登录失败(%d): %s", resp.StatusCode, strings.TrimSpace(string(data))))
}
var login struct {
AccessToken string `json:"AccessToken"`
@@ -599,8 +599,16 @@ func (r *EmbyRemoteService) doGet(ctx context.Context, acct *model.StrmAccount,
}
if lastErr != nil {
if r.log != nil && acct != nil {
r.log.Warn("remote emby request failed",
zap.String("account", acct.Name), zap.String("path", path), zap.Error(lastErr))
fields := []zap.Field{
zap.String("account", acct.Name),
zap.String("path", path),
zap.Error(redactSensitiveError(lastErr)),
}
if errors.Is(lastErr, context.Canceled) {
r.log.Debug("remote emby request canceled", fields...)
} else {
r.log.Warn("remote emby request failed", fields...)
}
}
return lastErr
}
@@ -629,7 +637,7 @@ func (r *EmbyRemoteService) doGetOnLine(ctx context.Context, acct *model.StrmAcc
req.Header.Set("X-Emby-Token", cfg.Token)
resp, err := r.http.Do(req)
if err != nil {
return fmt.Errorf("请求远程 Emby 失败: %w", err)
return redactSensitiveError(fmt.Errorf("请求远程 Emby 失败: %w", err))
}
// 读 8MB+1 以区分"刚好 8MB"与"被截断":截断的 JSON 会让
// Unmarshal 报 unexpected end,难以定位;这里显式报错。
@@ -656,7 +664,7 @@ func (r *EmbyRemoteService) doGetOnLine(ctx context.Context, acct *model.StrmAcc
continue
}
if resp.StatusCode >= 300 {
return fmt.Errorf("远程 Emby 请求失败(%d): %s", resp.StatusCode, strings.TrimSpace(string(data)))
return redactSensitiveError(fmt.Errorf("远程 Emby 请求失败(%d): %s", resp.StatusCode, strings.TrimSpace(string(data))))
}
if out == nil {
return nil
@@ -1121,7 +1129,7 @@ func (r *EmbyRemoteService) proxyVideoStreamOnLine(ctx context.Context, w http.R
defer resp.Body.Close()
if resp.StatusCode >= 400 {
data, _ := io.ReadAll(io.LimitReader(resp.Body, 256))
return fmt.Errorf("远程 Emby 视频流失败(%d): %s", resp.StatusCode, strings.TrimSpace(string(data)))
return redactSensitiveError(fmt.Errorf("远程 Emby 视频流失败(%d): %s", resp.StatusCode, strings.TrimSpace(string(data))))
}
for _, header := range []string{"Content-Type", "Content-Length", "Content-Range", "Accept-Ranges", "ETag", "Cache-Control"} {
if value := resp.Header.Get(header); value != "" {
@@ -1281,12 +1289,12 @@ func (r *EmbyRemoteService) doMutateOnLine(ctx context.Context, cfg *EmbyRemoteC
req.Header.Set("X-Emby-Token", cfg.Token)
resp, err := r.http.Do(req)
if err != nil {
return fmt.Errorf("请求远程 Emby 失败: %w", err)
return redactSensitiveError(fmt.Errorf("请求远程 Emby 失败: %w", err))
}
defer resp.Body.Close()
if resp.StatusCode >= 300 {
data, _ := io.ReadAll(io.LimitReader(resp.Body, 256))
return fmt.Errorf("远程 Emby 状态同步失败(%d): %s", resp.StatusCode, strings.TrimSpace(string(data)))
return redactSensitiveError(fmt.Errorf("远程 Emby 状态同步失败(%d): %s", resp.StatusCode, strings.TrimSpace(string(data))))
}
return nil
}
+9 -5
View File
@@ -180,18 +180,20 @@ func downloadFFmpegArchive(ctx context.Context, log *zap.Logger, urls []string,
var lastErr error
for i, u := range urls {
if i > 0 && log != nil {
log.Warn("ffmpeg 主下载源不可用,切换备用源", zap.String("url", u))
log.Warn("ffmpeg 主下载源不可用,切换备用源", zap.String("url", redactSensitiveURL(u)))
}
if err := downloadFFmpegFile(ctx, log, u, dest); err != nil {
lastErr = err
if log != nil {
log.Warn("ffmpeg 下载失败", zap.String("url", u), zap.Error(err))
log.Warn("ffmpeg 下载失败",
zap.String("url", redactSensitiveURL(u)),
zap.Error(redactSensitiveError(err)))
}
continue
}
return nil
}
return fmt.Errorf("所有下载源均失败:%v", lastErr)
return redactSensitiveError(fmt.Errorf("所有下载源均失败:%v", lastErr))
}
// downloadFFmpegFile 下载单个归档文件(最多 10 分钟,限制大小上限)。
@@ -221,10 +223,12 @@ func downloadFFmpegFile(ctx context.Context, log *zap.Logger, url, dest string)
return err
}
if n > 500<<20 {
return fmt.Errorf("归档文件过大(>500MB): %s", url)
return fmt.Errorf("归档文件过大(>500MB): %s", redactSensitiveURL(url))
}
if log != nil {
log.Info("ffmpeg 归档下载完成", zap.String("url", url), zap.Int64("bytes", n))
log.Info("ffmpeg 归档下载完成",
zap.String("url", redactSensitiveURL(url)),
zap.Int64("bytes", n))
}
return nil
}
+2 -2
View File
@@ -196,14 +196,14 @@ func (p *ImageProxy) fetchAndCacheRemoteImage(ctx context.Context, raw, host, ca
p.writeImageCache(cachePath, failPath, "img-*.tmp", data)
return data, ctype, contentLength, nil
}
p.log.Warn("imageproxy: curl fallback failed", zap.String("host", host), zap.Error(err))
logImageFetchError(p.log, "imageproxy: curl fallback failed", host, "curl", err)
lastErr = err
}
p.markImageFetchFailed(failPath)
if lastErr == nil {
lastErr = errors.New("upstream image fetch failed")
}
return nil, "", "", lastErr
return nil, "", "", redactSensitiveError(lastErr)
}
type sharedRemoteImageResult struct {
+24 -8
View File
@@ -48,14 +48,14 @@ func (p *ImageProxy) canUseExternalImageFallback() bool {
func (p *ImageProxy) fetchRemoteImageOnce(ctx context.Context, raw, host string, candidate remoteImageFetchClient) ([]byte, string, string, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, raw, nil)
if err != nil {
p.log.Warn("imageproxy: build request failed", zap.String("url", raw), zap.Error(err))
p.log.Warn("imageproxy: build request failed", zap.String("url", redactSensitiveURL(raw)), zap.Error(redactSensitiveError(err)))
return nil, "", "", errImageProxyRequestSetup
}
applyRemoteImageHeaders(req, host, raw)
resp, err := candidate.client.Do(req)
if err != nil {
p.log.Warn("imageproxy: upstream fetch failed", zap.String("host", host), zap.String("client", candidate.name), zap.Error(err))
logImageFetchError(p.log, "imageproxy: upstream fetch failed", host, candidate.name, err)
return nil, "", "", err
}
defer resp.Body.Close()
@@ -65,7 +65,7 @@ func (p *ImageProxy) fetchRemoteImageOnce(ctx context.Context, raw, host string,
}
data, err := io.ReadAll(io.LimitReader(resp.Body, 32<<20))
if err != nil || len(data) == 0 {
p.log.Warn("imageproxy: read upstream body failed", zap.String("host", host), zap.String("client", candidate.name), zap.Error(err))
p.log.Warn("imageproxy: read upstream body failed", zap.String("host", host), zap.String("client", candidate.name), zap.Error(redactSensitiveError(err)))
if err == nil {
err = errors.New("upstream image body is empty")
}
@@ -79,6 +79,22 @@ func (p *ImageProxy) fetchRemoteImageOnce(ctx context.Context, raw, host string,
return data, ctype, resp.Header.Get("Content-Length"), nil
}
func logImageFetchError(log *zap.Logger, message, host, client string, err error) {
if log == nil || err == nil {
return
}
fields := []zap.Field{
zap.String("host", host),
zap.String("client", client),
zap.Error(redactSensitiveError(err)),
}
if errors.Is(err, context.Canceled) {
log.Debug(message, fields...)
return
}
log.Warn(message, fields...)
}
func applyRemoteImageHeaders(req *http.Request, host, raw string) {
req.Header.Set("User-Agent", "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/125.0 Safari/537.36")
req.Header.Set("Accept", "image/avif,image/webp,image/apng,image/svg+xml,image/*,*/*;q=0.8")
@@ -162,9 +178,9 @@ func fetchRemoteImageWithCurl(ctx context.Context, raw, host string) ([]byte, st
"--header", "Cache-Control: no-cache",
"--header", "Pragma: no-cache",
}
if referer := remoteImageReferer(host, raw); referer != "" {
args = append(args, "--referer", referer)
}
if referer := remoteImageReferer(host, raw); referer != "" {
args = append(args, "--referer", referer)
}
if cookie := remoteImageCookie(host); cookie != "" {
args = append(args, "--cookie", cookie)
}
@@ -188,9 +204,9 @@ func fetchRemoteImageWithCurl(ctx context.Context, raw, host string) ([]byte, st
if waitErr != nil {
message := strings.TrimSpace(stderr.String())
if message != "" {
return nil, "", "", errors.New(message)
return nil, "", "", redactSensitiveError(errors.New(message))
}
return nil, "", "", waitErr
return nil, "", "", redactSensitiveError(waitErr)
}
if len(data) == 0 {
return nil, "", "", errors.New("curl image body is empty")
+89
View File
@@ -0,0 +1,89 @@
package service
import (
"errors"
"net/url"
"regexp"
"strings"
)
const sensitiveLogKeyPattern = `api[_-]?key|apikey|access[_-]?token|refresh[_-]?token|id[_-]?token|token|password|passwd|pwd|authorization|client[_-]?secret|secret|signature|sig|x[_-]?emby[_-]?token|x[_-]?api[_-]?key|x[_-]?amz[_-]?signature|x[_-]?amz[_-]?credential|x[_-]?amz[_-]?security[_-]?token|awsaccesskeyid|session[_-]?id`
var (
sensitiveLogQueryRE = regexp.MustCompile(`(?i)([?&;])(` + sensitiveLogKeyPattern + `)=([^&#\s"']+)`)
sensitiveLogJSONRE = regexp.MustCompile(`(?i)("(?:` + sensitiveLogKeyPattern + `)"\s*:\s*")((?:[^"\\]|\\.)*)(")`)
sensitiveLogAssignmentRE = regexp.MustCompile(`(?im)((?:^|[\s,{])(?:` + sensitiveLogKeyPattern + `)\s*[:=]\s*(?:(?:bearer|basic)\s+)?)([^\s,;"'}\]]+)`)
)
// redactSensitiveURL removes credentials from URLs before they reach logs or
// error messages. Query values are replaced with REDACTED and URL userinfo
// passwords are removed.
func redactSensitiveURL(raw string) string {
trimmed := strings.TrimSpace(raw)
u, err := url.Parse(trimmed)
if err != nil || u.Scheme == "" || u.Host == "" {
return redactSensitiveText(raw)
}
if u.User != nil {
if _, hasPassword := u.User.Password(); hasPassword {
u.User = url.UserPassword(u.User.Username(), "REDACTED")
}
}
query := u.Query()
changed := false
for key := range query {
if isSensitiveLogKey(key) {
query.Set(key, "REDACTED")
changed = true
}
}
if changed {
u.RawQuery = query.Encode()
}
return u.String()
}
func redactSensitiveText(value string) string {
value = sensitiveLogQueryRE.ReplaceAllString(value, `${1}${2}=REDACTED`)
value = sensitiveLogJSONRE.ReplaceAllString(value, `${1}REDACTED${3}`)
return sensitiveLogAssignmentRE.ReplaceAllString(value, `${1}REDACTED`)
}
type logRedactedError struct {
err error
}
func (e logRedactedError) Error() string {
return redactSensitiveText(e.err.Error())
}
func (e logRedactedError) Unwrap() error {
return e.err
}
// redactSensitiveError wraps an error with a redacted Error string while
// preserving errors.Is and errors.As behavior.
func redactSensitiveError(err error) error {
if err == nil {
return nil
}
var alreadyRedacted logRedactedError
if errors.As(err, &alreadyRedacted) {
return err
}
return logRedactedError{err: err}
}
func isSensitiveLogKey(key string) bool {
normalized := strings.ToLower(strings.TrimSpace(key))
normalized = strings.ReplaceAll(normalized, "-", "_")
switch normalized {
case "api_key", "apikey", "access_token", "refresh_token", "id_token", "token",
"password", "passwd", "pwd", "authorization", "client_secret", "secret",
"signature", "sig", "x_emby_token", "x_api_key", "x_amz_signature",
"x_amz_credential", "x_amz_security_token", "awsaccesskeyid", "session_id":
return true
default:
return false
}
}
+52
View File
@@ -0,0 +1,52 @@
package service
import (
"context"
"errors"
"strings"
"testing"
)
func TestRedactSensitiveURL(t *testing.T) {
raw := "https://user:secret@media.example/emby/Items/1/Images/primary?api_key=abc123&X-Emby-Token=xyz&X-Amz-Signature=sig123&quality=90"
got := redactSensitiveURL(raw)
for _, secret := range []string{"secret", "abc123", "xyz", "sig123"} {
if strings.Contains(got, secret) {
t.Fatalf("redacted URL still contains %q: %s", secret, got)
}
}
if !strings.Contains(got, "quality=90") {
t.Fatalf("non-sensitive query value was removed: %s", got)
}
}
func TestRedactSensitiveTextCoversJSONAndHeaders(t *testing.T) {
raw := `{"AccessToken":"json-token","quality":"90"} Authorization: Bearer bearer-token; X-Emby-Token=header-token`
got := redactSensitiveText(raw)
for _, secret := range []string{"json-token", "bearer-token", "header-token"} {
if strings.Contains(got, secret) {
t.Fatalf("redacted text still contains %q: %s", secret, got)
}
}
if !strings.Contains(got, `"quality":"90"`) {
t.Fatalf("non-sensitive JSON value was removed: %s", got)
}
}
func TestRedactSensitiveErrorPreservesWrapping(t *testing.T) {
base := errors.New(`Get "https://media.example/item?token=secret&quality=90": context canceled`)
err := redactSensitiveError(context.Canceled)
if !errors.Is(err, context.Canceled) {
t.Fatal("redacted error lost context.Canceled")
}
if strings.Contains(err.Error(), "secret") {
t.Fatalf("redacted error leaked token: %s", err.Error())
}
redactedBase := redactSensitiveError(base)
if strings.Contains(redactedBase.Error(), "secret") {
t.Fatalf("redacted base error leaked token: %s", redactedBase.Error())
}
if !errors.Is(redactedBase, base) {
t.Fatal("redacted error lost original error identity")
}
}
+96 -17
View File
@@ -25,7 +25,10 @@ const (
type RuntimeCacheService struct {
log *zap.Logger
client *redis.Client
prefix string
// redisGet is separated from client so cache ordering can be tested without
// requiring a live Redis instance.
redisGet func(ctx context.Context, key string) ([]byte, error)
prefix string
mu sync.RWMutex
memory map[string]runtimeCacheItem
@@ -36,10 +39,11 @@ type RuntimeCacheService struct {
}
type runtimeCacheItem struct {
raw []byte
expiresAt time.Time
lastUsed time.Time
size int64
raw []byte
expiresAt time.Time
staleUntil time.Time
lastUsed time.Time
size int64
}
// runtimeObjectItem 直存 Go 对象,跳过 JSON 编解码。热点路径(整库行、
@@ -95,6 +99,9 @@ func NewRuntimeCacheService(cfg *config.Config, log *zap.Logger) *RuntimeCacheSe
return c
}
c.client = client
c.redisGet = func(ctx context.Context, key string) ([]byte, error) {
return client.Get(ctx, key).Bytes()
}
if log != nil {
log.Info("redis runtime cache enabled with in-process L1", zap.String("addr", opts.Addr), zap.String("prefix", c.prefix))
}
@@ -120,8 +127,8 @@ func (c *RuntimeCacheService) GetJSON(ctx context.Context, key string, out any)
if raw, ok := c.getMemory(fullKey); ok {
return json.Unmarshal(raw, out) == nil
}
if c.client != nil {
raw, err := c.client.Get(ctx, fullKey).Bytes()
if c.redisGet != nil {
raw, err := c.redisGet(ctx, fullKey)
if err == nil {
if json.Unmarshal(raw, out) != nil {
return false
@@ -133,6 +140,36 @@ func (c *RuntimeCacheService) GetJSON(ctx context.Context, key string, out any)
return false
}
// GetJSONStale returns a fresh value when available and otherwise a retained
// stale value stored with SetJSONWithStale. The stale flag lets
// callers serve immediately while refreshing in the background.
func (c *RuntimeCacheService) GetJSONStale(ctx context.Context, key string, out any) (found bool, stale bool) {
if !c.Enabled() || strings.TrimSpace(key) == "" || out == nil {
return false, false
}
fullKey := c.key(key)
raw, ok, isStale := c.getMemoryWithStale(fullKey)
if ok && !isStale {
return json.Unmarshal(raw, out) == nil, false
}
// A stale L1 entry must not hide a newer value written by another
// instance. Prefer Redis whenever the local copy is stale, then fall back
// to it only when Redis is unavailable or its fresh value has expired.
if c.redisGet != nil {
redisRaw, err := c.redisGet(ctx, fullKey)
if err == nil {
if json.Unmarshal(redisRaw, out) == nil {
c.setMemoryOwned(fullKey, redisRaw, 2*time.Second)
return true, false
}
}
}
if ok {
return json.Unmarshal(raw, out) == nil, isStale
}
return false, false
}
func (c *RuntimeCacheService) SetJSON(ctx context.Context, key string, value any, ttl time.Duration) {
if !c.Enabled() || strings.TrimSpace(key) == "" || value == nil || ttl <= 0 {
return
@@ -148,6 +185,27 @@ func (c *RuntimeCacheService) SetJSON(ctx context.Context, key string, value any
}
}
// SetJSONWithStale stores a fresh value for freshTTL and keeps an in-process
// stale copy for staleTTL. Redis keeps only the fresh window so multi-instance
// deployments retain the existing consistency semantics.
func (c *RuntimeCacheService) SetJSONWithStale(ctx context.Context, key string, value any, freshTTL, staleTTL time.Duration) {
if !c.Enabled() || strings.TrimSpace(key) == "" || value == nil || freshTTL <= 0 {
return
}
if staleTTL < freshTTL {
staleTTL = freshTTL
}
raw, err := json.Marshal(value)
if err != nil {
return
}
fullKey := c.key(key)
c.setMemoryBytesWithStale(fullKey, raw, freshTTL, staleTTL, true)
if c.client != nil {
_ = c.client.Set(ctx, fullKey, raw, freshTTL).Err()
}
}
// GetObject 返回缓存中的对象。返回值不可变:调用方需要修改时必须先自行拷贝。
func (c *RuntimeCacheService) GetObject(key string) (any, bool) {
if !c.Enabled() || strings.TrimSpace(key) == "" {
@@ -255,20 +313,29 @@ func (c *RuntimeCacheService) key(key string) string {
}
func (c *RuntimeCacheService) getMemory(key string) ([]byte, bool) {
raw, ok, stale := c.getMemoryWithStale(key)
return raw, ok && !stale
}
func (c *RuntimeCacheService) getMemoryWithStale(key string) ([]byte, bool, bool) {
now := time.Now()
c.mu.Lock()
defer c.mu.Unlock()
item, ok := c.memory[key]
if !ok {
return nil, false
return nil, false, false
}
if !now.Before(item.expiresAt) {
staleUntil := item.staleUntil
if staleUntil.IsZero() {
staleUntil = item.expiresAt
}
if !now.Before(staleUntil) {
c.removeMemoryLocked(key)
return nil, false
return nil, false, false
}
item.lastUsed = now
c.memory[key] = item
return item.raw, true
return item.raw, true, !now.Before(item.expiresAt)
}
func (c *RuntimeCacheService) setMemory(key string, raw []byte, ttl time.Duration) {
@@ -280,9 +347,16 @@ func (c *RuntimeCacheService) setMemoryOwned(key string, raw []byte, ttl time.Du
}
func (c *RuntimeCacheService) setMemoryBytes(key string, raw []byte, ttl time.Duration, owned bool) {
if ttl <= 0 || len(raw) == 0 {
c.setMemoryBytesWithStale(key, raw, ttl, ttl, owned)
}
func (c *RuntimeCacheService) setMemoryBytesWithStale(key string, raw []byte, freshTTL, staleTTL time.Duration, owned bool) {
if freshTTL <= 0 || len(raw) == 0 {
return
}
if staleTTL < freshTTL {
staleTTL = freshTTL
}
size := int64(len(key)+len(raw)) + runtimeCacheEntryOverheadBytes
now := time.Now()
c.mu.Lock()
@@ -298,10 +372,11 @@ func (c *RuntimeCacheService) setMemoryBytes(key string, raw []byte, ttl time.Du
raw = append([]byte(nil), raw...)
}
c.memory[key] = runtimeCacheItem{
raw: raw,
expiresAt: now.Add(ttl),
lastUsed: now,
size: size,
raw: raw,
expiresAt: now.Add(freshTTL),
staleUntil: now.Add(staleTTL),
lastUsed: now,
size: size,
}
c.bytesUsed += size
}
@@ -334,7 +409,11 @@ func (c *RuntimeCacheService) entryCountLocked() int {
func (c *RuntimeCacheService) evictExpiredLocked(now time.Time) {
for key, item := range c.memory {
if !now.Before(item.expiresAt) {
staleUntil := item.staleUntil
if staleUntil.IsZero() {
staleUntil = item.expiresAt
}
if !now.Before(staleUntil) {
c.removeMemoryLocked(key)
}
}
+71
View File
@@ -116,3 +116,74 @@ func TestRuntimeCacheSetMaxSizeEvictsImmediately(t *testing.T) {
t.Fatal("newest entry should remain after lowering cache limit")
}
}
func TestRuntimeCacheReturnsStaleValueAfterFreshTTL(t *testing.T) {
cache := newRuntimeCacheForTest(t, 1)
key := "latest:stale"
cache.SetJSONWithStale(context.Background(), key, "cached-value", time.Minute, time.Hour)
fullKey := cache.key(key)
cache.mu.Lock()
item := cache.memory[fullKey]
item.expiresAt = time.Now().Add(-time.Second)
cache.memory[fullKey] = item
cache.mu.Unlock()
var fresh string
if cache.GetJSON(context.Background(), key, &fresh) {
t.Fatal("expired fresh entry must not be returned by GetJSON")
}
var stale string
found, isStale := cache.GetJSONStale(context.Background(), key, &stale)
if !found || !isStale {
t.Fatalf("GetJSONStale found=%t stale=%t, want true/true", found, isStale)
}
if stale != "cached-value" {
t.Fatalf("stale value=%q, want cached-value", stale)
}
}
func TestRuntimeCacheDoesNotReturnOrdinaryEntryAsStale(t *testing.T) {
cache := newRuntimeCacheForTest(t, 1)
key := "latest:ordinary"
cache.SetJSON(context.Background(), key, "cached-value", time.Minute)
fullKey := cache.key(key)
cache.mu.Lock()
item := cache.memory[fullKey]
expiredAt := time.Now().Add(-time.Second)
item.expiresAt = expiredAt
item.staleUntil = expiredAt
cache.memory[fullKey] = item
cache.mu.Unlock()
var out string
if found, isStale := cache.GetJSONStale(context.Background(), key, &out); found || isStale {
t.Fatalf("ordinary expired entry found=%t stale=%t, want false/false", found, isStale)
}
}
func TestRuntimeCacheStalePrefersFreshRedisValue(t *testing.T) {
cache := newRuntimeCacheForTest(t, 1)
key := "latest:redis-fresh"
cache.SetJSONWithStale(context.Background(), key, "local-stale", time.Minute, time.Hour)
fullKey := cache.key(key)
cache.mu.Lock()
item := cache.memory[fullKey]
item.expiresAt = time.Now().Add(-time.Second)
cache.memory[fullKey] = item
cache.mu.Unlock()
cache.redisGet = func(context.Context, string) ([]byte, error) {
return []byte(`"redis-fresh"`), nil
}
var out string
found, stale := cache.GetJSONStale(context.Background(), key, &out)
if !found || stale {
t.Fatalf("GetJSONStale found=%t stale=%t, want true/false", found, stale)
}
if out != "redis-fresh" {
t.Fatalf("value=%q, want redis-fresh", out)
}
}
+11 -9
View File
@@ -40,9 +40,9 @@ func (s *ScraperService) prepareScrapedArtworkURL(ctx context.Context, mediaID,
s.log.Warn("MetaTube artwork processing failed; using local face-aware crop",
zap.String("media_id", mediaID),
zap.String("field", field),
zap.String("candidate", candidate),
zap.String("source", originalSource),
zap.Error(err))
zap.String("candidate", redactSensitiveURL(candidate)),
zap.String("source", redactSensitiveURL(originalSource)),
zap.Error(redactSensitiveError(err)))
if current != "" && isHTTPish(current) {
return originalSource, current
}
@@ -53,16 +53,16 @@ func (s *ScraperService) prepareScrapedArtworkURL(ctx context.Context, mediaID,
s.log.Warn("scrape artwork prefetch failed; keeping existing artwork",
zap.String("media_id", mediaID),
zap.String("field", field),
zap.String("candidate", candidate),
zap.String("existing", current),
zap.Error(err))
zap.String("candidate", redactSensitiveURL(candidate)),
zap.String("existing", redactSensitiveURL(current)),
zap.Error(redactSensitiveError(err)))
return current, ""
}
s.log.Warn("scrape artwork prefetch failed; keeping new artwork URL for retry",
zap.String("media_id", mediaID),
zap.String("field", field),
zap.String("candidate", candidate),
zap.Error(err))
zap.String("candidate", redactSensitiveURL(candidate)),
zap.Error(redactSensitiveError(err)))
return candidate, ""
}
if current != "" && isHTTPish(current) {
@@ -98,7 +98,9 @@ func (s *ScraperService) removeCachedScrapedArtwork(urls ...string) {
}
seen[raw] = struct{}{}
if err := s.images.RemoveCached(raw); err != nil {
s.log.Debug("remove old scraped artwork cache failed", zap.String("url", raw), zap.Error(err))
s.log.Debug("remove old scraped artwork cache failed",
zap.String("url", redactSensitiveURL(raw)),
zap.Error(redactSensitiveError(err)))
}
}
}
@@ -117,8 +117,8 @@ func (s *ScraperService) downloadArtworkToPathWithOptions(ctx context.Context, d
if err != nil || len(data) == 0 {
s.log.Warn("scrape artwork download failed",
zap.String("name", name),
zap.String("url", raw),
zap.Error(err))
zap.String("url", redactSensitiveURL(raw)),
zap.Error(redactSensitiveError(err)))
return ""
}
if !isImageContentType(ctype) || isTransparentPlaceholderData(data) {
+3 -1
View File
@@ -88,7 +88,9 @@ func (b *serviceContainerBuilder) configureMediaSearchBackend() {
}
b.repos.Media.SetSearchBackend(searchBackend)
if b.log != nil {
b.log.Info("opensearch media search enabled", zap.String("index", b.cfg.Search.Index), zap.String("url", b.cfg.Search.OpenSearchURL))
b.log.Info("opensearch media search enabled",
zap.String("index", b.cfg.Search.Index),
zap.String("url", redactSensitiveURL(b.cfg.Search.OpenSearchURL)))
}
}