Stream report and trees instead of loading every record (closes #14)
check / check (push) Successful in 1m49s
check / check (push) Successful in 1m49s
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. A new test checks that both commands give the same output whatever order the records were inserted in. Model: opus-5-5
This commit is contained in:
@@ -51,6 +51,12 @@ CREATE TABLE files (
|
||||
) 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 = `
|
||||
@@ -256,6 +262,11 @@ func createSchema(ctx context.Context, db *sql.DB) error {
|
||||
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 {
|
||||
@@ -277,18 +288,18 @@ func userVersion(ctx context.Context, db *sql.DB) (int, error) {
|
||||
return v, nil
|
||||
}
|
||||
|
||||
// loadFileRows reads every record from the files table.
|
||||
func loadFileRows(ctx context.Context, db *sql.DB) ([]scanRec, error) {
|
||||
// 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")
|
||||
"SELECT path, size, mtime, head, tail, content FROM files "+
|
||||
"ORDER BY path")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read records: %w", err)
|
||||
return fmt.Errorf("read records: %w", err)
|
||||
}
|
||||
|
||||
defer func() { _ = rows.Close() }()
|
||||
|
||||
var recs []scanRec
|
||||
|
||||
for rows.Next() {
|
||||
var (
|
||||
path []byte
|
||||
@@ -298,19 +309,92 @@ func loadFileRows(ctx context.Context, db *sql.DB) ([]scanRec, error) {
|
||||
err = rows.Scan(&path, &r.size, &r.mtime, &r.head, &r.tail,
|
||||
&r.content)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read record: %w", err)
|
||||
return fmt.Errorf("read record: %w", err)
|
||||
}
|
||||
|
||||
r.path = string(path)
|
||||
recs = append(recs, r)
|
||||
fn(r)
|
||||
}
|
||||
|
||||
err = rows.Err()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read records: %w", err)
|
||||
return fmt.Errorf("read records: %w", err)
|
||||
}
|
||||
|
||||
return recs, nil
|
||||
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
|
||||
|
||||
Reference in New Issue
Block a user