Files
sfdupes/trees.go
sneak 1399249957
All checks were successful
check / check (push) Successful in 1m2s
Unwind the hash worker pool instead of abandoning it (closes #6)
hashPhase returned the moment recordRun failed and left the pool
running: the feeder parked forever on a full jobs channel and every
worker on a full results channel. Until #4 landed the process exited
before that mattered; now that runScan returns an error and unwinds,
the goroutines are a real leak.

The pool is now an owned hashPool. Its context is derived from the
scan's, every blocking send in the feeder and the workers selects on
ctx.Done(), the feeder closes jobs on every path out so the workers'
range always terminates, and hashPhase defers pool.stop(), which
cancels and then drains results until the last goroutine has exited.
Draining is the half that matters: a worker already parked on a send
cannot observe the cancellation until a receiver frees it.

ctx comes from cmd.Context() and is threaded through runScan,
syncScan, both worker pools and the database layer as the first
parameter throughout, so graceful interrupt handling has a path to
hook into rather than a pool to rewrite.

The walk pool never leaked, because walkPhase always drains its
events to close, but it has the same unbounded-send shape and gets
the same treatment, together with a ctx.Err() guard after the walk: a
cancelled walk leaves a partial size census, and the update phase
would read every file it never reached as vanished and delete its
record.

Tests drive the scan entry point against a database whose insert
trigger aborts, with a fixture large enough that the failure lands
partway through the hash phase with more runs queued than either pool
channel can hold, and assert that the scan fails instead of hanging
and that runtime.NumGoroutine polls back to its pre-scan baseline.
2026-08-09 03:00:01 +00:00

244 lines
5.9 KiB
Go

package main
import (
"bufio"
"context"
"crypto/sha256"
"fmt"
"os"
"slices"
"strconv"
"strings"
)
// fileSig is a file's duplicate signature; mtime is excluded.
type fileSig struct {
size int64
head string
tail string
}
// treeNode is one directory reconstructed from the scan stream.
type treeNode struct {
path string
parent *treeNode
dirs map[string]*treeNode
files map[string]fileSig
digest [sha256.Size]byte
fileCount int64
totalSize int64
}
// runTrees implements the trees subcommand: it reads every record from
// the database, reconstructs the directory hierarchy from the record
// paths, computes a Merkle-style digest per directory, and prints
// maximal duplicate-tree groups as TSV on stdout. It never touches the
// scanned filesystem; its only I/O is the database, stdout, and
// stderr.
func runTrees(ctx context.Context) error {
recs, err := loadRecords(ctx)
if err != nil {
return err
}
super, allDirs := buildHierarchy(recs)
super.compute()
dupes := collectTreeGroups(allDirs, super)
out := bufio.NewWriterSize(os.Stdout, ioBufSize)
_, err = fmt.Fprintln(out, "first\tdupe\tfiles\tsize")
if err != nil {
return fmt.Errorf("write stdout: %w", err)
}
dupeTrees := 0
var reclaimable int64
for _, g := range dupes {
first := g[0]
for _, n := range g[1:] {
_, err = fmt.Fprintf(out, "%s\t%s\t%d\t%d\n",
first.path, n.path, first.fileCount, first.totalSize)
if err != nil {
return fmt.Errorf("write stdout: %w", err)
}
dupeTrees++
reclaimable += first.totalSize
}
}
err = out.Flush()
if err != nil {
return fmt.Errorf("write stdout: %w", err)
}
fmt.Fprintf(os.Stderr,
"trees: %d records read, %d duplicate tree groups, %d dupe trees, "+
"%s reclaimable\n",
len(recs), len(dupes), dupeTrees, humanBytes(reclaimable))
return nil
}
// buildHierarchy reconstructs the directory hierarchy from the record
// paths under a synthetic super-root. Paths are split on "/"; for
// absolute paths the first component is empty, which simply becomes a
// top-level node representing "/". It returns the super-root and every
// directory node created.
func buildHierarchy(recs []scanRec) (*treeNode, []*treeNode) {
super := &treeNode{}
var allDirs []*treeNode
for _, r := range recs {
comps := strings.Split(r.path, "/")
node := super
for _, c := range comps[:len(comps)-1] {
child := node.dirs[c]
if child == nil {
childPath := c
if node != super {
childPath = node.path + "/" + c
}
child = &treeNode{path: childPath, parent: node}
if node.dirs == nil {
node.dirs = make(map[string]*treeNode)
}
node.dirs[c] = child
allDirs = append(allDirs, child)
}
node = child
}
if node.files == nil {
node.files = make(map[string]fileSig)
}
sig := fileSig{size: r.size, head: r.head, tail: r.tail}
// An unhashed record (its size was unique when last scanned)
// has unknown content: give it a signature no other file can
// share, so trees containing it never compare equal. Real
// heads are hex, so the NUL-prefixed form cannot collide.
if sig.head == "" {
sig.head = "unhashed\x00" + r.path
}
node.files[comps[len(comps)-1]] = sig
}
return super, allDirs
}
// collectTreeGroups groups directories by digest and returns every
// maximal group with two or more members, each group's members sorted
// by path, groups ordered by tree size descending then by first path
// ascending.
func collectTreeGroups(allDirs []*treeNode, super *treeNode) [][]*treeNode {
groups := make(map[[sha256.Size]byte][]*treeNode)
for _, d := range allDirs {
groups[d.digest] = append(groups[d.digest], d)
}
var dupes [][]*treeNode
for _, g := range groups {
if len(g) < minGroupSize || suppressed(g, super) {
continue
}
slices.SortFunc(g, func(a, b *treeNode) int {
return strings.Compare(a.path, b.path)
})
dupes = append(dupes, g)
}
// Biggest reclaimable space first; ties broken by first path.
slices.SortFunc(dupes, func(a, b []*treeNode) int {
if a[0].totalSize != b[0].totalSize {
if a[0].totalSize > b[0].totalSize {
return -1
}
return 1
}
return strings.Compare(a[0].path, b[0].path)
})
return dupes
}
// compute fills in digest, fileCount, and totalSize for n and all of
// its descendants. A directory's digest is the SHA-256 of its child
// entries — files serialized with name and signature, subdirectories
// with name and recursive digest — sorted byte-lexicographically.
// Filenames cannot contain NUL or "/", so NUL delimiters are
// unambiguous.
func (n *treeNode) compute() {
entries := make([]string, 0, len(n.dirs)+len(n.files))
for name, sig := range n.files {
entries = append(entries,
"f\x00"+name+"\x00"+strconv.FormatInt(sig.size, 10)+
"\x00"+sig.head+"\x00"+sig.tail)
n.fileCount++
n.totalSize += sig.size
}
for name, child := range n.dirs {
child.compute()
entries = append(entries, "d\x00"+name+"\x00"+string(child.digest[:]))
n.fileCount += child.fileCount
n.totalSize += child.totalSize
}
slices.Sort(entries)
h := sha256.New()
for _, e := range entries {
h.Write([]byte(e))
h.Write([]byte{0})
}
copy(n.digest[:], h.Sum(nil))
}
// suppressed reports whether a duplicate-tree group is non-maximal: its
// members' parents are pairwise distinct real directories that all
// share a single digest, so the group is wholly implied by its parents'
// (or a further ancestor's) group. Groups containing siblings (shared
// parent) or members whose parents differ are always reported.
func suppressed(g []*treeNode, super *treeNode) bool {
seen := make(map[*treeNode]bool, len(g))
var parentDigest [sha256.Size]byte
for i, n := range g {
p := n.parent
if p == nil || p == super {
return false
}
if seen[p] {
return false // siblings: not implied by any parent group
}
seen[p] = true
if i == 0 {
parentDigest = p.digest
} else if p.digest != parentDigest {
return false
}
}
return true
}