Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
51 changes: 0 additions & 51 deletions internal/fsnotifyext/fsnotify_dedup.go

This file was deleted.

1 change: 1 addition & 0 deletions testdata/watch_newdir/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
src/*
10 changes: 10 additions & 0 deletions testdata/watch_newdir/Taskfile.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
# https://taskfile.dev

version: '3'

tasks:
default:
sources:
- "src/**/*"
cmds:
- echo "Task running!"
1 change: 1 addition & 0 deletions testdata/watch_remove/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
src/*
10 changes: 10 additions & 0 deletions testdata/watch_remove/Taskfile.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
# https://taskfile.dev

version: '3'

tasks:
default:
sources:
- "src/*"
cmds:
- echo "Task running!"
1 change: 1 addition & 0 deletions testdata/watch_sources/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
src/*
10 changes: 10 additions & 0 deletions testdata/watch_sources/Taskfile.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
# https://taskfile.dev

version: '3'

tasks:
default:
sources:
- "src/*.txt"
cmds:
- echo "Task running!"
184 changes: 133 additions & 51 deletions watch.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package task
import (
"context"
"fmt"
"io/fs"
"os"
"os/signal"
"path/filepath"
Expand All @@ -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"
Expand Down Expand Up @@ -56,65 +55,69 @@ 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())

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()
Expand All @@ -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)
Expand All @@ -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
}

Expand Down
Loading
Loading