From 91df7521ceccec57bfd5123f733061a193bea387 Mon Sep 17 00:00:00 2001 From: ShukeBta Date: Fri, 12 Jun 2026 04:31:34 +0000 Subject: [PATCH 1/4] =?UTF-8?q?fix:=20=E8=B5=84=E6=BA=90=E5=8D=A0=E7=94=A8?= =?UTF-8?q?/=E7=99=BB=E5=BD=95=E7=A8=B3=E5=AE=9A=E6=80=A7/QB=E6=95=B4?= =?UTF-8?q?=E7=90=86=E5=85=A5=E5=BA=93/=E7=AC=AC=E4=B8=89=E6=96=B9?= =?UTF-8?q?=E6=92=AD=E6=94=BE404=20=E7=BB=BC=E5=90=88=E4=BF=AE=E5=A4=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 资源占用(Docker 部署 CPU/内存长期居高): - 云盘探测预算改为按尝试扣减,杜绝队列满时对每个文件反复入队 并刷出数万条 WARN(实测日志 41165 条) - 探测队列满时给文件挂 30 分钟退避 + 告警限速为每分钟一条 - 扫描时每个文件的海报/背景图由同步下载(单张最长 20s)改为 后台预取队列,云盘大库扫描不再串行拉图数小时 - PlaybackInfo 的云盘 ffprobe 探测改异步(原同步最长 8s, 既拖慢起播又放大云盘流量),带单飞去重 - 访问日志跳过 /api/health 与静态资源;logging.level/format 配置真正生效(此前是死配置) 登录稳定性(经常登录报错): - refresh token 未及时落库期间,刷新请求可识别「待落库令牌」, 不再把用户踢回登录页;轮换/登出后取消后台补写,防止旧令牌复活 QB 下载整理入库: - 新增 download.path_mappings 设置:自定义下载器→本程序路径映射 (每行 客户端路径=本地路径),并复用 compose 环境变量映射规则 - 应用重启后补整理最近 24h 内完成的种子(此前重启即永久漏掉) - 下载客户端初始化失败仍注册并惰性重连(容器启动顺序免疫) - 硬链接跨文件系统(EXDEV)自动降级为复制,保种语义不变 第三方播放器 404: - 播放处理器不再把所有错误吞成 404:媒体不存在→404, 云盘解析失败/STRM 关闭→502+原因 - 存库的云盘播放 URL 规范化为相对路径,免疫扫描时固化的旧 host - 云盘媒体 SupportsDirectPlay=false,强制走带鉴权的 DirectStream --- cmd/server/main.go | 17 ++++- internal/handler/emby.go | 19 ++++- internal/middleware/middleware.go | 16 +++- internal/service/download_manager_svc.go | 9 ++- internal/service/downloads.go | 86 ++++++++++++++++++++- internal/service/downloads_test.go | 56 +++++++++++++- internal/service/emby_compat.go | 87 +++++++++++++++------- internal/service/emby_compat_test.go | 43 ++++++++--- internal/service/qbittorrent.go | 4 + internal/service/scanner.go | 47 ++++++++++-- internal/service/stream.go | 36 ++++++++- internal/service/stream_normalize_test.go | 34 +++++++++ internal/service/stream_test.go | 4 +- internal/service/token_svc.go | 59 ++++++++++++--- internal/service/token_svc_pending_test.go | 78 +++++++++++++++++++ internal/service/transfer.go | 7 ++ 16 files changed, 530 insertions(+), 72 deletions(-) create mode 100644 internal/service/stream_normalize_test.go create mode 100644 internal/service/token_svc_pending_test.go diff --git a/cmd/server/main.go b/cmd/server/main.go index 92c8a50..9722e8c 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -278,11 +278,26 @@ func isFrontendLibraryRoute(path string) bool { return true } +// newLogger 根据 cfg.Logging 构建 Zap。此前 logging.level / logging.format +// 配置完全没有生效(固定 NewProduction),用户无法在生产环境降低日志量; +// 配合每请求一条 INFO 访问日志,几小时即可产生几十 MB 日志,在 Docker +// json-file 驱动下持续消耗磁盘 IO。 func newLogger(cfg *config.Config) (*zap.Logger, error) { if cfg.App.Debug { return zap.NewDevelopment() } - return zap.NewProduction() + zapCfg := zap.NewProductionConfig() + if level, err := zap.ParseAtomicLevel(strings.TrimSpace(cfg.Logging.Level)); err == nil && cfg.Logging.Level != "" { + zapCfg.Level = level + } + if strings.EqualFold(strings.TrimSpace(cfg.Logging.Format), "console") { + zapCfg.Encoding = "console" + } + if out := strings.TrimSpace(cfg.Logging.OutputPath); out != "" { + zapCfg.OutputPaths = append(zapCfg.OutputPaths, out) + zapCfg.ErrorOutputPaths = append(zapCfg.ErrorOutputPaths, out) + } + return zapCfg.Build() } // getLocalIP returns the first non-loopback IPv4 address of the machine. diff --git a/internal/handler/emby.go b/internal/handler/emby.go index cff79b2..5225e82 100644 --- a/internal/handler/emby.go +++ b/internal/handler/emby.go @@ -840,14 +840,27 @@ func embyVideoStreamHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { uid := embyUserID(c) item, err := svc.Emby.Item(c.Request.Context(), c.Param("id"), uid) - if err != nil || item == nil { + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + if item == nil { c.Status(http.StatusNotFound) return } - // 直接调用 Stream service 写入 response + // 直接调用 Stream service 写入 response。 + // 此前这里把所有错误一律吞成 404:云盘 Cookie 过期、直链解析失败、 + // STRM 播放被关闭……在第三方播放器上全部表现为「404 不存在」, + // 无法排查。现在区分:行不存在→404;云盘播放不可用/上游故障→502+原因。 err = svc.Stream.ServeFile(c.Writer, c.Request, c.Param("id")) - if err != nil { + switch { + case err == nil: + case errors.Is(err, service.ErrMediaNotFound): c.Status(http.StatusNotFound) + default: + if !c.Writer.Written() { + c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()}) + } } } } diff --git a/internal/middleware/middleware.go b/internal/middleware/middleware.go index fd9be73..75c18a6 100644 --- a/internal/middleware/middleware.go +++ b/internal/middleware/middleware.go @@ -22,14 +22,26 @@ const ( ) // RequestLogger logs one structured line per request. +// +// 健康检查与静态资源的成功请求被跳过:healthcheck 每 30s 一次、SPA 静态 +// 文件每页几十个请求,全部记 INFO 会让日志在几小时内膨胀到几十 MB, +// 在 Docker json-file 日志驱动下白白消耗磁盘 IO。 func RequestLogger(log *zap.Logger) gin.HandlerFunc { return func(c *gin.Context) { start := time.Now() c.Next() + path := c.Request.URL.Path + status := c.Writer.Status() + if status < 400 { + if path == "/api/health" || strings.HasPrefix(path, "/assets/") || + path == "/favicon.ico" || path == "/favicon.svg" { + return + } + } log.Info("http", zap.String("method", c.Request.Method), - zap.String("path", c.Request.URL.Path), - zap.Int("status", c.Writer.Status()), + zap.String("path", path), + zap.Int("status", status), zap.Duration("dur", time.Since(start)), zap.String("ip", c.ClientIP()), ) diff --git a/internal/service/download_manager_svc.go b/internal/service/download_manager_svc.go index 8c86812..cea39d7 100644 --- a/internal/service/download_manager_svc.go +++ b/internal/service/download_manager_svc.go @@ -73,17 +73,20 @@ func (m *DownloadManager) LoadAll(ctx context.Context) error { } if initErr := adapter.Initialize(ctx, cfg); initErr != nil { - m.log.Warn("failed to initialize download client", + // 初始化失败通常是 Docker 启动顺序问题(qBittorrent 还没就绪)。 + // 适配器内部支持按需重新登录(403/未登录时透明重试),所以 + // 这里仍然注册适配器,等下载器上线后自动恢复;此前直接 continue + // 会让该客户端在应用重启前永久不可用,下载完成也无法整理入库。 + m.log.Warn("download client init failed; registered for lazy reconnect", zap.String("id", dc.ID), zap.String("name", dc.Name), zap.Error(initErr), ) - continue } m.clients[dc.ID] = adapter m.configs[dc.ID] = cfg - m.log.Info("download client initialized", + m.log.Info("download client registered", zap.String("id", dc.ID), zap.String("name", dc.Name), zap.String("type", dc.Type), diff --git a/internal/service/downloads.go b/internal/service/downloads.go index 67f2b90..0366163 100644 --- a/internal/service/downloads.go +++ b/internal/service/downloads.go @@ -19,6 +19,7 @@ import ( "context" "errors" "math" + "os" "net/url" "path" "path/filepath" @@ -814,7 +815,15 @@ func (d *DownloadService) processDownloadSnapshot(ctx context.Context, live []QB wasComplete, wasSeen := d.prevStates[stateKey] switch { case complete && (firstSnapshot || !wasSeen): + // 首次快照里已完成的种子:此前一律标记「已见过」并跳过整理, + // 导致「下载完成时应用恰好不在线/正在重启」的种子永远不会被 + // 自动整理入库。现在对最近完成的种子补一次整理 + // (onTorrentComplete 内部仍受 organize.auto 开关约束,且 + // 整理对已存在的目标文件幂等跳过)。 d.prevStates[stateKey] = true + if recentlyCompletedTorrent(torrent, time.Now()) { + shouldQueue = true + } case complete && !wasComplete: shouldQueue = true case complete: @@ -915,6 +924,20 @@ func (d *DownloadService) markCompletedTorrentOrganizeDone(torrent QBitTorrent) d.mu.Unlock() } +// completedTorrentCatchupWindow 限定重启补整理只覆盖最近完成的种子, +// 防止每次启动都把全部历史种子重新过一遍整理流程。 +const completedTorrentCatchupWindow = 24 * time.Hour + +// recentlyCompletedTorrent 报告该种子是否在补整理时间窗内完成。 +// qBittorrent 未提供 completion_on 时保守地返回 false。 +func recentlyCompletedTorrent(torrent QBitTorrent, now time.Time) bool { + if torrent.CompletionOn <= 0 { + return false + } + completed := time.Unix(torrent.CompletionOn, 0) + return now.Sub(completed) <= completedTorrentCatchupWindow +} + func completedTorrentQueueKey(torrent QBitTorrent) string { hash := strings.ToLower(strings.TrimSpace(torrent.Hash)) if hash != "" { @@ -1010,7 +1033,7 @@ func (d *DownloadService) onTorrentComplete(ctx context.Context, torrent QBitTor d.log.Info("download completed, auto-organize disabled", zap.String("hash", torrent.Hash)) return } - source := d.completedTorrentSource(torrent) + source := d.completedTorrentSource(ctx, torrent) if source == "" { d.log.Warn("download completed but payload path is not accessible", zap.String("hash", torrent.Hash), @@ -1045,13 +1068,23 @@ func (d *DownloadService) onTorrentComplete(ctx context.Context, torrent QBitTor zap.Int("errors", len(res.Errors))) } -func (d *DownloadService) completedTorrentSource(torrent QBitTorrent) string { +// DownloadPathMappingsSettingKey 允许用户自定义「下载器路径 → 本程序路径」 +// 映射,每行一条,格式 `客户端路径=本地路径`(也接受 `=>` 或单个 `:` 分隔)。 +// qBittorrent 与本程序常在不同容器/主机里,对同一份数据看到的路径不同; +// 此前映射表是写死的三条猜测,对不上时整理静默失败。 +const DownloadPathMappingsSettingKey = "download.path_mappings" + +func (d *DownloadService) completedTorrentSource(ctx context.Context, torrent QBitTorrent) string { // 常见路径映射:qBittorrent容器路径 -> MediaStationGo容器路径 mappings := map[string]string{ "/var/apps/qBittorrent/shares/qBittorrent/Download": "/downloads", "/data/qBittorrent/downloads": "/downloads", "/downloads/qBittorrent": "/downloads", } + // 用户自定义映射优先(可覆盖内置猜测)。 + for clientPrefix, localPrefix := range d.userPathMappings(ctx) { + mappings[clientPrefix] = localPrefix + } for _, candidate := range []string{ torrent.ContentPath, filepath.Join(torrent.SavePath, torrent.Name), @@ -1064,6 +1097,55 @@ func (d *DownloadService) completedTorrentSource(torrent QBitTorrent) string { if translated := translateClientPath(clean, mappings); translated != "" { return translated } + // 复用 compose 注入的 MEDIASTATION_DOWNLOAD_DIR/MEDIA_DIR 宿主机↔容器 + // 映射(与媒体库路径换算同一套规则),覆盖「qB 跑在宿主机、 + // 本程序在容器里」的最常见部署形态。 + for _, mapped := range mappedPathCandidates(clean) { + if mapped == clean { + continue + } + if _, err := os.Stat(mapped); err == nil { + return mapped + } + } } return "" } + +// userPathMappings 解析用户配置的下载器路径映射。 +func (d *DownloadService) userPathMappings(ctx context.Context) map[string]string { + out := map[string]string{} + if d == nil || d.repo == nil || d.repo.Setting == nil { + return out + } + raw, err := d.repo.Setting.Get(ctx, DownloadPathMappingsSettingKey) + if err != nil { + return out + } + for _, line := range strings.Split(raw, "\n") { + line = strings.TrimSpace(line) + if line == "" || strings.HasPrefix(line, "#") { + continue + } + var from, to string + switch { + case strings.Contains(line, "=>"): + parts := strings.SplitN(line, "=>", 2) + from, to = parts[0], parts[1] + case strings.Contains(line, "="): + parts := strings.SplitN(line, "=", 2) + from, to = parts[0], parts[1] + case strings.Count(line, ":") == 1: + parts := strings.SplitN(line, ":", 2) + from, to = parts[0], parts[1] + default: + continue + } + from = strings.TrimSpace(from) + to = strings.TrimSpace(to) + if from != "" && to != "" { + out[from] = to + } + } + return out +} diff --git a/internal/service/downloads_test.go b/internal/service/downloads_test.go index dfdeec9..568fe9c 100644 --- a/internal/service/downloads_test.go +++ b/internal/service/downloads_test.go @@ -136,7 +136,7 @@ func TestCompletedTorrentSourceDoesNotFallbackToSavePath(t *testing.T) { } svc := NewDownloadService(zap.NewNop(), newOrganizerTestRepo(t), NewHub(zap.NewNop()), nil) - got := svc.completedTorrentSource(QBitTorrent{ + got := svc.completedTorrentSource(t.Context(), QBitTorrent{ Hash: "done123", Name: "Missing.Payload.S01", SavePath: savePath, @@ -190,6 +190,60 @@ func TestDownloadPollBaselinesAlreadyCompletedTorrents(t *testing.T) { } } +func TestDownloadPollCatchesUpRecentlyCompletedTorrents(t *testing.T) { + repos := newOrganizerTestRepo(t) + svc := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), nil) + + svc.processDownloadSnapshot(t.Context(), []QBitTorrent{ + {Hash: "fresh-complete", Name: "Fresh Complete S01E01", Progress: 1, CompletionOn: time.Now().Add(-time.Hour).Unix()}, + {Hash: "stale-complete", Name: "Stale Complete S01E01", Progress: 1, CompletionOn: time.Now().Add(-48 * time.Hour).Unix()}, + {Hash: "no-timestamp", Name: "No Timestamp S01E01", Progress: 1}, + }, nil) + + // 只有补整理时间窗内完成的种子会被补整理;无 completion_on 的保守跳过。 + if got := len(svc.organizeQueue); got != 1 { + t.Fatalf("first poll queued %d organize jobs, want 1 (recent completion only)", got) + } +} + +func TestCompletedTorrentSourceUsesConfiguredMapping(t *testing.T) { + root := t.TempDir() + localRoot := filepath.Join(root, "localdl") + payload := filepath.Join(localRoot, "Show.S01") + if err := os.MkdirAll(payload, 0o755); err != nil { + t.Fatal(err) + } + repos := newOrganizerTestRepo(t) + if err := repos.Setting.Set(t.Context(), DownloadPathMappingsSettingKey, "/qb/downloads="+localRoot); err != nil { + t.Fatal(err) + } + svc := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), nil) + + got := svc.completedTorrentSource(t.Context(), QBitTorrent{ContentPath: "/qb/downloads/Show.S01"}) + if got != payload { + t.Fatalf("completedTorrentSource = %q, want %q", got, payload) + } +} + +func TestUserPathMappingsParsing(t *testing.T) { + repos := newOrganizerTestRepo(t) + raw := "# comment\n/a=/b\n/c => /d\n/e:/f\nbad-line\n" + if err := repos.Setting.Set(t.Context(), DownloadPathMappingsSettingKey, raw); err != nil { + t.Fatal(err) + } + svc := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), nil) + got := svc.userPathMappings(t.Context()) + want := map[string]string{"/a": "/b", "/c": "/d", "/e": "/f"} + if len(got) != len(want) { + t.Fatalf("userPathMappings = %v, want %v", got, want) + } + for k, v := range want { + if got[k] != v { + t.Fatalf("mapping %q = %q, want %q", k, got[k], v) + } + } +} + func TestPublicDownloadTitleUsesMagnetDisplayName(t *testing.T) { got := publicDownloadTitle("magnet:?xt=urn:btih:abc&dn=%E6%B5%8B%E8%AF%95%E5%BD%B1%E7%89%87") if got != "测试影片" { diff --git a/internal/service/emby_compat.go b/internal/service/emby_compat.go index af7b4a9..ffc2325 100644 --- a/internal/service/emby_compat.go +++ b/internal/service/emby_compat.go @@ -69,6 +69,9 @@ type EmbyService struct { visibilityMu sync.RWMutex visibilityCache map[string]embyVisibilityCacheEntry + + cloudProbeMu sync.Mutex + cloudProbeInFlight map[string]struct{} } type cloudPlaybackResolver interface { @@ -1476,6 +1479,13 @@ func (e *EmbyService) PlaybackInfo(ctx context.Context, mediaID, userID string) }, nil } +// ensureCloudTrackMetadata 在后台补齐云盘媒体的轨道元数据。 +// +// 注意必须是异步的:此前这里在 PlaybackInfo 请求路径上同步执行 +// CloudResolve + ffprobe(HTTP)(最长 8 秒),既把第三方播放器的起播时间 +// 拖长到秒级,又让每一次点开详情/起播都可能触发一次云盘数据下载,是 +// Docker 部署下 CPU/带宽长期居高的来源之一。探测结果落库后,下一次 +// 请求自然能读到完整元数据。 func (e *EmbyService) ensureCloudTrackMetadata(ctx context.Context, m *model.Media) { if e == nil || m == nil || e.storage == nil || e.probe == nil || !mediaTrackMetadataMissing(m) { return @@ -1484,30 +1494,48 @@ func (e *EmbyService) ensureCloudTrackMetadata(ctx context.Context, m *model.Med if !ok { return } - probeCtx, cancel := context.WithTimeout(ctx, 8*time.Second) - defer cancel() - link, err := e.storage.CloudResolve(probeCtx, typ, ref, "") - if err != nil { - if e.log != nil { - e.log.Debug("resolve cloud media for playback probe failed", zap.String("media_id", m.ID), zap.Error(err)) + mediaID := m.ID + e.cloudProbeMu.Lock() + if e.cloudProbeInFlight == nil { + e.cloudProbeInFlight = make(map[string]struct{}) + } + if _, busy := e.cloudProbeInFlight[mediaID]; busy { + e.cloudProbeMu.Unlock() + return + } + e.cloudProbeInFlight[mediaID] = struct{}{} + e.cloudProbeMu.Unlock() + + go func() { + defer func() { + e.cloudProbeMu.Lock() + delete(e.cloudProbeInFlight, mediaID) + e.cloudProbeMu.Unlock() + }() + probeCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + link, err := e.storage.CloudResolve(probeCtx, typ, ref, "") + if err != nil { + if e.log != nil { + e.log.Debug("resolve cloud media for playback probe failed", zap.String("media_id", mediaID), zap.Error(err)) + } + return } - return - } - probe, err := e.probe.ProbeHTTP(probeCtx, link.URL, link.Headers) - if err != nil { - if e.log != nil { - e.log.Debug("playback cloud ffprobe failed", zap.String("media_id", m.ID), zap.Error(err)) + probe, err := e.probe.ProbeHTTP(probeCtx, link.URL, link.Headers) + if err != nil { + if e.log != nil { + e.log.Debug("playback cloud ffprobe failed", zap.String("media_id", mediaID), zap.Error(err)) + } + return } - return - } - updates := probeResultUpdates(probe) - if len(updates) == 0 { - return - } - if err := e.repo.DB.WithContext(ctx).Model(&model.Media{}).Where("id = ?", m.ID).Updates(updates).Error; err != nil && e.log != nil { - e.log.Debug("persist playback cloud probe failed", zap.String("media_id", m.ID), zap.Error(err)) - } - applyProbeResultToMediaValue(m, probe) + updates := probeResultUpdates(probe) + if len(updates) == 0 { + return + } + if err := e.repo.DB.WithContext(probeCtx).Model(&model.Media{}).Where("id = ?", mediaID).Updates(updates).Error; err != nil && e.log != nil { + e.log.Debug("persist playback cloud probe failed", zap.String("media_id", mediaID), zap.Error(err)) + } + }() } func mediaTrackMetadataMissing(m *model.Media) bool { @@ -1630,11 +1658,16 @@ func (e *EmbyService) mediaSource(m *model.Media, asEmbedded, directOnly bool) m "RequiresClosing": false, "ReadAtNativeFramerate": false, "SupportsTranscoding": !directOnly, - "SupportsDirectStream": true, - "SupportsDirectPlay": true, - "SupportsProbing": true, - "RunTimeTicks": int64(m.DurationSec) * 10_000_000, - "MediaStreams": e.mediaStreams(m), + // 云盘媒体禁用 DirectPlay:DirectPlay 语义是「客户端直接访问 + // Path」,而云盘媒体的 Path 是不带鉴权 token 的内部 /api/cloud/play + // 路径,Infuse/VidHub 等播放器直接请求会得到 401/404。强制它们走 + // DirectStream(/Videos/{id}/stream?api_key=...),由服务端校验后 + // 302 到云盘直链。 + "SupportsDirectStream": true, + "SupportsDirectPlay": !isCloud, + "SupportsProbing": true, + "RunTimeTicks": int64(m.DurationSec) * 10_000_000, + "MediaStreams": e.mediaStreams(m), } if !asEmbedded { src["DirectStreamUrl"] = embyDirectStreamURL(m.ID, container) diff --git a/internal/service/emby_compat_test.go b/internal/service/emby_compat_test.go index 6b93b38..91d5958 100644 --- a/internal/service/emby_compat_test.go +++ b/internal/service/emby_compat_test.go @@ -3,6 +3,7 @@ package service import ( "context" "testing" + "time" "github.com/glebarez/sqlite" "go.uber.org/zap" @@ -405,30 +406,45 @@ func TestEmbyPlaybackInfoProbesMissingCloudTrackMetadata(t *testing.T) { } svc.SetCloudProbe(resolver, prober) - pb, err := svc.PlaybackInfo(t.Context(), "cloud-probe-1", "user-1") - if err != nil { + if _, err := svc.PlaybackInfo(t.Context(), "cloud-probe-1", "user-1"); err != nil { t.Fatalf("playback info: %v", err) } + + // 探测现在是异步的(同步探测曾把起播拖慢最多 8 秒并放大云盘流量)。 + // 轮询等待后台探测结果落库。 + var persisted model.Media + deadline := time.Now().Add(3 * time.Second) + for { + if err := svc.repo.DB.First(&persisted, "id = ?", "cloud-probe-1").Error; err != nil { + t.Fatalf("reload media: %v", err) + } + if persisted.DurationSec > 0 || time.Now().After(deadline) { + break + } + time.Sleep(10 * time.Millisecond) + } + if persisted.DurationSec != 3661 || persisted.Width != 3840 || persisted.Height != 2160 || persisted.VideoCodec != "hevc" || persisted.AudioCodec != "eac3" { + t.Fatalf("probe metadata not persisted: %#v", persisted) + } if resolver.typ != "openlist" || resolver.ref != "/Movies/Movie.mkv" { t.Fatalf("resolver called with typ=%q ref=%q", resolver.typ, resolver.ref) } if prober.rawURL != "http://cdn.example.test/Movie.mkv" || prober.headers["Authorization"] != "Bearer probe-token" { t.Fatalf("probe called with url=%q headers=%#v", prober.rawURL, prober.headers) } + + // 落库之后,再次请求 PlaybackInfo 应当带上完整轨道元数据。 + pb, err := svc.PlaybackInfo(t.Context(), "cloud-probe-1", "user-1") + if err != nil { + t.Fatalf("playback info (second): %v", err) + } src := pb["MediaSources"].([]map[string]any)[0] if src["RunTimeTicks"] != int64(3661)*10_000_000 { - t.Fatalf("runtime ticks not filled from probe: %#v", src) + t.Fatalf("runtime ticks not filled after async probe: %#v", src) } streams := src["MediaStreams"].([]map[string]any) if len(streams) != 2 || streams[0]["Codec"] != "hevc" || streams[1]["Codec"] != "eac3" { - t.Fatalf("media streams not filled from probe: %#v", streams) - } - var persisted model.Media - if err := svc.repo.DB.First(&persisted, "id = ?", "cloud-probe-1").Error; err != nil { - t.Fatalf("reload media: %v", err) - } - if persisted.DurationSec != 3661 || persisted.Width != 3840 || persisted.Height != 2160 || persisted.VideoCodec != "hevc" || persisted.AudioCodec != "eac3" { - t.Fatalf("probe metadata not persisted: %#v", persisted) + t.Fatalf("media streams not filled after async probe: %#v", streams) } } @@ -438,6 +454,11 @@ func newTestEmbyService(t *testing.T) *EmbyService { if err != nil { t.Fatalf("open db: %v", err) } + // 内存库 + 异步探测协程:限制为单连接,避免连接池新建连接时 + // 拿到一个空白的 :memory: 实例(no such table)。 + if sqlDB, err := db.DB(); err == nil { + sqlDB.SetMaxOpenConns(1) + } if err := db.AutoMigrate(&model.Library{}, &model.Series{}, &model.Media{}, &model.Favorite{}, &model.PlaybackHistory{}, &model.User{}, &model.Setting{}); err != nil { t.Fatalf("migrate: %v", err) } diff --git a/internal/service/qbittorrent.go b/internal/service/qbittorrent.go index ce2e8c3..c397316 100644 --- a/internal/service/qbittorrent.go +++ b/internal/service/qbittorrent.go @@ -59,6 +59,10 @@ type QBitTorrent struct { // root folder. Prefer it for automatic organize so we do not scan the whole // download category. ContentPath string `json:"content_path"` + // CompletionOn 是 qBittorrent 报告的完成时间(Unix 秒,未完成为 0 或负值)。 + // 用于应用重启后的「补整理」判断:只补最近完成的种子,避免每次启动 + // 都重新触发全部历史种子的整理。 + CompletionOn int64 `json:"completion_on"` } // QBitClient is a thread-safe qBittorrent v2 API client. diff --git a/internal/service/scanner.go b/internal/service/scanner.go index 701b8a8..968f838 100644 --- a/internal/service/scanner.go +++ b/internal/service/scanner.go @@ -71,6 +71,8 @@ type ScannerService struct { cloudMediaProbeMu sync.Mutex cloudMediaProbing map[string]struct{} cloudMediaProbeBackoff map[string]time.Time + cloudMediaProbeWarnMu sync.Mutex + cloudMediaProbeLastWarn time.Time } // NewScannerService is the constructor. @@ -269,6 +271,10 @@ const maxCloudMediaProbeQueuePerScan = 32 const cloudMediaProbeFailureBackoff = 6 * time.Hour +// cloudMediaProbeQueueFullBackoff 是探测队列饱和时给单个文件挂的短退避, +// 防止后续扫描轮次对同一批文件反复尝试入队。 +const cloudMediaProbeQueueFullBackoff = 30 * time.Minute + // CloudScanStatus is the operator-facing state for long-running cloud scans. type CloudScanStatus struct { LibraryID string `json:"library_id"` @@ -359,9 +365,29 @@ func (s *ScannerService) queueCloudMediaProbe(typ, ref, path string) bool { default: s.cloudMediaProbeMu.Lock() delete(s.cloudMediaProbing, path) + // 队列满说明探测工人已饱和;给该文件挂一个短退避,避免下一轮 + // 扫描立刻重复尝试同一批文件。 + if s.cloudMediaProbeBackoff == nil { + s.cloudMediaProbeBackoff = make(map[string]time.Time) + } + s.cloudMediaProbeBackoff[path] = time.Now().Add(cloudMediaProbeQueueFullBackoff) s.cloudMediaProbeMu.Unlock() if s.log != nil { - s.log.Warn("cloud media probe queue full", zap.String("provider", typ), zap.String("path", path)) + // 限速告警:队列满在大库扫描中是常态而非异常,逐条 WARN 会 + // 在几小时内产生数万行日志(真实环境出现过 41165 条)。 + now := time.Now() + s.cloudMediaProbeWarnMu.Lock() + shouldWarn := now.Sub(s.cloudMediaProbeLastWarn) >= time.Minute + if shouldWarn { + s.cloudMediaProbeLastWarn = now + } + s.cloudMediaProbeWarnMu.Unlock() + if shouldWarn { + s.log.Warn("cloud media probe queue full; deferring remaining probes (logged at most once per minute)", + zap.String("provider", typ), zap.String("path", path)) + } else { + s.log.Debug("cloud media probe queue full", zap.String("provider", typ), zap.String("path", path)) + } } return false } @@ -372,14 +398,13 @@ func (s *ScannerService) queueCloudMediaProbeWithBudget(typ, ref, path string, b if *budget <= 0 { return false } - } - if !s.queueCloudMediaProbe(typ, ref, path) { - return false - } - if budget != nil { + // 预算按「尝试」扣减而不是按「成功入队」扣减。否则当探测队列被 + // 其他扫描填满时,本次扫描会对剩下的每一个文件都尝试入队并各打 + // 一条日志——真实环境里曾因此产生过 4 万多条 "queue full" WARN, + // 这本身就是一笔可观的 CPU/磁盘开销。 *budget-- } - return true + return s.queueCloudMediaProbe(typ, ref, path) } func (s *ScannerService) beginCloudScan(ctx context.Context, lib *model.Library, mount CloudMountInfo) (context.Context, func(*ScanResult, error), error) { @@ -917,7 +942,13 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar displayPath := joinCloudDisplayPath(displayDir, entry.Name) path := cloudMediaPath(typ, displayPath) localMeta := s.cloudFileMetadata(ctx, typ, displayPath, entry.Name, sidecars, dirMeta, librarySupportsSeasons(lib)) - s.cacheCloudMetadataArtworkNow(ctx, localMeta) + // 每个文件的海报/背景图改走后台预取队列。此前这里是同步 + // CloudResolve+下载(每张最多 20s 超时),几千个文件的云盘库 + // 扫描会变成持续数小时的串行下载,把 CPU/带宽长期吃满。 + if localMeta != nil { + s.queueCloudArtworkPrefetch(localMeta.PosterURL) + s.queueCloudArtworkPrefetch(localMeta.BackdropURL) + } candidate := cloudCandidate{ ref: ref, name: entry.Name, diff --git a/internal/service/stream.go b/internal/service/stream.go index 0fdc537..c05e77a 100644 --- a/internal/service/stream.go +++ b/internal/service/stream.go @@ -57,6 +57,32 @@ func NewStreamService(cfg *config.Config, log *zap.Logger, repo *repository.Cont // ErrMediaNotFound is returned when the media row or its file is missing. var ErrMediaNotFound = errors.New("media not found") +// ErrCloudPlaybackUnavailable 表示媒体行存在但属于云盘媒体、且当前无法 +// 构造可用的播放重定向(例如 STRM 播放被关闭或 STRMURL 缺失)。调用方 +// 应把它与「媒体不存在」区分开,避免把配置类故障当成 404 返回给播放器。 +var ErrCloudPlaybackUnavailable = errors.New("cloud media playback unavailable: strm playback disabled or media missing play url; re-scan the library or enable strm playback") + +// normalizeCloudPlayTarget 把存库的云盘播放 URL 规范化为相对路径。 +// +// STRMURL 是扫描时根据当时的 server_url/请求地址生成并固化进数据库的。 +// 在 Windows 开发机上扫描、再部署到 Docker(或更换了内网 IP/域名)后, +// 这些绝对 URL 会指向已失效的旧地址,第三方播放器跟随 302 就会拿到 +// 连接失败/404。这里只要能从 URL 中解析出 provider+ref,就重建为相对 +// /api/cloud/play 路径,由 absoluteInternalRedirect 基于「当前请求」补全 +// host,从而对历史脏数据免疫。 +func normalizeCloudPlayTarget(raw string) string { + typ, ref, ok := parseCloudMediaPlaybackURL(raw) + if !ok { + return raw + } + return BuildRelativeCloudPlayURL(typ, ref) +} + +// BuildRelativeCloudPlayURL 构造相对的云盘播放 API 路径。 +func BuildRelativeCloudPlayURL(typ, ref string) string { + return "/api/cloud/play/" + url.PathEscape(strings.TrimSpace(typ)) + "?" + url.Values{"ref": []string{ref}}.Encode() +} + // directPlayOnly reports whether the admin enabled「客户端直连解码」mode, // in which the host never transcodes (HLS is refused) and all playback is // handled by the client (direct play / 302 redirect). @@ -209,10 +235,18 @@ func (s *StreamService) ServeFile(w http.ResponseWriter, r *http.Request, mediaI return ErrMediaNotFound } if strings.TrimSpace(m.STRMURL) != "" && STRMPlaybackEnabled(r.Context(), s.repo) { - target := withAuthTokenForInternalRedirect(m.STRMURL, r, PublicServerURL(r.Context(), s.repo, s.cfg)) + // 云盘播放 URL 先规范化为相对路径,免疫扫描时固化的旧 host。 + target := normalizeCloudPlayTarget(m.STRMURL) + target = withAuthTokenForInternalRedirect(target, r, PublicServerURL(r.Context(), s.repo, s.cfg)) http.Redirect(w, r, absoluteInternalRedirect(target, r), http.StatusFound) return nil } + if strings.HasPrefix(strings.ToLower(strings.TrimSpace(m.Path)), "cloud://") { + // 云盘媒体没有本地文件可回退;走到这里说明 STRM 播放被关闭或 + // STRMURL 缺失。返回明确错误而不是笼统的「文件不存在」, + // 处理器据此回 502 + 原因,方便用户在播放器/日志里定位。 + return ErrCloudPlaybackUnavailable + } f, err := os.Open(m.Path) if err != nil { return ErrMediaNotFound diff --git a/internal/service/stream_normalize_test.go b/internal/service/stream_normalize_test.go new file mode 100644 index 0000000..52ab92d --- /dev/null +++ b/internal/service/stream_normalize_test.go @@ -0,0 +1,34 @@ +package service + +import ( + "net/url" + "testing" +) + +// TestNormalizeCloudPlayTarget 验证存库的云盘播放 URL(可能携带扫描时的 +// 旧 host)被规范化为相对路径,使 302 始终基于当前请求地址构造。 +func TestNormalizeCloudPlayTarget(t *testing.T) { + ref := "/电影/某部影片 (2024)/movie.mkv" + stale := "http://192.168.1.4:9011/api/cloud/play/openlist?ref=" + url.QueryEscape(ref) + got := normalizeCloudPlayTarget(stale) + want := BuildRelativeCloudPlayURL("openlist", ref) + if got != want { + t.Fatalf("normalizeCloudPlayTarget = %q, want %q", got, want) + } + parsed, err := url.Parse(got) + if err != nil { + t.Fatal(err) + } + if parsed.IsAbs() || parsed.Host != "" { + t.Fatalf("normalized target should be relative, got %q", got) + } + if parsed.Query().Get("ref") != ref { + t.Fatalf("ref round-trip failed: %q", parsed.Query().Get("ref")) + } + + // 非云盘播放 URL 保持原样(WebDAV/直链等)。 + passthrough := "https://dav.example.com/media/file.mkv" + if got := normalizeCloudPlayTarget(passthrough); got != passthrough { + t.Fatalf("non-cloud target should pass through, got %q", got) + } +} diff --git a/internal/service/stream_test.go b/internal/service/stream_test.go index 2810406..44333eb 100644 --- a/internal/service/stream_test.go +++ b/internal/service/stream_test.go @@ -103,7 +103,9 @@ func TestServeFileHonorsSTRMPlaybackDisabled(t *testing.T) { w := httptest.NewRecorder() err := svc.ServeFile(w, req, "cloud-1") - if err != ErrMediaNotFound { + // 云盘媒体在 STRM 播放关闭时返回明确的「云盘播放不可用」错误, + // 而不是和「媒体不存在」混在一起(后者会让播放器显示 404)。 + if err != ErrCloudPlaybackUnavailable { t.Fatalf("disabled STRM should not redirect cloud media, err=%v status=%d location=%q", err, w.Code, w.Header().Get("Location")) } if loc := w.Header().Get("Location"); loc != "" { diff --git a/internal/service/token_svc.go b/internal/service/token_svc.go index 7791b08..760842e 100644 --- a/internal/service/token_svc.go +++ b/internal/service/token_svc.go @@ -42,12 +42,22 @@ type TokenService struct { log *zap.Logger repo *repository.Container delayedStoreMu sync.Mutex - delayedStores map[string]struct{} + // delayedStores 记录「已发给客户端但还没写进库」的 refresh token。 + // 键是 token 哈希;值携带签发信息,让 Refresh 在落库完成前也能识别 + // 这些令牌——否则用户登录成功、一小时后 access token 过期,刷新时 + // 因为 refresh token 从未落库而被判定无效,被强制踢回登录页, + // 表现就是「经常登录报错」。 + delayedStores map[string]pendingRefreshToken +} + +type pendingRefreshToken struct { + UserID string + ExpiresAt time.Time } // NewTokenService 创建令牌服务实例。 func NewTokenService(cfg *config.Config, log *zap.Logger, repo *repository.Container) *TokenService { - return &TokenService{cfg: cfg, log: log, repo: repo, delayedStores: make(map[string]struct{})} + return &TokenService{cfg: cfg, log: log, repo: repo, delayedStores: make(map[string]pendingRefreshToken)} } // TokenPair 包含访问令牌和刷新令牌。 @@ -113,7 +123,7 @@ func (s *TokenService) issuePair(ctx context.Context, userID, role, tier string, zap.String("user_id", userID), zap.Error(err)) } - if s.trackDelayedStore(userID, tokenHash) { + if s.trackDelayedStore(userID, tokenHash, rt.ExpiresAt) { go s.storeRefreshTokenEventually(userID, tokenHash, rt.ExpiresAt) } } @@ -142,6 +152,11 @@ func (s *TokenService) storeRefreshTokenEventually(userID, tokenHash string, exp for attempt := 1; attempt <= 8; attempt++ { timer := time.NewTimer(delay) <-timer.C + // 令牌可能已在等待期间被轮换/登出(从 pending 表移除), + // 此时绝不能再写库,否则会复活一个已被替换的旧令牌。 + if _, stillPending := s.pendingDelayedStore(tokenHash); !stillPending { + return + } ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) err := s.storeRefreshToken(ctx, &model.RefreshToken{ UserID: userID, @@ -173,20 +188,19 @@ func (s *TokenService) storeRefreshTokenEventually(userID, tokenHash string, exp } } -func (s *TokenService) trackDelayedStore(userID, tokenHash string) bool { +func (s *TokenService) trackDelayedStore(userID, tokenHash string, expiresAt time.Time) bool { if s == nil { return false } - key := userID + "\x00" + tokenHash s.delayedStoreMu.Lock() defer s.delayedStoreMu.Unlock() if s.delayedStores == nil { - s.delayedStores = make(map[string]struct{}) + s.delayedStores = make(map[string]pendingRefreshToken) } - if _, ok := s.delayedStores[key]; ok { + if _, ok := s.delayedStores[tokenHash]; ok { return false } - s.delayedStores[key] = struct{}{} + s.delayedStores[tokenHash] = pendingRefreshToken{UserID: userID, ExpiresAt: expiresAt} return true } @@ -194,12 +208,22 @@ func (s *TokenService) untrackDelayedStore(userID, tokenHash string) { if s == nil { return } - key := userID + "\x00" + tokenHash s.delayedStoreMu.Lock() - delete(s.delayedStores, key) + delete(s.delayedStores, tokenHash) s.delayedStoreMu.Unlock() } +// pendingDelayedStore 返回尚未落库的 refresh token 信息(如果存在)。 +func (s *TokenService) pendingDelayedStore(tokenHash string) (pendingRefreshToken, bool) { + if s == nil { + return pendingRefreshToken{}, false + } + s.delayedStoreMu.Lock() + defer s.delayedStoreMu.Unlock() + pending, ok := s.delayedStores[tokenHash] + return pending, ok +} + func (s *TokenService) maxActiveRefreshTokens(ctx context.Context) int { cfg := loadBotConfig(ctx, s.repo) if cfg.MaxLoggedClients < 1 { @@ -244,7 +268,17 @@ func (s *TokenService) Refresh(ctx context.Context, refreshToken string) (*Token return nil, err } if rt == nil { - return nil, ErrInvalidRefreshToken + // 登录高峰/扫描写压力下,refresh token 可能还在后台补写队列里 + // 没来得及落库。此时令牌对客户端而言是合法的,不能判无效。 + pending, ok := s.pendingDelayedStore(tokenHash) + if !ok || time.Now().After(pending.ExpiresAt) { + return nil, ErrInvalidRefreshToken + } + rt = &model.RefreshToken{ + UserID: pending.UserID, + TokenHash: tokenHash, + ExpiresAt: pending.ExpiresAt, + } } // 检查是否已撤销 @@ -272,10 +306,11 @@ func (s *TokenService) Refresh(ctx context.Context, refreshToken string) (*Token return nil, ErrUserExpired } - // 撤销旧的 Refresh Token + // 撤销旧的 Refresh Token(包括可能仍在后台补写队列里的副本)。 if err := s.repo.RefreshToken.Revoke(ctx, tokenHash); err != nil { s.log.Warn("failed to revoke old refresh token", zap.Error(err)) } + s.untrackDelayedStore(rt.UserID, tokenHash) // 签发新的令牌对 return s.IssuePair(ctx, user.ID, user.Role, user.Tier) diff --git a/internal/service/token_svc_pending_test.go b/internal/service/token_svc_pending_test.go new file mode 100644 index 0000000..6fb2e77 --- /dev/null +++ b/internal/service/token_svc_pending_test.go @@ -0,0 +1,78 @@ +package service + +import ( + "testing" + "time" + + "github.com/glebarez/sqlite" + "go.uber.org/zap" + "gorm.io/gorm" + + "github.com/ShukeBta/MediaStationGo/internal/config" + "github.com/ShukeBta/MediaStationGo/internal/model" + "github.com/ShukeBta/MediaStationGo/internal/repository" +) + +func newTokenTestRepo(t *testing.T) *repository.Container { + t.Helper() + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&model.User{}, &model.RefreshToken{}, &model.Setting{}); err != nil { + t.Fatal(err) + } + return repository.New(db) +} + +// TestRefreshAcceptsPendingDelayedToken 验证:登录时因 SQLite 写压力未及时 +// 落库的 refresh token(仍在后台补写队列中)在刷新时被接受,而不是把用户 +// 踢回登录页(历史上「经常登录报错」的来源之一)。 +func TestRefreshAcceptsPendingDelayedToken(t *testing.T) { + repos := newTokenTestRepo(t) + cfg := &config.Config{} + cfg.Secrets.JWTSecret = "test-secret" + svc := NewTokenService(cfg, zap.NewNop(), repos) + + u := &model.User{Username: "u1", PasswordHash: "x", Role: "user", Tier: "free", IsActive: true} + if err := repos.User.Create(t.Context(), u); err != nil { + t.Fatal(err) + } + + refreshToken := "pending-token-value" + hash := repository.HashToken(refreshToken) + if !svc.trackDelayedStore(u.ID, hash, time.Now().Add(time.Hour)) { + t.Fatal("trackDelayedStore returned false") + } + + pair, err := svc.Refresh(t.Context(), refreshToken) + if err != nil { + t.Fatalf("Refresh rejected pending delayed token: %v", err) + } + if pair == nil || pair.AccessToken == "" || pair.RefreshToken == "" { + t.Fatalf("Refresh returned incomplete pair: %+v", pair) + } + // 轮换后旧令牌应从 pending 表移除,不能再次使用。 + if _, still := svc.pendingDelayedStore(hash); still { + t.Fatal("rotated pending token still tracked") + } + if _, err := svc.Refresh(t.Context(), refreshToken); err == nil { + t.Fatal("rotated pending token should not refresh twice") + } +} + +// TestRefreshRejectsExpiredPendingToken 验证过期的待落库令牌不会被接受。 +func TestRefreshRejectsExpiredPendingToken(t *testing.T) { + repos := newTokenTestRepo(t) + cfg := &config.Config{} + cfg.Secrets.JWTSecret = "test-secret" + svc := NewTokenService(cfg, zap.NewNop(), repos) + + refreshToken := "expired-pending" + hash := repository.HashToken(refreshToken) + svc.trackDelayedStore("user-x", hash, time.Now().Add(-time.Minute)) + + if _, err := svc.Refresh(t.Context(), refreshToken); err == nil { + t.Fatal("expired pending token should be rejected") + } +} diff --git a/internal/service/transfer.go b/internal/service/transfer.go index 8f47b8f..df1cc08 100644 --- a/internal/service/transfer.go +++ b/internal/service/transfer.go @@ -63,6 +63,13 @@ func transferFile(src, dst string, mode TransferMode) error { return copyFile(src, dst) case TransferHardlink: if err := os.Link(src, dst); err != nil { + // Docker 部署里下载目录和媒体目录往往是两个独立的 bind mount, + // 即使在宿主机上同属一块盘,容器内 os.Link 也会因跨文件系统 + // (EXDEV) 失败。此前直接报错导致 PT 下载完成后整理静默中断; + // 现在自动降级为复制(保留源文件继续做种,语义一致)。 + if copyErr := copyFile(src, dst); copyErr == nil { + return nil + } return fmt.Errorf("hardlink failed: %w; source and target must be on the same filesystem, choose copy if you want to duplicate data", err) } return nil From 8d1d1f18a873e007b8d79a161aabc7f3606ef722 Mon Sep 17 00:00:00 2001 From: ShukeBta Date: Fri, 12 Jun 2026 04:32:47 +0000 Subject: [PATCH 2/4] =?UTF-8?q?chore:=20=E9=99=90=E5=88=B6=E5=AE=B9?= =?UTF-8?q?=E5=99=A8=E6=97=A5=E5=BF=97=E4=BD=93=E7=A7=AF=E5=B9=B6=E8=A1=A5?= =?UTF-8?q?=E5=85=A8=20.gitignore?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - docker-compose 增加 json-file 日志上限(10m x 3),防止日志无限增长 持续消耗宿主机磁盘与 IO - .gitignore 补充 .tmp-* 运行产物、downloads/、media/、*.pid - 清理仓库目录中遗留的临时日志、诊断脚本与旧二进制(约 60MB,均未被 git 跟踪) --- .gitignore | 7 +++++++ docker-compose.yml | 8 ++++++++ 2 files changed, 15 insertions(+) diff --git a/.gitignore b/.gitignore index 5bbbb17..fbb7b7d 100644 --- a/.gitignore +++ b/.gitignore @@ -60,3 +60,10 @@ config.yaml # Editor backups *~ .tmp_* + +# Runtime / local-only artifacts (清理补充) +.tmp-live-backups/ +.tmp-* +downloads/ +media/ +*.pid diff --git a/docker-compose.yml b/docker-compose.yml index 46afb26..d77f696 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -76,3 +76,11 @@ services: timeout: 10s retries: 5 start_period: 30s + + # 限制容器日志体积:默认 json-file 驱动不设上限, + # 长期运行会持续吃宿主机磁盘与 IO。 + logging: + driver: json-file + options: + max-size: "10m" + max-file: "3" From 1568ae127a6a13b391230eefaec6f99d786a7093 Mon Sep 17 00:00:00 2001 From: ShukeBta Date: Fri, 12 Jun 2026 06:41:14 +0000 Subject: [PATCH 3/4] =?UTF-8?q?perf:=20=E9=87=8D=E6=9E=84=20FTS=20?= =?UTF-8?q?=E7=B4=A2=E5=BC=95=E4=B8=BA=20rowid=20=E5=AF=BB=E5=9D=80+?= =?UTF-8?q?=E8=A7=A6=E5=8F=91=E5=99=A8=E7=BB=B4=E6=8A=A4=EF=BC=8C=E6=A0=B9?= =?UTF-8?q?=E6=B2=BB=E9=87=8D=E5=90=AF=20CPU=20=E5=8D=A0=E6=BB=A1=E4=B8=8E?= =?UTF-8?q?=E6=89=AB=E6=8F=8F=E5=8D=A1=E6=AD=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - FTS5 普通列(含 UNINDEXED)不支持索引查找,旧版按 media_id 做 NOT EXISTS/DELETE 全是整表扫描:启动回填 O(N^2) 烧 CPU 数小时, 扫描时每个 upsert 一次全表扫描,均隔着全局写锁拖死登录(125s 超时) - v2 布局:FTS rowid 与 media.rowid 对齐,由 INSERT/UPDATE/DELETE 触发器实时维护;回填只剩 rowid 点查兜底,且去掉 ORDER BY - 顺带修复刮削直写 Updates() 后新标题搜不到的问题(触发器覆盖) - 写门改为 context 感知:长写语句不再让登录请求无限期排队 - warmMediaSearchIndex 延迟 30s 错峰启动 - library_scan 周期任务跳过云盘库(由 cloud_sync 夜间窗口负责), 首轮延迟从 15s 改为等满一个周期,重启不再立即全量扫描风暴 - prune 改为只取 id/path 并按 500 条批量删除,缩短写锁占用 --- internal/database/database.go | 119 ++++++++++++++++++++----- internal/repository/repository.go | 40 +++------ internal/repository/repository_test.go | 28 ++++-- internal/service/scanner.go | 54 +++++++---- internal/service/scheduler.go | 16 +++- internal/service/service.go | 7 ++ 6 files changed, 187 insertions(+), 77 deletions(-) diff --git a/internal/database/database.go b/internal/database/database.go index 4180acb..a8a1e8f 100644 --- a/internal/database/database.go +++ b/internal/database/database.go @@ -3,9 +3,9 @@ package database import ( + "context" "fmt" "path/filepath" - "sync" "github.com/glebarez/sqlite" "go.uber.org/zap" @@ -57,12 +57,23 @@ func installSQLiteWriteGate(db *gorm.DB) { if db == nil { return } - gate := &sqliteWriteGate{} + const lockedKey = "mediastation:sqlite_write_locked" + gate := newSQLiteWriteGate() lock := func(tx *gorm.DB) { - gate.Lock() + ctx := context.Background() + if tx.Statement != nil && tx.Statement.Context != nil { + ctx = tx.Statement.Context + } + if err := gate.Lock(ctx); err != nil { + _ = tx.AddError(err) + return + } + tx.InstanceSet(lockedKey, struct{}{}) } unlock := func(tx *gorm.DB) { - gate.Unlock() + if _, ok := tx.InstanceGet(lockedKey); ok { + gate.Unlock() + } } _ = db.Callback().Create().Before("gorm:create").Register("mediastation:sqlite_write_lock", lock) _ = db.Callback().Create().After("gorm:create").Register("mediastation:sqlite_write_unlock", unlock) @@ -74,16 +85,41 @@ func installSQLiteWriteGate(db *gorm.DB) { _ = db.Callback().Raw().After("gorm:raw").Register("mediastation:sqlite_write_unlock", unlock) } +// sqliteWriteGate 串行化进程内的 SQLite 写操作,避免多连接写竞争触发 +// SQLITE_BUSY。Lock 尊重语句自身的 context:此前用 sync.Mutex 时,一条 +// 长写语句(如 FTS 回填批次)会让登录等关键写操作无限期排队——客户端 +// 早已超时断开,goroutine 还挂在互斥锁上。现在等待方可随 context 取消 +// 及时失败,不再把整个进程的写路径拖死。 type sqliteWriteGate struct { - mu sync.Mutex + ch chan struct{} } -func (g *sqliteWriteGate) Lock() { - g.mu.Lock() +func newSQLiteWriteGate() *sqliteWriteGate { + return &sqliteWriteGate{ch: make(chan struct{}, 1)} +} + +func (g *sqliteWriteGate) Lock(ctx context.Context) error { + select { + case g.ch <- struct{}{}: + return nil + default: + } + if ctx == nil { + ctx = context.Background() + } + select { + case g.ch <- struct{}{}: + return nil + case <-ctx.Done(): + return ctx.Err() + } } func (g *sqliteWriteGate) Unlock() { - g.mu.Unlock() + select { + case <-g.ch: + default: + } } func buildDSN(cfg *config.Config) string { @@ -139,9 +175,31 @@ func ensurePerformanceIndexes(db *gorm.DB) error { return nil } +// mediaSearchIndexSchemaVersion 标识 FTS 索引的物理布局版本。 +// v2:FTS 行的 rowid 与 media.rowid 对齐,并由触发器实时维护。 +const mediaSearchIndexSchemaVersion = 2 + func ensureMediaSearchIndex(db *gorm.DB) error { - if mediaSearchIndexNeedsRebuild(db) { - _ = db.Exec(`DROP TABLE IF EXISTS media_search_fts`).Error + if err := db.Exec(`CREATE TABLE IF NOT EXISTS media_search_meta (id INTEGER PRIMARY KEY CHECK (id = 1), version INTEGER NOT NULL)`).Error; err != nil { + return nil + } + var version int + _ = db.Raw(`SELECT version FROM media_search_meta WHERE id = 1`).Scan(&version).Error + if version != mediaSearchIndexSchemaVersion { + // 旧版(v1)FTS 表按 UNINDEXED 的 media_id 寻址。FTS5 的普通列 + // 不支持索引查找,按 media_id 的 DELETE / NOT EXISTS 都是整表 + // 扫描:十几万行的库每次启动回填要做上百亿次行访问,纯 Go + // sqlite 直接把 CPU 钉满数小时,并隔着全局写锁拖死登录。 + // v2 起 FTS 行的 rowid 与 media.rowid 对齐,所有寻址走 rowid + // 点查,索引一致性交给下方触发器维护。 + for _, stmt := range []string{ + `DROP TRIGGER IF EXISTS media_search_fts_ai`, + `DROP TRIGGER IF EXISTS media_search_fts_au`, + `DROP TRIGGER IF EXISTS media_search_fts_ad`, + `DROP TABLE IF EXISTS media_search_fts`, + } { + _ = db.Exec(stmt).Error + } } if err := db.Exec(`CREATE VIRTUAL TABLE IF NOT EXISTS media_search_fts USING fts5(media_id UNINDEXED, title, original_name, path, genres, tokenize='trigram')`).Error; err != nil { if fallbackErr := db.Exec(`CREATE VIRTUAL TABLE IF NOT EXISTS media_search_fts USING fts5(media_id UNINDEXED, title, original_name, path, genres, tokenize='unicode61')`).Error; fallbackErr != nil { @@ -151,22 +209,35 @@ func ensureMediaSearchIndex(db *gorm.DB) error { return nil } } - return nil -} - -func mediaSearchIndexNeedsRebuild(db *gorm.DB) bool { - var cols []struct { - Name string - } - if err := db.Raw(`PRAGMA table_info(media_search_fts)`).Scan(&cols).Error; err != nil || len(cols) == 0 { - return false - } - for _, col := range cols { - if col.Name == "genres" { - return false + // 触发器让 FTS 与 media 行保持同步(新增/标题刮削改写/软删/恢复/ + // 硬删全覆盖),应用层不再需要按 media_id 手工刷新索引——也顺带 + // 修复了刮削直写 Updates() 后新标题搜不到的问题。 + for _, stmt := range []string{ + `CREATE TRIGGER IF NOT EXISTS media_search_fts_ai AFTER INSERT ON media WHEN new.deleted_at IS NULL BEGIN + DELETE FROM media_search_fts WHERE rowid = new.rowid; + INSERT INTO media_search_fts(rowid, media_id, title, original_name, path, genres) + VALUES (new.rowid, new.id, COALESCE(new.title, ''), COALESCE(new.original_name, ''), COALESCE(new.path, ''), COALESCE(new.genres, '')); + END`, + `CREATE TRIGGER IF NOT EXISTS media_search_fts_au AFTER UPDATE OF title, original_name, path, genres, deleted_at ON media BEGIN + DELETE FROM media_search_fts WHERE rowid = old.rowid; + INSERT INTO media_search_fts(rowid, media_id, title, original_name, path, genres) + SELECT new.rowid, new.id, COALESCE(new.title, ''), COALESCE(new.original_name, ''), COALESCE(new.path, ''), COALESCE(new.genres, '') + WHERE new.deleted_at IS NULL; + END`, + `CREATE TRIGGER IF NOT EXISTS media_search_fts_ad AFTER DELETE ON media BEGIN + DELETE FROM media_search_fts WHERE rowid = old.rowid; + END`, + } { + if err := db.Exec(stmt).Error; err != nil { + return err } } - return true + if version != mediaSearchIndexSchemaVersion { + if err := db.Exec(`INSERT INTO media_search_meta(id, version) VALUES (1, ?) ON CONFLICT(id) DO UPDATE SET version = excluded.version`, mediaSearchIndexSchemaVersion).Error; err != nil { + return err + } + } + return nil } func enforceTelegramBindingOneToOne(db *gorm.DB) error { diff --git a/internal/repository/repository.go b/internal/repository/repository.go index 898be8d..4ee91d6 100644 --- a/internal/repository/repository.go +++ b/internal/repository/repository.go @@ -323,7 +323,6 @@ func (r *MediaRepository) upsert(ctx context.Context, m *model.Media) error { m.ScrapeStatus = "pending" } if createErr := r.db.WithContext(ctx).Create(m).Error; createErr == nil { - _ = r.refreshSearchIndex(ctx, m.ID) return nil } else if retryErr := r.db.WithContext(ctx).Unscoped().Where("path = ?", m.Path).First(&existing).Error; retryErr != nil { return createErr @@ -426,7 +425,6 @@ func (r *MediaRepository) upsert(ctx context.Context, m *model.Media) error { Where("id = ?", existing.ID).Updates(updates).Error; err != nil { return err } - _ = r.refreshSearchIndex(ctx, existing.ID) // 回写 ID / 不可变字段,让 caller 拿到完整的现有行。 *m = existing return nil @@ -518,7 +516,7 @@ func (r *MediaRepository) searchFilteredFTS(ctx context.Context, query string, o var items []model.Media q := r.db.WithContext(ctx). Table("media"). - Joins("JOIN media_search_fts ON media_search_fts.media_id = media.id"). + Joins("JOIN media_search_fts ON media_search_fts.rowid = media.rowid"). Where("media.deleted_at IS NULL"). Where("media_search_fts MATCH ?", ftsQuery) q = applyQualifiedMediaQueryFilter(q, filter) @@ -625,23 +623,6 @@ func escapeLike(value string) string { return value } -func (r *MediaRepository) refreshSearchIndex(ctx context.Context, mediaID string) error { - if strings.TrimSpace(mediaID) == "" { - return nil - } - if !r.searchIndexEnabled(ctx) { - return nil - } - tx := r.db.WithContext(ctx) - _ = tx.Exec(`DELETE FROM media_search_fts WHERE media_id = ?`, mediaID).Error - return tx.Exec(` -INSERT INTO media_search_fts(media_id, title, original_name, path, genres) -SELECT id, COALESCE(title, ''), COALESCE(original_name, ''), COALESCE(path, ''), COALESCE(genres, '') -FROM media -WHERE id = ? AND deleted_at IS NULL -`, mediaID).Error -} - func (r *MediaRepository) BackfillSearchIndex(ctx context.Context, batchLimit int) (int64, error) { if batchLimit <= 0 { batchLimit = 1000 @@ -649,15 +630,19 @@ func (r *MediaRepository) BackfillSearchIndex(ctx context.Context, batchLimit in if !r.searchIndexEnabled(ctx) { return 0, nil } + // 关键性能点:FTS5 普通列(含 UNINDEXED)不支持索引查找,按 + // media_id 做 NOT EXISTS 是对 FTS 表的整表扫描,再叠加 ORDER BY + // 后每个批次都要对全部 media 行探测一遍——大库一次启动回填等于 + // 上百亿次行访问,曾把 CPU 钉满数小时。v2 布局下 FTS 行 rowid 与 + // media.rowid 对齐,NOT EXISTS 走 rowid 点查,且无需排序。 res := r.db.WithContext(ctx).Exec(` -INSERT INTO media_search_fts(media_id, title, original_name, path, genres) -SELECT m.id, COALESCE(m.title, ''), COALESCE(m.original_name, ''), COALESCE(m.path, ''), COALESCE(m.genres, '') +INSERT INTO media_search_fts(rowid, media_id, title, original_name, path, genres) +SELECT m.rowid, m.id, COALESCE(m.title, ''), COALESCE(m.original_name, ''), COALESCE(m.path, ''), COALESCE(m.genres, '') FROM media AS m WHERE m.deleted_at IS NULL AND NOT EXISTS ( - SELECT 1 FROM media_search_fts AS f WHERE f.media_id = m.id + SELECT 1 FROM media_search_fts AS f WHERE f.rowid = m.rowid ) -ORDER BY m.created_at DESC LIMIT ? `, batchLimit) return res.RowsAffected, res.Error @@ -679,18 +664,13 @@ func (r *MediaRepository) searchIndexEnabled(ctx context.Context) bool { // DeleteByLibrary purges all media tied to a library. func (r *MediaRepository) DeleteByLibrary(ctx context.Context, libraryID string) error { - if r.searchIndexEnabled(ctx) { - _ = r.db.WithContext(ctx).Exec(`DELETE FROM media_search_fts WHERE media_id IN (SELECT id FROM media WHERE library_id = ?)`, libraryID).Error - } + // FTS 行由 media 表上的触发器同步清理(软删/硬删都覆盖)。 return r.db.WithContext(ctx).Where("library_id = ?", libraryID).Delete(&model.Media{}).Error } // PurgeByLibrary permanently removes media tied to a library. Used for virtual // cloud mounts where "remove mount" must not populate the recycle bin. func (r *MediaRepository) PurgeByLibrary(ctx context.Context, libraryID string) error { - if r.searchIndexEnabled(ctx) { - _ = r.db.WithContext(ctx).Exec(`DELETE FROM media_search_fts WHERE media_id IN (SELECT id FROM media WHERE library_id = ?)`, libraryID).Error - } return r.db.WithContext(ctx).Unscoped().Where("library_id = ?", libraryID).Delete(&model.Media{}).Error } diff --git a/internal/repository/repository_test.go b/internal/repository/repository_test.go index cd471e9..7a2e37c 100644 --- a/internal/repository/repository_test.go +++ b/internal/repository/repository_test.go @@ -153,12 +153,17 @@ func TestMediaSearchIndexBackfillRunsInBatches(t *testing.T) { }).Error; err != nil { t.Fatal(err) } - var before int64 - if err := repos.DB.Raw(`SELECT COUNT(*) FROM media_search_fts`).Scan(&before).Error; err != nil { + // 插入触发器应当同步维护 FTS 行。 + var indexed int64 + if err := repos.DB.Raw(`SELECT COUNT(*) FROM media_search_fts`).Scan(&indexed).Error; err != nil { t.Fatal(err) } - if before != 0 { - t.Fatalf("startup migrate should not synchronously backfill FTS, got %d rows", before) + if indexed != 1 { + t.Fatalf("insert trigger should index new media, got %d rows", indexed) + } + // 清空 FTS 模拟旧库升级后索引缺失,回填应按批补齐且 rowid 对齐。 + if err := repos.DB.Exec(`DELETE FROM media_search_fts`).Error; err != nil { + t.Fatal(err) } n, err := repos.Media.BackfillSearchIndex(t.Context(), 1) if err != nil { @@ -167,11 +172,22 @@ func TestMediaSearchIndexBackfillRunsInBatches(t *testing.T) { if n != 1 { t.Fatalf("backfilled rows = %d, want 1", n) } + var aligned int64 + if err := repos.DB.Raw(`SELECT COUNT(*) FROM media_search_fts f JOIN media m ON f.rowid = m.rowid AND f.media_id = m.id`).Scan(&aligned).Error; err != nil { + t.Fatal(err) + } + if aligned != 1 { + t.Fatalf("fts rows aligned with media rowid = %d, want 1", aligned) + } + // 软删除后触发器应清理对应 FTS 行,避免搜索命中已删媒体。 + if err := repos.DB.Delete(&model.Media{}, "id = ?", "m-backfill").Error; err != nil { + t.Fatal(err) + } var after int64 if err := repos.DB.Raw(`SELECT COUNT(*) FROM media_search_fts`).Scan(&after).Error; err != nil { t.Fatal(err) } - if after != 1 { - t.Fatalf("fts rows = %d, want 1", after) + if after != 0 { + t.Fatalf("soft delete should drop fts row, got %d", after) } } diff --git a/internal/service/scanner.go b/internal/service/scanner.go index 968f838..a574dea 100644 --- a/internal/service/scanner.go +++ b/internal/service/scanner.go @@ -1491,13 +1491,20 @@ func (s *ScannerService) mediaPathExists(ctx context.Context, path string) bool } func (s *ScannerService) pruneMissingMedia(ctx context.Context, libraryID string, seen map[string]struct{}) (int64, error) { - var rows []model.Media + // 只取 id/path,并把删除按批提交:此前整表载入完整 Media 结构体、 + // 每行一条 DELETE,大库 prune 既费内存又长期占用写锁。 + var rows []struct { + ID string + Path string + } if err := s.repo.DB.WithContext(ctx). + Model(&model.Media{}). + Select("id, path"). Where("library_id = ?", libraryID). Find(&rows).Error; err != nil { return 0, err } - var removed int64 + stale := make([]string, 0) for _, row := range rows { if row.Path == "" { continue @@ -1510,9 +1517,26 @@ func (s *ScannerService) pruneMissingMedia(ctx context.Context, libraryID string } else if !os.IsNotExist(err) { continue } - res := s.repo.DB.WithContext(ctx). - Where("id = ?", row.ID). - Delete(&model.Media{}) + stale = append(stale, row.ID) + } + return s.deleteMediaByIDs(ctx, stale, false) +} + +// deleteMediaByIDs removes media rows in fixed-size batches so each write +// transaction stays short and the global write gate is released frequently. +func (s *ScannerService) deleteMediaByIDs(ctx context.Context, ids []string, hard bool) (int64, error) { + const batch = 500 + var removed int64 + for i := 0; i < len(ids); i += batch { + end := i + batch + if end > len(ids) { + end = len(ids) + } + q := s.repo.DB.WithContext(ctx) + if hard { + q = q.Unscoped() + } + res := q.Where("id IN ?", ids[i:end]).Delete(&model.Media{}) if res.Error != nil { return removed, res.Error } @@ -1522,27 +1546,25 @@ func (s *ScannerService) pruneMissingMedia(ctx context.Context, libraryID string } func (s *ScannerService) pruneMissingCloudMedia(ctx context.Context, libraryID string, seen map[string]struct{}) (int64, error) { - var rows []model.Media + var rows []struct { + ID string + Path string + } if err := s.repo.DB.WithContext(ctx). + Model(&model.Media{}). + Select("id, path"). Where("library_id = ? AND path LIKE ?", libraryID, "cloud://%"). Find(&rows).Error; err != nil { return 0, err } - var removed int64 + stale := make([]string, 0) for _, row := range rows { if _, ok := seen[row.Path]; ok { continue } - res := s.repo.DB.WithContext(ctx). - Unscoped(). - Where("id = ?", row.ID). - Delete(&model.Media{}) - if res.Error != nil { - return removed, res.Error - } - removed += res.RowsAffected + stale = append(stale, row.ID) } - return removed, nil + return s.deleteMediaByIDs(ctx, stale, true) } func parseCloudLibraryPath(raw string) (typ, dirID string, ok bool) { diff --git a/internal/service/scheduler.go b/internal/service/scheduler.go index d47a9d9..f1e635d 100644 --- a/internal/service/scheduler.go +++ b/internal/service/scheduler.go @@ -134,7 +134,14 @@ func (s *SchedulerService) Start(ctx context.Context) { }, } for _, j := range s.jobs { - go s.loop(ctx, j) + initialDelay := 15 * time.Second + if j.name == "library_scan" { + // 重启后不立即整库重扫:更新/重启窗口恰是登录高峰,启动 + // 15 秒即全量扫描曾把 CPU/磁盘打满导致无法登录。首轮等满 + // 一个完整周期再跑,平时的每小时节奏不变。 + initialDelay = j.interval + } + go s.loopWithInitialDelay(ctx, j, initialDelay) } } @@ -259,6 +266,13 @@ func (s *SchedulerService) jobScanLibraries(ctx context.Context) error { if !l.Enabled { continue } + if _, ok := ParseCloudLibraryMount(l.Path); ok { + // 云盘库由 cloud_sync 任务在夜间窗口低频同步;周期性整库 + // 重扫只面向本地磁盘库。否则十几个云盘库每小时全量遍历 + // 会把 CPU/网络长期吃满,还会占住唯一的云扫描槽位,让 + // 手动扫描看起来一直"卡死"在排队。 + continue + } if _, err := s.scanner.ScanLibrary(ctx, l.ID); err != nil { s.log.Warn("scheduled scan failed", zap.String("library", l.ID), zap.Error(err)) diff --git a/internal/service/service.go b/internal/service/service.go index fee4398..78d1c3e 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -268,6 +268,13 @@ func (c *Container) warmMediaSearchIndex(ctx context.Context) { if c == nil || c.Repo == nil || c.Repo.Media == nil { return } + // 错峰:FTS 正常由 media 表触发器实时维护,回填只是升级或异常后的 + // 兜底。先让登录、首页等关键路径跑起来,再开始后台补索引。 + select { + case <-ctx.Done(): + return + case <-time.After(30 * time.Second): + } const batchSize = 1000 const pause = 100 * time.Millisecond total := int64(0) From d270e9034f879907723a10c7c4f847c5f433e377 Mon Sep 17 00:00:00 2001 From: ShukeBta Date: Fri, 12 Jun 2026 19:06:47 +0800 Subject: [PATCH 4/4] =?UTF-8?q?fix:=20=E9=98=B2=E6=AD=A2=E5=90=AF=E5=8A=A8?= =?UTF-8?q?=E6=95=B4=E7=90=86=E6=89=AB=E6=8F=8F=E9=98=BB=E5=A1=9E=E7=99=BB?= =?UTF-8?q?=E5=BD=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/handler/auth.go | 2 +- internal/service/audit.go | 14 ++++ internal/service/auth.go | 15 ++++- internal/service/auth_user_limits_test.go | 5 ++ internal/service/downloads.go | 58 +++++++++++++++- internal/service/downloads_test.go | 69 +++++++++++++++++++ internal/service/organizer_scan.go | 8 +++ internal/service/path_translation.go | 11 +++- internal/service/scheduler.go | 18 +++-- internal/service/service.go | 66 ++++++++++++++++++- internal/service/token_svc.go | 80 +++++++++++++++++------ 11 files changed, 312 insertions(+), 34 deletions(-) diff --git a/internal/handler/auth.go b/internal/handler/auth.go index 9e85fcc..658c7f7 100644 --- a/internal/handler/auth.go +++ b/internal/handler/auth.go @@ -45,7 +45,7 @@ func loginHandler(svc *service.Container) gin.HandlerFunc { "user": resp.User, "tokens": resp.Tokens, }) - svc.Audit.Record(c.Request.Context(), resp.User.ID, "auth.login", resp.User.Username, c.ClientIP(), "") + svc.Audit.RecordBestEffort(resp.User.ID, "auth.login", resp.User.Username, c.ClientIP(), "") } } diff --git a/internal/service/audit.go b/internal/service/audit.go index 80a80f8..cd13bfc 100644 --- a/internal/service/audit.go +++ b/internal/service/audit.go @@ -8,6 +8,7 @@ package service import ( "context" + "time" "go.uber.org/zap" @@ -39,3 +40,16 @@ func (a *AuditService) Record(ctx context.Context, userID, action, target, ip, d a.log.Debug("audit write failed", zap.Error(err)) } } + +// RecordBestEffort writes an audit row off the request path. Login must not be +// held open by SQLite write pressure from scans or background maintenance. +func (a *AuditService) RecordBestEffort(userID, action, target, ip, detail string) { + if a == nil || a.repo == nil || a.repo.Log == nil { + return + } + go func() { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + a.Record(ctx, userID, action, target, ip, detail) + }() +} diff --git a/internal/service/auth.go b/internal/service/auth.go index 1293c34..f9c4435 100644 --- a/internal/service/auth.go +++ b/internal/service/auth.go @@ -163,10 +163,23 @@ func (s *AuthService) Login(ctx context.Context, username, password string) (*Lo if err != nil { return nil, err } - _ = s.repo.User.TouchLogin(ctx, u.ID) + s.touchLoginBestEffort(u.ID) return &LoginResponse{User: u, Tokens: tokens}, nil } +func (s *AuthService) touchLoginBestEffort(userID string) { + if s == nil || s.repo == nil || s.repo.User == nil || strings.TrimSpace(userID) == "" { + return + } + go func() { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + if err := s.repo.User.TouchLogin(ctx, userID); err != nil && s.log != nil { + s.log.Debug("touch login delayed", zap.String("user_id", userID), zap.Error(err)) + } + }() +} + // ChangePassword updates the user password if the old one matches. func (s *AuthService) ChangePassword(ctx context.Context, userID, oldPwd, newPwd string) error { if strings.TrimSpace(newPwd) == "" || len(newPwd) < 6 { diff --git a/internal/service/auth_user_limits_test.go b/internal/service/auth_user_limits_test.go index 9690c39..43f3be8 100644 --- a/internal/service/auth_user_limits_test.go +++ b/internal/service/auth_user_limits_test.go @@ -30,6 +30,11 @@ func newAuthTestServices(t *testing.T) (*repository.Container, *AuthService, *Pr if err := db.AutoMigrate(&model.User{}, &model.UserPermission{}, &model.RefreshToken{}, &model.TelegramBinding{}, &model.Setting{}); err != nil { t.Fatal(err) } + sqlDB, err := db.DB() + if err != nil { + t.Fatal(err) + } + sqlDB.SetMaxOpenConns(1) repos := repository.New(db) cfg := &config.Config{} cfg.Secrets.JWTSecret = "test-secret" diff --git a/internal/service/downloads.go b/internal/service/downloads.go index 0366163..84489b9 100644 --- a/internal/service/downloads.go +++ b/internal/service/downloads.go @@ -17,10 +17,12 @@ package service import ( "context" + "crypto/sha1" "errors" + "fmt" "math" - "os" "net/url" + "os" "path" "path/filepath" "regexp" @@ -821,7 +823,7 @@ func (d *DownloadService) processDownloadSnapshot(ctx context.Context, live []QB // (onTorrentComplete 内部仍受 organize.auto 开关约束,且 // 整理对已存在的目标文件幂等跳过)。 d.prevStates[stateKey] = true - if recentlyCompletedTorrent(torrent, time.Now()) { + if recentlyCompletedTorrent(torrent, time.Now()) && !d.completedTorrentCatchupRecorded(ctx, torrent) { shouldQueue = true } case complete && !wasComplete: @@ -928,6 +930,8 @@ func (d *DownloadService) markCompletedTorrentOrganizeDone(torrent QBitTorrent) // 防止每次启动都把全部历史种子重新过一遍整理流程。 const completedTorrentCatchupWindow = 24 * time.Hour +const completedTorrentCatchupSettingPrefix = "download.auto_organized." + // recentlyCompletedTorrent 报告该种子是否在补整理时间窗内完成。 // qBittorrent 未提供 completion_on 时保守地返回 false。 func recentlyCompletedTorrent(torrent QBitTorrent, now time.Time) bool { @@ -938,6 +942,46 @@ func recentlyCompletedTorrent(torrent QBitTorrent, now time.Time) bool { return now.Sub(completed) <= completedTorrentCatchupWindow } +func (d *DownloadService) completedTorrentCatchupRecorded(ctx context.Context, torrent QBitTorrent) bool { + if d == nil || d.repo == nil || d.repo.Setting == nil { + return false + } + key := completedTorrentCatchupSettingKey(torrent) + if key == "" { + return false + } + value, err := d.repo.Setting.Get(ctx, key) + if err != nil { + return false + } + return parseBoolSetting(value, false) +} + +func (d *DownloadService) markCompletedTorrentCatchupRecorded(ctx context.Context, torrent QBitTorrent) { + if d == nil || d.repo == nil || d.repo.Setting == nil { + return + } + key := completedTorrentCatchupSettingKey(torrent) + if key == "" { + return + } + if err := d.repo.Setting.Set(ctx, key, "true"); err != nil && d.log != nil { + d.log.Debug("mark completed torrent catchup failed", + zap.String("hash", torrent.Hash), + zap.String("name", torrent.Name), + zap.Error(err)) + } +} + +func completedTorrentCatchupSettingKey(torrent QBitTorrent) string { + key := completedTorrentQueueKey(torrent) + if key == "" { + return "" + } + sum := sha1.Sum([]byte(key)) + return completedTorrentCatchupSettingPrefix + fmt.Sprintf("%x", sum[:]) +} + func completedTorrentQueueKey(torrent QBitTorrent) string { hash := strings.ToLower(strings.TrimSpace(torrent.Hash)) if hash != "" { @@ -1054,9 +1098,17 @@ func (d *DownloadService) onTorrentComplete(ctx context.Context, torrent QBitTor zap.Error(err)) return } - if d.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" { + if d.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" && OrganizeResultHasChanges(res) { res.Scans, res.Scrapes = d.scanner.ScanAndScrapeLibrariesForPath(ctx, res.DestPath, "", OrganizeScrapeAfterEnabled(ctx, d.repo)) + } else if d.log != nil && res != nil && !OrganizeResultHasChanges(res) { + d.log.Info("auto organize completed torrent skipped scan; no destination changes", + zap.String("hash", torrent.Hash), + zap.String("source", source), + zap.Int("organized", res.Organized), + zap.Int("replaced", res.Replaced), + zap.Int("skipped", res.Skipped)) } + d.markCompletedTorrentCatchupRecorded(context.Background(), torrent) d.log.Info("auto organize completed torrent finished", zap.String("hash", torrent.Hash), zap.String("source", source), diff --git a/internal/service/downloads_test.go b/internal/service/downloads_test.go index 568fe9c..1d332a4 100644 --- a/internal/service/downloads_test.go +++ b/internal/service/downloads_test.go @@ -206,6 +206,75 @@ func TestDownloadPollCatchesUpRecentlyCompletedTorrents(t *testing.T) { } } +func TestDownloadPollSkipsRecordedCompletedTorrentCatchup(t *testing.T) { + repos := newOrganizerTestRepo(t) + torrent := QBitTorrent{ + Hash: "fresh-complete", + Name: "Fresh Complete S01E01", + Progress: 1, + CompletionOn: time.Now().Add(-time.Hour).Unix(), + } + if err := repos.Setting.Set(t.Context(), completedTorrentCatchupSettingKey(torrent), "true"); err != nil { + t.Fatal(err) + } + svc := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), nil) + + svc.processDownloadSnapshot(t.Context(), []QBitTorrent{torrent}, nil) + + if got := len(svc.organizeQueue); got != 0 { + t.Fatalf("recorded completed torrent queued %d organize jobs, want 0", got) + } +} + +func TestAutoOrganizeSkipsScanWhenNoFilesChanged(t *testing.T) { + root := t.TempDir() + src := filepath.Join(root, "downloads", "国产剧", "狂飙.S01E01.2023.1080p.mkv") + dest := filepath.Join(root, "media") + writeOrgFile(t, src, "episode") + + repos := newOrganizerTestRepo(t) + for key, value := range map[string]string{ + "organizer.auto_after_download": "true", + "organize.target_dir": dest, + "organize.transfer_mode": "copy", + } { + if err := repos.Setting.Set(t.Context(), key, value); err != nil { + t.Fatal(err) + } + } + lib := model.Library{Name: "国产剧", Path: filepath.Join(dest, "电视剧", "国产剧"), Type: "tv", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + if _, err := org.OrganizeDirectory(t.Context(), OrganizeOptions{ + SourcePath: src, + DestPath: dest, + TransferMode: TransferCopy, + }); err != nil { + t.Fatalf("seed organized destination: %v", err) + } + scanner := NewScannerService(&config.Config{}, zap.NewNop(), repos, NewHub(zap.NewNop()), nil, nil) + svc := NewDownloadService(zap.NewNop(), repos, NewHub(zap.NewNop()), org) + svc.SetScanner(scanner) + + svc.onTorrentComplete(t.Context(), QBitTorrent{ + Hash: "done123", + Name: "狂飙.S01E01.2023.1080p", + Progress: 1, + SavePath: filepath.Dir(src), + ContentPath: src, + }) + + var count int64 + if err := repos.DB.Model(&model.Media{}).Count(&count).Error; err != nil { + t.Fatal(err) + } + if count != 0 { + t.Fatalf("no-op auto organize triggered scan and created %d media rows, want 0", count) + } +} + func TestCompletedTorrentSourceUsesConfiguredMapping(t *testing.T) { root := t.TempDir() localRoot := filepath.Join(root, "localdl") diff --git a/internal/service/organizer_scan.go b/internal/service/organizer_scan.go index 360029f..17f3bdb 100644 --- a/internal/service/organizer_scan.go +++ b/internal/service/organizer_scan.go @@ -47,6 +47,14 @@ func OrganizeScrapeAfterEnabled(ctx context.Context, repo *repository.Container) return false } +// OrganizeResultHasChanges reports whether an organize run actually changed +// files in the destination library. Skipped duplicates are intentionally not a +// change: scanning after a no-op organize can turn a harmless restart into a +// full library ffprobe sweep. +func OrganizeResultHasChanges(res *OrganizeResult) bool { + return res != nil && (res.Organized > 0 || res.Replaced > 0) +} + // ScanLibrariesForPath recursively scans libraries affected by an organize // destination. If preferredLibraryID is set, only that library is scanned. // Otherwise every enabled library whose path intersects destRoot is scanned; diff --git a/internal/service/path_translation.go b/internal/service/path_translation.go index d2e446a..cd22332 100644 --- a/internal/service/path_translation.go +++ b/internal/service/path_translation.go @@ -18,9 +18,16 @@ func translateClientPath(clientPath string, mappings map[string]string) string { return clean } // 尝试路径映射 + cleanForMatch := filepath.ToSlash(clean) for clientPrefix, localPrefix := range mappings { - if strings.HasPrefix(clean, clientPrefix) { - translated := filepath.Join(localPrefix, strings.TrimPrefix(clean, clientPrefix)) + prefix := strings.TrimRight(filepath.ToSlash(filepath.Clean(clientPrefix)), "/") + if prefix == "" || prefix == "." { + continue + } + if cleanForMatch == prefix || strings.HasPrefix(cleanForMatch, prefix+"/") { + rel := strings.TrimPrefix(cleanForMatch, prefix) + rel = strings.TrimPrefix(rel, "/") + translated := filepath.Join(localPrefix, filepath.FromSlash(rel)) if _, err := os.Stat(translated); err == nil { return translated } diff --git a/internal/service/scheduler.go b/internal/service/scheduler.go index f1e635d..4a90913 100644 --- a/internal/service/scheduler.go +++ b/internal/service/scheduler.go @@ -135,10 +135,10 @@ func (s *SchedulerService) Start(ctx context.Context) { } for _, j := range s.jobs { initialDelay := 15 * time.Second - if j.name == "library_scan" { - // 重启后不立即整库重扫:更新/重启窗口恰是登录高峰,启动 - // 15 秒即全量扫描曾把 CPU/磁盘打满导致无法登录。首轮等满 - // 一个完整周期再跑,平时的每小时节奏不变。 + if j.name == "library_scan" || j.name == "organize_source" { + // 重启后不立即整库重扫/整理下载目录:更新窗口恰是登录高峰, + // 15 秒即全量 walk + ffprobe 曾把 CPU/磁盘打满导致无法登录。 + // 首轮等满一个完整周期再跑,平时节奏不变。 initialDelay = j.interval } go s.loopWithInitialDelay(ctx, j, initialDelay) @@ -483,8 +483,16 @@ func (s *SchedulerService) jobOrganizeSource(ctx context.Context) error { if err != nil { return err } - if s.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" { + if s.scanner != nil && res != nil && strings.TrimSpace(res.DestPath) != "" && OrganizeResultHasChanges(res) { res.Scans, res.Scrapes = s.scanner.ScanAndScrapeLibrariesForPath(ctx, res.DestPath, "", OrganizeScrapeAfterEnabled(ctx, s.repo)) + } else if s.log != nil && res != nil && !OrganizeResultHasChanges(res) { + s.log.Info("scheduled source organize skipped scan; no destination changes", + zap.String("source", res.SourcePath), + zap.String("dest", res.DestPath), + zap.Int("organized", res.Organized), + zap.Int("replaced", res.Replaced), + zap.Int("skipped", res.Skipped), + ) } if s.log != nil && res != nil { s.log.Info("scheduled source organize finished", diff --git a/internal/service/service.go b/internal/service/service.go index 78d1c3e..916967d 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -5,6 +5,7 @@ package service import ( "context" + "strconv" "strings" "time" @@ -268,15 +269,21 @@ func (c *Container) warmMediaSearchIndex(ctx context.Context) { if c == nil || c.Repo == nil || c.Repo.Media == nil { return } + if !mediaSearchWarmupEnabled(ctx, c.Repo) { + if c.Log != nil { + c.Log.Info("media search index warmup disabled") + } + return + } // 错峰:FTS 正常由 media 表触发器实时维护,回填只是升级或异常后的 // 兜底。先让登录、首页等关键路径跑起来,再开始后台补索引。 select { case <-ctx.Done(): return - case <-time.After(30 * time.Second): + case <-time.After(mediaSearchWarmupDelay(ctx, c.Repo)): } - const batchSize = 1000 - const pause = 100 * time.Millisecond + batchSize := mediaSearchWarmupBatchSize(ctx, c.Repo) + pause := mediaSearchWarmupPause(ctx, c.Repo) total := int64(0) for { select { @@ -304,6 +311,59 @@ func (c *Container) warmMediaSearchIndex(ctx context.Context) { } } +func mediaSearchWarmupEnabled(ctx context.Context, repo *repository.Container) bool { + if repo == nil || repo.Setting == nil { + return true + } + value, err := repo.Setting.Get(ctx, "search.index_warmup_enabled") + if err != nil || strings.TrimSpace(value) == "" { + return true + } + return parseBoolSetting(value, true) +} + +func mediaSearchWarmupDelay(ctx context.Context, repo *repository.Container) time.Duration { + seconds := mediaSearchWarmupIntSetting(ctx, repo, "search.index_warmup_delay_seconds", 120) + if seconds < 30 { + seconds = 30 + } + return time.Duration(seconds) * time.Second +} + +func mediaSearchWarmupBatchSize(ctx context.Context, repo *repository.Container) int { + size := mediaSearchWarmupIntSetting(ctx, repo, "search.index_warmup_batch_size", 100) + if size < 10 { + size = 10 + } + if size > 1000 { + size = 1000 + } + return size +} + +func mediaSearchWarmupPause(ctx context.Context, repo *repository.Container) time.Duration { + ms := mediaSearchWarmupIntSetting(ctx, repo, "search.index_warmup_pause_ms", 2000) + if ms < 250 { + ms = 250 + } + return time.Duration(ms) * time.Millisecond +} + +func mediaSearchWarmupIntSetting(ctx context.Context, repo *repository.Container, key string, fallback int) int { + if repo == nil || repo.Setting == nil { + return fallback + } + value, err := repo.Setting.Get(ctx, key) + if err != nil { + return fallback + } + n, err := strconv.Atoi(strings.TrimSpace(value)) + if err != nil || n <= 0 { + return fallback + } + return n +} + func (c *Container) NormalizeCloudLibraryTypes(ctx context.Context) error { if c == nil || c.Repo == nil || c.Repo.Library == nil || c.Repo.DB == nil { return nil diff --git a/internal/service/token_svc.go b/internal/service/token_svc.go index 760842e..fb2f029 100644 --- a/internal/service/token_svc.go +++ b/internal/service/token_svc.go @@ -107,25 +107,17 @@ func (s *TokenService) issuePair(ctx context.Context, userID, role, tier string, TokenHash: tokenHash, ExpiresAt: time.Now().Add(RefreshTokenDuration), } - storeCtx := ctx - cancel := func() {} if bestEffort { - storeCtx, cancel = context.WithTimeout(context.Background(), loginRefreshTokenStoreTimeout) + s.storeRefreshTokenBestEffort(userID, tokenHash, rt.ExpiresAt) + return &TokenPair{ + AccessToken: accessToken, + RefreshToken: refreshToken, + ExpiresIn: int64(AccessTokenDuration.Seconds()), + TokenType: "Bearer", + }, nil } - err = s.storeRefreshToken(storeCtx, rt) - cancel() - if err != nil { - if !bestEffort { - return nil, err - } - if s.log != nil { - s.log.Warn("refresh token store delayed; login will continue", - zap.String("user_id", userID), - zap.Error(err)) - } - if s.trackDelayedStore(userID, tokenHash, rt.ExpiresAt) { - go s.storeRefreshTokenEventually(userID, tokenHash, rt.ExpiresAt) - } + if err := s.storeRefreshToken(ctx, rt); err != nil { + return nil, err } return &TokenPair{ @@ -136,6 +128,56 @@ func (s *TokenService) issuePair(ctx context.Context, userID, role, tier string, }, nil } +func (s *TokenService) storeRefreshTokenBestEffort(userID, tokenHash string, expiresAt time.Time) { + if !s.trackDelayedStore(userID, tokenHash, expiresAt) { + return + } + done := make(chan error, 1) + go func() { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + done <- s.storeRefreshToken(ctx, &model.RefreshToken{ + UserID: userID, + TokenHash: tokenHash, + ExpiresAt: expiresAt, + }) + }() + select { + case err := <-done: + s.finishBestEffortRefreshTokenStore(userID, tokenHash, expiresAt, err) + case <-time.After(loginRefreshTokenStoreTimeout): + if s.log != nil { + s.log.Warn("refresh token store delayed; login will continue", + zap.String("user_id", userID), + zap.Error(context.DeadlineExceeded)) + } + go func() { + err := <-done + s.finishBestEffortRefreshTokenStore(userID, tokenHash, expiresAt, err) + }() + } +} + +func (s *TokenService) finishBestEffortRefreshTokenStore(userID, tokenHash string, expiresAt time.Time, err error) { + if err == nil { + s.untrackDelayedStore(userID, tokenHash) + return + } + if repository.IsSQLiteBusyError(err) || errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) { + if s.log != nil { + s.log.Warn("refresh token store delayed; login will continue", + zap.String("user_id", userID), + zap.Error(err)) + } + s.storeRefreshTokenEventually(userID, tokenHash, expiresAt) + return + } + s.untrackDelayedStore(userID, tokenHash) + if s.log != nil { + s.log.Warn("refresh token delayed store failed permanently", zap.String("user_id", userID), zap.Error(err)) + } +} + func (s *TokenService) storeRefreshToken(ctx context.Context, rt *model.RefreshToken) error { if err := s.repo.RefreshToken.Create(ctx, rt); err != nil { return err @@ -148,7 +190,7 @@ func (s *TokenService) storeRefreshToken(ctx context.Context, rt *model.RefreshT func (s *TokenService) storeRefreshTokenEventually(userID, tokenHash string, expiresAt time.Time) { defer s.untrackDelayedStore(userID, tokenHash) - delay := 5 * time.Second + delay := time.Second for attempt := 1; attempt <= 8; attempt++ { timer := time.NewTimer(delay) <-timer.C @@ -313,7 +355,7 @@ func (s *TokenService) Refresh(ctx context.Context, refreshToken string) (*Token s.untrackDelayedStore(rt.UserID, tokenHash) // 签发新的令牌对 - return s.IssuePair(ctx, user.ID, user.Role, user.Tier) + return s.IssuePairBestEffort(ctx, user.ID, user.Role, user.Tier) } // RevokeAll 撤销用户的所有 Refresh Token(用于登出)。