diff --git a/go.mod b/go.mod index cd382b86c5..770ca883b3 100644 --- a/go.mod +++ b/go.mod @@ -14,11 +14,11 @@ require ( github.com/dominikbraun/graph v0.23.0 github.com/elliotchance/orderedmap/v3 v3.1.1 github.com/fatih/color v1.19.0 - github.com/fsnotify/fsnotify v1.10.1 github.com/go-task/slim-sprig/v3 v3.0.0 github.com/go-task/template v0.2.0 github.com/google/uuid v1.6.0 github.com/hashicorp/go-getter v1.8.9 + github.com/helshabini/fsbroker v1.0.3 github.com/joho/godotenv v1.5.1 github.com/mitchellh/hashstructure/v2 v2.0.2 github.com/puzpuzpuz/xsync/v4 v4.5.0 @@ -82,6 +82,7 @@ require ( github.com/envoyproxy/go-control-plane/envoy v1.39.0 // indirect github.com/envoyproxy/protoc-gen-validate v1.3.3 // indirect github.com/felixge/httpsnoop v1.1.0 // indirect + github.com/fsnotify/fsnotify v1.10.1 // indirect github.com/go-jose/go-jose/v4 v4.1.4 // indirect github.com/go-logr/logr v1.4.4 // indirect github.com/go-logr/stdr v1.2.2 // indirect diff --git a/go.sum b/go.sum index 5d8a7e998e..11a45a6617 100644 --- a/go.sum +++ b/go.sum @@ -173,6 +173,8 @@ github.com/hashicorp/go-getter v1.8.9 h1:1AOTMUmz/S/GuqOTE6+bg6nwPXTEbKr3Tc/T/lV github.com/hashicorp/go-getter v1.8.9/go.mod h1:qI7vH/m552bXutKe4UkWptOJl4bo2KeO/CSOoh4s0ps= github.com/hashicorp/go-version v1.9.0 h1:CeOIz6k+LoN3qX9Z0tyQrPtiB1DFYRPfCIBtaXPSCnA= github.com/hashicorp/go-version v1.9.0/go.mod h1:fltr4n8CU8Ke44wwGCBoEymUuxUHl09ZGVZPK5anwXA= +github.com/helshabini/fsbroker v1.0.3 h1:O+MhAhRIfIBvM4uVuuwyjU3LlhUsnymo94PsP0w+4CM= +github.com/helshabini/fsbroker v1.0.3/go.mod h1:ZL9RVJ3zrRJS7Fi5UtyJcpz16Otzpwzq74N6qBrnoZE= github.com/hexops/gotextdiff v1.0.3 h1:gitA9+qJrrTCsiCl7+kh75nPqQt1cx4ZkudSTLoUqJM= github.com/hexops/gotextdiff v1.0.3/go.mod h1:pSWU5MAI3yDq+fZBTazCSJysOMbxWL1BSow5/V2vxeg= github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0= diff --git a/internal/fsnotifyext/fsnotify_dedup.go b/internal/fsnotifyext/fsnotify_dedup.go deleted file mode 100644 index d081842380..0000000000 --- a/internal/fsnotifyext/fsnotify_dedup.go +++ /dev/null @@ -1,51 +0,0 @@ -package fsnotifyext - -import ( - "math" - "time" - - "github.com/fsnotify/fsnotify" -) - -type Deduper struct { - w *fsnotify.Watcher - waitTime time.Duration -} - -func NewDeduper(w *fsnotify.Watcher, waitTime time.Duration) *Deduper { - return &Deduper{ - w: w, - waitTime: waitTime, - } -} - -// GetChan returns a chan of deduplicated [fsnotify.Event]. -// -// [fsnotify.Chmod] operations will be skipped. -func (d *Deduper) GetChan() <-chan fsnotify.Event { - channel := make(chan fsnotify.Event) - - go func() { - timers := make(map[string]*time.Timer) - for { - event, ok := <-d.w.Events - switch { - case !ok: - return - case event.Has(fsnotify.Chmod): - continue - } - - timer, ok := timers[event.String()] - if !ok { - timer = time.AfterFunc(math.MaxInt64, func() { channel <- event }) - timer.Stop() - timers[event.String()] = timer - } - - timer.Reset(d.waitTime) - } - }() - - return channel -} diff --git a/testdata/watch_newdir/.gitignore b/testdata/watch_newdir/.gitignore new file mode 100644 index 0000000000..71bcfc6c90 --- /dev/null +++ b/testdata/watch_newdir/.gitignore @@ -0,0 +1 @@ +src/* diff --git a/testdata/watch_newdir/Taskfile.yaml b/testdata/watch_newdir/Taskfile.yaml new file mode 100644 index 0000000000..7ebfc55128 --- /dev/null +++ b/testdata/watch_newdir/Taskfile.yaml @@ -0,0 +1,10 @@ +# https://taskfile.dev + +version: '3' + +tasks: + default: + sources: + - "src/**/*" + cmds: + - echo "Task running!" diff --git a/testdata/watch_remove/.gitignore b/testdata/watch_remove/.gitignore new file mode 100644 index 0000000000..71bcfc6c90 --- /dev/null +++ b/testdata/watch_remove/.gitignore @@ -0,0 +1 @@ +src/* diff --git a/testdata/watch_remove/Taskfile.yaml b/testdata/watch_remove/Taskfile.yaml new file mode 100644 index 0000000000..813ff45d74 --- /dev/null +++ b/testdata/watch_remove/Taskfile.yaml @@ -0,0 +1,10 @@ +# https://taskfile.dev + +version: '3' + +tasks: + default: + sources: + - "src/*" + cmds: + - echo "Task running!" diff --git a/testdata/watch_sources/.gitignore b/testdata/watch_sources/.gitignore new file mode 100644 index 0000000000..71bcfc6c90 --- /dev/null +++ b/testdata/watch_sources/.gitignore @@ -0,0 +1 @@ +src/* diff --git a/testdata/watch_sources/Taskfile.yaml b/testdata/watch_sources/Taskfile.yaml new file mode 100644 index 0000000000..8398c38cbe --- /dev/null +++ b/testdata/watch_sources/Taskfile.yaml @@ -0,0 +1,10 @@ +# https://taskfile.dev + +version: '3' + +tasks: + default: + sources: + - "src/*.txt" + cmds: + - echo "Task running!" diff --git a/watch.go b/watch.go index 80c08258d4..4af0e1a206 100644 --- a/watch.go +++ b/watch.go @@ -3,6 +3,7 @@ package task import ( "context" "fmt" + "io/fs" "os" "os/signal" "path/filepath" @@ -11,13 +12,11 @@ import ( "syscall" "time" - "github.com/fsnotify/fsnotify" + "github.com/helshabini/fsbroker" "github.com/puzpuzpuz/xsync/v4" "github.com/go-task/task/v3/errors" - "github.com/go-task/task/v3/internal/filepathext" "github.com/go-task/task/v3/internal/fingerprint" - "github.com/go-task/task/v3/internal/fsnotifyext" "github.com/go-task/task/v3/internal/logger" "github.com/go-task/task/v3/internal/slicesext" "github.com/go-task/task/v3/taskfile/ast" @@ -56,27 +55,52 @@ func (e *Executor) watchTasks(calls ...*Call) error { waitTime = defaultWaitTime } - w, err := fsnotify.NewWatcher() + config := fsbroker.DefaultFSConfig() + config.Timeout = waitTime + config.IgnoreHiddenFiles = false + + broker, err := fsbroker.NewFSBroker(config) if err != nil { cancel() return err } - defer w.Close() + defer broker.Stop() + + closeOnInterrupt(broker) - deduper := fsnotifyext.NewDeduper(w, waitTime) - eventsChan := deduper.GetChan() + e.watchedDirs = xsync.NewMap[string, bool]() - closeOnInterrupt(w) + broker.Start() go func() { for { select { - case event, ok := <-eventsChan: + case action, ok := <-broker.Next(): if !ok { cancel() return } - e.Logger.VerboseErrf(logger.Magenta, "task: received watch event: %v\n", event) + + actions := e.filterWatchActions(broker, drainActions(broker, action)) + if len(actions) == 0 { + continue + } + + files, err := e.collectSources(calls) + if err != nil { + e.Logger.Errf(logger.Red, "%v\n", err) + continue + } + + if !slices.ContainsFunc(actions, func(a *fsbroker.FSAction) bool { + return isRelevantWatchAction(a, files) + }) { + for _, a := range actions { + relPath, _ := filepath.Rel(e.Dir, a.Subject.Path) + e.Logger.VerboseErrf(logger.Magenta, "task: skipped for file not in sources: %s\n", relPath) + } + continue + } cancel() ctx, cancel = context.WithCancel(context.Background()) @@ -84,37 +108,16 @@ func (e *Executor) watchTasks(calls ...*Call) error { e.Compiler.ResetCache() for _, c := range calls { - go func() { - if ShouldIgnore(event.Name) { - e.Logger.VerboseErrf(logger.Magenta, "task: event skipped for being an ignored dir: %s\n", event.Name) - return - } - t, err := e.GetTask(c) - if err != nil { - e.Logger.Errf(logger.Red, "%v\n", err) - return - } - baseDir := filepathext.SmartJoin(e.Dir, t.Dir) - files, err := e.collectSources(calls) - if err != nil { - e.Logger.Errf(logger.Red, "%v\n", err) - return - } - - if !event.Has(fsnotify.Remove) && !slices.Contains(files, filepath.ToSlash(event.Name)) { - relPath, _ := filepath.Rel(baseDir, event.Name) - e.Logger.VerboseErrf(logger.Magenta, "task: skipped for file not in sources: %s\n", relPath) - return - } - err = e.RunTask(ctx, c) + go func(ctx context.Context, c *Call) { + err := e.RunTask(ctx, c) if err == nil { e.Logger.Errf(logger.Green, "task: task \"%s\" finished running\n", c.Task) } else if !isContextError(err) { e.Logger.Errf(logger.Red, "%v\n", err) } - }() + }(ctx, c) } - case err, ok := <-w.Errors: + case err, ok := <-broker.Error(): switch { case !ok: cancel() @@ -126,14 +129,12 @@ func (e *Executor) watchTasks(calls ...*Call) error { } }() - e.watchedDirs = xsync.NewMap[string, bool]() - go func() { // NOTE(@andreynering): New files can be created in directories // that were previously empty, so we need to check for new dirs // from time to time. for { - if err := e.registerWatchedDirs(w, calls...); err != nil { + if err := e.registerWatchedDirs(broker, calls...); err != nil { e.Logger.Errf(logger.Red, "%v\n", err) } time.Sleep(5 * time.Second) @@ -152,36 +153,117 @@ func isContextError(err error) bool { return errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) } -func closeOnInterrupt(w *fsnotify.Watcher) { +func closeOnInterrupt(broker *fsbroker.FSBroker) { ch := make(chan os.Signal, 1) signal.Notify(ch, os.Interrupt, syscall.SIGTERM) go func() { <-ch - w.Close() + broker.Stop() os.Exit(0) }() } -func (e *Executor) registerWatchedDirs(w *fsnotify.Watcher, calls ...*Call) error { +// drainActions returns the given action along with any other action that is +// already available, so that a burst of changes (e.g. a git checkout) results +// in a single task run instead of one per changed file. +func drainActions(broker *fsbroker.FSBroker, action *fsbroker.FSAction) []*fsbroker.FSAction { + actions := []*fsbroker.FSAction{action} + for { + select { + case next, ok := <-broker.Next(): + if !ok { + return actions + } + actions = append(actions, next) + default: + return actions + } + } +} + +// filterWatchActions applies the side effects of the given actions and returns +// the ones that are candidates for triggering a task run. +func (e *Executor) filterWatchActions(broker *fsbroker.FSBroker, actions []*fsbroker.FSAction) []*fsbroker.FSAction { + filtered := make([]*fsbroker.FSAction, 0, len(actions)) + for _, action := range actions { + if action.Subject == nil || action.Type == fsbroker.NoOp || action.Type == fsbroker.Chmod { + continue + } + + e.Logger.VerboseErrf(logger.Magenta, "task: received watch event: %s: %s\n", action.Type, action.Subject.Path) + + if action.Type == fsbroker.Remove || action.Type == fsbroker.Rename { + e.watchedDirs.Delete(action.Subject.Path) + if oldPath, ok := action.Properties["OldPath"].(string); ok { + e.watchedDirs.Delete(oldPath) + } + } + + if ShouldIgnore(action.Subject.Path) { + e.Logger.VerboseErrf(logger.Magenta, "task: event skipped for being an ignored dir: %s\n", action.Subject.Path) + continue + } + + if action.Subject.IsDir() && action.Type != fsbroker.Remove { + if err := e.registerWatchedTree(broker, action.Subject.Path); err != nil { + e.Logger.VerboseErrf(logger.Magenta, "task: failed to watch dir %s: %v\n", action.Subject.Path, err) + } + } + + filtered = append(filtered, action) + } + return filtered +} + +func isRelevantWatchAction(action *fsbroker.FSAction, files []string) bool { + if action.Subject.IsDir() || action.Type == fsbroker.Remove || action.Type == fsbroker.Rename { + return true + } + return slices.Contains(files, filepath.ToSlash(action.Subject.Path)) +} + +func (e *Executor) registerWatchedDirs(broker *fsbroker.FSBroker, calls ...*Call) error { files, err := e.collectSources(calls) if err != nil { return err } for _, f := range files { - d := filepath.Dir(f) - if isSet, ok := e.watchedDirs.Load(d); ok && isSet { - continue + if err := e.registerWatchedDir(broker, filepath.Dir(f)); err != nil { + return err } - if ShouldIgnore(d) { - continue + } + return nil +} + +func (e *Executor) registerWatchedTree(broker *fsbroker.FSBroker, dir string) error { + return filepath.WalkDir(dir, func(path string, d fs.DirEntry, err error) error { + if err != nil { + e.Logger.VerboseErrf(logger.Magenta, "task: failed to walk dir %s: %v\n", path, err) + return filepath.SkipDir } - if err := w.Add(d); err != nil { - return err + if !d.IsDir() { + return nil } - e.watchedDirs.Store(d, true) - relPath, _ := filepath.Rel(e.Dir, d) - e.Logger.VerboseOutf(logger.Green, "task: watching new dir: %v\n", relPath) + if ShouldIgnore(path) { + return filepath.SkipDir + } + return e.registerWatchedDir(broker, path) + }) +} + +func (e *Executor) registerWatchedDir(broker *fsbroker.FSBroker, d string) error { + if isSet, ok := e.watchedDirs.Load(d); ok && isSet { + return nil + } + if ShouldIgnore(d) { + return nil + } + if err := broker.AddWatch(d); err != nil { + return err } + e.watchedDirs.Store(d, true) + relPath, _ := filepath.Rel(e.Dir, d) + e.Logger.VerboseOutf(logger.Green, "task: watching new dir: %v\n", relPath) return nil } diff --git a/watch_test.go b/watch_test.go index 6d3c89632e..43cf43f5d8 100644 --- a/watch_test.go +++ b/watch_test.go @@ -8,6 +8,7 @@ import ( "fmt" "os" "strings" + "sync" "testing" "time" @@ -18,6 +19,23 @@ import ( "github.com/go-task/task/v3/internal/filepathext" ) +type syncBuffer struct { + mu sync.Mutex + buf bytes.Buffer +} + +func (b *syncBuffer) Write(p []byte) (int, error) { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.Write(p) +} + +func (b *syncBuffer) String() string { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.String() +} + func TestFileWatch(t *testing.T) { t.Parallel() @@ -80,6 +98,141 @@ task: task "default" finished running assert.Equal(t, expectedOutput, strings.TrimSpace(buff.String())) } +// TestFileWatchNewDir checks that a directory created while watching is +// watched without having to wait for the periodic sources rescan. +func TestFileWatchNewDir(t *testing.T) { + t.Parallel() + + const dir = "testdata/watch_newdir" + _ = os.RemoveAll(filepathext.SmartJoin(dir, ".task")) + _ = os.RemoveAll(filepathext.SmartJoin(dir, "src")) + + var buff syncBuffer + e := task.NewExecutor( + task.WithDir(dir), + task.WithStdout(&buff), + task.WithStderr(&buff), + task.WithWatch(true), + ) + require.NoError(t, e.Setup()) + + dirPath := filepathext.SmartJoin(dir, "src") + require.NoError(t, os.MkdirAll(dirPath, 0o755)) + require.NoError(t, os.WriteFile(filepathext.SmartJoin(dirPath, "a"), []byte("test"), 0o644)) + + ctx, cancel := context.WithCancel(context.Background()) + go func() { + if err := e.Run(ctx, &task.Call{Task: "default"}); err != nil { + panic(err) + } + }() + + time.Sleep(300 * time.Millisecond) + + newDirPath := filepathext.SmartJoin(dirPath, "sub") + require.NoError(t, os.MkdirAll(newDirPath, 0o755)) + require.NoError(t, os.WriteFile(filepathext.SmartJoin(newDirPath, "b"), []byte("test"), 0o644)) + + require.Eventually(t, func() bool { + return strings.Count(buff.String(), `"default" finished running`) >= 2 + }, 2*time.Second, 50*time.Millisecond) + + cancel() +} + +// TestFileWatchSources checks that changes to files that are not task sources +// do not trigger a run. +func TestFileWatchSources(t *testing.T) { + t.Parallel() + + const dir = "testdata/watch_sources" + _ = os.RemoveAll(filepathext.SmartJoin(dir, ".task")) + _ = os.RemoveAll(filepathext.SmartJoin(dir, "src")) + + expectedOutput := strings.TrimSpace(` +task: Started watching for tasks: default +task: [default] echo "Task running!" +Task running! +task: task "default" finished running + `) + + var buff syncBuffer + e := task.NewExecutor( + task.WithDir(dir), + task.WithStdout(&buff), + task.WithStderr(&buff), + task.WithWatch(true), + ) + require.NoError(t, e.Setup()) + + dirPath := filepathext.SmartJoin(dir, "src") + require.NoError(t, os.MkdirAll(dirPath, 0o755)) + require.NoError(t, os.WriteFile(filepathext.SmartJoin(dirPath, "a.txt"), []byte("test"), 0o644)) + + ctx, cancel := context.WithCancel(context.Background()) + go func() { + if err := e.Run(ctx, &task.Call{Task: "default"}); err != nil { + panic(err) + } + }() + + time.Sleep(300 * time.Millisecond) + + require.NoError(t, os.WriteFile(filepathext.SmartJoin(dirPath, "notes.md"), []byte("other"), 0o644)) + + time.Sleep(500 * time.Millisecond) + cancel() + assert.Equal(t, expectedOutput, strings.TrimSpace(buff.String())) +} + +// TestFileWatchRemove checks that removing a source file triggers a run. +func TestFileWatchRemove(t *testing.T) { + t.Parallel() + + const dir = "testdata/watch_remove" + _ = os.RemoveAll(filepathext.SmartJoin(dir, ".task")) + _ = os.RemoveAll(filepathext.SmartJoin(dir, "src")) + + expectedOutput := strings.TrimSpace(` +task: Started watching for tasks: default +task: [default] echo "Task running!" +Task running! +task: task "default" finished running +task: [default] echo "Task running!" +Task running! +task: task "default" finished running + `) + + var buff syncBuffer + e := task.NewExecutor( + task.WithDir(dir), + task.WithStdout(&buff), + task.WithStderr(&buff), + task.WithWatch(true), + ) + require.NoError(t, e.Setup()) + + dirPath := filepathext.SmartJoin(dir, "src") + filePath := filepathext.SmartJoin(dirPath, "a") + require.NoError(t, os.MkdirAll(dirPath, 0o755)) + require.NoError(t, os.WriteFile(filePath, []byte("test"), 0o644)) + + ctx, cancel := context.WithCancel(context.Background()) + go func() { + if err := e.Run(ctx, &task.Call{Task: "default"}); err != nil { + panic(err) + } + }() + + time.Sleep(300 * time.Millisecond) + + require.NoError(t, os.Remove(filePath)) + + time.Sleep(500 * time.Millisecond) + cancel() + assert.Equal(t, expectedOutput, strings.TrimSpace(buff.String())) +} + func TestShouldIgnore(t *testing.T) { t.Parallel()