Skip to content
Open
421 changes: 403 additions & 18 deletions cmd/analyze/analyze_test.go

Large diffs are not rendered by default.

20 changes: 17 additions & 3 deletions cmd/analyze/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -644,10 +644,14 @@ func loadStaleCacheFromDisk(path string) (*cacheEntry, error) {
}

func saveCacheToDisk(path string, result scanResult) error {
return saveCacheToDiskWithOptions(path, result, false)
ctx := context.Background()
return saveCacheToDiskWithOptions(newScanPublication(ctx, nil), path, result, false)
}

func saveCacheToDiskWithOptions(path string, result scanResult, needsRefresh bool) error {
func saveCacheToDiskWithOptions(publication *scanPublication, path string, result scanResult, needsRefresh bool) error {
if err := publication.ctx.Err(); err != nil {
return err
}
cachePath, err := getCachePath(path)
if err != nil {
return err
Expand Down Expand Up @@ -691,13 +695,23 @@ func saveCacheToDiskWithOptions(path string, result scanResult, needsRefresh boo
_ = os.Remove(tmpPath)
return err
}
if err := os.Rename(tmpPath, cachePath); err != nil {
err = publication.commit(func() error {
return os.Rename(tmpPath, cachePath)
})
if err != nil {
_ = os.Remove(tmpPath)
return err
}
return nil
}

func removeCacheEntryForScan(publication *scanPublication, path string) error {
return publication.commit(func() error {
removeCacheEntry(path)
return nil
})
}

// peekCacheTotalFiles reads the total file count from cache, ignoring
// expiration, for initial scan progress estimates. It shares
// loadRawCacheFromDisk so a schema-stale or corrupt entry is rejected and
Expand Down
3 changes: 2 additions & 1 deletion cmd/analyze/json.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
package main

import (
"context"
"encoding/json"
"fmt"
"os"
Expand Down Expand Up @@ -60,7 +61,7 @@ func performDirectoryScanForJSON(path string) jsonOutput {
currentPath := &atomic.Value{}
currentPath.Store("")

result, err := scanPathConcurrentAllEntries(path, &filesScanned, &dirsScanned, &bytesScanned, currentPath)
result, err := scanPathConcurrentAllEntries(context.Background(), path, &filesScanned, &dirsScanned, &bytesScanned, currentPath)
if err != nil {
fmt.Fprintf(os.Stderr, "failed to scan directory: %v\n", err)
os.Exit(1)
Expand Down
115 changes: 80 additions & 35 deletions cmd/analyze/live_scan.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,19 +34,73 @@ type liveScanTarget struct {
kind liveScanTargetKind
}

// liveScanEventStream uses the scan publication boundary for event delivery.
// Progress is lossy and may use only the non-reserved portion of the buffer;
// one slot per target plus the final completion slot stays available.
type liveScanEventStream struct {
publication *scanPublication
events chan liveScanEventMsg
requiredSlots int
closed bool
}

func newLiveScanEventStream(publication *scanPublication, targetCount int) *liveScanEventStream {
return &liveScanEventStream{
publication: publication,
events: make(chan liveScanEventMsg, max(targetCount*4, 1)),
requiredSlots: targetCount + 1,
}
}

func (s *liveScanEventStream) cancel() {
s.publication.cancel()
}

func (s *liveScanEventStream) publish(msg liveScanEventMsg) {
_ = s.publication.commit(func() error {
if s.closed {
return nil
}
// Progress publication reserves enough capacity that every target can
// emit one result or failure and the coordinator can emit completion.
s.events <- msg
return nil
})
}

func (s *liveScanEventStream) publishProgress(msg liveScanEventMsg) {
_ = s.publication.commit(func() error {
if s.closed || len(s.events) >= cap(s.events)-s.requiredSlots {
return nil
}
s.events <- msg
return nil
})
}

func (s *liveScanEventStream) close() {
s.publication.finish(func() {
if s.closed {
return
}
s.closed = true
close(s.events)
})
}

func startLiveScanCmd(path string, filesScanned, dirsScanned, bytesScanned *int64, currentPath *atomic.Value) tea.Cmd {
return startLiveScanCmdWithPolicy(path, filesScanned, dirsScanned, bytesScanned, currentPath, scanCacheReuse)
}

func startLiveScanCmdWithPolicy(path string, filesScanned, dirsScanned, bytesScanned *int64, currentPath *atomic.Value, cachePolicy scanCachePolicy) tea.Cmd {
return func() tea.Msg {
id := nextLiveScanID.Add(1)
ctx, cancel := context.WithCancel(context.Background())
ctx, cancelContext := context.WithCancel(context.Background())

limiter := newScanLimiter(0)
entries, targets, totalSize, totalFiles, largeFiles, err := readLiveScanInitialEntries(path, limiter)
if err != nil {
cancel()
cancelContext()
return liveScanStartMsg{id: id, path: path, err: err}
}

Expand All @@ -57,8 +111,9 @@ func startLiveScanCmdWithPolicy(path string, filesScanned, dirsScanned, bytesSca
atomic.AddInt64(bytesScanned, totalSize)
}

events := make(chan liveScanEventMsg, max(len(targets)*4, 1))
go runLiveScan(ctx, id, path, entries, targets, totalSize, totalFiles, largeFiles, limiter, filesScanned, dirsScanned, bytesScanned, currentPath, events, cachePolicy)
publication := newScanPublication(ctx, cancelContext)
stream := newLiveScanEventStream(publication, len(targets))
go runLiveScan(ctx, id, path, entries, targets, totalSize, totalFiles, largeFiles, limiter, filesScanned, dirsScanned, bytesScanned, currentPath, stream, cachePolicy)

scanningPaths := make([]string, 0, len(targets))
for _, target := range targets {
Expand All @@ -73,8 +128,8 @@ func startLiveScanCmdWithPolicy(path string, filesScanned, dirsScanned, bytesSca
totalFiles: totalFiles,
largeFiles: largeFiles,
scanningPaths: scanningPaths,
events: events,
cancel: cancel,
events: stream.events,
cancel: stream.cancel,
}
}
}
Expand Down Expand Up @@ -188,10 +243,10 @@ func runLiveScan(
limiter *scanLimiter,
filesScanned, dirsScanned, bytesScanned *int64,
currentPath *atomic.Value,
events chan<- liveScanEventMsg,
stream *liveScanEventStream,
cachePolicy scanCachePolicy,
) {
defer close(events)
defer stream.close()

entriesByPath := make(map[string]dirEntry, len(initialEntries))
for _, entry := range initialEntries {
Expand Down Expand Up @@ -219,9 +274,9 @@ func runLiveScan(
target := target
scanTarget := func() {
defer wg.Done()
result, err := scanLiveTargetWithProgress(ctx, id, root, target, largeFileChan, &largeFileMinSize, limiter, currentPath, events, cachePolicy)
result, err := scanLiveTargetWithProgress(ctx, id, root, target, largeFileChan, &largeFileMinSize, limiter, currentPath, stream, cachePolicy)
if err != nil && !errors.Is(err, context.Canceled) {
sendLiveScanEvent(ctx, events, liveScanEventMsg{id: id, path: root, kind: liveScanFailed, entry: dirEntry{Name: target.name, Path: target.path, IsDir: true}, err: err})
stream.publish(liveScanEventMsg{id: id, path: root, kind: liveScanFailed, entry: dirEntry{Name: target.name, Path: target.path, IsDir: true}, err: err})
return
}
if ctx.Err() != nil {
Expand Down Expand Up @@ -253,7 +308,7 @@ func runLiveScan(
atomic.AddInt64(bytesScanned, result.TotalSize)
}

sendLiveScanEvent(ctx, events, liveScanEventMsg{
stream.publish(liveScanEventMsg{
id: id,
path: root,
kind: liveScanChildDone,
Expand All @@ -278,7 +333,6 @@ func runLiveScan(
largeFiles := <-largeFilesDone

if ctx.Err() != nil {
sendLiveScanEvent(context.Background(), events, liveScanEventMsg{id: id, path: root, kind: liveScanCanceled, err: ctx.Err()})
return
}

Expand All @@ -301,10 +355,10 @@ func runLiveScan(
dedupedHardlink: dedupedHardlink.Load(),
}

sendLiveScanEvent(ctx, events, liveScanEventMsg{id: id, path: root, kind: liveScanComplete, result: result})
stream.publish(liveScanEventMsg{id: id, path: root, kind: liveScanComplete, result: result})
}

func scanLiveTargetWithProgress(ctx context.Context, id int64, root string, target liveScanTarget, largeFileChan chan<- fileEntry, largeFileMinSize *int64, limiter *scanLimiter, currentPath *atomic.Value, events chan<- liveScanEventMsg, cachePolicy scanCachePolicy) (scanResult, error) {
func scanLiveTargetWithProgress(ctx context.Context, id int64, root string, target liveScanTarget, largeFileChan chan<- fileEntry, largeFileMinSize *int64, limiter *scanLimiter, currentPath *atomic.Value, stream *liveScanEventStream, cachePolicy scanCachePolicy) (scanResult, error) {
var filesScanned int64
var dirsScanned int64
var bytesScanned int64
Expand Down Expand Up @@ -336,7 +390,7 @@ func scanLiveTargetWithProgress(ctx context.Context, id int64, root string, targ
currentPath.Store(path)
}
}
sendLiveScanProgress(ctx, events, liveScanEventMsg{
stream.publishProgress(liveScanEventMsg{
id: id,
path: root,
kind: liveScanChildProgress,
Expand All @@ -351,7 +405,7 @@ func scanLiveTargetWithProgress(ctx context.Context, id int64, root string, targ
}
}()

result, err := scanLiveTarget(ctx, target, largeFileChan, largeFileMinSize, limiter, &filesScanned, &dirsScanned, &bytesScanned, localCurrentPath, cachePolicy)
result, err := scanLiveTarget(ctx, target, largeFileChan, largeFileMinSize, limiter, &filesScanned, &dirsScanned, &bytesScanned, localCurrentPath, cachePolicy, stream.publication)
close(done)
<-progressDone
if result.TotalFiles == 0 {
Expand All @@ -363,7 +417,7 @@ func scanLiveTargetWithProgress(ctx context.Context, id int64, root string, targ
return result, err
}

func scanLiveTarget(ctx context.Context, target liveScanTarget, largeFileChan chan<- fileEntry, largeFileMinSize *int64, limiter *scanLimiter, filesScanned, dirsScanned, bytesScanned *int64, currentPath *atomic.Value, cachePolicy scanCachePolicy) (scanResult, error) {
func scanLiveTarget(ctx context.Context, target liveScanTarget, largeFileChan chan<- fileEntry, largeFileMinSize *int64, limiter *scanLimiter, filesScanned, dirsScanned, bytesScanned *int64, currentPath *atomic.Value, cachePolicy scanCachePolicy, publication *scanPublication) (scanResult, error) {
if err := ctx.Err(); err != nil {
return scanResult{}, err
}
Expand All @@ -376,20 +430,26 @@ func scanLiveTarget(ctx context.Context, target liveScanTarget, largeFileChan ch
}
}
case liveScanTargetFoldedDirectory:
size, err := getDirectorySizeFromDu(target.path)
size, err := getDirectorySizeFromDu(ctx, target.path)
if ctx.Err() != nil {
return scanResult{}, ctx.Err()
}
if err != nil || size <= 0 {
size = calculateDirSizeFastWithLimiter(target.path, limiter, filesScanned, dirsScanned, bytesScanned, currentPath)
size = calculateDirSizeFastWithLimiter(ctx, target.path, limiter, filesScanned, dirsScanned, bytesScanned, currentPath)
} else {
atomic.AddInt64(bytesScanned, size)
}
if ctx.Err() != nil {
return scanResult{}, ctx.Err()
}
return scanResult{TotalSize: size}, nil
}

if err := ctx.Err(); err != nil {
return scanResult{}, err
}

result := scanSubdirWithCache(target.path, largeFileChan, largeFileMinSize, limiter, limiter.dirSem, limiter.duSem, limiter.duQueueSem, filesScanned, dirsScanned, bytesScanned, currentPath, cachePolicy)
result := scanSubdirWithCache(ctx, target.path, largeFileChan, largeFileMinSize, limiter, limiter.dirSem, limiter.duSem, limiter.duQueueSem, filesScanned, dirsScanned, bytesScanned, currentPath, cachePolicy, publication)
return result, ctx.Err()
}

Expand Down Expand Up @@ -424,21 +484,6 @@ func pushLiveLargeFile(h *largeFileHeap, file fileEntry, largeFileMinSize *int64
}
}

func sendLiveScanEvent(ctx context.Context, events chan<- liveScanEventMsg, msg liveScanEventMsg) {
select {
case <-ctx.Done():
case events <- msg:
}
}

func sendLiveScanProgress(ctx context.Context, events chan<- liveScanEventMsg, msg liveScanEventMsg) {
select {
case <-ctx.Done():
case events <- msg:
default:
}
}

func waitLiveScanEventCmd(events <-chan liveScanEventMsg) tea.Cmd {
return func() tea.Msg {
msg, ok := <-events
Expand Down
1 change: 0 additions & 1 deletion cmd/analyze/model.go
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,6 @@ const (
liveScanChildDone
liveScanComplete
liveScanFailed
liveScanCanceled
)

type liveScanEventMsg struct {
Expand Down
Loading
Loading