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 }