From f8f2f5a1cc38a9e6f50e8d1b46562b6d84c58b4a Mon Sep 17 00:00:00 2001 From: isink17 <39876158+isink17@users.noreply.github.com> Date: Fri, 24 Apr 2026 17:55:24 +0200 Subject: [PATCH 1/2] watch: apply repo include/exclude and harden dirty_files draining --- internal/indexer/indexer.go | 8 ++++++ internal/store/store.go | 12 ++++++++ internal/store/store_test.go | 51 +++++++++++++++++++++++++++++++++ internal/watcher/watcher.go | 55 +++++++++++++++++++++++------------- 4 files changed, 106 insertions(+), 20 deletions(-) diff --git a/internal/indexer/indexer.go b/internal/indexer/indexer.go index d8626c6..468ad99 100644 --- a/internal/indexer/indexer.go +++ b/internal/indexer/indexer.go @@ -749,6 +749,14 @@ func ShouldSkipDir(rel string, excludes []string) bool { return shouldSkipDir(rel, excludes) } +func ShouldSkipFile(rel string, includes, excludes []string) bool { + return shouldSkipFile(rel, includes, excludes) +} + +func ShouldIgnorePath(rel string, excludes []string) bool { + return shouldIgnorePath(rel, excludes) +} + func shouldSkipDir(rel string, excludes []string) bool { rel = filepath.ToSlash(rel) base := filepath.Base(rel) diff --git a/internal/store/store.go b/internal/store/store.go index 8d250e1..b7e915b 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -3184,6 +3184,18 @@ func (s *Store) QueueDirtyFile(ctx context.Context, repoID int64, path, reason s return err } +func (s *Store) HasDirtyFiles(ctx context.Context, repoID int64) (bool, error) { + var exists int + err := s.db.QueryRowContext(ctx, `SELECT 1 FROM dirty_files WHERE repo_id = ? LIMIT 1`, repoID).Scan(&exists) + if err == sql.ErrNoRows { + return false, nil + } + if err != nil { + return false, err + } + return true, nil +} + func (s *Store) DrainDirtyFiles(ctx context.Context, repoID int64) ([]string, error) { // Prefer an atomic drain. `DELETE ... RETURNING` guarantees we only remove rows // that are returned to the caller (no SELECT+DELETE race). diff --git a/internal/store/store_test.go b/internal/store/store_test.go index 35b2791..179f8d0 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -12,6 +12,57 @@ import ( "github.com/isink17/codegraph/internal/store" ) +func TestDirtyFilesQueueAndDrain(t *testing.T) { + ctx := context.Background() + dbPath := filepath.Join(t.TempDir(), "graph.sqlite") + s, err := store.Open(dbPath) + if err != nil { + t.Fatalf("store.Open() error = %v", err) + } + defer s.Close() + + repo, err := s.UpsertRepo(ctx, t.TempDir()) + if err != nil { + t.Fatalf("UpsertRepo() error = %v", err) + } + + if ok, err := s.HasDirtyFiles(ctx, repo.ID); err != nil { + t.Fatalf("HasDirtyFiles() error = %v", err) + } else if ok { + t.Fatalf("expected no dirty files at start") + } + + if err := s.QueueDirtyFile(ctx, repo.ID, "a.go", "test"); err != nil { + t.Fatalf("QueueDirtyFile(a.go) error = %v", err) + } + if err := s.QueueDirtyFile(ctx, repo.ID, "b.go", "test"); err != nil { + t.Fatalf("QueueDirtyFile(b.go) error = %v", err) + } + if err := s.QueueDirtyFile(ctx, repo.ID, "a.go", "test2"); err != nil { + t.Fatalf("QueueDirtyFile(a.go update) error = %v", err) + } + + if ok, err := s.HasDirtyFiles(ctx, repo.ID); err != nil { + t.Fatalf("HasDirtyFiles() error = %v", err) + } else if !ok { + t.Fatalf("expected dirty files after queueing") + } + + paths, err := s.DrainDirtyFiles(ctx, repo.ID) + if err != nil { + t.Fatalf("DrainDirtyFiles() error = %v", err) + } + if len(paths) != 2 { + t.Fatalf("expected 2 paths, got %d (%v)", len(paths), paths) + } + + if ok, err := s.HasDirtyFiles(ctx, repo.ID); err != nil { + t.Fatalf("HasDirtyFiles() error = %v", err) + } else if ok { + t.Fatalf("expected no dirty files after drain") + } +} + func TestListScansIncludesLanguageCoverage(t *testing.T) { ctx := context.Background() repoRoot := t.TempDir() diff --git a/internal/watcher/watcher.go b/internal/watcher/watcher.go index 1cbf7de..0494d86 100644 --- a/internal/watcher/watcher.go +++ b/internal/watcher/watcher.go @@ -12,6 +12,7 @@ import ( "github.com/fsnotify/fsnotify" + "github.com/isink17/codegraph/internal/config" "github.com/isink17/codegraph/internal/indexer" "github.com/isink17/codegraph/internal/store" ) @@ -112,6 +113,13 @@ func (w *Watcher) Run(ctx context.Context, repoRoot string, repoID int64, deboun debounce = 750 * time.Millisecond } + repoCfg, err := config.LoadRepo(repoRoot) + if err != nil { + return err + } + includes := repoCfg.Include + excludes := repoCfg.Exclude + fsw, err := fsnotify.NewWatcher() if err != nil { return err @@ -129,7 +137,7 @@ func (w *Watcher) Run(ctx context.Context, repoRoot string, repoID int64, deboun return relErr } rel = filepath.Clean(rel) - if d.IsDir() && indexer.ShouldSkipDir(rel, nil) { + if d.IsDir() && indexer.ShouldSkipDir(rel, excludes) { return filepath.SkipDir } } @@ -160,11 +168,29 @@ func (w *Watcher) Run(ctx context.Context, repoRoot string, repoID int64, deboun Paths: paths, ScanKind: "watch", }) - if err == nil { - w.updateRuns.Add(1) - w.updatePaths.Add(int64(len(paths))) + if err != nil { + // `DrainDirtyFiles` is destructive; re-queue the drained paths on failure so + // work isn't silently dropped. + for _, path := range paths { + _ = w.store.QueueDirtyFile(ctx, repoID, path, "watch_retry") + } + return err } + w.updateRuns.Add(1) + w.updatePaths.Add(int64(len(paths))) + return err + } + + if hasDirty, err := w.store.HasDirtyFiles(ctx, repoID); err != nil { return err + } else if hasDirty { + // Ensure any queued work from previous runs is processed even if no new + // fsnotify events occur. + w.flushRuns.Add(1) + if err := flush(); err != nil { + w.flushErrors.Add(1) + return err + } } flushSignalCh := make(chan struct{}, 1) @@ -242,7 +268,7 @@ func (w *Watcher) Run(ctx context.Context, repoRoot string, repoID int64, deboun continue } rel = filepath.Clean(rel) - if shouldIgnorePath(rel) { + if indexer.ShouldIgnorePath(rel, excludes) { w.eventsIgnored.Add(1) continue } @@ -257,6 +283,10 @@ func (w *Watcher) Run(ctx context.Context, repoRoot string, repoID int64, deboun continue } } + if indexer.ShouldSkipFile(rel, includes, excludes) { + w.eventsIgnored.Add(1) + continue + } seenMu.Lock() _, alreadyQueued := seenSinceFlush[rel] @@ -307,18 +337,3 @@ func (w *Watcher) queueDirtyWithRetry(ctx context.Context, repoID int64, path, r } return fmt.Errorf("queue dirty file %s: %w", path, lastErr) } - -func shouldIgnorePath(rel string) bool { - current := rel - for current != "." && current != "" { - if indexer.ShouldSkipDir(current, nil) { - return true - } - next := filepath.Dir(current) - if next == current { - break - } - current = next - } - return false -} From 1bc492852146364e6fd25f6bda34203cf3dbf480 Mon Sep 17 00:00:00 2001 From: isink17 <39876158+isink17@users.noreply.github.com> Date: Fri, 24 Apr 2026 18:19:22 +0200 Subject: [PATCH 2/2] watch: batch re-queue dirty_files on update failure --- internal/store/store.go | 38 ++++++++++++++++++++++++++++++++++++ internal/store/store_test.go | 12 ++++++++++++ internal/watcher/watcher.go | 4 ++-- 3 files changed, 52 insertions(+), 2 deletions(-) diff --git a/internal/store/store.go b/internal/store/store.go index b7e915b..40f1ea5 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -3184,6 +3184,44 @@ func (s *Store) QueueDirtyFile(ctx context.Context, repoID int64, path, reason s return err } +func (s *Store) QueueDirtyFiles(ctx context.Context, repoID int64, paths []string, reason string) error { + if len(paths) == 0 { + return nil + } + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return err + } + committed := false + defer func() { + if committed { + return + } + _ = tx.Rollback() + }() + + stmt, err := tx.PrepareContext(ctx, ` + INSERT INTO dirty_files(repo_id, path, reason, queued_at) + VALUES(?, ?, ?, ?) + ON CONFLICT(repo_id, path) DO UPDATE SET reason=excluded.reason, queued_at=excluded.queued_at + `) + if err != nil { + return err + } + defer stmt.Close() + + for _, path := range paths { + if _, err := stmt.ExecContext(ctx, repoID, path, reason, time.Now().UTC().Format(time.RFC3339)); err != nil { + return err + } + } + if err := tx.Commit(); err != nil { + return err + } + committed = true + return nil +} + func (s *Store) HasDirtyFiles(ctx context.Context, repoID int64) (bool, error) { var exists int err := s.db.QueryRowContext(ctx, `SELECT 1 FROM dirty_files WHERE repo_id = ? LIMIT 1`, repoID).Scan(&exists) diff --git a/internal/store/store_test.go b/internal/store/store_test.go index 179f8d0..4310377 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -61,6 +61,18 @@ func TestDirtyFilesQueueAndDrain(t *testing.T) { } else if ok { t.Fatalf("expected no dirty files after drain") } + + if err := s.QueueDirtyFiles(ctx, repo.ID, []string{"c.go", "d.go"}, "batch"); err != nil { + t.Fatalf("QueueDirtyFiles() error = %v", err) + } + + paths, err = s.DrainDirtyFiles(ctx, repo.ID) + if err != nil { + t.Fatalf("DrainDirtyFiles() (2) error = %v", err) + } + if len(paths) != 2 { + t.Fatalf("expected 2 paths from batch, got %d (%v)", len(paths), paths) + } } func TestListScansIncludesLanguageCoverage(t *testing.T) { diff --git a/internal/watcher/watcher.go b/internal/watcher/watcher.go index 0494d86..fa739f1 100644 --- a/internal/watcher/watcher.go +++ b/internal/watcher/watcher.go @@ -171,8 +171,8 @@ func (w *Watcher) Run(ctx context.Context, repoRoot string, repoID int64, deboun if err != nil { // `DrainDirtyFiles` is destructive; re-queue the drained paths on failure so // work isn't silently dropped. - for _, path := range paths { - _ = w.store.QueueDirtyFile(ctx, repoID, path, "watch_retry") + if requeueErr := w.store.QueueDirtyFiles(ctx, repoID, paths, "watch_retry"); requeueErr != nil { + return fmt.Errorf("update failed: %w (retry re-queue failed: %v)", err, requeueErr) } return err }