feat(backend): directory scanner — diff, error state, covers, cache sweep, redis lock
This commit is contained in:
@@ -9,6 +9,7 @@ import (
|
||||
"booklib/internal/auth"
|
||||
"booklib/internal/config"
|
||||
"booklib/internal/redispkg"
|
||||
"booklib/internal/scanner"
|
||||
"booklib/internal/store"
|
||||
)
|
||||
|
||||
@@ -16,6 +17,7 @@ type api struct {
|
||||
cfg *config.Config
|
||||
st *store.Store
|
||||
rdb *redispkg.R
|
||||
sc *scanner.Scanner
|
||||
}
|
||||
|
||||
func err(c *gin.Context, status int, code, msg string) {
|
||||
|
||||
@@ -15,6 +15,7 @@ import (
|
||||
"booklib/internal/auth"
|
||||
"booklib/internal/db"
|
||||
"booklib/internal/redispkg"
|
||||
"booklib/internal/scanner"
|
||||
"booklib/internal/store"
|
||||
)
|
||||
|
||||
@@ -43,7 +44,8 @@ func setupAPI(t *testing.T) (*store.Store, http.Handler) {
|
||||
cfg := testCfg()
|
||||
cfg.BooksDir = filepath.Clean(os.TempDir())
|
||||
cfg.CacheDir = t.TempDir()
|
||||
r := NewRouter(cfg, st, redispkg.New(os.Getenv("REDIS_URL")))
|
||||
rdb := redispkg.New(os.Getenv("REDIS_URL"))
|
||||
r := NewRouter(cfg, st, rdb, scanner.New(st, cfg, rdb))
|
||||
if u := os.Getenv("REDIS_URL"); u != "" { // 测试卫生: 共享 redis 上重置登录限流桶, 防跨测试累计 429
|
||||
if opt, e := redis.ParseURL(u); e == nil {
|
||||
rc := redis.NewClient(opt)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
@@ -82,11 +83,15 @@ func (a *api) getLibrary(c *gin.Context) (store.Library, bool) {
|
||||
}
|
||||
|
||||
func (a *api) scanLibrary(c *gin.Context) {
|
||||
if _, ok := a.getLibrary(c); !ok {
|
||||
lib, ok := a.getLibrary(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
// Task 9 接线: go a.sc.ScanLibraryByID(...)
|
||||
err(c, http.StatusNotImplemented, "not_ready", "scanner not wired yet")
|
||||
if _, ok := a.libRoot(c, lib); !ok {
|
||||
return
|
||||
}
|
||||
go a.sc.ScanLibraryByID(context.WithoutCancel(c), lib.ID)
|
||||
c.JSON(http.StatusAccepted, gin.H{"accepted": true})
|
||||
}
|
||||
|
||||
func (a *api) upload(c *gin.Context) {
|
||||
|
||||
@@ -7,12 +7,13 @@ import (
|
||||
|
||||
"booklib/internal/config"
|
||||
"booklib/internal/redispkg"
|
||||
"booklib/internal/scanner"
|
||||
"booklib/internal/store"
|
||||
)
|
||||
|
||||
func NewRouter(cfg *config.Config, st *store.Store, rdb *redispkg.R) *gin.Engine {
|
||||
func NewRouter(cfg *config.Config, st *store.Store, rdb *redispkg.R, sc *scanner.Scanner) *gin.Engine {
|
||||
gin.SetMode(gin.ReleaseMode)
|
||||
a := &api{cfg: cfg, st: st, rdb: rdb}
|
||||
a := &api{cfg: cfg, st: st, rdb: rdb, sc: sc}
|
||||
r := gin.New()
|
||||
r.Use(gin.Recovery())
|
||||
g := r.Group("/api")
|
||||
|
||||
@@ -15,7 +15,7 @@ func testCfg() *config.Config {
|
||||
}
|
||||
|
||||
func TestHealthz(t *testing.T) {
|
||||
r := NewRouter(testCfg(), nil, redispkg.New(""))
|
||||
r := NewRouter(testCfg(), nil, redispkg.New(""), nil)
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/healthz", nil)
|
||||
w := httptest.NewRecorder()
|
||||
r.ServeHTTP(w, req)
|
||||
|
||||
@@ -0,0 +1,251 @@
|
||||
package scanner
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/fs"
|
||||
"log"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"booklib/internal/bookfile"
|
||||
"booklib/internal/config"
|
||||
"booklib/internal/redispkg"
|
||||
"booklib/internal/store"
|
||||
)
|
||||
|
||||
type Scanner struct {
|
||||
st *store.Store
|
||||
cfg *config.Config
|
||||
rdb *redispkg.R
|
||||
}
|
||||
|
||||
func New(st *store.Store, cfg *config.Config, rdb *redispkg.R) *Scanner {
|
||||
return &Scanner{st: st, cfg: cfg, rdb: rdb}
|
||||
}
|
||||
|
||||
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:
|
||||
libs, err := s.st.ListLibraries(ctx)
|
||||
if err != nil {
|
||||
log.Printf("scan: list libraries: %v", err)
|
||||
continue
|
||||
}
|
||||
for _, l := range libs {
|
||||
s.ScanLibrary(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.ScanLibrary(ctx, lib)
|
||||
}
|
||||
|
||||
func (s *Scanner) ScanLibrary(ctx context.Context, lib store.Library) {
|
||||
// ponytail: 5min lock TTL; a scan longer than this lets another replica join — refresh mid-walk if libs ever outgrow it
|
||||
unlock, ok := s.rdb.Lock(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 || !inside(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.add(ctx, lib.ID, root, rel, ds)
|
||||
case old.Size != ds.size || old.ModTS != ds.modTS:
|
||||
s.update(ctx, lib.ID, old.ID, root, rel, ds)
|
||||
}
|
||||
}
|
||||
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 inside(booksDir, root string) bool {
|
||||
b := filepath.Clean(booksDir)
|
||||
return root == b || strings.HasPrefix(root, b+string(os.PathSeparator))
|
||||
}
|
||||
|
||||
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)), "_", " "))
|
||||
}
|
||||
|
||||
// cbz 完整性判定集中在 add/update:PageIndex 失败 → state=error。
|
||||
// InsertBook/UpdateBookFile 的 SQL 已把 state 重置为 ready(Task 2),无需显式清 error。
|
||||
func (s *Scanner) add(ctx context.Context, libID int64, root, rel string, ds diskStat) {
|
||||
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
|
||||
}
|
||||
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
|
||||
}
|
||||
if idxErr != nil {
|
||||
s.st.SetBookState(ctx, id, "error", idxErr.Error())
|
||||
return
|
||||
}
|
||||
s.cover(ctx, id, root, rel, format, ds)
|
||||
}
|
||||
|
||||
func (s *Scanner) update(ctx context.Context, libID, bookID int64, root, rel string, ds diskStat) {
|
||||
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 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 {
|
||||
s.st.SetBookState(ctx, bookID, "error", idxErr.Error())
|
||||
return
|
||||
}
|
||||
s.cover(ctx, bookID, root, rel, format, ds)
|
||||
}
|
||||
|
||||
func (s *Scanner) zipIndex(root, rel string) ([]string, error) {
|
||||
f, err := os.Open(filepath.Join(root, filepath.FromSlash(rel)))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer f.Close()
|
||||
st, err := f.Stat()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return bookfile.PageIndex(f, st.Size())
|
||||
}
|
||||
|
||||
// cover 失败(坏 epub、无图等)只 log — 书的 state 由 PageIndex 判定,封面缺了有占位 SVG 兜底
|
||||
func (s *Scanner) cover(ctx context.Context, id int64, root, rel, format string, ds diskStat) {
|
||||
var img []byte
|
||||
var ext string
|
||||
var err error
|
||||
switch format {
|
||||
case "cbz":
|
||||
img, ext, err = s.readCover(root, rel, bookfile.CBZCover)
|
||||
case "epub":
|
||||
img, ext, err = s.readCover(root, rel, bookfile.EPUBCover)
|
||||
default:
|
||||
return // pdf/txt/md 用占位 SVG,不落盘
|
||||
}
|
||||
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 := os.MkdirAll(dir, 0o755); e != nil {
|
||||
log.Printf("scan: coverdir %s: %v", rel, e)
|
||||
return
|
||||
}
|
||||
tmp := filepath.Join(dir, "cover"+ext+".tmp")
|
||||
dst := filepath.Join(dir, "cover"+ext)
|
||||
if e := os.WriteFile(tmp, img, 0o644); e == nil {
|
||||
os.Rename(tmp, dst)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Scanner) readCover(root, rel string, fn func(io.ReaderAt, int64) ([]byte, string, error)) ([]byte, string, error) {
|
||||
f, err := os.Open(filepath.Join(root, filepath.FromSlash(rel)))
|
||||
if err != nil {
|
||||
return nil, "", err
|
||||
}
|
||||
defer f.Close()
|
||||
st, err := f.Stat()
|
||||
if err != nil {
|
||||
return nil, "", err
|
||||
}
|
||||
return fn(f, st.Size())
|
||||
}
|
||||
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,142 @@
|
||||
package scanner
|
||||
|
||||
import (
|
||||
"archive/zip"
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"booklib/internal/bookfile"
|
||||
"booklib/internal/config"
|
||||
"booklib/internal/db"
|
||||
"booklib/internal/redispkg"
|
||||
"booklib/internal/store"
|
||||
)
|
||||
|
||||
// 返回 scanner、library、解析过符号链接的 root、以及一个满足 progress FK 的 uid
|
||||
func setupLib(t *testing.T) (*Scanner, store.Library, string, int64) {
|
||||
t.Helper()
|
||||
url := os.Getenv("DATABASE_URL")
|
||||
if url == "" {
|
||||
t.Skip("DATABASE_URL not set")
|
||||
}
|
||||
ctx := context.Background()
|
||||
p, _ := db.Connect(ctx, url)
|
||||
db.Migrate(ctx, p)
|
||||
st := store.New(p)
|
||||
p.Exec(ctx, "DELETE FROM reading_progress; DELETE FROM books; DELETE FROM libraries; DELETE FROM users")
|
||||
books := t.TempDir()
|
||||
cache := t.TempDir()
|
||||
resolved := mustResolve(t, books)
|
||||
libID, err := st.CreateLibrary(ctx, "t", filepath.Join(resolved, "lib"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
uid, err := st.CreateUser(ctx, "scantest", "h", "member")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
root := filepath.Join(resolved, "lib")
|
||||
os.MkdirAll(filepath.Join(root, "series-a"), 0o755)
|
||||
cfg := &config.Config{BooksDir: resolved, CacheDir: cache, ScanInterval: time.Minute}
|
||||
lib, _ := st.GetLibrary(ctx, libID)
|
||||
return New(st, cfg, redispkg.New(os.Getenv("REDIS_URL"))), lib, root, uid
|
||||
}
|
||||
|
||||
func mustResolve(t *testing.T, p string) string {
|
||||
r, err := filepath.EvalSymlinks(p) // macOS 上 t.TempDir 是 /var→/private 符号链接
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return r
|
||||
}
|
||||
|
||||
func writeCBZ(t *testing.T, path string, pages int) {
|
||||
t.Helper()
|
||||
os.MkdirAll(filepath.Dir(path), 0o755)
|
||||
buf := &bytes.Buffer{}
|
||||
zw := zip.NewWriter(buf)
|
||||
for i := 1; i <= pages; i++ {
|
||||
w, _ := zw.Create(fmt.Sprintf("%02d.jpg", i))
|
||||
w.Write(bytes.Repeat([]byte("JPG"), 100))
|
||||
}
|
||||
zw.Close()
|
||||
os.WriteFile(path, buf.Bytes(), 0o644)
|
||||
}
|
||||
|
||||
func TestScanFullLifecycle(t *testing.T) {
|
||||
sc, lib, root, uid := setupLib(t)
|
||||
ctx := context.Background()
|
||||
writeCBZ(t, filepath.Join(root, "series-a", "vol_01.cbz"), 3)
|
||||
os.WriteFile(filepath.Join(root, "notes.txt"), []byte("hi"), 0o644)
|
||||
sc.ScanLibrary(ctx, lib)
|
||||
|
||||
meta, _ := sc.st.ListBookMeta(ctx, lib.ID)
|
||||
if len(meta) != 2 {
|
||||
t.Fatalf("want 2 books got %v", meta)
|
||||
}
|
||||
b, err := sc.st.GetBook(ctx, meta["series-a/vol_01.cbz"].ID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if b.Format != "cbz" || b.PageCount != 3 || b.Title != "vol 01" {
|
||||
t.Fatalf("bad book %+v", b)
|
||||
}
|
||||
key := bookfile.DirKey(b.ID, bookfile.Hash(b.FileSize, b.ModTS))
|
||||
covers, _ := os.ReadDir(bookfile.CoverDir(sc.cfg.CacheDir, key))
|
||||
if len(covers) == 0 {
|
||||
t.Fatal("cover not built")
|
||||
}
|
||||
// 二次扫描:无变化 → 不动
|
||||
before := b.AddedAt
|
||||
sc.ScanLibrary(ctx, lib)
|
||||
b2, _ := sc.st.GetBook(ctx, b.ID)
|
||||
if !b2.AddedAt.Equal(before) || b2.FileSize != b.FileSize {
|
||||
t.Fatal("unchanged file must be untouched")
|
||||
}
|
||||
// 修改:page_count 变、hash 变、旧缓存被扫尾清掉
|
||||
writeCBZ(t, filepath.Join(root, "series-a", "vol_01.cbz"), 5)
|
||||
os.Chtimes(filepath.Join(root, "series-a", "vol_01.cbz"), time.Now(), time.Now())
|
||||
sc.ScanLibrary(ctx, lib)
|
||||
b3, _ := sc.st.GetBook(ctx, b.ID)
|
||||
if b3.PageCount != 5 {
|
||||
t.Fatalf("not updated: %+v", b3)
|
||||
}
|
||||
if _, err := os.Stat(bookfile.CoverDir(sc.cfg.CacheDir, key)); !os.IsNotExist(err) {
|
||||
t.Fatal("stale cover cache remains")
|
||||
}
|
||||
// 删除:行没了,进度按 path 还在(spec §4)
|
||||
if err := sc.st.UpsertProgress(ctx, uid, lib.ID, "notes.txt", []byte("{}"), 0.5); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.Remove(filepath.Join(root, "notes.txt")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sc.ScanLibrary(ctx, lib)
|
||||
if _, err := sc.st.GetBook(ctx, b.ID); err != nil {
|
||||
t.Fatal("cbz vanished wrongly")
|
||||
}
|
||||
meta2, _ := sc.st.ListBookMeta(ctx, lib.ID)
|
||||
if _, ok := meta2["notes.txt"]; ok {
|
||||
t.Fatal("deleted file still in db")
|
||||
}
|
||||
if pr, err := sc.st.GetProgress(ctx, uid, lib.ID, "notes.txt"); err != nil || pr.Percent != 0.5 {
|
||||
t.Fatalf("progress must survive book removal: %+v %v", pr, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestScanBrokenCBZStateError(t *testing.T) {
|
||||
sc, lib, root, _ := setupLib(t)
|
||||
ctx := context.Background()
|
||||
os.WriteFile(filepath.Join(root, "bad.cbz"), []byte("not a zip"), 0o644)
|
||||
sc.ScanLibrary(ctx, lib)
|
||||
meta, _ := sc.st.ListBookMeta(ctx, lib.ID)
|
||||
b, _ := sc.st.GetBook(ctx, meta["bad.cbz"].ID)
|
||||
if b.State != "error" || b.ErrMsg == "" {
|
||||
t.Fatalf("want error state, got %+v", b)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user