From df885c76d858b38ef254d8642cd76c107ee56e9b Mon Sep 17 00:00:00 2001 From: Jesse Herrick Date: Sat, 3 Oct 2026 16:34:27 -0400 Subject: [PATCH 1/2] Make large warm reconciles and prunes fast, and cancel them on shutdown A warm start that found a large change set (for example, nested checkouts of the whole tree) wrote one transaction per file. Each commit wrote every index page that the file touched. A large prune ran three DELETE statements per path. Both took minutes on a large monorepo. During that time, stop and SIGTERM could not end the daemon, and it kept the workspace lock without serving its socket. Changes: - The sweep reads all stored mtimes with one query, not one query for each file. The walk is the same: one lstat for each file. - Changed files are parsed on all cores. A small change set is written in batches of many files for each transaction, with multi-row INSERTs. When a batch fails, its files are written again one at a time, so only a file that fails by itself is lost. - A change set of at least 4096 files and a quarter of the index uses Store.BeginRebuild. One transaction writes new, unindexed copies of definitions and refs, copies in the rows of all other files, swaps the copies in, and builds each index with one sort. Readers keep their snapshot until the commit. The rebuild takes the write lock with its first statement, so it waits for busy_timeout like other writers. If the rebuild fails, the change set is written in batches. - Removal deletes by sets of file ids. When the removed files are a quarter of the index, it rebuilds the tables without them, and falls back to the deletes if the rebuild fails. - IndexCoordinator.CancelWork stops the sweep, the prune, every removal, the rebuild, and a cold build (before and during its walk and parse). Runtime.Close calls it first. A file and its rows are always in one transaction, so a canceled pass writes no partial file. The next start finishes the pass from the stored mtimes. - In an embedded server, single-file writes from saves and watched-file events take the coordinator's mutation lock, as daemon events do. They wait behind a pass instead of racing its batches for the write lock. Measured on a synthetic project: 56k Elixir files from the Hex cache, plus four nested copies (223k files). - Warm add of 223k files: 3m48s -> 44s to 53s - Prune of 223k files: 66s -> 23s to 24s - Prune of 56k files (half of the index): 16.5s -> 16s - Daemon exit after stop or SIGTERM during these passes: 14s to 225s -> 0.07s to 0.9s - Warm start, no changes: 0.87s -> 0.69s; five changed files: 1.04s -> 0.79s Co-Authored-By: Claude Opus 5.5 --- docs/architecture.md | 7 +- internal/indexer/indexer.go | 49 ++- internal/lsp/reconcile.go | 426 +++++++++++++++++++++ internal/lsp/reconcile_test.go | 358 ++++++++++++++++++ internal/lsp/server.go | 120 +++--- internal/store/bulk_test.go | 402 ++++++++++++++++++++ internal/store/store.go | 576 +++++++++++++++++++++-------- internal/workspace/runtime.go | 18 +- internal/workspace/runtime_test.go | 133 +++++++ 9 files changed, 1868 insertions(+), 221 deletions(-) create mode 100644 internal/lsp/reconcile.go create mode 100644 internal/lsp/reconcile_test.go create mode 100644 internal/store/bulk_test.go diff --git a/docs/architecture.md b/docs/architecture.md index 780b385..aed361b 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -212,7 +212,7 @@ The cold index is bound by the single SQLite writer, not by parsing. On a 70k-fi Consequences that are easy to undo by accident: - **Refs are deduplicated in the parser** (`dedupeRefs`), not in the store. Refs are line-granular, so `@spec f(String.t(), String.t())` produces identical rows; ~60k of them on a large monorepo. Identical rows cannot change a result — no query counts refs, and the References handler dedupes by file+line — so the parse workers drop them before they reach the writer. -- **The bulk path batches inserts** into multi-row `INSERT`s (`multiRowInsert`, 900 bound parameters per statement, which is under even the legacy `SQLITE_MAX_VARIABLE_NUMBER` of 999). Incremental reindex keeps the row-at-a-time path, where a file's `DELETE` must stay ordered ahead of its `INSERT`s. +- **Every batch buffers inserts** into multi-row `INSERT`s (`multiRowInsert`, 900 bound parameters per statement, which is under even the legacy `SQLITE_MAX_VARIABLE_NUMBER` of 999). In an incremental batch a file's `DELETE` targets only its own id, so buffered rows of other files are safe; a file written twice in one batch flushes the buffer before its second `DELETE`. - **Prefix queries use a range, never `LIKE`.** `LIKE` is case-insensitive by default, so SQLite cannot turn `module LIKE 'Prefix.%'` into an index range and scans all refs. `module >= 'Prefix.' AND module < 'Prefix/'` ('/' is '.'+1) uses `idx_refs_module_function` and turns that scan into a range search: on a 3.9M-row index, 11-14x faster with a warm page cache and ~190x faster cold. - **`idx_refs_function_kind` was retired.** No query leads with `function`; the two that filter on function/kind both lead with `file_path`. It cost 80 MB and a share of every index rebuild. Check `EXPLAIN QUERY PLAN` before adding an index here — index build time is ~40% of a cold index. @@ -223,8 +223,9 @@ The largest remaining win is interning `file_path`: every ref row stores a ~122- - **Tokenizer instead of tree-sitter for indexing** — a hand-rolled tokenizer + walker replaced the original regex-based parser for both file indexing and runtime `__using__` parsing. The tokenizer handles heredocs, sigils, multi-line expressions, and comments as opaque tokens, eliminating fragile line-joining heuristics. Tree-sitter is only used for scope-aware variable operations in files already opened by the editor. - **SQLite for storage** — single file, fast reads, incremental updates via mtime tracking. - **Parallel indexing** — the cold build uses all CPU cores for parsing, single writer for SQLite. Both callers share `indexer.FullBuild`; the server used to have a second, serial implementation that parsed one file at a time and committed a transaction per file, which measured ~4x slower on a 10k-file corpus. -- **Two write paths** — `indexer.FullBuild` is insert-only and allocates file ids from a counter, so it is only correct on an empty index and only with no other writer active. Everything else — incremental sweeps included — goes through the per-file path, which deletes a file's rows before reinserting them. -- **`IndexCoordinator.writes` covers every writer, without exception.** A full build holds it for writing; every other writer — a save, a watched-file event, a rename, and the incremental sweep's own walk — holds it for reading. A new writer that skips it can run inside an insert-only build, whose emptiness precondition it then violates. The coordinator belongs to the workspace rather than to a session, so every LSP session attached to one store shares it, and the daemon's ownership lock extends the precondition across processes: `dexter reindex` in a shell now enters the same queue instead of opening a second writer the lock could not see. +- **Three write paths** — `indexer.FullBuild` is insert-only and allocates file ids from a counter, so it is only correct on an empty index and only with no other writer active. A single-file write deletes a file's rows before reinserting them. The warm sweep reads every stored mtime with one query (`Store.FileStates`), walks, and then parses the changed files on every core: a small change set goes through incremental batches (many files per transaction), and a change set of at least 4096 files and a quarter of the index goes through `Store.BeginRebuild`. A rebuild writes new, unindexed copies of `definitions` and `refs`, copies in the rows of every other file, swaps the copies in, and builds each index with one sort, all in one transaction. Random updates of the name-sorted indexes are what made one transaction per file take minutes on a large change set. Removal follows the same rule: `Store.RemoveFileIDs` deletes by id sets, or rebuilds without the removed files when they are a quarter of the index. A failed rebuild falls back to batches (or to id-set deletes), and a failed batch is written again one file at a time, so only a file that fails by itself is lost. A rebuild takes SQLite's write lock with its first statement, because SQLite does not call the busy handler when a deferred transaction upgrades a read to a write. +- **Shutdown cancels index work.** `workspace.Runtime.Close` calls `IndexCoordinator.CancelWork` first. The sweep, the prune, every removal, the rebuild, and a cold build stop; an open transaction rolls back, and a canceled SQL statement is interrupted. A file and its rows are always in one transaction, so nothing is left half written, and the files a canceled pass did not reach keep their old mtime, so the next start does them. +- **`IndexCoordinator.writes` covers every writer, without exception.** A full build holds it for writing; every other writer — a save, a watched-file event, a rename, and the incremental sweep's own walk — holds it for reading. Single-file writes and removals also take `IndexCoordinator.reindexing`, so they wait behind a sweep instead of racing its batches for SQLite's write lock. A new writer that skips it can run inside an insert-only build, whose emptiness precondition it then violates. The coordinator belongs to the workspace rather than to a session, so every LSP session attached to one store shares it, and the daemon's ownership lock extends the precondition across processes: `dexter reindex` in a shell now enters the same queue instead of opening a second writer the lock could not see. - **Emptiness is decided under the lock, never sampled.** `Server.fullBuild` tests `IsEmpty` while holding the write lock and reports through its `ran` return value, because one save arriving between a sample and the lock invalidates the answer. `Store.IsEmpty` answers `false` when its own query fails, so a database too broken to count is never taken for an empty one. - **The sweep's prune re-checks the filesystem.** `pruneMissingFiles` deletes only stored paths that are absent from the walk *and* fail to stat. The walk is one traversal, so a file saved after it passed that directory is legitimately missing from `seen`; the stat is also what stops a walker that yields nothing from deleting the whole index. - **A cold build is one WAL transaction**, so the `-wal` file grows to roughly the size of the index — a few hundred MiB on a large monorepo — until the post-build checkpoint reclaims it. `wal_autocheckpoint` cannot touch frames belonging to an open transaction, and `InProcess` deliberately leaves `journal_mode` alone, so the server cannot avoid this the way `dexter init` does. diff --git a/internal/indexer/indexer.go b/internal/indexer/indexer.go index 35c11cd..f79c5d1 100644 --- a/internal/indexer/indexer.go +++ b/internal/indexer/indexer.go @@ -10,6 +10,7 @@ package indexer import ( + "context" "errors" "fmt" "os" @@ -40,6 +41,11 @@ type Options struct { // store.SetBulkPragmas. InProcess bool + // Context, when set, cancels the build: parsing stops, the bulk + // transaction rolls back, and FullBuild returns the context's error. The + // index is left empty, so the next start builds it again. Optional. + Context context.Context + // Warn reports a recoverable per-file failure. Optional. // // It is called from every parse worker, so it must be safe for concurrent @@ -146,6 +152,10 @@ func statFilesParallel(paths []string) []fileEntry { func FullBuild(s *store.Store, projectRoot string, opts Options) (Stats, error) { var stats Stats start := time.Now() + ctx := opts.Context + if ctx == nil { + ctx = context.Background() + } warn := opts.serialWarn() @@ -157,6 +167,12 @@ func FullBuild(s *store.Store, projectRoot string, opts Options) (Stats, error) } } + // Nothing is written before the bulk transaction, so a cancel up to that + // point returns with the database untouched. + if err := ctx.Err(); err != nil { + return stats, err + } + // Phase 1: collect file paths and mtimes. Both halves run on all cores: // the traversal fans out per directory, and stat costs one syscall per // file — ~70k on a large monorepo, which dominated this phase. @@ -165,13 +181,19 @@ func FullBuild(s *store.Store, projectRoot string, opts Options) (Stats, error) if opts.StdlibRoot != "" { stdlibPaths = dedupeAgainst(parser.CollectElixirFilesParallel(opts.StdlibRoot), filePaths) } + if err := ctx.Err(); err != nil { + return stats, err + } files := statFilesParallel(filePaths) stdlibFiles := statFilesParallel(stdlibPaths) stats.Walk = time.Since(start) // Phase 2a: parse stdlib files in parallel. Definitions only — refs are // not indexed for stdlib. - stdlibResults := parseStdlib(stdlibFiles) + stdlibResults := parseStdlib(ctx, stdlibFiles) + if err := ctx.Err(); err != nil { + return stats, err + } // Phase 2b: parse project files in parallel, streaming to the writer. workers := runtime.NumCPU() @@ -199,8 +221,13 @@ func FullBuild(s *store.Store, projectRoot string, opts Options) (Stats, error) } go func() { + feed: for _, f := range files { - fileCh <- f + select { + case fileCh <- f: + case <-ctx.Done(): + break feed + } } close(fileCh) wg.Wait() @@ -242,6 +269,9 @@ func FullBuild(s *store.Store, projectRoot string, opts Options) (Stats, error) var writeNanos time.Duration for res := range resultCh { + if err := ctx.Err(); err != nil { + return stats, abortBuild(s, batch, resultCh, err) + } writeStart := time.Now() err := batch.IndexFileWithMtimeAndRefs(res.path, res.mtimeNano, res.defs, res.refs) writeNanos += time.Since(writeStart) @@ -255,6 +285,9 @@ func FullBuild(s *store.Store, projectRoot string, opts Options) (Stats, error) stats.Write = writeNanos stats.Parse = time.Duration(parseNanos.Load()) + if err := ctx.Err(); err != nil { + return stats, abortBuild(s, batch, resultCh, err) + } commitStart := time.Now() if err := batch.Commit(); err != nil { return stats, restoreIndexes(s, fmt.Errorf("commit: %w", err)) @@ -290,7 +323,7 @@ type stdlibResult struct { defs []parser.Definition } -func parseStdlib(stdlibFiles []fileEntry) []stdlibResult { +func parseStdlib(ctx context.Context, stdlibFiles []fileEntry) []stdlibResult { if len(stdlibFiles) == 0 { return nil } @@ -305,6 +338,9 @@ func parseStdlib(stdlibFiles []fileEntry) []stdlibResult { go func() { defer wg.Done() for f := range fileCh { + if ctx.Err() != nil { + continue + } defs, _, err := parser.ParseFile(f.path) if err != nil { continue @@ -314,8 +350,13 @@ func parseStdlib(stdlibFiles []fileEntry) []stdlibResult { }() } go func() { + feed: for _, f := range stdlibFiles { - fileCh <- f + select { + case fileCh <- f: + case <-ctx.Done(): + break feed + } } close(fileCh) wg.Wait() diff --git a/internal/lsp/reconcile.go b/internal/lsp/reconcile.go new file mode 100644 index 0000000..1d20411 --- /dev/null +++ b/internal/lsp/reconcile.go @@ -0,0 +1,426 @@ +package lsp + +import ( + "context" + "io/fs" + "log" + "os" + "runtime" + "sync" + "time" + + "github.com/remoteoss/dexter/internal/parser" + "github.com/remoteoss/dexter/internal/store" +) + +// Write batching for the warm reconciliation. One transaction per file made +// each commit write every index page the file touched: on a 280k-file index, a +// warm pass over 223k new files spent 59% of its CPU time in COMMIT and took +// almost four minutes. A batch shares those pages across many files. The time +// bound keeps the write lock short for any other writer: a batch commits at +// whichever limit it reaches first. +var ( + reconcileBatchFiles = 2048 + reconcileBatchTime = 500 * time.Millisecond +) + +// A change set this large, and at least 1/rebuildShare of the files the index +// will hold, is written with store.BeginRebuild instead: the live indexes then +// take no row-at-a-time updates at all. On a 280k-file index, a warm pass over +// 223k new files took 1m58s through batches and 44s as a rebuild (3m48s with +// one transaction per file). A batch costs about 0.5ms per changed file there, +// a rebuild about 0.16ms per file in the whole index, so they break even near +// a third; a quarter leaves margin for the larger index a rebuild re-sorts. +var ( + rebuildMinFiles = 4096 + rebuildShare = 4 +) + +// testHookParse, when set by a test, runs before each changed file is parsed. +var testHookParse func(path string) + +// changedFile is one file the walk found new or changed. +type changedFile struct { + path string + mtimeNano int64 // from the walk's lstat; 0 for a symlink, statted on parse + symlink bool + refs bool +} + +type parsedFile struct { + changedFile + defs []parser.Definition + refsList []parser.Reference +} + +// reconcileChangedFiles walks the stdlib and project roots and indexes every +// file whose mtime differs from the stored one. It returns the set of paths the +// walk saw (for the prune), how many files it wrote, and false when the index +// is unavailable or the pass was canceled, in which case nothing may be pruned. +// +// The stored mtimes are read with one query before the walk instead of one +// query per file, and the walk itself is unchanged: one lstat per file, from +// the directory entry. Changed files are parsed on every core while the walk +// continues, and one writer stores them in batched transactions. A file and its +// rows are always in one transaction, so a canceled pass leaves each file either +// fully old or fully new, and the next pass redoes the rest from their mtimes. +func (s *Server) reconcileChangedFiles() (map[string]struct{}, int, bool) { + ctx := s.index.work + // The walk writes, so it takes indexWrites for reading, the same as every + // other single-file write. That is what keeps it from overlapping a cold + // build. + s.index.writes.RLock() + defer s.index.writes.RUnlock() + if s.index.unavailable || ctx.Err() != nil { + return nil, 0, false + } + + stored, err := s.store.FileStates() + if err != nil { + log.Printf("Warning: reading stored file states: %v", err) + return nil, 0, false + } + + seen := make(map[string]struct{}, len(stored)) + var changed []changedFile + walk := func(root string, indexRefs bool) { + _ = parser.WalkElixirFiles(root, func(path string, d fs.DirEntry) error { + if err := ctx.Err(); err != nil { + return err + } + // The stdlib root can lie inside the project; the second visit of + // a path must not queue it again. + if _, dup := seen[path]; dup { + return nil + } + seen[path] = struct{}{} + + info, err := d.Info() + if err != nil { + return nil + } + mtime := info.ModTime().UnixNano() + if st, found := stored[path]; found && st.Mtime == mtime { + return nil + } + f := changedFile{path: path, mtimeNano: mtime, refs: indexRefs} + if d.Type()&fs.ModeSymlink != 0 { + // The stored mtime is the target's, as the cold build and + // single-file writes record it. + f.symlink, f.mtimeNano = true, 0 + } + changed = append(changed, f) + return nil + }) + } + + if stdlibRoot := s.StdlibRoot(); stdlibRoot != "" { + walk(stdlibRoot, false) + } + walk(s.projectRoot, true) + + if ctx.Err() != nil { + return seen, 0, false + } + + written := 0 + switch { + case len(changed) == 0: + case len(changed) >= rebuildMinFiles && len(changed)*rebuildShare >= len(stored)+len(changed): + // The rebuild must be the only writer for its whole transaction, as a + // cold build is, so it trades the read lock for the write lock. A + // single-file write that lands in between is either rewritten here or + // copied over as it is. + s.index.writes.RUnlock() + n, ok := s.rebuildChanged(ctx, changed) + s.index.writes.RLock() + written = n + if !ok && ctx.Err() == nil && !s.index.unavailable { + // Batches are slower but need nothing the rebuild did, so a + // failed rebuild does not leave the change set unindexed. + log.Printf("Warning: index rebuild failed, writing %d changed files in batches", len(changed)) + written = s.writeChangedInBatches(ctx, changed) + } + default: + written = s.writeChangedInBatches(ctx, changed) + } + if ctx.Err() != nil { + return seen, written, false + } + return seen, written, true +} + +// writeChangedInBatches parses changed files on every core and writes them +// through batched transactions. It returns how many files it wrote. +func (s *Server) writeChangedInBatches(ctx context.Context, changed []changedFile) int { + pipe := s.startReconcilePipeline(ctx) + for _, f := range changed { + if pipe.send(f) != nil { + break + } + } + return pipe.finish() +} + +// parsePool parses changed files on every core. +type parsePool struct { + ctx context.Context + work chan changedFile + results chan parsedFile +} + +func newParsePool(ctx context.Context) *parsePool { + workers := runtime.NumCPU() + p := &parsePool{ + ctx: ctx, + work: make(chan changedFile, workers*4), + results: make(chan parsedFile, workers*4), + } + var wg sync.WaitGroup + for i := 0; i < workers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for f := range p.work { + if ctx.Err() != nil { + continue // drain without parsing + } + if testHookParse != nil { + testHookParse(f.path) + } + if f.symlink { + info, err := os.Stat(f.path) + if err != nil { + continue + } + f.mtimeNano = info.ModTime().UnixNano() + } + defs, refs, err := parser.ParseFile(f.path) + if err != nil { + continue + } + if !f.refs { + refs = nil + } + select { + case p.results <- parsedFile{changedFile: f, defs: defs, refsList: refs}: + case <-ctx.Done(): + } + } + }() + } + go func() { + wg.Wait() + close(p.results) + }() + return p +} + +// send queues one changed file for parsing. It returns the context error when +// the pass is canceled. +func (p *parsePool) send(f changedFile) error { + select { + case p.work <- f: + return nil + case <-p.ctx.Done(): + return p.ctx.Err() + } +} + +// reconcilePipeline is a parse pool with one writer that stores its results in +// batched transactions. +type reconcilePipeline struct { + *parsePool + written chan int +} + +func (s *Server) startReconcilePipeline(ctx context.Context) *reconcilePipeline { + p := &reconcilePipeline{parsePool: newParsePool(ctx), written: make(chan int, 1)} + go func() { p.written <- s.writeReconciled(ctx, p.results) }() + return p +} + +// finish waits for every queued file to be parsed and written, and returns how +// many files were written. +func (p *reconcilePipeline) finish() int { + close(p.work) + return <-p.written +} + +// rebuildChanged parses the change set on every core and writes it through one +// rebuild transaction. It returns how many files it wrote, and false when the +// rebuild failed or was canceled, in which case it wrote none. +func (s *Server) rebuildChanged(ctx context.Context, changed []changedFile) (int, bool) { + s.index.writes.Lock() + defer s.index.writes.Unlock() + if s.index.unavailable { + return 0, false + } + start := time.Now() + batch, err := s.store.BeginRebuild(ctx) + if err != nil { + if ctx.Err() == nil { + log.Printf("Warning: starting index rebuild: %v", err) + } + return 0, false + } + results := s.parseChanged(ctx, changed) + files := 0 + var writeErr error + for res := range results { + if writeErr != nil || ctx.Err() != nil { + continue // drain so the parse workers can exit + } + if writeErr = writeParsed(batch, res); writeErr == nil { + files++ + } + } + if writeErr == nil { + writeErr = ctx.Err() + } + if writeErr != nil { + _ = batch.Rollback() + if ctx.Err() == nil { + log.Printf("Warning: index rebuild: %v", writeErr) + } + return 0, false + } + written := time.Now() + if err := batch.Commit(); err != nil { + if ctx.Err() == nil { + log.Printf("Warning: index rebuild: %v", err) + } + return 0, false + } + log.Printf("Rebuilt the index with %d changed files (parse and write %s, swap and index %s)", files, + written.Sub(start).Round(time.Millisecond), time.Since(written).Round(time.Millisecond)) + return files, true +} + +// testHookWrite, when set by a test, runs before each parsed file is written; +// an error it returns fails that write. +var testHookWrite func(path string) error + +func writeParsed(b *store.Batch, res parsedFile) error { + if testHookWrite != nil { + if err := testHookWrite(res.path); err != nil { + return err + } + } + return b.IndexFileWithMtimeAndRefs(res.path, res.mtimeNano, res.defs, res.refsList) +} + +// parseChanged parses files on every core and streams the results. The channel +// closes when every file is parsed, or soon after ctx is canceled. +func (s *Server) parseChanged(ctx context.Context, changed []changedFile) <-chan parsedFile { + p := newParsePool(ctx) + go func() { + for _, f := range changed { + if p.send(f) != nil { + break + } + } + close(p.work) + }() + return p.results +} + +// writeReconciled stores parsed files in batched transactions until results +// closes. When a batch fails for a reason other than cancellation, it is rolled +// back and its files are written again one by one, each in its own +// transaction, so only a file that fails by itself is lost; it keeps its old +// mtime, so the next pass retries it. After cancellation it rolls back the open +// batch and only drains. +func (s *Server) writeReconciled(ctx context.Context, results <-chan parsedFile) int { + var ( + batch *store.Batch + group []parsedFile // the files in batch, kept for a one-by-one retry + started time.Time + written int + ) + retryOneByOne := func(cause error) { + log.Printf("Warning: reindex batch of %d files failed, writing them one by one: %v", len(group), cause) + for _, res := range group { + if ctx.Err() != nil { + return + } + if err := s.writeOne(ctx, res); err != nil { + if ctx.Err() == nil { + log.Printf("Warning: reindex %s: %v", res.path, err) + } + continue + } + written++ + } + } + commit := func() { + if batch == nil { + return + } + err := batch.Commit() + switch { + case err == nil: + written += len(group) + case ctx.Err() == nil: + retryOneByOne(err) + } + batch, group = nil, group[:0] + } + abort := func() { + if batch != nil { + _ = batch.Rollback() + batch, group = nil, group[:0] + } + } + + for res := range results { + if ctx.Err() != nil { + abort() + continue + } + if batch == nil { + b, err := s.store.BeginBatchContext(ctx) + if err != nil { + if ctx.Err() == nil { + log.Printf("Warning: reindex %s: %v", res.path, err) + } + continue + } + batch, started = b, time.Now() + } + group = append(group, res) + if err := writeParsed(batch, res); err != nil { + // The failed file may be partly written inside the transaction, + // so the batch rolls back, and its files are written one by one. + _ = batch.Rollback() + batch = nil + if ctx.Err() == nil { + retryOneByOne(err) + } + group = group[:0] + continue + } + if len(group) >= reconcileBatchFiles || time.Since(started) >= reconcileBatchTime { + commit() + } + } + if ctx.Err() != nil { + abort() + } else { + commit() + } + return written +} + +// writeOne writes one parsed file in its own transaction. +func (s *Server) writeOne(ctx context.Context, res parsedFile) error { + b, err := s.store.BeginBatchContext(ctx) + if err != nil { + return err + } + if err := writeParsed(b, res); err != nil { + _ = b.Rollback() + return err + } + return b.Commit() +} diff --git a/internal/lsp/reconcile_test.go b/internal/lsp/reconcile_test.go new file mode 100644 index 0000000..6dd5e38 --- /dev/null +++ b/internal/lsp/reconcile_test.go @@ -0,0 +1,358 @@ +package lsp + +import ( + "bytes" + "context" + "fmt" + "log" + "os" + "strings" + "sync" + "testing" + "time" + + "go.lsp.dev/protocol" + "go.lsp.dev/uri" +) + +func moduleSource(i, version int) string { + return fmt.Sprintf(`defmodule MyApp.Gen%d do + def run_v%d(id), do: SharedLib.Worker.perform(id) +end +`, i, version) +} + +func setReconcileVars(t *testing.T, batchFiles, minFiles, share int) { + t.Helper() + oldBatch, oldMin, oldShare := reconcileBatchFiles, rebuildMinFiles, rebuildShare + reconcileBatchFiles, rebuildMinFiles, rebuildShare = batchFiles, minFiles, share + t.Cleanup(func() { + reconcileBatchFiles, rebuildMinFiles, rebuildShare = oldBatch, oldMin, oldShare + }) +} + +func captureLog(t *testing.T) *syncBuffer { + t.Helper() + buf := &syncBuffer{} + prev := log.Writer() + log.SetOutput(buf) + t.Cleanup(func() { log.SetOutput(prev) }) + return buf +} + +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() +} + +// reindexOnce runs one background pass to completion. +func reindexOnce(t *testing.T, s *Server) { + t.Helper() + select { + case <-s.startBackgroundReindex(): + case <-time.After(30 * time.Second): + t.Fatal("reindex did not finish") + } +} + +// A warm pass writes a large change set through the rebuild and a small one +// through batches. Both must leave changed files with only their new rows and +// add every new file. +func TestReconcile_WarmChangeSetPaths(t *testing.T) { + for _, tc := range []struct { + name string + batchFiles int + minFiles int + wantLog string + }{ + {name: "batched", batchFiles: 3, minFiles: 1 << 30, wantLog: ""}, + {name: "rebuild", batchFiles: 3, minFiles: 1, wantLog: "Rebuilt the index with 21 changed files"}, + } { + t.Run(tc.name, func(t *testing.T) { + setReconcileVars(t, tc.batchFiles, tc.minFiles, 4) + logs := captureLog(t) + server, cleanup := setupTestServer(t) + defer cleanup() + + writeTestFile(t, server.projectRoot, "lib/gen0.ex", moduleSource(0, 0)) + writeTestFile(t, server.projectRoot, "lib/kept.ex", "defmodule MyApp.Kept do\n def here, do: :ok\nend\n") + reindexOnce(t, server) // cold: full build + + // Change gen0 (with a different mtime) and add 20 files. + path := writeTestFile(t, server.projectRoot, "lib/gen0.ex", moduleSource(0, 1)) + future := time.Now().Add(time.Hour) + if err := os.Chtimes(path, future, future); err != nil { + t.Fatal(err) + } + for i := 1; i <= 20; i++ { + writeTestFile(t, server.projectRoot, fmt.Sprintf("lib/gen%d.ex", i), moduleSource(i, 0)) + } + reindexOnce(t, server) + + if r, _ := server.store.LookupFunction("MyApp.Gen0", "run_v1"); len(r) != 1 { + t.Error("changed file does not have its new definition") + } + if r, _ := server.store.LookupFunction("MyApp.Gen0", "run_v0"); len(r) != 0 { + t.Error("changed file kept its old definition") + } + for i := 1; i <= 20; i++ { + if r, _ := server.store.LookupFunction(fmt.Sprintf("MyApp.Gen%d", i), "run_v0"); len(r) != 1 { + t.Errorf("new file %d: %d definitions, want 1", i, len(r)) + } + } + if r, _ := server.store.LookupFunction("MyApp.Kept", "here"); len(r) != 1 { + t.Error("an unchanged file lost its definition") + } + if refs, _ := server.store.LookupReferences("SharedLib.Worker", "perform"); len(refs) != 21 { + t.Errorf("references = %d, want 21", len(refs)) + } + if tc.wantLog != "" && !strings.Contains(logs.String(), tc.wantLog) { + t.Errorf("log does not say %q:\n%s", tc.wantLog, logs.String()) + } + + // A pass with nothing changed writes nothing. + before := logs.String() + reindexOnce(t, server) + if !strings.Contains(strings.TrimPrefix(logs.String(), before), "Background reindex: 0 files updated") { + t.Errorf("a pass without changes wrote files:\n%s", strings.TrimPrefix(logs.String(), before)) + } + }) + } +} + +// CancelWork must end a warm pass in flight promptly, leave no file half +// written, and leave the rest for the next start, which finishes it from the +// stored mtimes. This is what shutdown relies on. +func TestReconcile_CancelWorkStopsPassAndNextStartFinishes(t *testing.T) { + const files = 40 + for _, tc := range []struct { + name string + minFiles int + }{ + {name: "batched", minFiles: 1 << 30}, + {name: "rebuild", minFiles: 1}, + } { + t.Run(tc.name, func(t *testing.T) { + setReconcileVars(t, 4, tc.minFiles, 4) + server, cleanup := setupTestServer(t) + defer cleanup() + + writeTestFile(t, server.projectRoot, "lib/gen0.ex", moduleSource(0, 0)) + reindexOnce(t, server) + for i := 1; i <= files; i++ { + writeTestFile(t, server.projectRoot, fmt.Sprintf("lib/gen%d.ex", i), moduleSource(i, 0)) + } + + // Let some files through, then hold every parse until the pass is + // canceled, as a long pass would be at shutdown. + var parsed sync.WaitGroup + parsed.Add(10) + var mu sync.Mutex + count := 0 + testHookParse = func(string) { + mu.Lock() + count++ + n := count + mu.Unlock() + if n <= 10 { + parsed.Done() + return + } + <-server.index.work.Done() + } + t.Cleanup(func() { testHookParse = nil }) + + done := server.startBackgroundReindex() + parsed.Wait() + canceled := time.Now() + server.index.CancelWork() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("the pass did not stop after CancelWork") + } + if waited := time.Since(canceled); waited > 2*time.Second { + t.Errorf("the pass took %s to stop", waited) + } + testHookParse = nil + + // No half-written file: every stored file has its definition. + paths, err := server.store.ListFilePaths() + if err != nil { + t.Fatal(err) + } + for _, p := range paths { + mods, err := server.store.LookupModulesInFile(p) + if err != nil || len(mods) != 1 { + t.Errorf("%s is stored without its module (%v, %v)", p, mods, err) + } + } + if len(paths) > files { + t.Errorf("a canceled pass stored all %d files", len(paths)) + } + + // The next start, with a new coordinator, finishes the work. + next := NewServer(server.store, server.projectRoot) + reindexOnce(t, next) + for i := 0; i <= files; i++ { + if r, _ := next.store.LookupFunction(fmt.Sprintf("MyApp.Gen%d", i), "run_v0"); len(r) != 1 { + t.Errorf("file %d not indexed after the next start", i) + } + } + }) + } +} + +// A canceled cold build leaves an empty index, which the next start builds. +func TestReconcile_CancelWorkBeforeColdBuild(t *testing.T) { + server, cleanup := setupTestServer(t) + defer cleanup() + writeTestFile(t, server.projectRoot, "lib/gen1.ex", moduleSource(1, 0)) + + server.index.CancelWork() + reindexOnce(t, server) + if !server.store.IsEmpty() { + t.Error("a canceled cold build stored files") + } + if names, err := server.store.IndexNames(); err != nil || len(names) != 6 { + t.Errorf("indexes after a canceled cold build: %v (%v)", names, err) + } + + next := NewServer(server.store, server.projectRoot) + reindexOnce(t, next) + if r, _ := next.store.LookupFunction("MyApp.Gen1", "run_v0"); len(r) != 1 { + t.Error("the next start did not build the index") + } +} + +// Removals take the coordinator's context, so CancelWork stops them as it +// stops a pass. Before, they ran with no context and could hold shutdown for +// a whole rebuild of the symbol tables. +func TestServer_RemoveFilesStopsAfterCancelWork(t *testing.T) { + server, cleanup := setupTestServer(t) + defer cleanup() + var paths []string + for i := 0; i < 5; i++ { + paths = append(paths, writeTestFile(t, server.projectRoot, fmt.Sprintf("lib/gen%d.ex", i), moduleSource(i, 0))) + } + reindexOnce(t, server) + + server.index.CancelWork() + server.RemoveFiles(paths) + server.RemoveFilesUnderRoot(server.projectRoot) + if stored, _ := server.store.ListFilePaths(); len(stored) != len(paths) { + t.Errorf("%d files left after canceled removals, want %d", len(stored), len(paths)) + } +} + +// One file that fails to write must not take the other files of its batch +// with it: the batch is written again one file at a time. +func TestReconcile_FailedWriteLosesOnlyThatFile(t *testing.T) { + for _, tc := range []struct { + name string + minFiles int + wantLog []string + }{ + {name: "batched", minFiles: 1 << 30, wantLog: []string{"writing them one by one"}}, + // A failed rebuild falls back to batches, which then retry one by one. + {name: "rebuild", minFiles: 1, wantLog: []string{"writing 10 changed files in batches", "writing them one by one"}}, + } { + t.Run(tc.name, func(t *testing.T) { testFailedWrite(t, tc.minFiles, tc.wantLog) }) + } +} + +func testFailedWrite(t *testing.T, minFiles int, wantLog []string) { + setReconcileVars(t, 1000, minFiles, 4) + logs := captureLog(t) + server, cleanup := setupTestServer(t) + defer cleanup() + writeTestFile(t, server.projectRoot, "lib/gen0.ex", moduleSource(0, 0)) + reindexOnce(t, server) + + var bad string + for i := 1; i <= 10; i++ { + p := writeTestFile(t, server.projectRoot, fmt.Sprintf("lib/gen%d.ex", i), moduleSource(i, 0)) + if i == 5 { + bad = p + } + } + testHookWrite = func(path string) error { + if path == bad { + return fmt.Errorf("injected write failure") + } + return nil + } + t.Cleanup(func() { testHookWrite = nil }) + reindexOnce(t, server) + testHookWrite = nil + + for i := 1; i <= 10; i++ { + want := 1 + if i == 5 { + want = 0 + } + if r, _ := server.store.LookupFunction(fmt.Sprintf("MyApp.Gen%d", i), "run_v0"); len(r) != want { + t.Errorf("file %d: %d definitions, want %d", i, len(r), want) + } + } + for _, want := range wantLog { + if !strings.Contains(logs.String(), want) { + t.Errorf("log does not say %q:\n%s", want, logs.String()) + } + } + + // The failed file kept no stored mtime, so the next pass writes it. + reindexOnce(t, server) + if r, _ := server.store.LookupFunction("MyApp.Gen5", "run_v0"); len(r) != 1 { + t.Error("the next pass did not write the file that failed") + } +} + +// In an embedded server, a watched-file change waits for the coordinator's +// mutation lock, as the daemon's events do. Before, it took only the read side +// of the write lock and raced the pass's batches for SQLite's write lock: it +// could fail after busy_timeout and stay unindexed until the next pass. +func TestServer_EmbeddedFileChangeWaitsForReconcile(t *testing.T) { + server, cleanup := setupTestServer(t) + defer cleanup() + writeTestFile(t, server.projectRoot, "lib/gen0.ex", moduleSource(0, 0)) + reindexOnce(t, server) + + path := writeTestFile(t, server.projectRoot, "lib/saved.ex", "defmodule MyApp.Saved do\n def ok, do: :ok\nend\n") + server.index.reindexing.Lock() // a pass in flight + if err := server.DidChangeWatchedFiles(context.Background(), &protocol.DidChangeWatchedFilesParams{ + Changes: []*protocol.FileEvent{{URI: uri.File(path), Type: protocol.FileChangeTypeChanged}}, + }); err != nil { + server.index.reindexing.Unlock() + t.Fatal(err) + } + time.Sleep(100 * time.Millisecond) + early, _ := server.store.LookupFunction("MyApp.Saved", "ok") + server.index.reindexing.Unlock() + if len(early) != 0 { + t.Error("the change was written while a pass held the mutation lock") + } + + deadline := time.Now().Add(10 * time.Second) + for { + if r, _ := server.store.LookupFunction("MyApp.Saved", "ok"); len(r) == 1 { + break + } + if time.Now().After(deadline) { + t.Fatal("the change was not indexed after the pass") + } + time.Sleep(10 * time.Millisecond) + } +} diff --git a/internal/lsp/server.go b/internal/lsp/server.go index 3e5cbfc..ee9bdd6 100644 --- a/internal/lsp/server.go +++ b/internal/lsp/server.go @@ -103,8 +103,21 @@ type IndexCoordinator struct { unavailable bool // guarded by writes backgroundWork sync.WaitGroup + + // work is canceled by CancelWork when the workspace shuts down. A + // reconciliation pass, a prune, and a cold build check it and stop. + work context.Context + cancelWork context.CancelFunc } +// CancelWork stops the reconciliation in flight and makes every later one +// return at once. The workspace calls it first when it shuts down, so a pass +// over a large change set cannot keep the process (and its workspace lock) +// alive for minutes. A stopped pass leaves no file half written: each write +// transaction holds whole files, and a canceled one rolls back. The files it +// did not reach keep their old mtime, so the next start picks them up. +func (c *IndexCoordinator) CancelWork() { c.cancelWork() } + func (c *IndexCoordinator) setStdlibRoot(root string) (string, bool) { c.stdlibMu.Lock() defer c.stdlibMu.Unlock() @@ -124,7 +137,8 @@ func (c *IndexCoordinator) getStdlibRoot() string { // NewIndexCoordinator returns write coordination for one workspace store. func NewIndexCoordinator() *IndexCoordinator { - return &IndexCoordinator{} + work, cancel := context.WithCancel(context.Background()) + return &IndexCoordinator{work: work, cancelWork: cancel} } // ServerOptions configures a Server attached to a daemon-owned workspace. @@ -378,6 +392,10 @@ func (s *Server) warmUsingCache() { // removed is still pruned, which is why the check is per candidate rather than // a test on seen being empty. func (s *Server) pruneMissingFiles(seen map[string]struct{}) { + ctx := s.index.work + if ctx.Err() != nil { + return + } s.index.writes.Lock() defer s.index.writes.Unlock() if s.index.unavailable { @@ -398,10 +416,13 @@ func (s *Server) pruneMissingFiles(seen map[string]struct{}) { inWorktree := make(map[string]bool) var recorded map[string]struct{} recordedRead := false - for _, storedPath := range storedPaths { + for i, storedPath := range storedPaths { if _, ok := seen[storedPath]; ok { continue } + if i&1023 == 0 && ctx.Err() != nil { + return + } _, err := os.Lstat(storedPath) if err != nil && !errors.Is(err, fs.ErrNotExist) { continue @@ -423,7 +444,14 @@ func (s *Server) pruneMissingFiles(seen map[string]struct{}) { toRemove = append(toRemove, storedPath) } if len(toRemove) > 0 { - _ = s.store.RemoveFiles(toRemove) + start := time.Now() + if err := s.store.RemoveFilesContext(ctx, toRemove); err != nil { + if ctx.Err() == nil { + log.Printf("Warning: removing %d files from the index: %v", len(toRemove), err) + } + return + } + log.Printf("Removed %d files from the index (%s)", len(toRemove), time.Since(start).Round(time.Millisecond)) } } @@ -487,6 +515,7 @@ func (s *Server) fullBuild() (stats indexer.Stats, ran bool, err error) { stats, err = indexer.FullBuild(s.store, s.projectRoot, indexer.Options{ StdlibRoot: s.StdlibRoot(), InProcess: true, + Context: s.index.work, Warn: func(format string, args ...interface{}) { log.Printf("Warning: "+format, args...) }, @@ -539,6 +568,9 @@ func (s *Server) startBackgroundReindex() <-chan struct{} { defer s.index.backgroundWork.Done() s.index.reindexing.Lock() defer s.index.reindexing.Unlock() + if s.index.work.Err() != nil { + return // the workspace is shutting down + } start := time.Now() reindexed := 0 @@ -575,6 +607,9 @@ func (s *Server) startBackgroundReindex() <-chan struct{} { log.Printf("Warning: WAL checkpoint: %v", err) } return + case s.index.work.Err() != nil: + log.Printf("Index build canceled") + return case err != nil: // The incremental walk below needs nothing to be true of the // database, so it is the safe thing to fall back to. It is @@ -591,65 +626,23 @@ func (s *Server) startBackgroundReindex() <-chan struct{} { } } - // Re-read rather than reusing coldStart. A full build, a failed build - // and a concurrent write all change the answer, and reading a stale - // true here would skip the mtime short-circuit for every file. - isEmpty := s.store.IsEmpty() - - seen := make(map[string]struct{}) - walkAndIndex := func(root string, indexRefs bool) { - _ = parser.WalkElixirFiles(root, func(path string, d fs.DirEntry) error { - seen[path] = struct{}{} - - if !isEmpty { - info, err := d.Info() - if err != nil { - return nil - } - storedMtime, found := s.store.GetFileMtime(path) - currentMtime := info.ModTime().UnixNano() - if found && storedMtime == currentMtime { - return nil - } - } - - defs, refs, err := parser.ParseFile(path) - if err != nil { - return nil - } - if !indexRefs { - refs = nil - } - if err := s.store.IndexFileWithRefs(path, defs, refs); err != nil { - log.Printf("Warning: reindex %s: %v", path, err) - } - reindexed++ - return nil - }) - } - // A full build already indexed every file on disk from the traversal - // this walk would repeat, so skipping it saves a second traversal and a - // stored-mtime query per file. The prune lives in the same branch and - // so cannot run without the walk that fills `seen`. + // this walk would repeat, so skipping it saves a second traversal. The + // prune lives in the same branch and so cannot run without the walk + // that fills `seen`. if !fullBuilt { - // The walk writes, so it takes indexWrites for reading, the same as - // every other single-file write. That is what keeps it from - // overlapping a cold build. - s.index.writes.RLock() - if s.index.unavailable { - s.index.writes.RUnlock() + seen, n, ok := s.reconcileChangedFiles() + reindexed += n + if ok { + s.pruneMissingFiles(seen) + } + if s.index.work.Err() != nil { + log.Printf("Background reindex canceled after %d files (%s)", reindexed, time.Since(start).Round(time.Millisecond)) return } - // Index stdlib first (definitions only). - if stdlibRoot := s.StdlibRoot(); stdlibRoot != "" { - walkAndIndex(stdlibRoot, false) + if !ok { + return } - - walkAndIndex(s.projectRoot, true) - s.index.writes.RUnlock() - - s.pruneMissingFiles(seen) } // Collapse the WAL back to disk now that the (potentially large) reindex @@ -707,7 +700,10 @@ func (s *Server) RemoveFiles(paths []string) { if s.index.unavailable { return } - if err := s.store.RemoveFiles(paths); err != nil { + // A large removal can rebuild the symbol tables, which takes seconds on a + // large index, so shutdown must be able to stop it. A canceled removal + // leaves the files in the index, and the next start's prune removes them. + if err := s.store.RemoveFilesContext(s.index.work, paths); err != nil && s.index.work.Err() == nil { log.Printf("Error removing %d files from index: %v", len(paths), err) } } @@ -1092,7 +1088,11 @@ func (s *Server) DidSave(ctx context.Context, params *protocol.DidSaveTextDocume if s.workspaceEvents != nil { s.workspaceEvents.ReconcileFile(path) } else { - go s.indexOneFile(path) + // ReconcileFile waits behind a reconciliation pass in flight. With + // only the write coordinator's read lock, the save would race the + // pass's batches for SQLite's write lock and could fail after + // busy_timeout, and then stay unindexed until the next pass. + go s.ReconcileFile(path) } return nil @@ -4271,7 +4271,7 @@ func (s *Server) DidChangeWatchedFiles(ctx context.Context, params *protocol.Did if s.workspaceEvents != nil { s.workspaceEvents.ReconcileFile(path) } else { - go s.indexOneFile(path) + go s.ReconcileFile(path) // see DidSave } case protocol.FileChangeTypeDeleted: if s.workspaceEvents != nil { diff --git a/internal/store/bulk_test.go b/internal/store/bulk_test.go new file mode 100644 index 0000000..3c48574 --- /dev/null +++ b/internal/store/bulk_test.go @@ -0,0 +1,402 @@ +package store + +import ( + "context" + "database/sql" + "errors" + "fmt" + "path/filepath" + "reflect" + "sort" + "strings" + "testing" + "time" + + "github.com/remoteoss/dexter/internal/parser" +) + +// synthRows returns the rows of a generated file. version changes every row, +// so a test can tell an old copy of a file from a new one. +func synthRows(i, version int) ([]parser.Definition, []parser.Reference) { + mod := fmt.Sprintf("MyApp.Mod%d", i) + defs := []parser.Definition{ + {Module: mod, Kind: "module", Line: 1}, + {Module: mod, Function: fmt.Sprintf("run_v%d", version), Arity: 1, Kind: "def", Line: 2, Params: "id"}, + {Module: mod, Function: "helper", Arity: 0, Kind: "defp", Line: 3 + version}, + } + refs := []parser.Reference{ + {Module: "SharedLib.Worker", Function: fmt.Sprintf("call%d", i%7), Line: 2, Kind: "call"}, + {Module: "SharedLib.Worker", Function: "perform", Line: 3 + version, Kind: "call"}, + {Module: "MyApp.Accounts", Line: 4, Kind: "alias"}, + } + return defs, refs +} + +func synthPath(dir string, i int) string { + return filepath.Join(dir, "lib", fmt.Sprintf("mod%d.ex", i)) +} + +// populate writes files [0, n) at version 0 through one incremental batch. +func populate(t *testing.T, s *Store, dir string, n int) { + t.Helper() + b, err := s.BeginBatch() + if err != nil { + t.Fatal(err) + } + for i := 0; i < n; i++ { + defs, refs := synthRows(i, 0) + if err := b.IndexFileWithMtimeAndRefs(synthPath(dir, i), int64(i+1), defs, refs); err != nil { + t.Fatal(err) + } + } + if err := b.Commit(); err != nil { + t.Fatal(err) + } +} + +// dump returns every row of the index keyed by path rather than file id, so two +// stores that hold the same content compare equal. +func dump(t *testing.T, s *Store, root string) []string { + t.Helper() + var out []string + collect := func(q string) { + rows, err := s.db.Query(q) + if err != nil { + t.Fatal(err) + } + defer func() { _ = rows.Close() }() + cols, _ := rows.Columns() + vals := make([]interface{}, len(cols)) + ptrs := make([]interface{}, len(cols)) + for i := range vals { + ptrs[i] = &vals[i] + } + for rows.Next() { + if err := rows.Scan(ptrs...); err != nil { + t.Fatal(err) + } + parts := make([]string, len(vals)) + for i, v := range vals { + if b, ok := v.([]byte); ok { + v = string(b) + } + parts[i] = strings.TrimPrefix(fmt.Sprint(v), root) + } + out = append(out, strings.Join(parts, "|")) + } + if err := rows.Err(); err != nil { + t.Fatal(err) + } + } + collect("SELECT 'file', path, mtime FROM files") + collect("SELECT 'def', f.path, d.module, d.function, d.arity, d.kind, d.line, d.delegate_to, d.delegate_as, d.params FROM definitions d JOIN files f ON f.id = d.file_id") + collect("SELECT 'ref', f.path, r.module, r.function, r.line, r.kind FROM refs r JOIN files f ON f.id = r.file_id") + collect("SELECT 'orphan-def', file_id FROM definitions WHERE file_id NOT IN (SELECT id FROM files)") + collect("SELECT 'orphan-ref', file_id FROM refs WHERE file_id NOT IN (SELECT id FROM files)") + sort.Strings(out) + return out +} + +func schemaObjects(t *testing.T, s *Store) []string { + t.Helper() + rows, err := s.db.Query("SELECT type || ' ' || name FROM sqlite_master WHERE name NOT LIKE 'sqlite_%' ORDER BY 1") + if err != nil { + t.Fatal(err) + } + defer func() { _ = rows.Close() }() + var out []string + for rows.Next() { + var v string + if err := rows.Scan(&v); err != nil { + t.Fatal(err) + } + out = append(out, v) + } + return out +} + +func withRebuildThreshold(t *testing.T, minFiles, share int) { + t.Helper() + oldMin, oldShare := rebuildMinFiles, rebuildShare + rebuildMinFiles, rebuildShare = minFiles, share + t.Cleanup(func() { rebuildMinFiles, rebuildShare = oldMin, oldShare }) +} + +// Both removal strategies must leave exactly the same index: the rows of every +// other file, no orphans, and every table and index in place. +func TestRemoveFiles_DeleteAndRebuildAgree(t *testing.T) { + const n = 2000 + var results [][]string + for _, mode := range []string{"delete", "rebuild"} { + t.Run(mode, func(t *testing.T) { + if mode == "delete" { + withRebuildThreshold(t, 1<<30, 4) + } else { + withRebuildThreshold(t, 1, 4) + } + s, dir := setupTestStore(t) + populate(t, s, dir, n) + before := schemaObjects(t, s) + + // More than one id chunk, plus a path that was never indexed. + var remove []string + for i := 1; i < n; i += 2 { + remove = append(remove, synthPath(dir, i)) + } + remove = append(remove, filepath.Join(dir, "lib", "never_indexed.ex")) + if err := s.RemoveFiles(remove); err != nil { + t.Fatal(err) + } + + got := dump(t, s, dir) + for _, row := range got { + if strings.HasPrefix(row, "orphan") { + t.Fatalf("removal left an orphan row: %s", row) + } + } + if after := schemaObjects(t, s); !reflect.DeepEqual(after, before) { + t.Errorf("schema changed:\nbefore %v\nafter %v", before, after) + } + if r, _ := s.LookupFunction("MyApp.Mod3", "run_v0"); len(r) != 0 { + t.Error("a removed file's definition is still found") + } + if r, _ := s.LookupFunction("MyApp.Mod4", "run_v0"); len(r) != 1 { + t.Errorf("a kept file's definition: got %d results, want 1", len(r)) + } + results = append(results, got) + }) + } + if len(results) == 2 && !reflect.DeepEqual(results[0], results[1]) { + t.Error("the delete and rebuild strategies left different rows") + } +} + +// A rebuild batch must produce the same index as the incremental batch that +// it replaces for large change sets: changed files replaced, new files added, +// untouched files copied, and existing files keeping their id. +func TestRebuildMatchesIncremental(t *testing.T) { + const n = 120 + apply := func(t *testing.T, rebuild bool) ([]string, *Store, string) { + s, dir := setupTestStore(t) + populate(t, s, dir, n) + var b *Batch + var err error + if rebuild { + b, err = s.BeginRebuild(context.Background()) + } else { + b, err = s.BeginBatch() + } + if err != nil { + t.Fatal(err) + } + // Change every other existing file and add as many new ones. + for i := 0; i < n*2; i += 2 { + defs, refs := synthRows(i, 1) + if err := b.IndexFileWithMtimeAndRefs(synthPath(dir, i), int64(1000+i), defs, refs); err != nil { + t.Fatal(err) + } + } + if err := b.Commit(); err != nil { + t.Fatal(err) + } + return dump(t, s, dir), s, dir + } + + incremental, _, _ := apply(t, false) + rebuilt, s, dir := apply(t, true) + if !reflect.DeepEqual(incremental, rebuilt) { + t.Fatalf("rebuild and incremental results differ (%d vs %d rows)", len(rebuilt), len(incremental)) + } + + names, err := s.IndexNames() + if err != nil { + t.Fatal(err) + } + if len(names) != 6 { + t.Errorf("indexes after rebuild: %v", names) + } + var id int64 + if err := s.db.QueryRow("SELECT id FROM files WHERE path = ?", synthPath(dir, 1)).Scan(&id); err != nil || id != 2 { + t.Errorf("an untouched file changed id: %d (%v)", id, err) + } + if r, _ := s.LookupFunction("MyApp.Mod2", "run_v1"); len(r) != 1 { + t.Error("changed file does not have its new rows") + } + if r, _ := s.LookupFunction("MyApp.Mod2", "run_v0"); len(r) != 0 { + t.Error("changed file kept its old rows") + } +} + +func TestRebuildRejectsSecondWriteOfAFile(t *testing.T) { + s, dir := setupTestStore(t) + b, err := s.BeginRebuild(context.Background()) + if err != nil { + t.Fatal(err) + } + defer func() { _ = b.Rollback() }() + defs, refs := synthRows(1, 0) + if err := b.IndexFileWithMtimeAndRefs(synthPath(dir, 1), 1, defs, refs); err != nil { + t.Fatal(err) + } + if err := b.IndexFileWithMtimeAndRefs(synthPath(dir, 1), 2, defs, refs); err == nil { + t.Fatal("a second write of one path in a rebuild was accepted") + } +} + +// A canceled rebuild rolls back completely: the old rows, tables, and indexes +// stay, and no scratch table is left behind. +func TestRebuildCanceledLeavesIndexUnchanged(t *testing.T) { + s, dir := setupTestStore(t) + populate(t, s, dir, 50) + before := dump(t, s, dir) + schema := schemaObjects(t, s) + + ctx, cancel := context.WithCancel(context.Background()) + b, err := s.BeginRebuild(ctx) + if err != nil { + t.Fatal(err) + } + defs, refs := synthRows(1, 9) + if err := b.IndexFileWithMtimeAndRefs(synthPath(dir, 1), 99, defs, refs); err != nil { + t.Fatal(err) + } + cancel() + if err := b.Commit(); err == nil { + t.Fatal("commit of a canceled rebuild succeeded") + } + + if got := dump(t, s, dir); !reflect.DeepEqual(got, before) { + t.Error("a canceled rebuild changed the index") + } + if got := schemaObjects(t, s); !reflect.DeepEqual(got, schema) { + t.Errorf("a canceled rebuild changed the schema: %v", got) + } +} + +func TestRemoveFilesContextCanceledRemovesNothing(t *testing.T) { + withRebuildThreshold(t, 1, 4) + s, dir := setupTestStore(t) + populate(t, s, dir, 20) + before := dump(t, s, dir) + ctx, cancel := context.WithCancel(context.Background()) + cancel() + err := s.RemoveFilesContext(ctx, []string{synthPath(dir, 1), synthPath(dir, 2)}) + if !errors.Is(err, context.Canceled) { + t.Fatalf("err = %v, want context.Canceled", err) + } + if got := dump(t, s, dir); !reflect.DeepEqual(got, before) { + t.Error("a canceled removal changed the index") + } +} + +// Incremental batches buffer rows across files. A file written twice in one +// batch must still end with only its last version. +func TestIncrementalBatchSecondWriteReplacesBufferedRows(t *testing.T) { + s, dir := setupTestStore(t) + b, err := s.BeginBatch() + if err != nil { + t.Fatal(err) + } + for v := 0; v < 2; v++ { + defs, refs := synthRows(1, v) + if err := b.IndexFileWithMtimeAndRefs(synthPath(dir, 1), int64(v+1), defs, refs); err != nil { + t.Fatal(err) + } + } + if err := b.Commit(); err != nil { + t.Fatal(err) + } + if r, _ := s.LookupFunction("MyApp.Mod1", "run_v0"); len(r) != 0 { + t.Error("the first version's rows survived the second write") + } + var defs, refs int + if err := s.db.QueryRow("SELECT (SELECT COUNT(*) FROM definitions), (SELECT COUNT(*) FROM refs)").Scan(&defs, &refs); err != nil { + t.Fatal(err) + } + if defs != 3 || refs != 3 { + t.Errorf("rows = %d defs, %d refs; want 3 and 3", defs, refs) + } +} + +func TestFileStates(t *testing.T) { + s, dir := setupTestStore(t) + populate(t, s, dir, 3) + states, err := s.FileStates() + if err != nil { + t.Fatal(err) + } + if len(states) != 3 { + t.Fatalf("got %d states, want 3", len(states)) + } + st := states[synthPath(dir, 2)] + if st.Mtime != 3 || st.ID != 3 { + t.Errorf("state = %+v, want id 3, mtime 3", st) + } +} + +// holdWriteLock opens a second connection to the same database and holds a +// write transaction on it for d. It returns when the lock is held. +func holdWriteLock(t *testing.T, dir string, d time.Duration) <-chan struct{} { + t.Helper() + other, err := sql.Open(driverName, DBPath(dir)+"?_journal_mode=WAL&_busy_timeout=5000") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = other.Close() }) + tx, err := other.Begin() + if err != nil { + t.Fatal(err) + } + if _, err := tx.Exec("INSERT OR REPLACE INTO metadata (key, value) VALUES ('other_writer', 'x')"); err != nil { + t.Fatal(err) + } + released := make(chan struct{}) + go func() { + time.Sleep(d) + _ = tx.Commit() + close(released) + }() + return released +} + +// A rebuild that starts while another connection writes must wait for it, as +// every other writer does, and not fail at once with "database is locked". +func TestRebuildWaitsForConcurrentWriter(t *testing.T) { + s, dir := setupTestStore(t) + populate(t, s, dir, 20) + released := holdWriteLock(t, dir, 300*time.Millisecond) + + b, err := s.BeginRebuild(context.Background()) + if err != nil { + t.Fatalf("BeginRebuild with a concurrent writer: %v", err) + } + select { + case <-released: + default: + t.Error("BeginRebuild got the write lock while another connection held it") + } + defs, refs := synthRows(1, 5) + if err := b.IndexFileWithMtimeAndRefs(synthPath(dir, 1), 50, defs, refs); err != nil { + t.Fatal(err) + } + if err := b.Commit(); err != nil { + t.Fatal(err) + } + if r, _ := s.LookupFunction("MyApp.Mod1", "run_v5"); len(r) != 1 { + t.Error("the rebuild did not write its file") + } + + // A removal large enough to rebuild also waits and succeeds. + withRebuildThreshold(t, 1, 4) + holdWriteLock(t, dir, 300*time.Millisecond) + var remove []string + for i := 0; i < 15; i++ { + remove = append(remove, synthPath(dir, i)) + } + if err := s.RemoveFiles(remove); err != nil { + t.Fatalf("removal with a concurrent writer: %v", err) + } + if paths, _ := s.ListFilePaths(); len(paths) != 5 { + t.Errorf("%d files left after the removal, want 5", len(paths)) + } +} diff --git a/internal/store/store.go b/internal/store/store.go index ac6c873..9ccd0fc 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -1,8 +1,11 @@ package store import ( + "context" "database/sql" + "errors" "fmt" + "log" "os" "path/filepath" "sort" @@ -255,8 +258,33 @@ func migrate(db *sql.DB) error { path TEXT NOT NULL UNIQUE, mtime INTEGER NOT NULL ); + ` + createDefinitionsTable("definitions") + createRefsTable("refs") + ` + CREATE TABLE IF NOT EXISTS metadata ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL + ); + `) + if err != nil { + return err + } + // idx_refs_function_kind was retired: no query leads with `function`, so + // SQLite never chose it (the two queries that filter on function/kind both + // lead with file_path and use idx_refs_file_path). On a 3.9M-row index it + // cost 80 MB and a share of every index rebuild. Drop it from databases + // that still carry it; this changes no query plan. + if _, err := db.Exec(`DROP INDEX IF EXISTS idx_refs_function_kind`); err != nil { + return err + } + return createIndexes(db) +} - CREATE TABLE IF NOT EXISTS definitions ( +// createDefinitionsTable and createRefsTable return the DDL of the two symbol +// tables under the given name. migrate creates them, and rebuildWithout creates +// a copy with the same shape before it swaps the copy in. One definition keeps +// the two from drifting apart. +func createDefinitionsTable(name string) string { + return ` + CREATE TABLE IF NOT EXISTS ` + name + ` ( module TEXT NOT NULL, function TEXT NOT NULL DEFAULT '', arity INTEGER NOT NULL DEFAULT 0, @@ -268,8 +296,12 @@ func migrate(db *sql.DB) error { params TEXT NOT NULL DEFAULT '', FOREIGN KEY (file_id) REFERENCES files(id) ON DELETE CASCADE ); + ` +} - CREATE TABLE IF NOT EXISTS refs ( +func createRefsTable(name string) string { + return ` + CREATE TABLE IF NOT EXISTS ` + name + ` ( module TEXT NOT NULL, function TEXT NOT NULL DEFAULT '', line INTEGER NOT NULL, @@ -281,24 +313,7 @@ func migrate(db *sql.DB) error { -- Every path that removes a file deletes its refs explicitly first, -- so the cascade was never the thing keeping them consistent. ); - - CREATE TABLE IF NOT EXISTS metadata ( - key TEXT PRIMARY KEY, - value TEXT NOT NULL - ); - `) - if err != nil { - return err - } - // idx_refs_function_kind was retired: no query leads with `function`, so - // SQLite never chose it (the two queries that filter on function/kind both - // lead with file_path and use idx_refs_file_path). On a 3.9M-row index it - // cost 80 MB and a share of every index rebuild. Drop it from databases - // that still carry it; this changes no query plan. - if _, err := db.Exec(`DROP INDEX IF EXISTS idx_refs_function_kind`); err != nil { - return err - } - return createIndexes(db) + ` } // fileIDSubquery resolves a path argument to files.id inside a WHERE clause, so @@ -311,7 +326,11 @@ type dbExecer interface { } func createIndexes(db dbExecer) error { - _, err := db.Exec(` + _, err := db.Exec(createIndexesSQL) + return err +} + +const createIndexesSQL = ` CREATE INDEX IF NOT EXISTS idx_definitions_module_function ON definitions(module, function); CREATE INDEX IF NOT EXISTS idx_definitions_file_id_line ON definitions(file_id, line); CREATE INDEX IF NOT EXISTS idx_refs_module_function ON refs(module, function, file_id, line, kind); @@ -323,9 +342,7 @@ func createIndexes(db dbExecer) error { -- leads with module and cannot be seeked by function alone. The partial -- index holds only the __using__ rows, so the scan becomes a small range. CREATE INDEX IF NOT EXISTS idx_definitions_using ON definitions(module, file_id) WHERE function = '__using__'; - `) - return err -} + ` func (s *Store) DropIndexes() error { _, err := s.db.Exec(` @@ -515,28 +532,53 @@ const ( refChunkRows = maxBindVars / refColumns // 180 ) +// batchMode selects how a Batch writes. +type batchMode uint8 + +const ( + // batchIncremental replaces each file's rows in the live tables: an upsert + // of the files row, a DELETE by file id, and buffered INSERTs. + batchIncremental batchMode = iota + // batchInsertOnly fills empty tables whose indexes the caller has dropped. + batchInsertOnly + // batchRebuild writes into new, unindexed copies of the symbol tables and + // swaps them in at Commit. See BeginRebuild. + batchRebuild +) + // Batch wraps multiple IndexFile operations in a single SQLite transaction // with shared prepared statements. type Batch struct { + ctx context.Context tx *sql.Tx + mode batchMode defStmt *sql.Stmt refStmt *sql.Stmt fileStmt *sql.Stmt - delDefStmt *sql.Stmt // nil in insert-only mode - delRefStmt *sql.Stmt // nil in insert-only mode - insertOnly bool - - // Multi-row INSERT buffers, used in insert-only mode only. A cold index - // writes ~4.4M rows through one connection, and the writer is the - // bottleneck of the whole indexing pipeline; batching turns ~4.4M cgo - // crossings into a few tens of thousands. Incremental reindexing keeps the - // row-at-a-time path, where a file's DELETE must stay ordered ahead of its - // INSERTs and the row count is far too small to matter. + delDefStmt *sql.Stmt // incremental mode only + delRefStmt *sql.Stmt // incremental mode only + skipStmt *sql.Stmt // rebuild mode only + + // Multi-row INSERT buffers. A cold index writes ~4.4M rows through one + // connection, and the writer is the bottleneck of the whole indexing + // pipeline; batching turns ~4.4M cgo crossings into a few tens of + // thousands. A large warm change set (a new checkout of a whole tree, for + // example) writes as many rows through the other modes, so they buffer + // too. Rows buffered for one file never meet another file's DELETE, which + // targets only that file's id; a file seen twice in one incremental batch + // flushes the buffer first, so its DELETE stays ordered ahead of its earlier + // INSERTs. defChunkStmt *sql.Stmt refChunkStmt *sql.Stmt defArgs []interface{} refArgs []interface{} + // written holds the file ids this batch has written, in the incremental and + // rebuild modes. An incremental batch flushes before a second write of one + // file; a rebuild rejects it, because its new tables have no index to find + // the earlier rows by. + written map[int64]struct{} + // Bulk-path file id allocation. Insert-only mode assigns ids in Go from a // counter and writes files rows with an explicit id, so a cold index never // pays a round trip per file to learn what id it just wrote. @@ -563,97 +605,125 @@ func multiRowInsert(table, columns string, n, rows int) string { } func (s *Store) BeginBatch() (*Batch, error) { - return s.beginBatch(false) + return s.beginBatch(context.Background(), batchIncremental) +} + +// BeginBatchContext starts an incremental batch whose transaction rolls back +// when ctx is canceled. Writes after that fail, and Commit reports the +// cancellation, so none of the batch's files is left half written. +func (s *Store) BeginBatchContext(ctx context.Context) (*Batch, error) { + return s.beginBatch(ctx, batchIncremental) } // BeginBulkInsert starts a batch optimized for inserting into an empty table. // It skips DELETE statements before each insert. Callers should drop indexes // before calling this and recreate them after Commit. func (s *Store) BeginBulkInsert() (*Batch, error) { - return s.beginBatch(true) + return s.beginBatch(context.Background(), batchInsertOnly) } -func (s *Store) beginBatch(insertOnly bool) (*Batch, error) { - tx, err := s.db.Begin() +// BeginRebuild starts a batch that rewrites the symbol tables in one +// transaction, for a change set too large to apply row by row to the live +// indexes. Each index of the live tables sorts rows by name, so the rows of one +// file land on pages all over it, and a write of a large share of the index +// touches most of its pages, many times. +// +// Files written to the batch go into new tables that have no indexes. Commit +// then copies in the rows of every file the batch did not write or remove, +// swaps the new tables in for the old ones, and builds each index with one +// sort, as a cold build does. Readers keep their snapshot of the old tables +// until the commit, and a canceled or failed rebuild rolls back to exactly the +// old state. Each path may be written only once. +func (s *Store) BeginRebuild(ctx context.Context) (*Batch, error) { + return s.beginBatch(ctx, batchRebuild) +} + +const ( + defInsertColumns = "module, function, arity, kind, line, file_id, delegate_to, delegate_as, params" + refInsertColumns = "module, function, line, file_id, kind" +) + +func (s *Store) beginBatch(ctx context.Context, mode batchMode) (_ *Batch, err error) { + tx, err := s.db.BeginTx(ctx, nil) if err != nil { return nil, err } + b := &Batch{ctx: ctx, tx: tx, mode: mode} + defer func() { + if err != nil { + b.closeStmts() + _ = tx.Rollback() + } + }() + + defTable, refTable := "definitions", "refs" + if mode == batchRebuild { + defTable, refTable = "definitions_rebuild", "refs_rebuild" + // Take the write lock first. The transaction is deferred, so the + // DROP below would open it as a read, and the CREATE after it would + // have to upgrade that read to a write. SQLite does not call the busy + // handler for that upgrade: with another writer active, the rebuild + // would fail at once with "database is locked" instead of waiting + // for busy_timeout. A write statement as the first statement waits + // like every other writer. + if _, err = tx.ExecContext(ctx, "DELETE FROM metadata WHERE key = 'rebuild_lock'"); err != nil { + return nil, err + } + // A rolled-back rebuild leaves nothing behind, but a copy from a + // process that died mid-way would make CREATE ... IF NOT EXISTS keep + // its rows. + if _, err = tx.ExecContext(ctx, ` + DROP TABLE IF EXISTS temp.rebuild_skip; + CREATE TEMP TABLE rebuild_skip (id INTEGER PRIMARY KEY, removed INTEGER NOT NULL DEFAULT 0); + DROP TABLE IF EXISTS definitions_rebuild; + DROP TABLE IF EXISTS refs_rebuild; + `+createDefinitionsTable(defTable)+createRefsTable(refTable)); err != nil { + return nil, err + } + if b.skipStmt, err = tx.Prepare("INSERT OR IGNORE INTO temp.rebuild_skip (id) VALUES (?)"); err != nil { + return nil, err + } + } - defStmt, err := tx.Prepare("INSERT INTO definitions (module, function, arity, kind, line, file_id, delegate_to, delegate_as, params) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)") - if err != nil { - _ = tx.Rollback() + if b.defStmt, err = tx.Prepare("INSERT INTO " + defTable + " (" + defInsertColumns + ") VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)"); err != nil { return nil, err } - - refStmt, err := tx.Prepare("INSERT INTO refs (module, function, line, file_id, kind) VALUES (?, ?, ?, ?, ?)") - if err != nil { - _ = defStmt.Close() - _ = tx.Rollback() + if b.refStmt, err = tx.Prepare("INSERT INTO " + refTable + " (" + refInsertColumns + ") VALUES (?, ?, ?, ?, ?)"); err != nil { return nil, err } - fileSQL := "INSERT INTO files (path, mtime) VALUES (?, ?) ON CONFLICT(path) DO UPDATE SET mtime = excluded.mtime RETURNING id" - if insertOnly { + if mode == batchInsertOnly { fileSQL = "INSERT INTO files (id, path, mtime) VALUES (?, ?, ?)" } - fileStmt, err := tx.Prepare(fileSQL) - if err != nil { - _ = defStmt.Close() - _ = refStmt.Close() - _ = tx.Rollback() + if b.fileStmt, err = tx.Prepare(fileSQL); err != nil { return nil, err } - - b := &Batch{ - tx: tx, - defStmt: defStmt, - refStmt: refStmt, - fileStmt: fileStmt, - insertOnly: insertOnly, - } - - if !insertOnly { - b.delDefStmt, err = tx.Prepare("DELETE FROM definitions WHERE file_id = ?") - if err != nil { - b.closeStmts() - _ = tx.Rollback() + if mode == batchIncremental { + if b.delDefStmt, err = tx.Prepare("DELETE FROM definitions WHERE file_id = ?"); err != nil { return nil, err } - b.delRefStmt, err = tx.Prepare("DELETE FROM refs WHERE file_id = ?") - if err != nil { - b.closeStmts() - _ = tx.Rollback() + if b.delRefStmt, err = tx.Prepare("DELETE FROM refs WHERE file_id = ?"); err != nil { return nil, err } - } else { - b.defChunkStmt, err = tx.Prepare(multiRowInsert( - "definitions", - "module, function, arity, kind, line, file_id, delegate_to, delegate_as, params", - defColumns, defChunkRows)) - if err != nil { - b.closeStmts() - _ = tx.Rollback() - return nil, err - } - b.refChunkStmt, err = tx.Prepare(multiRowInsert( - "refs", "module, function, line, file_id, kind", refColumns, refChunkRows)) - if err != nil { - b.closeStmts() - _ = tx.Rollback() - return nil, err - } - b.defArgs = make([]interface{}, 0, defColumns*defChunkRows) - b.refArgs = make([]interface{}, 0, refColumns*refChunkRows) - + } + if mode != batchInsertOnly { + b.written = make(map[int64]struct{}) + } + if b.defChunkStmt, err = tx.Prepare(multiRowInsert(defTable, defInsertColumns, defColumns, defChunkRows)); err != nil { + return nil, err + } + if b.refChunkStmt, err = tx.Prepare(multiRowInsert(refTable, refInsertColumns, refColumns, refChunkRows)); err != nil { + return nil, err + } + b.defArgs = make([]interface{}, 0, defColumns*defChunkRows) + b.refArgs = make([]interface{}, 0, refColumns*refChunkRows) + if mode == batchInsertOnly { // Bulk mode is used on a freshly created database, but seed from the // table anyway so a non-empty one cannot collide on the primary key. - if err := tx.QueryRow("SELECT COALESCE(MAX(id), 0) FROM files").Scan(&b.nextFileID); err != nil { - b.closeStmts() - _ = tx.Rollback() + if err = tx.QueryRow("SELECT COALESCE(MAX(id), 0) FROM files").Scan(&b.nextFileID); err != nil { return nil, err } } - return b, nil } @@ -679,54 +749,75 @@ func (b *Batch) indexFile(path string, mtimeNano int64, defs []parser.Definition return err } - if !b.insertOnly { + switch b.mode { + case batchIncremental: + if _, again := b.written[fileID]; again { + if err := b.flushPending(); err != nil { + return err + } + } + b.written[fileID] = struct{}{} if _, err := b.delDefStmt.Exec(fileID); err != nil { return err } if _, err := b.delRefStmt.Exec(fileID); err != nil { return err } + case batchRebuild: + if _, again := b.written[fileID]; again { + return fmt.Errorf("rebuild: %s written twice", path) + } + b.written[fileID] = struct{}{} + // The old rows of this file are not copied into the new tables. + if _, err := b.skipStmt.Exec(fileID); err != nil { + return err + } } - if b.insertOnly { - // Reset the buffer before checking the error: the chunk boundary is an - // exact-multiple test, so a buffer left full would never match again and - // the rest of the batch would silently fall back to flushPending. - for _, d := range defs { - b.defArgs = append(b.defArgs, d.Module, d.Function, d.Arity, d.Kind, d.Line, fileID, d.DelegateTo, d.DelegateAs, d.Params) - if len(b.defArgs) == defColumns*defChunkRows { - _, err := b.defChunkStmt.Exec(b.defArgs...) - b.defArgs = b.defArgs[:0] - if err != nil { - return err - } + // Reset the buffer before checking the error: the chunk boundary is an + // exact-multiple test, so a buffer left full would never match again and + // the rest of the batch would silently fall back to flushPending. + for _, d := range defs { + b.defArgs = append(b.defArgs, d.Module, d.Function, d.Arity, d.Kind, d.Line, fileID, d.DelegateTo, d.DelegateAs, d.Params) + if len(b.defArgs) == defColumns*defChunkRows { + _, err := b.defChunkStmt.Exec(b.defArgs...) + b.defArgs = b.defArgs[:0] + if err != nil { + return err } } - for _, r := range refs { - b.refArgs = append(b.refArgs, r.Module, r.Function, r.Line, fileID, r.Kind) - if len(b.refArgs) == refColumns*refChunkRows { - _, err := b.refChunkStmt.Exec(b.refArgs...) - b.refArgs = b.refArgs[:0] - if err != nil { - return err - } + } + for _, r := range refs { + b.refArgs = append(b.refArgs, r.Module, r.Function, r.Line, fileID, r.Kind) + if len(b.refArgs) == refColumns*refChunkRows { + _, err := b.refChunkStmt.Exec(b.refArgs...) + b.refArgs = b.refArgs[:0] + if err != nil { + return err } } - return nil } + return nil +} - for _, d := range defs { - if _, err := b.defStmt.Exec(d.Module, d.Function, d.Arity, d.Kind, d.Line, fileID, d.DelegateTo, d.DelegateAs, d.Params); err != nil { - return err - } +// RemoveIDs marks files for removal in a rebuild: their rows are not copied +// into the new tables, and their files rows go at Commit. +func (b *Batch) RemoveIDs(ids []int64) error { + if b.mode != batchRebuild { + return errors.New("RemoveIDs needs a rebuild batch") } - - for _, r := range refs { - if _, err := b.refStmt.Exec(r.Module, r.Function, r.Line, fileID, r.Kind); err != nil { + for start := 0; start < len(ids); start += maxBindVars { + chunk := ids[start:min(start+maxBindVars, len(ids))] + args := make([]interface{}, len(chunk)) + for i, id := range chunk { + args[i] = id + } + if _, err := b.tx.ExecContext(b.ctx, + "INSERT OR REPLACE INTO temp.rebuild_skip (id, removed) VALUES "+strings.TrimSuffix(strings.Repeat("(?,1),", len(chunk)), ","), + args...); err != nil { return err } } - return nil } @@ -735,10 +826,10 @@ func (b *Batch) indexFile(path string, mtimeNano int64, defs []parser.Definition // The bulk path allocates ids from a counter and inserts them explicitly: a // cold index writes ~70k files, and asking SQLite to hand back each id would // add a round trip per file to the one thread that is already the bottleneck. -// The incremental path upserts, so an existing file keeps the id that its +// The other modes upsert, so an existing file keeps the id that its // definitions and refs already point at. func (b *Batch) fileID(path string, mtimeNano int64) (int64, error) { - if b.insertOnly { + if b.mode == batchInsertOnly { b.nextFileID++ if _, err := b.fileStmt.Exec(b.nextFileID, path, mtimeNano); err != nil { return 0, err @@ -769,35 +860,55 @@ func (b *Batch) flushPending() error { } func (b *Batch) Commit() error { - if err := b.flushPending(); err != nil { - b.closeStmts() + err := b.flushPending() + b.closeStmts() + if err == nil && b.mode == batchRebuild { + err = b.finishRebuild() + } + if err != nil { _ = b.tx.Rollback() return err } - b.closeStmts() return b.tx.Commit() } +// finishRebuild copies the rows of every file the rebuild did not write or +// remove into the new tables, swaps them in, builds the indexes, and deletes +// the removed files. Every statement runs under the batch context, so a cancel +// interrupts the one in progress. +// +// Dropping a table that only references files (definitions) or nothing (refs) +// does no foreign key work, and the indexes go with their tables. +func (b *Batch) finishRebuild() error { + const keep = " WHERE file_id NOT IN (SELECT id FROM temp.rebuild_skip)" + for _, q := range []string{ + "INSERT INTO definitions_rebuild (" + defInsertColumns + ") SELECT " + defInsertColumns + " FROM definitions" + keep, + "DROP TABLE definitions", + "ALTER TABLE definitions_rebuild RENAME TO definitions", + "INSERT INTO refs_rebuild (" + refInsertColumns + ") SELECT " + refInsertColumns + " FROM refs" + keep, + "DROP TABLE refs", + "ALTER TABLE refs_rebuild RENAME TO refs", + createIndexesSQL, + "DELETE FROM files WHERE id IN (SELECT id FROM temp.rebuild_skip WHERE removed = 1)", + "DROP TABLE temp.rebuild_skip", + } { + if _, err := b.tx.ExecContext(b.ctx, q); err != nil { + return err + } + } + return nil +} + func (b *Batch) Rollback() error { b.closeStmts() return b.tx.Rollback() } func (b *Batch) closeStmts() { - _ = b.defStmt.Close() - _ = b.refStmt.Close() - _ = b.fileStmt.Close() - if b.defChunkStmt != nil { - _ = b.defChunkStmt.Close() - } - if b.refChunkStmt != nil { - _ = b.refChunkStmt.Close() - } - if b.delDefStmt != nil { - _ = b.delDefStmt.Close() - } - if b.delRefStmt != nil { - _ = b.delRefStmt.Close() + for _, st := range []*sql.Stmt{b.defStmt, b.refStmt, b.fileStmt, b.defChunkStmt, b.refChunkStmt, b.delDefStmt, b.delRefStmt, b.skipStmt} { + if st != nil { + _ = st.Close() + } } } @@ -865,30 +976,191 @@ func (s *Store) RemoveFile(path string) error { return s.RemoveFiles([]string{path}) } +// RemoveFiles removes the given paths and every row that belongs to them. +// Paths that are not indexed are ignored. func (s *Store) RemoveFiles(paths []string) error { + return s.RemoveFilesContext(context.Background(), paths) +} + +// RemoveFilesContext is RemoveFiles with cancellation. It resolves the paths to +// file ids with one query per chunk and removes the ids as sets; see +// RemoveFileIDs. +func (s *Store) RemoveFilesContext(ctx context.Context, paths []string) error { if len(paths) == 0 { return nil } - tx, err := s.db.Begin() + ids := make([]int64, 0, len(paths)) + for start := 0; start < len(paths); start += maxBindVars { + chunk := paths[min(start, len(paths)):min(start+maxBindVars, len(paths))] + args := make([]interface{}, len(chunk)) + for i, p := range chunk { + args[i] = p + } + rows, err := s.db.QueryContext(ctx, "SELECT id FROM files WHERE path IN ("+placeholders(len(chunk))+")", args...) + if err != nil { + return err + } + for rows.Next() { + var id int64 + if err := rows.Scan(&id); err != nil { + _ = rows.Close() + return err + } + ids = append(ids, id) + } + err = rows.Err() + _ = rows.Close() + if err != nil { + return err + } + } + return s.RemoveFileIDs(ctx, ids) +} + +// FileState is what the index records about one file: its row id and the +// modification time the stored symbols were parsed from. +type FileState struct { + ID int64 + Mtime int64 +} + +// FileStates returns every indexed file with its id and stored mtime, read in +// one query. A reconciliation pass compares the walk against it instead of +// issuing one query per file, and removes the files the walk did not see by id. +func (s *Store) FileStates() (map[string]FileState, error) { + var n int + if err := s.db.QueryRow("SELECT COUNT(*) FROM files").Scan(&n); err != nil { + return nil, err + } + rows, err := s.db.Query("SELECT path, id, mtime FROM files") if err != nil { - return err + return nil, err } - defer func() { _ = tx.Rollback() }() + defer func() { _ = rows.Close() }() - for _, path := range paths { - if _, err = tx.Exec("DELETE FROM definitions WHERE file_id = "+fileIDSubquery, path); err != nil { - return err + states := make(map[string]FileState, n) + for rows.Next() { + var path string + var st FileState + if err := rows.Scan(&path, &st.ID, &st.Mtime); err != nil { + return nil, err } - if _, err = tx.Exec("DELETE FROM refs WHERE file_id = "+fileIDSubquery, path); err != nil { + states[path] = st + } + return states, rows.Err() +} + +// Removal strategy. Deleting a file's rows costs one random B-tree update per +// row in every index on the table: idx_refs_module_function sorts by name, so +// the refs of one file are spread over the whole index. Removing a nested +// checkout of a large monorepo is millions of such updates. Copying the rows +// that stay into a new table and indexing it afresh costs one sequential read of +// the table and one sort per index instead, which is far cheaper once a large +// share of the rows goes. Both are measured on a 280k-file index, where removing +// 223k files took 31s as set-based DELETEs and 15s as a rebuild (and 66s as the +// former per-path DELETEs). +var ( + // rebuildMinFiles keeps small removals on the DELETE path, whatever the + // share: a rebuild always reads both tables in full. A variable so tests + // can reach the rebuild with a small index. + rebuildMinFiles = 4096 + // rebuildShare is the share of indexed files, as 1/rebuildShare, at or + // above which a removal rebuilds the symbol tables. + rebuildShare = 4 +) + +const ( + // deleteGroupIDs is how many files one DELETE transaction removes. Each + // group commits on its own, so a canceled removal keeps the groups already + // done, and a group is small enough to commit quickly. + deleteGroupIDs = 8 * maxBindVars +) + +// RemoveFileIDs removes the files with the given ids and all their rows. Each +// file goes in one transaction with its definitions and refs, so no file is +// ever left half removed, also when ctx is canceled part way through. +// +// A removal of a large share of the index rebuilds the symbol tables without +// the removed files (see rebuildWithout); a smaller one deletes by id sets. +func (s *Store) RemoveFileIDs(ctx context.Context, ids []int64) error { + if len(ids) == 0 { + return nil + } + if len(ids) >= rebuildMinFiles { + var total int + if err := s.db.QueryRowContext(ctx, "SELECT COUNT(*) FROM files").Scan(&total); err != nil { return err } - if _, err = tx.Exec("DELETE FROM files WHERE path = ?", path); err != nil { + if len(ids)*rebuildShare >= total { + err := s.rebuildWithout(ctx, ids) + if err == nil || ctx.Err() != nil { + return err + } + // The DELETE path is slower but needs nothing the rebuild did, + // so a failed rebuild does not leave the files in the index. + log.Printf("Warning: index rebuild for %d removed files failed, deleting them by id: %v", len(ids), err) + } + } + for start := 0; start < len(ids); start += deleteGroupIDs { + if err := s.deleteFileIDs(ctx, ids[start:min(start+deleteGroupIDs, len(ids))]); err != nil { return err } } + return nil +} + +// deleteFileIDs deletes one group of files in one transaction, by id sets of at +// most maxBindVars. Sorting lets each statement walk idx_refs_file_id and the +// definitions index in key order. +func (s *Store) deleteFileIDs(ctx context.Context, ids []int64) error { + sorted := append([]int64(nil), ids...) + sort.Slice(sorted, func(i, j int) bool { return sorted[i] < sorted[j] }) + + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return err + } + defer func() { _ = tx.Rollback() }() + + for start := 0; start < len(sorted); start += maxBindVars { + chunk := sorted[start:min(start+maxBindVars, len(sorted))] + args := make([]interface{}, len(chunk)) + for i, id := range chunk { + args[i] = id + } + in := "(" + placeholders(len(chunk)) + ")" + for _, q := range []string{ + "DELETE FROM definitions WHERE file_id IN " + in, + "DELETE FROM refs WHERE file_id IN " + in, + "DELETE FROM files WHERE id IN " + in, + } { + if _, err := tx.ExecContext(ctx, q, args...); err != nil { + return err + } + } + } return tx.Commit() } +// rebuildWithout removes the given files by rebuilding the symbol tables +// without them. See BeginRebuild. +func (s *Store) rebuildWithout(ctx context.Context, ids []int64) error { + b, err := s.BeginRebuild(ctx) + if err != nil { + return err + } + if err := b.RemoveIDs(ids); err != nil { + _ = b.Rollback() + return err + } + return b.Commit() +} + +// placeholders returns n comma-separated "?" markers. +func placeholders(n int) string { + return strings.TrimSuffix(strings.Repeat("?,", n), ",") +} + type CompletionResult struct { Module string Function string diff --git a/internal/workspace/runtime.go b/internal/workspace/runtime.go index f2c8eb3..9c5941b 100644 --- a/internal/workspace/runtime.go +++ b/internal/workspace/runtime.go @@ -595,7 +595,14 @@ func (r *Runtime) publishBatch(full bool, paths map[string]struct{}) { r.publish(c) } +// testHookReconcilePath, when set by a test, runs before each path event is +// reconciled. +var testHookReconcilePath func(path string) + func (r *Runtime) reconcilePath(path string) error { + if testHookReconcilePath != nil { + testHookReconcilePath(path) + } info, err := os.Stat(path) if err != nil { if os.IsNotExist(err) { @@ -640,10 +647,17 @@ func (r *Runtime) reconcilePath(path string) error { return nil } -// Close stops event sources, drains accepted mutations, waits for background -// index work, checkpoints, and closes the store. +// Close cancels a reconciliation in flight, stops event sources, drains +// accepted mutations, waits for background index work, checkpoints, and closes +// the store. +// +// The cancel comes first. A warm pass over a large change set can run for +// minutes, and the daemon holds the workspace lock without serving its socket +// until Close returns. A canceled pass leaves every file fully old or fully new, +// and the next start finishes it from the stored mtimes. func (r *Runtime) Close() error { r.closeOnce.Do(func() { + r.index.CancelWork() close(r.watcherStop) r.watcherWG.Wait() <-r.watcherReady diff --git a/internal/workspace/runtime_test.go b/internal/workspace/runtime_test.go index df014ab..7805274 100644 --- a/internal/workspace/runtime_test.go +++ b/internal/workspace/runtime_test.go @@ -13,6 +13,7 @@ import ( "go.lsp.dev/protocol" "github.com/remoteoss/dexter/internal/lsp" + "github.com/remoteoss/dexter/internal/store" "github.com/remoteoss/dexter/internal/version" ) @@ -717,3 +718,135 @@ func TestWatchCoverageTransitionsTriggerOneFullReconcileEach(t *testing.T) { case <-time.After(100 * time.Millisecond): } } + +// Close cancels the initial reconciliation instead of waiting for it, so a +// daemon that is told to stop releases its workspace lock at once. The pass it +// cancels leaves a consistent index, and the next open finishes it. +func TestCloseCancelsInitialReconcileAndReopenFinishes(t *testing.T) { + t.Setenv("PATH", t.TempDir()) + t.Setenv("SHELL", "/bin/false") + root := t.TempDir() + writeTestModule(t, root, "lib/seed.ex", "MyApp.Seed") + rt, err := OpenWithOptions(root, Options{NoWatch: true}) + if err != nil { + t.Fatal(err) + } + if err := rt.WaitReady(testContext(t, 30*time.Second)); err != nil { + t.Fatal(err) + } + if err := rt.Close(); err != nil { + t.Fatal(err) + } + + const files = 2000 + for i := 0; i < files; i++ { + writeTestModule(t, root, fmt.Sprintf("lib/gen/mod%d.ex", i), fmt.Sprintf("MyApp.Gen%d", i)) + } + + rt, err = OpenWithOptions(root, Options{NoWatch: true}) + if err != nil { + t.Fatal(err) + } + start := time.Now() + if err := rt.Close(); err != nil { + t.Fatal(err) + } + if took := time.Since(start); took > 5*time.Second { + t.Errorf("Close took %s during the initial reconciliation", took) + } + + rt, err = OpenWithOptions(root, Options{NoWatch: true}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = rt.Close() }) + if err := rt.WaitReady(testContext(t, 60*time.Second)); err != nil { + t.Fatal(err) + } + paths, err := rt.Store().ListFilePathsUnder(root) + if err != nil { + t.Fatal(err) + } + if len(paths) != files+1 { + t.Errorf("after reopen: %d project files indexed, want %d", len(paths), files+1) + } + for _, module := range []string{"MyApp.Seed", "MyApp.Gen0", fmt.Sprintf("MyApp.Gen%d", files-1)} { + if got := countModule(t, rt, module); got != 1 { + t.Errorf("%s: %d definitions after reopen, want 1", module, got) + } + } +} + +// A removal that is queued or running when Close starts is canceled with the +// rest of the index work, so a large removal (a rebuild of the symbol tables) +// cannot hold shutdown. The rows it did not remove stay consistent, and the +// next open prunes them. +func TestCloseCancelsQueuedRemoval(t *testing.T) { + rt, root := newTestRuntime(t) + gone := filepath.Join(root, "lib", "gone") + for i := 0; i < 20; i++ { + writeTestModule(t, root, fmt.Sprintf("lib/gone/mod%d.ex", i), fmt.Sprintf("MyApp.Gone%d", i)) + } + if err := rt.Reindex(testContext(t, 30*time.Second)); err != nil { + t.Fatal(err) + } + if got := countModule(t, rt, "MyApp.Gone3"); got != 1 { + t.Fatalf("MyApp.Gone3 indexed %d times before the removal, want 1", got) + } + if err := os.RemoveAll(gone); err != nil { + t.Fatal(err) + } + + // Hold the mutation loop at the removal until Close has started. + reached := make(chan struct{}) + release := make(chan struct{}) + testHookReconcilePath = func(path string) { + if path == gone { + close(reached) + <-release + } + } + t.Cleanup(func() { testHookReconcilePath = nil }) + rt.RemoveFile(gone) + <-reached + + closed := make(chan error, 1) + start := time.Now() + go func() { closed <- rt.Close() }() + time.Sleep(100 * time.Millisecond) // Close cancels index work first + close(release) + if err := <-closed; err != nil { + t.Fatal(err) + } + if took := time.Since(start); took > 2*time.Second { + t.Errorf("Close took %s", took) + } + testHookReconcilePath = nil + + // The canceled removal left the rows of the deleted files. + s, err := store.Open(root) + if err != nil { + t.Fatal(err) + } + under, err := s.ListFilePathsUnder(gone) + _ = s.Close() + if err != nil { + t.Fatal(err) + } + if len(under) != 20 { + t.Errorf("%d rows under the removed directory after Close, want 20 (the removal was not canceled)", len(under)) + } + + // The next open prunes them. + reopened, err := OpenWithOptions(root, Options{NoWatch: true}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = reopened.Close() }) + if err := reopened.WaitReady(testContext(t, 30*time.Second)); err != nil { + t.Fatal(err) + } + if got := countModule(t, reopened, "MyApp.Gone3"); got != 0 { + t.Errorf("MyApp.Gone3 still indexed after reopen: %d", got) + } +} From 2fcf5ecbc434a9e4c7a9f81297c027b463522cb6 Mon Sep 17 00:00:00 2001 From: Jesse Herrick Date: Sat, 3 Oct 2026 20:35:26 -0400 Subject: [PATCH 2/2] Skip a rebuild that has no parsed file to write When every changed file of a large change set failed to parse, the rebuild still copied the old tables, swapped them in and built every index, although nothing changed: a file that fails to parse is not written, so its old rows are copied and kept. The next warm start then did the same for the same files. The rebuild now rolls back when no file was written. Co-Authored-By: Claude Opus 5.5 --- internal/lsp/reconcile.go | 8 +++++++ internal/lsp/reconcile_test.go | 40 ++++++++++++++++++++++++++++++++++ 2 files changed, 48 insertions(+) diff --git a/internal/lsp/reconcile.go b/internal/lsp/reconcile.go index 1d20411..9a96166 100644 --- a/internal/lsp/reconcile.go +++ b/internal/lsp/reconcile.go @@ -285,6 +285,14 @@ func (s *Server) rebuildChanged(ctx context.Context, changed []changedFile) (int } return 0, false } + if files == 0 { + // Every changed file failed to parse, so nothing changes: their old + // rows stay. Copying and swapping the tables would only cost time, and + // the next start would do it again for the same files. + _ = batch.Rollback() + log.Printf("Index rebuild skipped: none of %d changed files could be parsed", len(changed)) + return 0, true + } written := time.Now() if err := batch.Commit(); err != nil { if ctx.Err() == nil { diff --git a/internal/lsp/reconcile_test.go b/internal/lsp/reconcile_test.go index 6dd5e38..a3dab16 100644 --- a/internal/lsp/reconcile_test.go +++ b/internal/lsp/reconcile_test.go @@ -132,6 +132,46 @@ func TestReconcile_WarmChangeSetPaths(t *testing.T) { } } +// When every changed file fails to parse, a rebuild has nothing to write. It +// must not copy and swap the tables, and the old rows of those files stay. +func TestReconcile_RebuildWithNoParsedFileChangesNothing(t *testing.T) { + setReconcileVars(t, 3, 1, 4) + logs := captureLog(t) + server, cleanup := setupTestServer(t) + defer cleanup() + + var paths []string + for i := 0; i < 3; i++ { + paths = append(paths, writeTestFile(t, server.projectRoot, fmt.Sprintf("lib/gen%d.ex", i), moduleSource(i, 0))) + } + reindexOnce(t, server) // cold: full build + + future := time.Now().Add(time.Hour) + for _, path := range paths { + if err := os.Chtimes(path, future, future); err != nil { + t.Fatal(err) + } + if err := os.Chmod(path, 0); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.Chmod(path, 0o644) }) + } + before := logs.String() + reindexOnce(t, server) + pass := strings.TrimPrefix(logs.String(), before) + if strings.Contains(pass, "Rebuilt the index") { + t.Errorf("a rebuild with no parsed file swapped the tables:\n%s", pass) + } + if !strings.Contains(pass, "Index rebuild skipped") { + t.Errorf("the log does not say that the rebuild was skipped:\n%s", pass) + } + for i := 0; i < 3; i++ { + if r, _ := server.store.LookupFunction(fmt.Sprintf("MyApp.Gen%d", i), "run_v0"); len(r) != 1 { + t.Errorf("file %d lost its old definition: %d rows", i, len(r)) + } + } +} + // CancelWork must end a warm pass in flight promptly, leave no file half // written, and leave the rest for the next start, which finishes it from the // stored mtimes. This is what shutdown relies on.