//go:build windows package main import ( "context" "flag" "fmt" "net" "os" "os/signal" "path/filepath" "sync" "syscall" "time" "github.com/fsnotify/fsnotify" ) var ( sourceDir = flag.String("dir", "", "Source directory to watch (required)") connect = flag.String("connect", "localhost:5151", "TCP address of the container to connect to") bufSizeKB = flag.Int("bufsize", 256, "ReadDirectoryChangesW buffer size in KB per directory (default 256)") ) func run() { flag.Parse() if *sourceDir == "" { fmt.Fprintln(os.Stderr, "Usage: sourcewatch.exe -dir [-connect host:port] [-bufsize KB]") os.Exit(1) } level, err := ParseLogLevel(*logLevel) if err != nil { fmt.Fprintln(os.Stderr, err) os.Exit(1) } SetLogLevel(level) logInfo("watcher: scanning %s ...", *sourceDir) stack := NewGitignoreStack(*sourceDir) manifest := scanDirectory(*sourceDir, stack) logInfo("watcher: found %d files", len(manifest.Files)) watcher, err := fsnotify.NewWatcher() if err != nil { logErrorf("watcher: %v", err) os.Exit(1) } defer watcher.Close() if err := addWatchersRecursively(watcher, *sourceDir, stack); err != nil { logErrorf("watcher: %v", err) os.Exit(1) } ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer cancel() for { if ctx.Err() != nil { return } logInfo("watcher: connecting to %s ...", *connect) var conn net.Conn for { if ctx.Err() != nil { return } conn, err = net.DialTimeout("tcp", *connect, 5*time.Second) if err == nil { break } logDebug("connect failed: %v, retrying...", err) select { case <-ctx.Done(): return case <-time.After(2 * time.Second): } } logInfo("watcher: connected, sending manifest ...") if err := sendManifest(conn, manifest); err != nil { logErrorf("send manifest: %v", err) conn.Close() continue } logInfo("watcher: monitoring for changes ...") // Heartbeat goroutine to detect disconnects go func() { ticker := time.NewTicker(10 * time.Second) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: msg := &ProtocolMessage{Type: "ping"} data, _ := EncodeMessage(msg) data = append(data, '\n') if _, err := conn.Write(data); err != nil { logErrorf("heartbeat failed: %v", err) conn.Close() return } } } }() if !watchLoop(ctx, conn, watcher, manifest, stack) { conn.Close() continue } conn.Close() return } } func watchLoop(ctx context.Context, conn net.Conn, watcher *fsnotify.Watcher, manifest *Manifest, stack *GitignoreStack) bool { var mu sync.Mutex debounceTimers := make(map[string]*time.Timer) debounceOps := make(map[string]fsnotify.Op) for { select { case <-ctx.Done(): return true case event, ok := <-watcher.Events: if !ok { return true } relPath, err := filepath.Rel(*sourceDir, event.Name) if err != nil { continue } if isGitDir(relPath) { continue } if stack.IsIgnored(relPath) { continue } logDebug("event: %s %s", event.Op, relPath) if event.Op&fsnotify.Create != 0 { srcPath := filepath.Join(*sourceDir, relPath) info, err := os.Lstat(srcPath) if err == nil && info.IsDir() { bufBytes := *bufSizeKB * 1024 if hasGitignore(srcPath) { stack.Push(srcPath) } if err := watcher.AddWith(srcPath, fsnotify.WithBufferSize(bufBytes)); err != nil { logErrorf("watch %s: %v", relPath, err) } var added int filepath.Walk(srcPath, func(path string, info os.FileInfo, err error) error { if err != nil || !info.IsDir() || path == srcPath { return nil } if err := watcher.AddWith(path, fsnotify.WithBufferSize(bufBytes)); err != nil { logDebug("watch sub %s: %v", path, err) } else { added++ } return nil }) if added > 0 { logDebug("added %d sub-watchers for %s", added, relPath) } } } mu.Lock() if existing, exists := debounceOps[relPath]; exists { debounceTimers[relPath].Stop() debounceOps[relPath] = existing | event.Op } else { debounceOps[relPath] = event.Op } path := relPath debounceTimers[relPath] = time.AfterFunc(200*time.Millisecond, func() { mu.Lock() op := debounceOps[path] delete(debounceOps, path) delete(debounceTimers, path) mu.Unlock() evt := FileEvent{Path: path} switch { case op&fsnotify.Create != 0: evt.Op = "create" case op&fsnotify.Write != 0: evt.Op = "write" case op&fsnotify.Remove != 0: evt.Op = "remove" case op&fsnotify.Rename != 0: evt.Op = "remove" default: return } srcPath := filepath.Join(*sourceDir, path) switch evt.Op { case "create": if info, err := os.Lstat(srcPath); err == nil { linkTarget := "" if info.Mode()&os.ModeSymlink != 0 { linkTarget, _ = os.Readlink(srcPath) } manifest.Files[path] = FileMeta{ Size: info.Size(), Mode: info.Mode(), ModTime: info.ModTime(), Symlink: linkTarget, } } case "write": if info, err := os.Lstat(srcPath); err == nil { linkTarget := "" if info.Mode()&os.ModeSymlink != 0 { linkTarget, _ = os.Readlink(srcPath) } manifest.Files[path] = FileMeta{ Size: info.Size(), Mode: info.Mode(), ModTime: info.ModTime(), Symlink: linkTarget, } } case "remove": delete(manifest.Files, path) } if err := sendEvent(conn, evt); err != nil { logErrorf("send event: %v", err) } }) mu.Unlock() case err, ok := <-watcher.Errors: if !ok { return true } logErrorf("watcher: %v", err) return false } } } func scanDirectory(root string, stack *GitignoreStack) *Manifest { manifest := &Manifest{Files: make(map[string]FileMeta)} filepath.Walk(root, func(path string, info os.FileInfo, err error) error { if err != nil { return nil } relPath, err := filepath.Rel(root, path) if err != nil { return nil } if relPath == "." { return nil } if isGitDir(relPath) { if info.IsDir() { return filepath.SkipDir } return nil } if info.IsDir() && hasGitignore(path) { stack.Push(path) } if stack.IsIgnored(relPath) { return nil } linkTarget := "" if info.Mode()&os.ModeSymlink != 0 { linkTarget, _ = os.Readlink(path) } manifest.Files[relPath] = FileMeta{ Size: info.Size(), Mode: info.Mode(), ModTime: info.ModTime(), Symlink: linkTarget, } return nil }) return manifest } func sendManifest(conn net.Conn, manifest *Manifest) error { const batchSize = 1000 files := make([]FileMetaJSON, 0, batchSize) paths := make([]string, 0, batchSize) for path, meta := range manifest.ToJSON() { paths = append(paths, path) files = append(files, meta) if len(files) >= batchSize { msg := &ProtocolMessage{ Type: "manifest_batch", Files: make(map[string]FileMetaJSON), } for i, p := range paths { msg.Files[p] = files[i] } data, err := EncodeMessage(msg) if err != nil { return err } data = append(data, '\n') if _, err := conn.Write(data); err != nil { return err } files = files[:0] paths = paths[:0] } } if len(files) > 0 { msg := &ProtocolMessage{ Type: "manifest_batch", Files: make(map[string]FileMetaJSON), } for i, p := range paths { msg.Files[p] = files[i] } data, err := EncodeMessage(msg) if err != nil { return err } data = append(data, '\n') if _, err := conn.Write(data); err != nil { return err } } done := &ProtocolMessage{Type: "manifest_done"} doneData, err := EncodeMessage(done) if err != nil { return err } doneData = append(doneData, '\n') _, err = conn.Write(doneData) return err } func sendEvent(conn net.Conn, evt FileEvent) error { msg := &ProtocolMessage{Path: evt.Path} switch evt.Op { case "create": msg.Type = "event_create" case "write": msg.Type = "event_write" case "remove": msg.Type = "event_remove" } data, err := EncodeMessage(msg) if err != nil { return err } data = append(data, '\n') _, err = conn.Write(data) return err } func addWatchersRecursively(watcher *fsnotify.Watcher, root string, stack *GitignoreStack) error { var count, failCount, skippedCount int bufBytes := *bufSizeKB * 1024 err := filepath.Walk(root, func(path string, info os.FileInfo, err error) error { if err != nil { logDebug("walk error: %v", err) return nil // continue walking even if we can't stat a path } relPath, err := filepath.Rel(root, path) if err != nil { logDebug("rel error: %v", err) return nil } if isGitDir(relPath) { if info.IsDir() { return filepath.SkipDir } return nil } if info.IsDir() { if hasGitignore(path) { stack.Push(path) } if !stack.IsIgnored(relPath) { if err := watcher.AddWith(path, fsnotify.WithBufferSize(bufBytes)); err != nil { failCount++ if failCount <= 10 { logErrorf("watch %s: %v", relPath, err) } } else { count++ } } else { skippedCount++ } } return nil }) if failCount > 10 { logErrorf("watcher: %d additional watch failures (suppressed)", failCount-10) } logInfo("watcher: watching %d directories, skipped %d ignored, %d failed (buffer=%dKB)", count, skippedCount, failCount, *bufSizeKB) return err }