Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1ed3df06f2 |
+110
@@ -4,12 +4,16 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
"errors"
|
"errors"
|
||||||
|
"fmt"
|
||||||
"os"
|
"os"
|
||||||
|
"os/signal"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"slices"
|
"slices"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
|
"syscall"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
@@ -384,6 +388,112 @@ func TestRunScanInterrupted(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 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,
|
// drainClosed counts the values received from ch until it closes,
|
||||||
// failing the test if it does not close within poolUnwind. A pool that
|
// failing the test if it does not close within poolUnwind. A pool that
|
||||||
// ignored its cancellation leaves its channel open with its goroutines
|
// ignored its cancellation leaves its channel open with its goroutines
|
||||||
|
|||||||
@@ -135,6 +135,9 @@ func newRootCommand(stdout, stderr io.Writer) *cobra.Command {
|
|||||||
Short: "Walk trees and synchronize the scan database",
|
Short: "Walk trees and synchronize the scan database",
|
||||||
Args: cobra.MinimumNArgs(1),
|
Args: cobra.MinimumNArgs(1),
|
||||||
RunE: runE(func(ctx context.Context, args []string) error {
|
RunE: runE(func(ctx context.Context, args []string) error {
|
||||||
|
ctx, stop := interruptContext(ctx)
|
||||||
|
defer stop()
|
||||||
|
|
||||||
return runScan(ctx, args, scanWorkers, scanOneFS)
|
return runScan(ctx, args, scanWorkers, scanOneFS)
|
||||||
}),
|
}),
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -95,9 +95,10 @@ type fileMeta struct {
|
|||||||
// scan fails before it walks the filesystem or opens the database.
|
// scan fails before it walks the filesystem or opens the database.
|
||||||
// Errors are returned rather than exiting, so that the deferred close —
|
// Errors are returned rather than exiting, so that the deferred close —
|
||||||
// which takes the database out of WAL mode — always runs, and the lock
|
// which takes the database out of WAL mode — always runs, and the lock
|
||||||
// is released after it. SIGINT or SIGTERM cancels ctx: the scan keeps
|
// is released after it. When ctx is cancelled, as by the SIGINT or
|
||||||
// what it has hashed (see syncScan), prints how many files its walk
|
// SIGTERM that interruptContext catches, the scan keeps what it has
|
||||||
// reached, and returns errInterrupted.
|
// hashed (see syncScan), prints how many files its walk reached, and
|
||||||
|
// returns errInterrupted.
|
||||||
func runScan(ctx context.Context, roots []string, workers int,
|
func runScan(ctx context.Context, roots []string, workers int,
|
||||||
oneFS bool,
|
oneFS bool,
|
||||||
) error {
|
) error {
|
||||||
@@ -105,20 +106,6 @@ func runScan(ctx context.Context, roots []string, workers int,
|
|||||||
workers = 1
|
workers = 1
|
||||||
}
|
}
|
||||||
|
|
||||||
// A SIGINT ignored from the start, as by a script's background job,
|
|
||||||
// stays ignored.
|
|
||||||
signals := []os.Signal{syscall.SIGTERM}
|
|
||||||
if !signal.Ignored(syscall.SIGINT) {
|
|
||||||
signals = append(signals, syscall.SIGINT)
|
|
||||||
}
|
|
||||||
|
|
||||||
ctx, stop := signal.NotifyContext(ctx, signals...)
|
|
||||||
defer stop()
|
|
||||||
|
|
||||||
// Stopping restores the default handling, so a second signal ends
|
|
||||||
// the process at once.
|
|
||||||
context.AfterFunc(ctx, stop)
|
|
||||||
|
|
||||||
roots, err := resolveRoots(roots)
|
roots, err := resolveRoots(roots)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -164,6 +151,26 @@ func runScan(ctx context.Context, roots []string, workers int,
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// interruptContext returns a copy of ctx that the first SIGINT or
|
||||||
|
// SIGTERM cancels; the scan command runs the scan under it. stop
|
||||||
|
// releases the signals.
|
||||||
|
func interruptContext(ctx context.Context) (context.Context, func()) {
|
||||||
|
// A SIGINT ignored from the start, as by a script's background job,
|
||||||
|
// stays ignored.
|
||||||
|
signals := []os.Signal{syscall.SIGTERM}
|
||||||
|
if !signal.Ignored(syscall.SIGINT) {
|
||||||
|
signals = append(signals, syscall.SIGINT)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx, stop := signal.NotifyContext(ctx, signals...)
|
||||||
|
|
||||||
|
// Stopping restores the default handling, so a second signal ends
|
||||||
|
// the process at once.
|
||||||
|
context.AfterFunc(ctx, stop)
|
||||||
|
|
||||||
|
return ctx, stop
|
||||||
|
}
|
||||||
|
|
||||||
// interrupted prints the line for a scan stopped by a signal after its
|
// interrupted prints the line for a scan stopped by a signal after its
|
||||||
// walk reached walked files, and returns errInterrupted.
|
// walk reached walked files, and returns errInterrupted.
|
||||||
func interrupted(walked int) error {
|
func interrupted(walked int) error {
|
||||||
|
|||||||
+3
-3
@@ -1500,9 +1500,9 @@ func injectWriteFailure(t *testing.T, path string) {
|
|||||||
func baselineGoroutines(t *testing.T) int {
|
func baselineGoroutines(t *testing.T) int {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
// The first scan in a process starts os/signal's goroutine, which
|
// The first scan command in a process starts os/signal's goroutine,
|
||||||
// never exits. Start it now, so the baseline counts it instead of
|
// which never exits. Start it now, so the baseline counts it
|
||||||
// the scan seeming to leave it behind.
|
// instead of the scan seeming to leave it behind.
|
||||||
ch := make(chan os.Signal, 1)
|
ch := make(chan os.Signal, 1)
|
||||||
signal.Notify(ch, syscall.SIGINT)
|
signal.Notify(ch, syscall.SIGINT)
|
||||||
signal.Stop(ch)
|
signal.Stop(ch)
|
||||||
|
|||||||
Reference in New Issue
Block a user