- add/update consolidated into ingest(isNew) — persist branch differs, error handling and cover path fully shared (Task 24) - zipIndex/cover use bookfile.OpenReaderAt instead of hand-rolled open+stat pairs - root containment check uses bookfile.Contains (local inside removed) - cover write goes through media.WriteAtomic (now exported, log-free — callers own context); media no longer logs inside the atomic helper
254 lines
6.7 KiB
Go
254 lines
6.7 KiB
Go
package scanner
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"io/fs"
|
|
"log"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"booklib/internal/bookfile"
|
|
"booklib/internal/config"
|
|
"booklib/internal/media"
|
|
"booklib/internal/redispkg"
|
|
"booklib/internal/store"
|
|
)
|
|
|
|
// Sweeper 是 scanner 每轮顺手调用的清理钩子;upload.U 满足它(B16)。
|
|
type Sweeper interface {
|
|
Sweep(ctx context.Context) error
|
|
}
|
|
|
|
type Scanner struct {
|
|
st *store.Store
|
|
cfg *config.Config
|
|
rdb *redispkg.R
|
|
sweepers []Sweeper
|
|
// B9-②: per-library single-flight — concurrent scan triggers for the same
|
|
// library are merged into one execution, even without redis.
|
|
flights sync.Map // map[int64]*sync.WaitGroup
|
|
}
|
|
|
|
func New(st *store.Store, cfg *config.Config, rdb *redispkg.R, sweepers ...Sweeper) *Scanner {
|
|
return &Scanner{st: st, cfg: cfg, rdb: rdb, sweepers: sweepers}
|
|
}
|
|
|
|
func (s *Scanner) Run(ctx context.Context) {
|
|
t := time.NewTicker(s.cfg.ScanInterval)
|
|
defer t.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-t.C:
|
|
for _, sw := range s.sweepers { // B16: 上传会话清扫随扫描周期跑
|
|
if err := sw.Sweep(ctx); err != nil {
|
|
log.Printf("scan: sweep: %v", err)
|
|
}
|
|
}
|
|
libs, err := s.st.ListLibraries(ctx)
|
|
if err != nil {
|
|
log.Printf("scan: list libraries: %v", err)
|
|
continue
|
|
}
|
|
for _, l := range libs {
|
|
s.scanOnce(ctx, l)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Scanner) ScanLibraryByID(ctx context.Context, id int64) {
|
|
lib, err := s.st.GetLibrary(ctx, id)
|
|
if err != nil {
|
|
log.Printf("scan: library %d: %v", id, err)
|
|
return
|
|
}
|
|
s.scanOnce(ctx, lib)
|
|
}
|
|
|
|
// scanOnce ensures only one scan per library runs concurrently in this process.
|
|
// Concurrent callers block until the in-flight scan completes (B9-②).
|
|
func (s *Scanner) scanOnce(ctx context.Context, lib store.Library) {
|
|
wg := &sync.WaitGroup{}
|
|
wg.Add(1)
|
|
if existing, loaded := s.flights.LoadOrStore(lib.ID, wg); loaded {
|
|
existing.(*sync.WaitGroup).Wait()
|
|
return
|
|
}
|
|
defer func() {
|
|
s.flights.Delete(lib.ID)
|
|
wg.Done()
|
|
}()
|
|
s.ScanLibrary(ctx, lib)
|
|
}
|
|
|
|
func (s *Scanner) ScanLibrary(ctx context.Context, lib store.Library) {
|
|
// B9-①: ScanLock auto-renews every TTL/2 during long scans.
|
|
unlock, ok := s.rdb.ScanLock(ctx, fmt.Sprintf("scan:%d", lib.ID), 5*time.Minute)
|
|
if !ok {
|
|
return // 别的副本在扫
|
|
}
|
|
defer unlock()
|
|
|
|
root, err := filepath.EvalSymlinks(filepath.Clean(lib.RootPath))
|
|
if err != nil || !bookfile.Contains(s.cfg.BooksDir, root) {
|
|
log.Printf("scan: library %d root %q rejected", lib.ID, lib.RootPath)
|
|
return
|
|
}
|
|
disk, err := walk(root)
|
|
if err != nil {
|
|
log.Printf("scan: walk %s: %v", root, err)
|
|
return
|
|
}
|
|
dbMeta, err := s.st.ListBookMeta(ctx, lib.ID)
|
|
if err != nil {
|
|
log.Printf("scan: list books: %v", err)
|
|
return
|
|
}
|
|
for rel, ds := range disk {
|
|
old, exists := dbMeta[rel]
|
|
delete(dbMeta, rel)
|
|
switch {
|
|
case !exists:
|
|
s.ingest(ctx, lib.ID, 0, root, rel, ds, true)
|
|
case old.Size != ds.size || old.ModTS != ds.modTS:
|
|
s.ingest(ctx, lib.ID, old.ID, root, rel, ds, false)
|
|
}
|
|
}
|
|
for rel := range dbMeta { // 只剩被删的文件
|
|
if err := s.st.DeleteBookByPath(ctx, lib.ID, rel); err != nil {
|
|
log.Printf("scan: delete %s: %v", rel, err)
|
|
}
|
|
}
|
|
s.sweepCache(ctx)
|
|
}
|
|
|
|
type diskStat struct{ size, modTS int64 }
|
|
|
|
func walk(root string) (map[string]diskStat, error) {
|
|
out := map[string]diskStat{}
|
|
err := filepath.WalkDir(root, func(p string, d fs.DirEntry, err error) error {
|
|
if err != nil {
|
|
log.Printf("scan: walk %s: %v", p, err)
|
|
return nil // 单点失败不中断
|
|
}
|
|
if d.IsDir() {
|
|
return nil
|
|
}
|
|
if bookfile.FormatFromExt(d.Name()) == "" {
|
|
return nil
|
|
}
|
|
info, err := d.Info()
|
|
if err != nil {
|
|
return nil
|
|
}
|
|
rel, err := filepath.Rel(root, p)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
out[filepath.ToSlash(rel)] = diskStat{info.Size(), info.ModTime().Unix()}
|
|
return nil
|
|
})
|
|
return out, err
|
|
}
|
|
|
|
func titleOf(rel string) string {
|
|
base := filepath.Base(rel)
|
|
return strings.TrimSpace(strings.ReplaceAll(strings.TrimSuffix(base, filepath.Ext(base)), "_", " "))
|
|
}
|
|
|
|
// ingest 是 add/update 的合一实现(Task 24):isNew 决定走 Insert 还是 UpdateBookFile,
|
|
// 之后的错误处理与封面生成完全共享。bookID 仅在 isNew=false 时有意义。
|
|
func (s *Scanner) ingest(ctx context.Context, libID, bookID int64, root, rel string, ds diskStat, isNew bool) {
|
|
format := bookfile.FormatFromExt(filepath.Base(rel))
|
|
pageCount := 0
|
|
var idxErr error
|
|
if format == "cbz" {
|
|
idx, err := s.zipIndex(root, rel)
|
|
pageCount = len(idx)
|
|
idxErr = err
|
|
if idxErr == nil && pageCount == 0 {
|
|
idxErr = errors.New("no images in archive")
|
|
}
|
|
}
|
|
if isNew {
|
|
id, err := s.st.InsertBook(ctx, libID, rel, titleOf(rel), format, ds.size, ds.modTS, pageCount)
|
|
if err != nil {
|
|
log.Printf("scan: insert %s: %v", rel, err)
|
|
return
|
|
}
|
|
bookID = id
|
|
} else if err := s.st.UpdateBookFile(ctx, bookID, ds.size, ds.modTS, pageCount); err != nil {
|
|
log.Printf("scan: update %s: %v", rel, err)
|
|
return
|
|
}
|
|
if idxErr != nil {
|
|
// B10: log SetBookState errors instead of discarding.
|
|
if e := s.st.SetBookState(ctx, bookID, "error", idxErr.Error()); e != nil {
|
|
log.Printf("scan: SetBookState %s: %v", rel, e)
|
|
}
|
|
return
|
|
}
|
|
s.cover(ctx, bookID, root, rel, format, ds)
|
|
}
|
|
|
|
func (s *Scanner) zipIndex(root, rel string) ([]string, error) {
|
|
f, size, err := bookfile.OpenReaderAt(root, rel)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer f.Close()
|
|
return bookfile.PageIndex(f, size)
|
|
}
|
|
|
|
// cover writes the cover image to the cache dir via media.WriteAtomic.
|
|
// B11: write failures are logged; orphan .tmp files are cleaned only on failure.
|
|
func (s *Scanner) cover(ctx context.Context, id int64, root, rel, format string, ds diskStat) {
|
|
var fn func(io.ReaderAt, int64) ([]byte, string, error)
|
|
switch format {
|
|
case "cbz":
|
|
fn = bookfile.CBZCover
|
|
case "epub":
|
|
fn = bookfile.EPUBCover
|
|
default:
|
|
return // pdf/txt/md 用占位 SVG,不落盘
|
|
}
|
|
f, size, err := bookfile.OpenReaderAt(root, rel)
|
|
if err != nil {
|
|
log.Printf("scan: cover %s: %v", rel, err)
|
|
return
|
|
}
|
|
defer f.Close()
|
|
img, ext, err := fn(f, size)
|
|
if err != nil {
|
|
log.Printf("scan: cover %s: %v", rel, err)
|
|
return
|
|
}
|
|
dir := bookfile.CoverDir(s.cfg.CacheDir, bookfile.DirKey(id, bookfile.Hash(ds.size, ds.modTS)))
|
|
if e := media.WriteAtomic(dir, "cover"+ext, img); e != nil {
|
|
log.Printf("scan: cover write %s: %v", rel, e)
|
|
}
|
|
}
|
|
|
|
func (s *Scanner) sweepCache(ctx context.Context) {
|
|
hashes, err := s.st.BookHashes(ctx)
|
|
if err != nil {
|
|
return
|
|
}
|
|
live := map[string]bool{}
|
|
for id, v := range hashes {
|
|
live[bookfile.DirKey(id, bookfile.Hash(v[0], v[1]))] = true
|
|
}
|
|
if n, err := bookfile.SweepStale(s.cfg.CacheDir, live); err != nil {
|
|
log.Printf("scan: sweep: %v", err)
|
|
} else if n > 0 {
|
|
log.Printf("scan: swept %d stale cache dirs", n)
|
|
}
|
|
}
|