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..9a96166 --- /dev/null +++ b/internal/lsp/reconcile.go @@ -0,0 +1,434 @@ +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 + } + 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 { + 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..a3dab16 --- /dev/null +++ b/internal/lsp/reconcile_test.go @@ -0,0 +1,398 @@ +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)) + } + }) + } +} + +// 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. +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 50f464f..b793a54 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) { @@ -642,10 +649,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) + } +}