Files
sfdupes/db.go
T
clawbot 33cf3dd29a
check / check (push) Successful in 1m46s
Stream report and trees instead of loading every record (closes #14)
report now has SQLite group the records and put the rows in report
order, helped by a new files_signature index on (size, head, tail,
content), and writes each row as it reads it. trees reads the records
in path order, where all the paths under a directory come together, so
it computes each directory's digest as soon as the stream leaves it and
keeps only its path, parent, digest and totals. Output is unchanged.

The tests that called the removed in-memory grouping functions now group
records stored in a database. New tests check that both commands give
the same output whatever order the records were inserted in, and that a
stdout failure partway through a long report is reported as one.

Model: opus-5-5
2026-10-04 06:47:25 +02:00

600 lines
16 KiB
Go

package main
import (
"context"
"database/sql"
"errors"
"fmt"
"io/fs"
"os"
"path/filepath"
"slices"
"strconv"
"golang.org/x/sys/unix"
// The pure-Go SQLite driver, registered as "sqlite"; keeps cgo
// disabled.
_ "modernc.org/sqlite"
)
// defaultDatabasePath is where the persistent scan database lives when
// SFDUPES_DATABASE is not set.
const defaultDatabasePath = "/var/lib/sfdupes/db.sqlite"
// databaseEnv is the environment variable that overrides the database
// path.
const databaseEnv = "SFDUPES_DATABASE"
// schemaVersion is the database schema version this build reads and
// writes, stored in PRAGMA user_version.
const schemaVersion = 1
// dbDirPerm is the mode for a database parent directory created by
// scan.
const dbDirPerm = 0o755
// lockFilePerm is the mode for the scan lock file. Anyone who can open
// the file can hold the lock and keep every scan from running, so it
// is open to its owner only.
const lockFilePerm = 0o600
// createTableSQL is the schema applied to a fresh database. Paths are
// BLOBs because Unix paths are raw bytes, not guaranteed UTF-8.
const createTableSQL = `
CREATE TABLE files (
path BLOB PRIMARY KEY,
size INTEGER NOT NULL,
mtime INTEGER NOT NULL,
head TEXT NOT NULL,
tail TEXT NOT NULL,
content TEXT NOT NULL
) WITHOUT ROWID
`
// createIndexSQL indexes the records by signature, so report can have
// SQLite group them without sorting the whole table.
const createIndexSQL = `
CREATE INDEX files_signature ON files (size, head, tail, content)
`
// upsertSQL inserts one file record, replacing any existing record for
// the same path.
const upsertSQL = `
INSERT INTO files (path, size, mtime, head, tail, content)
VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT (path) DO UPDATE SET
size = excluded.size, mtime = excluded.mtime,
head = excluded.head, tail = excluded.tail,
content = excluded.content
`
// errNoDatabase reports a missing database file for report/trees.
var errNoDatabase = errors.New(
"no database (run \"sfdupes scan\" first, or set " + databaseEnv + ")")
// errSchemaVersion reports a database whose schema version this build
// does not understand.
var errSchemaVersion = errors.New("unsupported database schema version")
// errScanRunning reports that another scan holds the lock on the
// database.
var errScanRunning = errors.New("another scan is running")
// databasePath resolves the database location: SFDUPES_DATABASE when
// set and non-empty, the compiled-in default otherwise.
func databasePath() string {
if p := os.Getenv(databaseEnv); p != "" {
return p
}
return defaultDatabasePath
}
// scanParams are the connection parameters for scan: read-write, with
// WAL journaling and a busy timeout, so a report can run while a cron
// scan is in progress. closeScanDatabase leaves WAL mode again.
const scanParams = "_pragma=busy_timeout(10000)" +
"&_pragma=journal_mode(WAL)" +
"&_pragma=synchronous(NORMAL)"
// reportParams are the connection parameters for report and trees:
// read-only, with the same busy timeout. They set no journal mode,
// because setting one is a write.
const reportParams = "mode=ro" +
"&_pragma=busy_timeout(10000)" +
"&_pragma=query_only(1)"
// openDB opens the SQLite database at path with the connection
// parameters params. It does not create or verify the schema.
func openDB(path, params string) (*sql.DB, error) {
db, err := sql.Open("sqlite", "file:"+path+"?"+params)
if err != nil {
return nil, fmt.Errorf("open database %s: %w", path, err)
}
// A single connection avoids SQLITE_BUSY between this process's
// own connections; concurrency lives in the worker pools, not in
// parallel database access.
db.SetMaxOpenConns(1)
return db, nil
}
// lockScanDatabase takes the lock that keeps a second scan off the
// database at path: an exclusive flock(2) on the file beside it named
// path with ".lock" appended, created along with the database's parent
// directory if missing. A lock held by another scan fails at once
// instead of waiting. The lock lasts until the returned file is closed
// or the process ends. The file is never deleted: a scan that deleted
// it would let the next scan lock a new file while another still holds
// the old one.
func lockScanDatabase(path string) (*os.File, error) {
err := os.MkdirAll(filepath.Dir(path), dbDirPerm)
if err != nil {
return nil, fmt.Errorf("create database directory: %w", err)
}
lockPath := path + ".lock"
//nolint:gosec // the operator chooses the database path
f, err := os.OpenFile(lockPath, os.O_RDWR|os.O_CREATE, lockFilePerm)
if err != nil {
return nil, err
}
err = unix.Flock(int(f.Fd()), unix.LOCK_EX|unix.LOCK_NB)
if err != nil {
_ = f.Close()
if errors.Is(err, unix.EWOULDBLOCK) {
return nil, fmt.Errorf("%w (lock held on %s)",
errScanRunning, lockPath)
}
return nil, fmt.Errorf("lock %s: %w", lockPath, err)
}
return f, nil
}
// openScanDatabase opens the database for the scan subcommand, creating
// the file, its parent directory, and the schema as needed.
func openScanDatabase(ctx context.Context, path string) (*sql.DB, error) {
err := os.MkdirAll(filepath.Dir(path), dbDirPerm)
if err != nil {
return nil, fmt.Errorf("create database directory: %w", err)
}
db, err := openDB(path, scanParams)
if err != nil {
return nil, err
}
err = initSchema(ctx, db)
if err != nil {
_ = db.Close()
return nil, fmt.Errorf("database %s: %w", path, err)
}
return db, nil
}
// closeScanDatabase switches the database at path from WAL back to
// rollback-journal mode and closes it. Out of WAL mode the database
// file alone holds the whole database, so a reader needs no -wal or
// -shm file beside it, nor write access to create them. The switch
// fails while a report has the database open; the database then stays
// in WAL mode, still readable, until a later scan closes it.
func closeScanDatabase(ctx context.Context, db *sql.DB, path string) {
// Runs on the way out of a cancelled scan too.
_, err := db.ExecContext(context.WithoutCancel(ctx),
"PRAGMA journal_mode = DELETE")
if err != nil {
fmt.Fprintf(os.Stderr, "scan: database %s left in WAL mode: %v\n",
path, err)
}
_ = db.Close()
}
// openReportDatabase opens an existing database for the report and
// trees subcommands. A missing database file is an error directing the
// user to run scan first; the schema version must match exactly.
func openReportDatabase(ctx context.Context,
path string,
) (*sql.DB, error) {
_, err := os.Stat(path)
if errors.Is(err, fs.ErrNotExist) {
return nil, fmt.Errorf("%s: %w", path, errNoDatabase)
}
if err != nil {
return nil, fmt.Errorf("database: %w", err)
}
db, err := openDB(path, reportParams)
if err != nil {
return nil, err
}
v, err := userVersion(ctx, db)
if err != nil {
_ = db.Close()
return nil, fmt.Errorf("database %s: %w", path, err)
}
if v != schemaVersion {
_ = db.Close()
return nil, fmt.Errorf("database %s: version %d, want %d: %w",
path, v, schemaVersion, errSchemaVersion)
}
return db, nil
}
// initSchema creates the schema on a fresh database and verifies the
// schema version on an existing one.
func initSchema(ctx context.Context, db *sql.DB) error {
v, err := userVersion(ctx, db)
if err != nil {
return err
}
switch v {
case 0:
return createSchema(ctx, db)
case schemaVersion:
return nil
default:
return fmt.Errorf("version %d, want %d: %w",
v, schemaVersion, errSchemaVersion)
}
}
// createSchema applies the schema to a fresh database and stamps the
// schema version.
func createSchema(ctx context.Context, db *sql.DB) error {
_, err := db.ExecContext(ctx, createTableSQL)
if err != nil {
return fmt.Errorf("create schema: %w", err)
}
_, err = db.ExecContext(ctx, createIndexSQL)
if err != nil {
return fmt.Errorf("create schema: %w", err)
}
_, err = db.ExecContext(ctx,
"PRAGMA user_version = "+strconv.Itoa(schemaVersion))
if err != nil {
return fmt.Errorf("set schema version: %w", err)
}
return nil
}
// userVersion reads the database's PRAGMA user_version.
func userVersion(ctx context.Context, db *sql.DB) (int, error) {
var v int
err := db.QueryRowContext(ctx, "PRAGMA user_version").Scan(&v)
if err != nil {
return 0, fmt.Errorf("read schema version: %w", err)
}
return v, nil
}
// loadFileRows streams every record to fn in path order: byte order,
// which is the order of the primary key, so SQLite does not sort.
func loadFileRows(ctx context.Context, db *sql.DB, fn func(r scanRec)) error {
rows, err := db.QueryContext(ctx,
"SELECT path, size, mtime, head, tail, content FROM files "+
"ORDER BY path")
if err != nil {
return fmt.Errorf("read records: %w", err)
}
defer func() { _ = rows.Close() }()
for rows.Next() {
var (
path []byte
r scanRec
)
err = rows.Scan(&path, &r.size, &r.mtime, &r.head, &r.tail,
&r.content)
if err != nil {
return fmt.Errorf("read record: %w", err)
}
r.path = string(path)
fn(r)
}
err = rows.Err()
if err != nil {
return fmt.Errorf("read records: %w", err)
}
return nil
}
// dupeRowsSQL selects every record in a duplicate group, with the
// group's first path. A group is the records with a content hash that
// share a size, head, tail, and content, when there are two or more of
// them. The rows come in report order: groups by size descending, then
// by first path, and each group's paths ascending.
const dupeRowsSQL = `
SELECT g.first, f.path, f.size
FROM files AS f
JOIN (
SELECT size, head, tail, content, MIN(path) AS first
FROM files
WHERE content <> ''
GROUP BY size, head, tail, content
HAVING COUNT(*) > 1
) AS g USING (size, head, tail, content)
ORDER BY f.size DESC, g.first, f.path
`
// loadDupeRows streams the rows of dupeRowsSQL to fn and returns the
// number of records in the database. The count and the rows are read
// in one transaction, so they agree while a scan is committing. An
// error from fn stops the reading and is returned as it is.
func loadDupeRows(ctx context.Context, db *sql.DB,
fn func(first, path string, size int64) error,
) (int, error) {
// Everything goes through tx: the report connection is the only
// one, so a query on db would wait for tx forever.
tx, err := db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true})
if err != nil {
return 0, fmt.Errorf("read records: %w", err)
}
defer func() { _ = tx.Rollback() }()
var records int
err = tx.QueryRowContext(ctx, "SELECT COUNT(*) FROM files").Scan(&records)
if err != nil {
return 0, fmt.Errorf("read records: %w", err)
}
rows, err := tx.QueryContext(ctx, dupeRowsSQL)
if err != nil {
return 0, fmt.Errorf("read records: %w", err)
}
defer func() { _ = rows.Close() }()
for rows.Next() {
var (
first, path []byte
size int64
)
err = rows.Scan(&first, &path, &size)
if err != nil {
return 0, fmt.Errorf("read record: %w", err)
}
err = fn(string(first), string(path), size)
if err != nil {
return 0, err
}
}
err = rows.Err()
if err != nil {
return 0, fmt.Errorf("read records: %w", err)
}
return records, nil
}
// loadFileMeta streams every record's path, size, mtime, and whether
// it carries hashes to fn. Scan change detection needs no hash
// values, and skipping the hash columns keeps the scan's in-memory
// index small on multi-million-file databases.
func loadFileMeta(ctx context.Context, db *sql.DB,
fn func(path string, size, mtime int64, hashed bool),
) error {
rows, err := db.QueryContext(ctx,
"SELECT path, size, mtime, head <> '' FROM files")
if err != nil {
return fmt.Errorf("read records: %w", err)
}
defer func() { _ = rows.Close() }()
for rows.Next() {
var (
path []byte
size, mtime int64
hashed int64
)
err = rows.Scan(&path, &size, &mtime, &hashed)
if err != nil {
return fmt.Errorf("read record: %w", err)
}
fn(string(path), size, mtime, hashed != 0)
}
err = rows.Err()
if err != nil {
return fmt.Errorf("read records: %w", err)
}
return nil
}
// contentCandidatesSQL selects every record of at least headTailMin
// bytes whose size, head, and tail equal another record's, in each
// group (the records sharing a size, head, and tail) where at least one
// record has no content hash, with whether each record has one. SQLite
// does the grouping, so no other record's hashes are loaded into
// memory; the rows come ordered by size, head, and tail, so each
// group's rows arrive together.
const contentCandidatesSQL = `
SELECT f.path, f.size, f.mtime, f.head, f.tail, f.content <> ''
FROM files AS f
JOIN (
SELECT size, head, tail
FROM files
WHERE size >= ? AND head <> ''
GROUP BY size, head, tail
HAVING COUNT(*) > 1 AND SUM(content = '') > 0
) AS g USING (size, head, tail)
ORDER BY size, head, tail
`
// loadContentCandidates streams the rows of contentCandidatesSQL to fn:
// each record, without its content hash, and whether it has one.
func loadContentCandidates(ctx context.Context, db *sql.DB,
fn func(r scanRec, hashed bool),
) error {
rows, err := db.QueryContext(ctx, contentCandidatesSQL, headTailMin)
if err != nil {
return fmt.Errorf("read records: %w", err)
}
defer func() { _ = rows.Close() }()
for rows.Next() {
var (
path []byte
r scanRec
hashed int64
)
err = rows.Scan(&path, &r.size, &r.mtime, &r.head, &r.tail, &hashed)
if err != nil {
return fmt.Errorf("read record: %w", err)
}
r.path = string(path)
fn(r, hashed != 0)
}
err = rows.Err()
if err != nil {
return fmt.Errorf("read records: %w", err)
}
return nil
}
// updateBatchSize is the number of record changes committed per
// transaction during the update pass. The filesystem is authoritative
// and the database an eventually-consistent reflection of it, so
// scan-level atomicity is not required; smaller transactions keep the
// WAL small and let concurrent reports observe progress.
const updateBatchSize = 10000
// applyChanges writes one scan's database changes — upserts for new and
// changed files, deletes for vanished ones — in batched transactions.
// Progress is rendered on prog (one increment per change).
func applyChanges(ctx context.Context, db *sql.DB, upserts []scanRec,
deletes []string, prog *progress,
) error {
for batch := range slices.Chunk(upserts, updateBatchSize) {
err := applyBatch(ctx, db, batch, nil, prog)
if err != nil {
return err
}
}
for batch := range slices.Chunk(deletes, updateBatchSize) {
err := applyBatch(ctx, db, nil, batch, prog)
if err != nil {
return err
}
}
return nil
}
// applyBatch commits one batch of upserts and deletes in a single
// transaction.
func applyBatch(ctx context.Context, db *sql.DB, upserts []scanRec,
deletes []string, prog *progress,
) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin transaction: %w", err)
}
defer func() { _ = tx.Rollback() }()
err = execUpserts(ctx, tx, upserts, prog)
if err != nil {
return err
}
err = execDeletes(ctx, tx, deletes, prog)
if err != nil {
return err
}
err = tx.Commit()
if err != nil {
return fmt.Errorf("commit: %w", err)
}
return nil
}
// execUpserts inserts or updates one record per new or changed file.
func execUpserts(ctx context.Context, tx *sql.Tx, upserts []scanRec,
prog *progress,
) error {
st, err := tx.PrepareContext(ctx, upsertSQL)
if err != nil {
return fmt.Errorf("prepare upsert: %w", err)
}
defer func() { _ = st.Close() }()
for _, r := range upserts {
_, err = st.ExecContext(ctx,
[]byte(r.path), r.size, r.mtime, r.head, r.tail, r.content)
if err != nil {
return fmt.Errorf("upsert %s: %w", r.path, err)
}
prog.increment()
}
return nil
}
// execDeletes removes the records for paths no longer present.
func execDeletes(ctx context.Context, tx *sql.Tx, deletes []string,
prog *progress,
) error {
st, err := tx.PrepareContext(ctx, "DELETE FROM files WHERE path = ?")
if err != nil {
return fmt.Errorf("prepare delete: %w", err)
}
defer func() { _ = st.Close() }()
for _, p := range deletes {
_, err = st.ExecContext(ctx, []byte(p))
if err != nil {
return fmt.Errorf("delete %s: %w", p, err)
}
prog.increment()
}
return nil
}