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)