check / check (push) Waiting to run
TestHashWorkerDropsQueuedRuns now passes hashWorker a hash function that records being called, so a worker that hashes a run after the scan is cancelled fails the test every time instead of only when it then chose to send its result. TestHashWorkerAbandonsBlockedSend cancels the scan from inside the hash function and leaves the result channel unread, so the worker can only return through the cancellation case beside its send. The scan tests could not show this, because stop drains results and frees a parked worker anyway. Model: opus-5-5
824 lines
25 KiB
Go
824 lines
25 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"os/signal"
|
|
"path/filepath"
|
|
"slices"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"syscall"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
// This file gathers the tests for scan cancellation and worker-pool
|
|
// unwinding. Everything it exercises lives in scan.go, so by the repo's
|
|
// convention of one test file per source file it would belong in
|
|
// scan_test.go. It is kept separate on purpose: cancellation behaviour
|
|
// cuts across both the walk pool and the hash pool as a single concern,
|
|
// and scan_test.go is already over 1,600 lines. That is the deliberate
|
|
// exception the convention otherwise expects to be stated.
|
|
|
|
// poolUnwind bounds how long a test waits for a cancellation to take
|
|
// effect: for a goroutine to return or a channel to close once its
|
|
// context is cancelled, or for a signal to cancel the scan's context.
|
|
// Only a failing run waits this long, and the bound is what makes that
|
|
// failure an assertion instead of a hang. A call made without it, as
|
|
// most of this file's scans are, has no bound: a regression that parks
|
|
// it is caught only as the test binary's own timeout.
|
|
const poolUnwind = 2 * time.Second
|
|
|
|
// walkClock is a context whose cancellation is driven by the scan's
|
|
// own progress rather than by the wall clock: it cancels itself the
|
|
// moment its Done method has been consulted n times. That is what
|
|
// makes "cancel in the middle of the walk" reproducible instead of a
|
|
// race against a timer.
|
|
//
|
|
// The accounting behind the n chosen by each test: every blocking
|
|
// channel operation in the walk selects on Done, so the walk spends
|
|
// one consultation per file event plus a couple per directory. The
|
|
// index load that runs ahead of it also consults Done, but a bounded
|
|
// number of times that does not grow with the record count. The tests
|
|
// depend on that property, not on the bound's exact value: each test
|
|
// sets n from the consultations of the walk, plus those of the hash
|
|
// phase when it cancels mid-hash, far from both ends of the phase it
|
|
// interrupts, so the cancellation lands inside that phase whatever the
|
|
// record count.
|
|
type walkClock struct {
|
|
n int64
|
|
seen atomic.Int64
|
|
once sync.Once
|
|
done chan struct{}
|
|
}
|
|
|
|
// newWalkClock returns a context that cancels itself on the nth
|
|
// consultation of its Done method.
|
|
func newWalkClock(n int64) *walkClock {
|
|
return &walkClock{n: n, done: make(chan struct{})}
|
|
}
|
|
|
|
// Done returns the cancellation channel, cancelling the context on the
|
|
// nth call and on every call after it. The same channel is returned
|
|
// throughout, so a caller that took it before the cancellation still
|
|
// observes the close.
|
|
func (c *walkClock) Done() <-chan struct{} {
|
|
if c.seen.Add(1) >= c.n {
|
|
c.once.Do(func() { close(c.done) })
|
|
}
|
|
|
|
return c.done
|
|
}
|
|
|
|
// Err reports the cancellation without consuming a consultation, which
|
|
// is what lets the post-walk guard read it without disturbing the
|
|
// count.
|
|
func (c *walkClock) Err() error {
|
|
select {
|
|
case <-c.done:
|
|
return context.Canceled
|
|
default:
|
|
return nil
|
|
}
|
|
}
|
|
|
|
// Deadline reports no deadline: this context is cancelled by progress,
|
|
// never by time.
|
|
func (c *walkClock) Deadline() (time.Time, bool) {
|
|
return time.Time{}, false
|
|
}
|
|
|
|
// Value carries nothing.
|
|
func (c *walkClock) Value(_ any) any {
|
|
return nil
|
|
}
|
|
|
|
// walkCancelDirs and walkCancelFilesPerDir shape the fixture for the
|
|
// mid-walk cancellation test. Spreading the files over directories is
|
|
// load-bearing: it is what bounds how much of the tree can still be
|
|
// walked after the cancellation, since the workers drop every
|
|
// directory still queued and only the handful already in flight can
|
|
// emit anything more.
|
|
const (
|
|
walkCancelDirs = 100
|
|
walkCancelFilesPerDir = 20
|
|
walkCancelFiles = walkCancelDirs * walkCancelFilesPerDir
|
|
walkCancelWorkers = 4
|
|
// The most files the walkCancelWorkers directories already in
|
|
// flight when the scan is cancelled can still emit, at
|
|
// walkCancelFilesPerDir each. A file count, not a directory count.
|
|
walkCancelInFlightFiles = walkCancelWorkers * walkCancelFilesPerDir
|
|
)
|
|
|
|
// walkCancelAtDone is the consultation on which the fixture's context
|
|
// cancels itself. A quarter of the file count is far past the index
|
|
// load's fixed handful and far short of the walk's total, so the
|
|
// cancellation lands deep inside the walk and nowhere near either end
|
|
// of it.
|
|
const walkCancelAtDone = walkCancelFiles / 4
|
|
|
|
// buildWalkCancelTree writes walkCancelFiles empty files spread over
|
|
// walkCancelDirs subdirectories. Zero-length files are never opened by
|
|
// the hasher, so the fixture costs directory entries and no read I/O
|
|
// while still giving the walk thousands of events to emit.
|
|
func buildWalkCancelTree(t *testing.T) string {
|
|
t.Helper()
|
|
|
|
dir := t.TempDir()
|
|
|
|
for i := range walkCancelDirs {
|
|
sub := filepath.Join(dir, "d"+strconv.Itoa(i))
|
|
|
|
err := os.Mkdir(sub, 0o750)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
writeEmptyFiles(t, sub, walkCancelFilesPerDir)
|
|
}
|
|
|
|
return dir
|
|
}
|
|
|
|
// assertRecordsIntact fails when the database no longer holds exactly
|
|
// the records it held before, reporting the first difference rather
|
|
// than dumping thousands of paths.
|
|
func assertRecordsIntact(t *testing.T, db *sql.DB, before []string) {
|
|
t.Helper()
|
|
|
|
got := recordPaths(dbRecords(t, db))
|
|
if len(got) != len(before) {
|
|
t.Fatalf("%d records after the cancelled scan, want %d",
|
|
len(got), len(before))
|
|
}
|
|
|
|
for i := range got {
|
|
if got[i] != before[i] {
|
|
t.Fatalf("record %d = %q after the cancelled scan, want %q",
|
|
i, got[i], before[i])
|
|
}
|
|
}
|
|
}
|
|
|
|
// TestSyncScanCancelledMidWalkKeepsRecords is the regression net under
|
|
// the post-walk guard. The scan is cancelled part-way through the
|
|
// walk, so it reaches the guard holding a genuinely partial size
|
|
// census and a still-populated index of records the walk never got to.
|
|
// Every one of those records would look vanished to the update phase.
|
|
// The guard is what stops the scan there, and this test is what
|
|
// notices if it stops doing so: deleting the guard, or making it
|
|
// unreachable, makes the scan carry its truncated view into the update
|
|
// phase, which counts every record the walk never reached for removal.
|
|
//
|
|
// The syncScan call here is not bounded by poolUnwind: a regression
|
|
// that left a worker pool parked would hang it, and that regression is
|
|
// caught only by the test binary's own timeout, not by a quick
|
|
// assertion.
|
|
//
|
|
//nolint:paralleltest // counts goroutines: must not run beside others
|
|
func TestSyncScanCancelledMidWalkKeepsRecords(t *testing.T) {
|
|
dir := buildWalkCancelTree(t)
|
|
db := openTestDB(t)
|
|
|
|
st := syncTree(t, db, dir)
|
|
if st.added != walkCancelFiles {
|
|
t.Fatalf("setup scan added %d records, want %d",
|
|
st.added, walkCancelFiles)
|
|
}
|
|
|
|
before := recordPaths(dbRecords(t, db))
|
|
base := baselineGoroutines(t)
|
|
|
|
st, err := syncScan(newWalkClock(walkCancelAtDone), db,
|
|
[]string{dir}, walkCancelWorkers, false)
|
|
|
|
assertWalkGuardAborted(t, st, err)
|
|
assertRecordsIntact(t, db, before)
|
|
|
|
if got := settledGoroutines(t, base); got > base {
|
|
t.Errorf("goroutines = %d after the cancelled scan, want %d back",
|
|
got, base)
|
|
}
|
|
}
|
|
|
|
// assertWalkGuardAborted checks that the scan stopped at the post-walk
|
|
// guard: with a census that is neither empty (the walk really ran)
|
|
// nor complete (it really was cut short), and with no record counted
|
|
// for removal. A removal count means the partial census was carried
|
|
// past the guard into the update phase, which is the failure this test
|
|
// exists to catch.
|
|
func assertWalkGuardAborted(t *testing.T, st scanStats, err error) {
|
|
t.Helper()
|
|
|
|
if !errors.Is(err, context.Canceled) {
|
|
t.Fatalf("syncScan cancelled mid-walk = %v, want %v",
|
|
err, context.Canceled)
|
|
}
|
|
|
|
if st.unchanged == 0 {
|
|
t.Fatalf("stats = %+v: the census is empty, so the walk never "+
|
|
"ran and the guard was reached for the wrong reason", st)
|
|
}
|
|
|
|
if st.unchanged >= walkCancelFiles {
|
|
t.Fatalf("stats = %+v: the census covers the whole tree, so the "+
|
|
"walk was not cut short", st)
|
|
}
|
|
|
|
// The workers drop every directory still queued once the scan is
|
|
// cancelled, so only the files in the directories already in flight
|
|
// can add to the census after the fact. A census beyond that bound
|
|
// would mean the cancellation was not observed where it should have
|
|
// been.
|
|
limit := walkCancelAtDone + walkCancelInFlightFiles
|
|
if st.unchanged > limit {
|
|
t.Errorf("census covers %d files, want at most %d: the walk kept "+
|
|
"taking directories off the queue after cancellation",
|
|
st.unchanged, limit)
|
|
}
|
|
|
|
if st.removed != 0 {
|
|
t.Errorf("stats = %+v: the scan counted records for removal from "+
|
|
"a partial census", st)
|
|
}
|
|
}
|
|
|
|
// TestSyncScanCancelledBeforeLoadIndex covers the trivial end of the
|
|
// cancellation path: a scan handed a context that is already cancelled
|
|
// fails in the index load, before the walk pool is ever started. It
|
|
// says nothing about the post-walk guard — nothing downstream of
|
|
// loadIndex runs at all — only that the failure surfaces as a
|
|
// cancellation and that no record is touched on the way out.
|
|
func TestSyncScanCancelledBeforeLoadIndex(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
dir := buildSmokeTree(t)
|
|
db := openTestDB(t)
|
|
|
|
syncTree(t, db, dir)
|
|
|
|
before := recordPaths(dbRecords(t, db))
|
|
|
|
ctx, cancel := context.WithCancel(t.Context())
|
|
cancel()
|
|
|
|
st, err := syncScan(ctx, db, []string{dir}, walkCancelWorkers, false)
|
|
if !errors.Is(err, context.Canceled) {
|
|
t.Fatalf("syncScan on a cancelled context = %v, want %v",
|
|
err, context.Canceled)
|
|
}
|
|
|
|
if st != (scanStats{}) {
|
|
t.Errorf("stats = %+v, want none: the scan gave up in the index "+
|
|
"load, before any phase ran", st)
|
|
}
|
|
|
|
assertRecordsIntact(t, db, before)
|
|
}
|
|
|
|
// hashCancelAtDone is the consultation on which the mid-hash test's
|
|
// context cancels itself. The walk of buildWalkCancelTree spends about
|
|
// one per file and three per directory, and the hash phase then one per
|
|
// file hashed, so this lands about half way through the hash phase.
|
|
const hashCancelAtDone = walkCancelFiles + 3*walkCancelDirs +
|
|
walkCancelFiles/2
|
|
|
|
// TestSyncScanCancelledMidHashKeepsHashedRecords cancels a first scan
|
|
// part-way through its hash phase. The fixture holds fewer files than a
|
|
// batch, so every file hashed is still waiting to be committed: the scan
|
|
// must commit them all before it returns, and the next scan must hash
|
|
// only the rest.
|
|
func TestSyncScanCancelledMidHashKeepsHashedRecords(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
dir := buildWalkCancelTree(t)
|
|
db := openTestDB(t)
|
|
|
|
st, err := syncScan(newWalkClock(hashCancelAtDone), db,
|
|
[]string{dir}, walkCancelWorkers, false)
|
|
if !errors.Is(err, context.Canceled) {
|
|
t.Fatalf("syncScan cancelled mid-hash = %v, want %v",
|
|
err, context.Canceled)
|
|
}
|
|
|
|
if st.walked != walkCancelFiles || st.added == 0 ||
|
|
st.added >= walkCancelFiles {
|
|
t.Fatalf("stats = %+v: want the walk complete and the hash phase "+
|
|
"cut short", st)
|
|
}
|
|
|
|
if got := len(dbRecords(t, db)); got != st.added {
|
|
t.Errorf("%d records after the cancelled scan, want the %d it hashed",
|
|
got, st.added)
|
|
}
|
|
|
|
hashed := st.added
|
|
|
|
st = syncTree(t, db, dir)
|
|
if st.added != walkCancelFiles-hashed || st.unchanged != hashed {
|
|
t.Errorf("next scan stats = %+v, want %d added %d unchanged",
|
|
st, walkCancelFiles-hashed, hashed)
|
|
}
|
|
}
|
|
|
|
// storedPaths opens the database at path as report does, which fails
|
|
// unless it is a valid database, and returns its records' paths.
|
|
func storedPaths(t *testing.T, path string) []string {
|
|
t.Helper()
|
|
|
|
db, err := openReportDatabase(t.Context(), path)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
defer func() { _ = db.Close() }()
|
|
|
|
return recordPaths(dbRecords(t, db))
|
|
}
|
|
|
|
// TestRunScanInterrupted calls the scan entrypoint with a context that
|
|
// is already cancelled, as when a signal arrives at once. It must return
|
|
// errInterrupted promptly with its one line on stderr, leave the
|
|
// database valid and as it was, and leave nothing in the way of the
|
|
// next scan, which must bring the database up to date.
|
|
func TestRunScanInterrupted(t *testing.T) {
|
|
path := testDBPath(t)
|
|
t.Setenv(databaseEnv, path)
|
|
|
|
stderr := captureStderr(t)
|
|
dir := buildSmokeTree(t)
|
|
|
|
err := runScan(t.Context(), []string{dir}, walkCancelWorkers, false)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
before := storedPaths(t, path)
|
|
|
|
// A vanished file and a new one: the interrupted scan records
|
|
// neither.
|
|
gone := filepath.Join(dir, "a", "unique.bin")
|
|
|
|
err = os.Remove(gone)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
added := writeFile(t, dir, "a/new.bin", pattern(50, 10))
|
|
shown := len(stderr())
|
|
done := make(chan struct{})
|
|
|
|
go func() {
|
|
defer close(done)
|
|
|
|
err = runScan(cancelledContext(t), []string{dir}, walkCancelWorkers,
|
|
false)
|
|
}()
|
|
|
|
awaitReturn(t, done, "runScan")
|
|
|
|
if !errors.Is(err, errInterrupted) {
|
|
t.Fatalf("runScan on a cancelled context = %v, want %v",
|
|
err, errInterrupted)
|
|
}
|
|
|
|
want := "scan: interrupted after 0 files\n"
|
|
if got := stderr()[shown:]; got != want {
|
|
t.Errorf("stderr = %q, want %q", got, want)
|
|
}
|
|
|
|
assertNoSidecars(t, path)
|
|
|
|
if got := storedPaths(t, path); !slices.Equal(got, before) {
|
|
t.Errorf("records = %q after the interrupted scan, want %q",
|
|
got, before)
|
|
}
|
|
|
|
err = runScan(t.Context(), []string{dir}, walkCancelWorkers, false)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
got := storedPaths(t, path)
|
|
if slices.Contains(got, gone) || !slices.Contains(got, added) {
|
|
t.Errorf("records = %q after the next scan, want %q gone and %q "+
|
|
"added", got, gone, added)
|
|
}
|
|
}
|
|
|
|
// TestRunScanInterruptedMidHash interrupts the scan entrypoint part-way
|
|
// through its hash phase, after the database is open. It must return
|
|
// errInterrupted, release the lock, end stderr with its line counting
|
|
// every file the walk reached, close the database out of WAL mode, and
|
|
// keep the records it hashed.
|
|
func TestRunScanInterruptedMidHash(t *testing.T) {
|
|
path := testDBPath(t)
|
|
t.Setenv(databaseEnv, path)
|
|
|
|
stderr := captureStderr(t)
|
|
dir := buildWalkCancelTree(t)
|
|
|
|
err := runScan(newWalkClock(hashCancelAtDone), []string{dir},
|
|
walkCancelWorkers, false)
|
|
if !errors.Is(err, errInterrupted) {
|
|
t.Fatalf("runScan interrupted mid-hash = %v, want %v",
|
|
err, errInterrupted)
|
|
}
|
|
|
|
holdScanLock(t, path)
|
|
|
|
want := fmt.Sprintf("scan: interrupted after %d files\n", walkCancelFiles)
|
|
if got := stderr(); !strings.HasSuffix(got, want) {
|
|
t.Errorf("stderr = %q, want it to end with %q", got, want)
|
|
}
|
|
|
|
assertNoSidecars(t, path)
|
|
|
|
db, err := openReportDatabase(t.Context(), path)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
defer func() { _ = db.Close() }()
|
|
|
|
// A plain close also removes the sidecars, but leaves WAL mode on.
|
|
var mode string
|
|
|
|
err = db.QueryRowContext(t.Context(), "PRAGMA journal_mode").Scan(&mode)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
if mode != "delete" {
|
|
t.Errorf("journal mode = %q after the interrupted scan, want %q",
|
|
mode, "delete")
|
|
}
|
|
|
|
kept := len(dbRecords(t, db))
|
|
if kept == 0 || kept >= walkCancelFiles {
|
|
t.Errorf("%d records after the interrupted scan, want those it "+
|
|
"hashed: some but not all of the %d files", kept, walkCancelFiles)
|
|
}
|
|
}
|
|
|
|
// TestInterruptContextCatchesSIGTERM sends SIGTERM to the test process
|
|
// while the scan's handler is installed, and checks that it cancels the
|
|
// scan's context.
|
|
//
|
|
//nolint:paralleltest // signals the whole process: must not run beside a scan
|
|
func TestInterruptContextCatchesSIGTERM(t *testing.T) {
|
|
// Caught here as well, so that a handler that misses SIGTERM fails
|
|
// this test instead of ending the test process.
|
|
caught := make(chan os.Signal, 1)
|
|
signal.Notify(caught, syscall.SIGTERM)
|
|
|
|
defer signal.Stop(caught)
|
|
|
|
ctx, stop := interruptContext(t.Context())
|
|
defer stop()
|
|
|
|
err := syscall.Kill(os.Getpid(), syscall.SIGTERM)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
case <-time.After(poolUnwind):
|
|
t.Fatal("SIGTERM did not cancel the scan's context")
|
|
}
|
|
}
|
|
|
|
// TestCommitFullBatchKeepsFailedBatch checks that a full batch whose
|
|
// commit fails, as it does once the scan is interrupted, stays in the
|
|
// batch, so that syncScan's final commit saves it.
|
|
func TestCommitFullBatchKeepsFailedBatch(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
s := &scanState{db: openTestDB(t)}
|
|
for i := range updateBatchSize {
|
|
s.batch = append(s.batch, scanRec{path: "/f" + strconv.Itoa(i)})
|
|
}
|
|
|
|
err := s.commitFullBatch(cancelledContext(t))
|
|
if !errors.Is(err, context.Canceled) {
|
|
t.Fatalf("commitFullBatch on a cancelled context = %v, want %v",
|
|
err, context.Canceled)
|
|
}
|
|
|
|
if len(s.batch) != updateBatchSize {
|
|
t.Errorf("batch holds %d records after the failed commit, want %d",
|
|
len(s.batch), updateBatchSize)
|
|
}
|
|
}
|
|
|
|
// drainClosed counts the values received from ch until it closes,
|
|
// failing the test if it does not close within poolUnwind. A pool that
|
|
// ignored its cancellation leaves its channel open with its goroutines
|
|
// parked, and this is what reports that as an assertion.
|
|
func drainClosed[T any](t *testing.T, ch <-chan T, what string) int {
|
|
t.Helper()
|
|
|
|
counted := make(chan int, 1)
|
|
|
|
go func() {
|
|
n := 0
|
|
for range ch {
|
|
n++
|
|
}
|
|
|
|
counted <- n
|
|
}()
|
|
|
|
select {
|
|
case n := <-counted:
|
|
return n
|
|
case <-time.After(poolUnwind):
|
|
t.Fatalf("%s stayed open after cancellation", what)
|
|
|
|
return 0
|
|
}
|
|
}
|
|
|
|
// awaitReturn fails the test if done is not closed within poolUnwind.
|
|
func awaitReturn(t *testing.T, done <-chan struct{}, what string) {
|
|
t.Helper()
|
|
|
|
select {
|
|
case <-done:
|
|
case <-time.After(poolUnwind):
|
|
t.Fatalf("%s did not return after cancellation", what)
|
|
}
|
|
}
|
|
|
|
// cancelledContext returns a context that is already cancelled.
|
|
func cancelledContext(t *testing.T) context.Context {
|
|
t.Helper()
|
|
|
|
ctx, cancel := context.WithCancel(t.Context())
|
|
cancel()
|
|
|
|
return ctx
|
|
}
|
|
|
|
// TestSendEventAbandonsBlockedSend checks that a walk goroutine with an
|
|
// event to deliver and nobody to deliver it to leaves on cancellation
|
|
// instead of holding the pool open. The channel here is unbuffered and
|
|
// unread, so the send can never complete.
|
|
func TestSendEventAbandonsBlockedSend(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
done := make(chan struct{})
|
|
events := make(chan walkEvent)
|
|
|
|
go func() {
|
|
defer close(done)
|
|
|
|
sendEvent(cancelledContext(t), events, walkEvent{})
|
|
}()
|
|
|
|
awaitReturn(t, done, "sendEvent")
|
|
}
|
|
|
|
// TestWalkOneDirStopsWhenCancelled checks that a cancelled scan stops
|
|
// reading a directory instead of going through the rest of its
|
|
// entries. A walk that kept going would return the subdirectory below
|
|
// to descend into. Unlike a file event, that return is not a send the
|
|
// cancellation can abandon, so the test catches the regression every
|
|
// time.
|
|
func TestWalkOneDirStopsWhenCancelled(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
dir := t.TempDir()
|
|
|
|
err := os.Mkdir(filepath.Join(dir, "sub"), 0o750)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
// Unbuffered and unread: on a cancelled scan every send gives up.
|
|
events := make(chan walkEvent)
|
|
|
|
subs := walkOneDir(cancelledContext(t), dirJob{path: dir}, false, events)
|
|
if len(subs) != 0 {
|
|
t.Errorf("cancelled walkOneDir returned %+v to descend into, "+
|
|
"want none", subs)
|
|
}
|
|
}
|
|
|
|
// TestWalkWorkersDropQueuedDirs checks that cancelled walk workers keep
|
|
// reading jobs and drop the directories rather than stopping their
|
|
// read: the range over jobs has to run out for the pool to tear down
|
|
// and close its event stream. The queued directory does not exist, so
|
|
// a worker that walked it anyway would send a warning before
|
|
// walkOneDir's own cancellation check could stop it. On a cancelled
|
|
// scan that send delivers or gives up at random, so with 64 jobs
|
|
// queued the regression has a one in 2^64 chance of passing.
|
|
func TestWalkWorkersDropQueuedDirs(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
missing := filepath.Join(t.TempDir(), "missing")
|
|
|
|
jobs, _, events := startWalkWorkers(cancelledContext(t), 2, false)
|
|
|
|
for range 64 {
|
|
jobs <- dirJob{path: missing}
|
|
}
|
|
|
|
close(jobs)
|
|
|
|
if n := drainClosed(t, events, "the walk event stream"); n != 0 {
|
|
t.Errorf("cancelled walk workers emitted %d events, want none", n)
|
|
}
|
|
}
|
|
|
|
// TestWalkWorkerAbandonsSubdirHandoff checks the other blocking send a
|
|
// walk worker makes: handing discovered subdirectories back to the
|
|
// dispatcher. Once the dispatcher has left, nothing drains that
|
|
// channel, and a worker parked on it would hold the pool open forever.
|
|
func TestWalkWorkerAbandonsSubdirHandoff(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
dir := t.TempDir()
|
|
writeEmptyFiles(t, dir, 1)
|
|
|
|
err := os.Mkdir(filepath.Join(dir, "sub"), 0o750)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
ctx, cancel := context.WithCancel(t.Context())
|
|
defer cancel()
|
|
|
|
jobs, subdirs, events := startWalkWorkers(ctx, 1, false)
|
|
|
|
// Fill the hand-back channel to its capacity — one slot per worker
|
|
// — so the worker's own hand-back is certain to block.
|
|
subdirs <- nil
|
|
|
|
jobs <- dirJob{path: dir}
|
|
|
|
// The file event proves the worker has read the directory and has
|
|
// nothing left to do but the blocked hand-back.
|
|
ev := <-events
|
|
if ev.fail {
|
|
t.Fatalf("walk event = %+v, want the fixture file", ev)
|
|
}
|
|
|
|
cancel()
|
|
close(jobs)
|
|
drainClosed(t, events, "the walk event stream")
|
|
}
|
|
|
|
// TestDispatchDirsClosesJobsWhenCancelled checks that a dispatcher
|
|
// leaving on cancellation closes the job channel on its way out. The
|
|
// workers range over that channel; a dispatcher that returned without
|
|
// closing it would strand every one of them.
|
|
func TestDispatchDirsClosesJobsWhenCancelled(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
// Unbuffered and unread: with no worker pool behind it, the
|
|
// dispatcher can only leave through its cancellation case.
|
|
jobs := make(chan dirJob)
|
|
subdirs := make(chan []dirJob)
|
|
initial := []dirJob{{path: "/a"}, {path: "/b"}}
|
|
|
|
dispatchDirs(cancelledContext(t), initial, jobs, subdirs)
|
|
|
|
if n := drainClosed(t, jobs, "the walk job queue"); n > len(initial) {
|
|
t.Errorf("dispatcher queued %d jobs, want at most %d",
|
|
n, len(initial))
|
|
}
|
|
}
|
|
|
|
// TestFeedHashJobsClosesJobsWhenCancelled checks that the hash feeder
|
|
// abandons the runs it has not queued yet and still closes the job
|
|
// channel, which is what lets the workers' range terminate. The
|
|
// receive on jobs below is not bounded: a feeder that returned without
|
|
// closing jobs would leave that receive with no sender and no close, so
|
|
// this regression is caught by the test binary's timeout rather than by
|
|
// a bounded assertion.
|
|
func TestFeedHashJobsClosesJobsWhenCancelled(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
done := make(chan struct{})
|
|
// Unbuffered and unread until the feeder has returned, so the only
|
|
// way out of the feeder is its cancellation case.
|
|
jobs := make(chan []fileRec)
|
|
runs := [][]fileRec{{{path: "a"}}, {{path: "b"}}}
|
|
|
|
go func() {
|
|
defer close(done)
|
|
|
|
feedHashJobs(cancelledContext(t), runs, jobs)
|
|
}()
|
|
|
|
awaitReturn(t, done, "feedHashJobs")
|
|
|
|
if _, ok := <-jobs; ok {
|
|
t.Error("the hash job channel was left open after cancellation")
|
|
}
|
|
}
|
|
|
|
// TestHashWorkerDropsQueuedRuns checks that a cancelled hash worker
|
|
// keeps reading jobs and drops the runs rather than reading files
|
|
// nobody wants the hashes of — while still letting the range run out
|
|
// so the pool tears down. The hash function records that it was
|
|
// called, so a worker that hashed the queued run anyway is caught
|
|
// every time.
|
|
func TestHashWorkerDropsQueuedRuns(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
done := make(chan struct{})
|
|
jobs := make(chan []fileRec, 1)
|
|
results := make(chan hashResult)
|
|
|
|
jobs <- []fileRec{{path: filepath.Join(t.TempDir(), "missing"), size: 1}}
|
|
|
|
close(jobs)
|
|
|
|
var hashed atomic.Bool
|
|
|
|
hash := func(path string, size int64) (string, string, string, error) {
|
|
hashed.Store(true)
|
|
|
|
return hashSignature(path, size)
|
|
}
|
|
|
|
go func() {
|
|
defer close(done)
|
|
|
|
hashWorker(cancelledContext(t), jobs, results, hash)
|
|
}()
|
|
|
|
awaitReturn(t, done, "hashWorker")
|
|
|
|
if hashed.Load() {
|
|
t.Error("cancelled hash worker hashed the queued run, want it dropped")
|
|
}
|
|
}
|
|
|
|
// TestHashWorkerAbandonsBlockedSend checks that a hash worker with a
|
|
// result to deliver and nobody to deliver it to leaves once the scan
|
|
// is cancelled, instead of holding the pool open. The scan tests do
|
|
// not catch this: stop drains results, which frees a parked worker
|
|
// anyway.
|
|
func TestHashWorkerAbandonsBlockedSend(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
ctx, cancel := context.WithCancel(t.Context())
|
|
defer cancel()
|
|
|
|
done := make(chan struct{})
|
|
jobs := make(chan []fileRec, 1)
|
|
// Unbuffered and unread, with jobs left open: the worker's only way
|
|
// out is the cancellation case beside its send.
|
|
results := make(chan hashResult)
|
|
|
|
jobs <- []fileRec{{path: filepath.Join(t.TempDir(), "missing"), size: 1}}
|
|
|
|
// The scan is cancelled while the worker hashes, so the worker has
|
|
// already passed the check that drops queued runs.
|
|
hash := func(path string, size int64) (string, string, string, error) {
|
|
cancel()
|
|
|
|
return hashSignature(path, size)
|
|
}
|
|
|
|
go func() {
|
|
defer close(done)
|
|
|
|
hashWorker(ctx, jobs, results, hash)
|
|
}()
|
|
|
|
awaitReturn(t, done, "hashWorker")
|
|
}
|
|
|
|
// TestHashPhaseCancelledReturnsContextError checks the result loop's
|
|
// own exit: with the pool cancelled, no result will ever arrive, and
|
|
// the loop must leave through the cancellation rather than wait for a
|
|
// receive that cannot happen. This call is not bounded by poolUnwind: a
|
|
// loop that dropped its cancellation case would block on that receive,
|
|
// so the regression surfaces as the test binary's timeout rather than
|
|
// as a bounded assertion.
|
|
func TestHashPhaseCancelledReturnsContextError(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
s := &scanState{
|
|
db: openTestDB(t),
|
|
toHash: []fileRec{{path: "a", size: 1, dev: 1, ino: 1}},
|
|
}
|
|
|
|
err := s.hashPhase(cancelledContext(t), 2)
|
|
if !errors.Is(err, context.Canceled) {
|
|
t.Fatalf("hashPhase on a cancelled context = %v, want %v",
|
|
err, context.Canceled)
|
|
}
|
|
}
|