Răsfoiți Sursa

perf(synthesis): fan dynamic-dispatch passes across the resolver pool, byte-identical graphs

The ~36 independent synthesis passes (callback/event/framework wiring) ran
sequentially on the indexer's main thread — 2.0s of a 4,402-file Java repo's
index, and the stage where kernel-class repos die (#1212). They now live in
an explicit registry (SYNTH_PASSES) and, when the resolver pool is alive
(>=150k-ref repos), fan out across its read-only workers: dubbo synthesis
2,024ms -> ~900ms (-55%), total fresh init 13.5s -> 11.9s. Graphs verified
byte-for-byte identical on both the pool path (dubbo) and the sequential
path (excalidraw).

Why this is safe: no pass's edges persist until the ordered merge, so every
pass sees the same committed post-resolution DB state in either mode, and
results merge in registry order regardless of completion order — the
first-seen dedup is unchanged. The pool now survives through synthesis
(destroy moved after it) instead of being torn down moments before the one
stage that could reuse it.

Robustness: a pass that fails on a worker (crash, OOM) is retried on the
main thread — a synthesizer blow-up now costs one worker instead of the
whole index, which is half the #1212 story on very large repos.

Also: ref-row cleanup deletes now run as one transaction with a cached
statement instead of one implicit commit per 500-row chunk (mechanically
fewer WAL commits; matters most on HDD-class storage). A set-based rewrite
of failed-ref parking was tried, measured ~zero on NVMe, and dropped — the
remaining persist cost is edge-index B-tree maintenance, not statement
dispatch.

SYNTH_PROGRESS_STEPS now derives from the registry (passes + fixed marks);
the pin test counts registry entries plus literal __mark sites.

Suite green (2444). Sequential-path timing unchanged on excalidraw.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Colby McHenry 1 lună în urmă
părinte
comite
27ac0c0126

+ 1 - 0
CHANGELOG.md

@@ -13,6 +13,7 @@ and adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
 
 - Reference resolution now runs in parallel on large projects. When a project has enough pending references to make it worthwhile (roughly 150k+, typical for big Java/Kotlin/Spring codebases), resolution fans out across worker threads while results are applied in the exact order the single-threaded path would have used — the graph comes out byte-for-byte identical, about twice as fast end-to-end on a 4,000-file Java project in our testing. Small projects keep the single-threaded path automatically (the fan-out costs more than it saves there). Set `CODEGRAPH_NO_PARALLEL_RESOLVE=1` to disable, or `CODEGRAPH_PARALLEL_RESOLVE_MIN=<count>` to tune when it engages.
 - Indexing large projects got another sizeable speedup — about a quarter less wall-clock on the same 4,000-file Java project, with the graph still byte-for-byte identical. Two changes: the database no longer interleaves expensive checkpoint housekeeping into the middle of resolution on a fresh index (it's folded once at the end instead), and while one batch's results are being written out, the worker threads are already resolving the next batch instead of sitting idle.
+- The dynamic-dispatch analysis that runs at the end of indexing (callback, event, and framework wiring) now runs its passes in parallel on large projects, cutting that stage roughly in half there — and a pass that crashes now retries safely instead of failing the whole index, which also makes very large codebases that previously died in this stage more likely to index to completion. Graphs remain byte-for-byte identical.
 - Indexing is significantly faster — a fresh `codegraph init` on a medium TypeScript project takes about a third less wall-clock time, with the same graph produced byte-for-byte. The gains come from batching database writes, storing files on a dedicated writer thread, memoizing repeated import-resolution lookups, skipping per-row search-index maintenance during the bulk build (rebuilt once at the end), and — on completely fresh databases only — deferring disk durability until the index completes, since an interrupted first index is simply re-run. Set `CODEGRAPH_NO_FAST_INIT=1` to keep full crash-durability during the initial build, or `CODEGRAPH_NO_STORE_WORKER=1` to store on the main thread.
 - `codegraph install` and `codegraph upgrade` now offer CodeGraph Pro beta access after finishing — answer yes, type your email, and you join the same waitlist as the getcodegraph.com homepage form. Strictly opt-in and asked at most once per machine total: nothing is sent unless you say yes and enter an email, either answer is remembered so no later install or upgrade ever re-asks, and non-interactive runs (`--yes`, scripts, CI) never see the question.
 - Every release is now cryptographically verifiable: npm packages publish with npm provenance (the "Provenance" badge on npmjs.com, proving each version was built by this repository's release workflow from a specific commit), and the GitHub Release bundles carry signed build attestations you can check with `gh attestation verify <file> -R colbymchenry/codegraph`.

+ 9 - 6
__tests__/synthesis-progress.test.ts

@@ -14,19 +14,22 @@ import * as fs from 'fs';
 import * as os from 'os';
 import * as path from 'path';
 import { CodeGraph, IndexProgress } from '../src/index';
-import { SYNTH_PROGRESS_STEPS } from '../src/resolution/callback-synthesizer';
+import { SYNTH_PASSES, SYNTH_PROGRESS_STEPS } from '../src/resolution/callback-synthesizer';
 
 describe('synthesis progress ("Linking dynamic dispatch" phase)', () => {
-  it('SYNTH_PROGRESS_STEPS matches the synthesizer’s actual __mark() step count', () => {
+  it('SYNTH_PROGRESS_STEPS matches the synthesizer’s actual step count', () => {
     // The constant is cosmetic (progress denominator), but drift makes the bar
-    // end early or jump to 100% — adding a pass must bump it. Every step site
-    // calls __mark('<label>') with a string literal, so count those.
+    // end early or jump to 100%. Steps = one per SYNTH_PASSES registry entry
+    // (each marks exactly once, run or gated-out, sequential or pooled) plus
+    // the fixed literal __mark('<label>') sites (the ordered Go pre-passes and
+    // the merge/insert tail). Adding a pass = adding a registry entry, so the
+    // constant tracks automatically; this pins the fixed-site count.
     const src = fs.readFileSync(
       path.join(__dirname, '../src/resolution/callback-synthesizer.ts'),
       'utf8'
     );
-    const stepSites = (src.match(/__mark\('/g) ?? []).length;
-    expect(SYNTH_PROGRESS_STEPS).toBe(stepSites);
+    const fixedSites = (src.match(/__mark\('/g) ?? []).length;
+    expect(SYNTH_PROGRESS_STEPS).toBe(SYNTH_PASSES.length + fixedSites);
   });
 
   it('indexing emits a monotonic linking phase ending at the full step count', async () => {

+ 23 - 5
src/db/queries.ts

@@ -230,6 +230,7 @@ export class QueryBuilder {
     getNodesByLowerName?: SqliteStatement;
     getUnresolvedCount?: SqliteStatement;
     getUnresolvedBatch?: SqliteStatement;
+    deleteRefsByRowIdsFull?: SqliteStatement;
     getAllFilePaths?: SqliteStatement;
     getAllNodeNames?: SqliteStatement;
     getDominantFile?: SqliteStatement;
@@ -2204,11 +2205,28 @@ export class QueryBuilder {
    */
   deleteReferencesByRowIds(rowIds: number[]): void {
     if (rowIds.length === 0) return;
-    for (let i = 0; i < rowIds.length; i += SQLITE_PARAM_CHUNK_SIZE) {
-      const chunk = rowIds.slice(i, i + SQLITE_PARAM_CHUNK_SIZE);
-      const placeholders = chunk.map(() => '?').join(',');
-      this.db.prepare(`DELETE FROM unresolved_refs WHERE id IN (${placeholders})`).run(...chunk);
-    }
+    // One transaction for all chunks (each chunk was previously its own
+    // implicit transaction = its own WAL commit — measurable on 100k+-ref
+    // resolution persists), and the full-size chunk statement is cached so
+    // repeat calls skip the re-prepare; only the final partial chunk (if any)
+    // prepares ad hoc.
+    this.db.transaction(() => {
+      for (let i = 0; i < rowIds.length; i += SQLITE_PARAM_CHUNK_SIZE) {
+        const chunk = rowIds.slice(i, i + SQLITE_PARAM_CHUNK_SIZE);
+        if (chunk.length === SQLITE_PARAM_CHUNK_SIZE) {
+          if (!this.stmts.deleteRefsByRowIdsFull) {
+            const placeholders = new Array(SQLITE_PARAM_CHUNK_SIZE).fill('?').join(',');
+            this.stmts.deleteRefsByRowIdsFull = this.db.prepare(
+              `DELETE FROM unresolved_refs WHERE id IN (${placeholders})`
+            );
+          }
+          this.stmts.deleteRefsByRowIdsFull.run(...chunk);
+        } else {
+          const placeholders = chunk.map(() => '?').join(',');
+          this.db.prepare(`DELETE FROM unresolved_refs WHERE id IN (${placeholders})`).run(...chunk);
+        }
+      }
+    })();
   }
 
   /**

+ 150 - 80
src/resolution/callback-synthesizer.ts

@@ -3451,14 +3451,109 @@ async function laravelEventEdges(ctx: ResolutionContext, onYield: MaybeYield): P
  * Number of progress steps synthesizeCallbackEdges reports: one per `__mark()`
  * call (every synthesis pass, plus the dedupe-merge and edge-insert steps).
  * Cosmetic only — drift just makes the progress bar end early or jump — and a
- * test pins it to the actual `__mark(` call count so adding a pass without
- * bumping this fails loudly instead of silently skewing the bar.
+ * test pins it to the actual step count (registry passes + the fixed
+ * pre/post marks) so adding a pass without bumping this fails loudly instead
+ * of silently skewing the bar.
  */
-export const SYNTH_PROGRESS_STEPS = 40;
+const JS_FAMILY = ['typescript', 'javascript', 'tsx', 'jsx'];
+
+/** `has(...)` shape passed to pass gates — true when the project contains any of the languages. */
+type HasLang = (...ls: string[]) => boolean;
+
+/**
+ * One independent synthesis pass. Every pass scans the COMMITTED graph (plus
+ * source via ctx) and returns an edge list; nothing it produces is persisted
+ * until the ordered merge in synthesizeCallbackEdges — which is what makes
+ * execution order free and the passes safe to fan out across the resolver
+ * pool's read-only workers. `gate` short-circuits a pass whose language never
+ * appears in the project (its result is provably empty — see #1212).
+ */
+export interface SynthPassDef {
+  name: string;
+  gate: (has: HasLang) => boolean;
+  run: (
+    queries: QueryBuilder,
+    ctx: ResolutionContext,
+    yieldToLoop: MaybeYield,
+    subProgress?: (fraction: number) => void
+  ) => Promise<Edge[]>;
+}
+
+const ALWAYS = (): boolean => true;
+
+/**
+ * The independent passes, in MERGE ORDER — the first-seen dedup in
+ * synthesizeCallbackEdges follows this array, so reordering entries changes
+ * which duplicate edge wins. The two Go pre-passes (cross-file method
+ * `contains`, implicit `implements`) are NOT here: they persist before these
+ * run because interfaceOverrideEdges reads their edges from the DB.
+ */
+export const SYNTH_PASSES: SynthPassDef[] = [
+  { name: 'fieldEdges', gate: ALWAYS, run: (q, c, y) => fieldChannelEdges(q, c, y) },
+  { name: 'closureCollEdges', gate: ALWAYS, run: (q, c, y) => closureCollectionEdges(q, c, y) },
+  { name: 'emitterEdges', gate: ALWAYS, run: (_q, c, y) => eventEmitterEdges(c, y) },
+  { name: 'renderEdges', gate: ALWAYS, run: (q, c, y) => reactRenderEdges(q, c, y) },
+  { name: 'jsxEdges', gate: ALWAYS, run: (_q, c, y) => reactJsxChildEdges(c, y) },
+  { name: 'vueEdges', gate: (has) => has('vue'), run: (_q, c, y) => vueTemplateEdges(c, y) },
+  { name: 'svelteKitEdges', gate: (has) => has('svelte'), run: (_q, c, y) => svelteKitLoadEdges(c, y) },
+  { name: 'pascalEdges', gate: ALWAYS, run: (_q, c, y) => pascalFormEdges(c, y) },
+  { name: 'flutterEdges', gate: (has) => has('dart'), run: (q, c, y) => flutterBuildEdges(q, c, y) },
+  { name: 'arkuiStateEdges', gate: (has) => has('arkts'), run: (q, c, y) => arkuiStateBuildEdges(q, c, y) },
+  { name: 'arkuiEmitter', gate: (has) => has('arkts'), run: (_q, c, y) => arkuiEmitterEdges(c, y) },
+  { name: 'arkuiRoutes', gate: (has) => has('arkts'), run: (_q, c, y) => arkuiRouterEdges(c, y) },
+  { name: 'cppEdges', gate: (has) => has('cpp'), run: (q, _c, y) => cppOverrideEdges(q, y) },
+  {
+    name: 'ifaceEdges',
+    gate: (has) => has('java', 'kotlin', 'csharp', 'swift', 'scala', 'go', 'rust', 'arkts', ...JS_FAMILY),
+    run: (q, _c, y) => interfaceOverrideEdges(q, y),
+  },
+  { name: 'kotlinExpectActual', gate: (has) => has('kotlin'), run: (q, _c, y) => kotlinExpectActualEdges(q, y) },
+  { name: 'goGrpcEdges', gate: (has) => has('go'), run: (q, _c, y) => goGrpcStubImplEdges(q, y) },
+  { name: 'rnEventEdgesList', gate: (has) => has(...JS_FAMILY), run: (_q, c, y) => rnEventEdges(c, y) },
+  { name: 'fabricNativeEdges', gate: ALWAYS, run: (_q, c, y) => fabricNativeImplEdges(c, y) },
+  { name: 'expoXPlatEdges', gate: ALWAYS, run: (q, _c, y) => expoCrossPlatformEdges(q, y) },
+  { name: 'rnXPlatEdges', gate: ALWAYS, run: (q, _c, y) => rnCrossPlatformEdges(q, y) },
+  {
+    name: 'mybatisEdges',
+    gate: (has) => has('java', 'kotlin') && has('xml'),
+    run: (q, _c, y) => mybatisJavaXmlEdges(q, y),
+  },
+  { name: 'ginEdges', gate: (has) => has('go'), run: (q, c, y) => ginMiddlewareChainEdges(q, c, y) },
+  { name: 'thunkEdges', gate: (has) => has(...JS_FAMILY), run: (q, c, y) => reduxThunkEdges(q, c, y) },
+  { name: 'registryEdges', gate: ALWAYS, run: (_q, c, y) => objectRegistryEdges(c, y) },
+  { name: 'rtkEdges', gate: (has) => has(...JS_FAMILY), run: (q, c, y) => rtkQueryEdges(q, c, y) },
+  { name: 'piniaEdges', gate: (has) => has('vue', ...JS_FAMILY), run: (_q, c, y) => piniaStoreEdges(c, y) },
+  { name: 'vuexEdges', gate: (has) => has('vue', ...JS_FAMILY), run: (_q, c, y) => vuexDispatchEdges(c, y) },
+  { name: 'celeryEdges', gate: (has) => has('python'), run: (_q, c, y) => celeryDispatchEdges(c, y) },
+  { name: 'springEdges', gate: (has) => has('java'), run: (_q, c, y) => springEventEdges(c, y) },
+  { name: 'mediatrEdges', gate: (has) => has('csharp'), run: (_q, c, y) => mediatrDispatchEdges(c, y) },
+  { name: 'sidekiqEdges', gate: (has) => has('ruby'), run: (_q, c, y) => sidekiqDispatchEdges(c, y) },
+  {
+    name: 'erlangBehaviourEdges',
+    gate: (has) => has('erlang'),
+    run: (q, c, y) => erlangBehaviourDispatchEdges(q, c, y),
+  },
+  { name: 'laravelEdges', gate: (has) => has('php'), run: (_q, c, y) => laravelEventEdges(c, y) },
+  {
+    name: 'cFnPtrEdges',
+    gate: (has) => has('c', 'cpp'),
+    run: (q, c, y, sub) => cFnPointerDispatchEdges(q, c, y, sub),
+  },
+  { name: 'goframeEdges', gate: (has) => has('go'), run: (_q, c, y) => goframeRouteEdges(c, y) },
+  { name: 'nixOptionEdges', gate: (has) => has('nix'), run: (q, _c, y) => nixOptionPathEdges(q, y) },
+];
+
+/** Fixed non-registry steps: goMethodContains, goImplements, dedupe-merge, insertMergedEdges. */
+const FIXED_SYNTH_STEPS = 4;
+export const SYNTH_PROGRESS_STEPS = SYNTH_PASSES.length + FIXED_SYNTH_STEPS;
 export async function synthesizeCallbackEdges(
   queries: QueryBuilder,
   ctx: ResolutionContext,
-  onProgress?: (done: number, total: number) => void
+  onProgress?: (done: number, total: number) => void,
+  // A live resolver pool to fan the independent passes across (structural type
+  // so this file never imports the pool — resolver-worker imports THIS file).
+  // Null/omitted → the sequential path, byte-identical to the pool path.
+  pool?: { runSynthPass(name: string): Promise<{ edges: Edge[]; ms: number }> } | null
 ): Promise<number> {
   // Each sub-pass below is a whole-graph scan, and there are ~30 of them, all
   // running synchronously on the indexer's main thread. Their AGGREGATE can run
@@ -3516,7 +3611,6 @@ export async function synthesizeCallbackEdges(
   // Passes without an explicit language filter always run.
   const langs = queries.getDistinctFileLanguages();
   const has = (...ls: string[]): boolean => ls.some((l) => langs.has(l));
-  const JS_FAMILY = ['typescript', 'javascript', 'tsx', 'jsx'];
   const NONE: Edge[] = [];
 
   // Cross-file Go method→type `contains` edges must be synthesized AND persisted
@@ -3542,84 +3636,60 @@ export async function synthesizeCallbackEdges(
   }
   await yieldToLoop(); __mark('goImplements');
 
-  const fieldEdges = await fieldChannelEdges(queries, ctx, yieldToLoop); await yieldToLoop(); __mark('fieldEdges');
-  const closureCollEdges = await closureCollectionEdges(queries, ctx, yieldToLoop); await yieldToLoop(); __mark('closureCollEdges');
-  const emitterEdges = await eventEmitterEdges(ctx, yieldToLoop); await yieldToLoop(); __mark('emitterEdges');
-  const renderEdges = await reactRenderEdges(queries, ctx, yieldToLoop); await yieldToLoop(); __mark('renderEdges');
-  const jsxEdges = await reactJsxChildEdges(ctx, yieldToLoop); await yieldToLoop(); __mark('jsxEdges');
-  const vueEdges = has('vue') ? await vueTemplateEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('vueEdges');
-  const svelteKitEdges = has('svelte') ? await svelteKitLoadEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('svelteKitEdges');
-  const pascalEdges = await pascalFormEdges(ctx, yieldToLoop); await yieldToLoop(); __mark('pascalEdges');
-  const flutterEdges = has('dart') ? await flutterBuildEdges(queries, ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('flutterEdges');
-  const arkuiStateEdges = has('arkts') ? await arkuiStateBuildEdges(queries, ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('arkuiStateEdges');
-  const arkuiEmitter = has('arkts') ? await arkuiEmitterEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('arkuiEmitter');
-  const arkuiRoutes = has('arkts') ? await arkuiRouterEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('arkuiRoutes');
-  const cppEdges = has('cpp') ? await cppOverrideEdges(queries, yieldToLoop) : NONE; await yieldToLoop(); __mark('cppEdges');
-  const ifaceEdges = has('java', 'kotlin', 'csharp', 'swift', 'scala', 'go', 'rust', 'arkts', ...JS_FAMILY)
-    ? await interfaceOverrideEdges(queries, yieldToLoop) : NONE; await yieldToLoop(); __mark('ifaceEdges');
-  const kotlinExpectActual = has('kotlin') ? await kotlinExpectActualEdges(queries, yieldToLoop) : NONE; await yieldToLoop(); __mark('kotlinExpectActual');
-  const goGrpcEdges = has('go') ? await goGrpcStubImplEdges(queries, yieldToLoop) : NONE; await yieldToLoop(); __mark('goGrpcEdges');
-  const rnEventEdgesList = has(...JS_FAMILY) ? await rnEventEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('rnEventEdgesList');
-  const fabricNativeEdges = await fabricNativeImplEdges(ctx, yieldToLoop); await yieldToLoop(); __mark('fabricNativeEdges');
-  const expoXPlatEdges = await expoCrossPlatformEdges(queries, yieldToLoop); await yieldToLoop(); __mark('expoXPlatEdges');
-  const rnXPlatEdges = await rnCrossPlatformEdges(queries, yieldToLoop); await yieldToLoop(); __mark('rnXPlatEdges');
-  const mybatisEdges = has('java', 'kotlin') && has('xml') ? await mybatisJavaXmlEdges(queries, yieldToLoop) : NONE; await yieldToLoop(); __mark('mybatisEdges');
-  const ginEdges = has('go') ? await ginMiddlewareChainEdges(queries, ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('ginEdges');
-  const thunkEdges = has(...JS_FAMILY) ? await reduxThunkEdges(queries, ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('thunkEdges');
-  const registryEdges = await objectRegistryEdges(ctx, yieldToLoop); await yieldToLoop(); __mark('registryEdges');
-  const rtkEdges = has(...JS_FAMILY) ? await rtkQueryEdges(queries, ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('rtkEdges');
-  const piniaEdges = has('vue', ...JS_FAMILY) ? await piniaStoreEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('piniaEdges');
-  const vuexEdges = has('vue', ...JS_FAMILY) ? await vuexDispatchEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('vuexEdges');
-  const celeryEdges = has('python') ? await celeryDispatchEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('celeryEdges');
-  const springEdges = has('java') ? await springEventEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('springEdges');
-  const mediatrEdges = has('csharp') ? await mediatrDispatchEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('mediatrEdges');
-  const sidekiqEdges = has('ruby') ? await sidekiqDispatchEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('sidekiqEdges');
-  const erlangBehaviourEdges = has('erlang') ? await erlangBehaviourDispatchEdges(queries, ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('erlangBehaviourEdges');
-  const laravelEdges = has('php') ? await laravelEventEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('laravelEdges');
-  const cFnPtrEdges = has('c', 'cpp') ? await cFnPointerDispatchEdges(queries, ctx, yieldToLoop, subProgress) : NONE; await yieldToLoop(); __mark('cFnPtrEdges');
-  const goframeEdges = has('go') ? await goframeRouteEdges(ctx, yieldToLoop) : NONE; await yieldToLoop(); __mark('goframeEdges');
-  const nixOptionEdges = has('nix') ? await nixOptionPathEdges(queries, yieldToLoop) : NONE; await yieldToLoop(); __mark('nixOptionEdges');
+  // Run the independent passes (see SYNTH_PASSES). Their results are merged in
+  // REGISTRY ORDER below regardless of execution order, and none of their edges
+  // persist until that merge — so every pass sees the same committed
+  // post-resolution DB state whether it runs sequentially here or on a resolver
+  // pool worker. With a live pool (already booted on ≥150k-ref repos), passes
+  // fan out across its read-only workers and the per-pass wall-clock comes from
+  // the worker; a pass that fails on a worker falls back to running on the main
+  // thread, so a worker crash isolates to a retry instead of failing synthesis.
+  const passEdges: Edge[][] = new Array<Edge[]>(SYNTH_PASSES.length).fill(NONE);
+  const markPass = (label: string, dt: number): void => {
+    if (process.env.CODEGRAPH_SYNTH_TIMINGS && (dt > 250 || process.env.CODEGRAPH_SYNTH_TIMINGS === 'all')) {
+      console.error(`[synth-timing] ${label}: ${dt}ms`);
+    }
+    passesDone++;
+    emit(passesDone);
+  };
+  const runPassOnMain = async (i: number): Promise<void> => {
+    const pass = SYNTH_PASSES[i]!;
+    const t0 = Date.now();
+    passEdges[i] = await pass.run(queries, ctx, yieldToLoop, subProgress);
+    await yieldToLoop();
+    markPass(pass.name, Date.now() - t0);
+  };
+
+  const gatedIn: number[] = [];
+  for (let i = 0; i < SYNTH_PASSES.length; i++) {
+    if (SYNTH_PASSES[i]!.gate(has)) gatedIn.push(i);
+    else markPass(SYNTH_PASSES[i]!.name, 0);
+  }
+
+  if (pool && gatedIn.length > 1) {
+    await Promise.all(
+      gatedIn.map(async (i) => {
+        const pass = SYNTH_PASSES[i]!;
+        try {
+          const out = await pool.runSynthPass(pass.name);
+          passEdges[i] = out.edges;
+          markPass(pass.name, out.ms);
+        } catch {
+          // Worker-side failure (crash, OOM, unknown pass after a version
+          // mismatch): retry this one pass on the main thread.
+          await runPassOnMain(i);
+        }
+      })
+    );
+  } else {
+    for (const i of gatedIn) {
+      await runPassOnMain(i);
+    }
+  }
 
   const merged: Edge[] = [];
   const seen = new Set<string>();
-  for (const e of [
-    ...fieldEdges,
-    ...closureCollEdges,
-    ...emitterEdges,
-    ...renderEdges,
-    ...jsxEdges,
-    ...vueEdges,
-    ...svelteKitEdges,
-    ...pascalEdges,
-    ...flutterEdges,
-    ...arkuiStateEdges,
-    ...arkuiEmitter,
-    ...arkuiRoutes,
-    ...cppEdges,
-    ...ifaceEdges,
-    ...kotlinExpectActual,
-    ...goGrpcEdges,
-    ...rnEventEdgesList,
-    ...fabricNativeEdges,
-    ...expoXPlatEdges,
-    ...rnXPlatEdges,
-    ...mybatisEdges,
-    ...ginEdges,
-    ...thunkEdges,
-    ...registryEdges,
-    ...rtkEdges,
-    ...piniaEdges,
-    ...vuexEdges,
-    ...celeryEdges,
-    ...springEdges,
-    ...mediatrEdges,
-    ...sidekiqEdges,
-    ...erlangBehaviourEdges,
-    ...laravelEdges,
-    ...cFnPtrEdges,
-    ...goframeEdges,
-    ...nixOptionEdges,
-  ]) {
+  for (const e of passEdges.flat()) {
     const key = `${e.source}>${e.target}`;
     if (seen.has(key)) continue;
     seen.add(key);

+ 17 - 5
src/resolution/index.ts

@@ -1305,6 +1305,14 @@ export class ReferenceResolver {
     };
   }
 
+  /**
+   * The resolver's live ResolutionContext — resolver-pool workers use it to
+   * run synthesis passes against their own read-only connection.
+   */
+  getResolutionContext(): ResolutionContext {
+    return this.context;
+  }
+
   /**
    * Re-queue deferred post-pass refs produced by resolver workers, preserving
    * their admission order so resolveChainedCallsViaConformance /
@@ -1551,25 +1559,29 @@ export class ReferenceResolver {
       batch = nextBatch;
       inFlight = nextInFlight;
     }
-    } finally {
-      if (pool) await pool.destroy().catch(() => undefined);
-    }
 
     // Dynamic-edge synthesis: now that all base `calls` edges are persisted,
     // synthesize observer/callback dispatch edges (dispatcher → registered
     // callbacks) that static parsing leaves out. Best-effort — never fail the
-    // index on it. See docs/design/callback-edge-synthesis.md.
+    // index on it. The pool (when it survived resolution) is REUSED to fan the
+    // independent passes across its read-only workers — that's why its destroy
+    // lives in the finally below, after synthesis, not at the end of the batch
+    // loop. See docs/design/callback-edge-synthesis.md.
     const tSynth = Date.now();
     try {
       aggregateStats.byMethod['callback-synthesis'] = await synthesizeCallbackEdges(
         this.queries,
         this.context,
-        onSynthesisProgress
+        onSynthesisProgress,
+        pool
       );
     } catch {
       // synthesis is additive and optional; ignore failures
     }
     if (process.env.CODEGRAPH_SYNTH_TIMINGS) console.error(`[phase-timing] callback-synthesis: ${Date.now() - tSynth}ms`);
+    } finally {
+      if (pool) await pool.destroy().catch(() => undefined);
+    }
 
     return {
       resolved: [],

+ 37 - 2
src/resolution/resolver-pool.ts

@@ -13,9 +13,15 @@ import { Worker } from 'worker_threads';
 import * as fs from 'fs';
 import * as path from 'path';
 import * as os from 'os';
-import type { UnresolvedReference } from '../types';
+import type { Edge, UnresolvedReference } from '../types';
 import type { ResolvedRef, UnresolvedRef } from './types';
 
+/** One synthesis pass's output: its edge list + worker-measured wall clock. */
+export interface SynthPassResult {
+  edges: Edge[];
+  ms: number;
+}
+
 export interface ChunkResult {
   resolved: ResolvedRef[];
   unresolved: UnresolvedRef[];
@@ -55,6 +61,7 @@ export class ResolverPool {
   private workers: PoolWorker[] = [];
   private nextId = 0;
   private waiters = new Map<number, { resolve: (r: ChunkResult) => void; reject: (e: Error) => void }>();
+  private synthWaiters = new Map<number, { resolve: (r: SynthPassResult) => void; reject: (e: Error) => void }>();
   private failed: Error | null = null;
 
   /**
@@ -85,7 +92,7 @@ export class ResolverPool {
         readyReject = reject;
       });
       const pw: PoolWorker = { worker, ready, busy: 0 };
-      worker.on('message', (msg: { type: string; id?: number; message?: string } & Partial<ChunkResult>) => {
+      worker.on('message', (msg: { type: string; id?: number; message?: string; edges?: Edge[]; ms?: number } & Partial<ChunkResult>) => {
         if (msg.type === 'ready') {
           readyResolve();
         } else if (msg.type === 'result' && msg.id !== undefined) {
@@ -99,6 +106,11 @@ export class ResolverPool {
             deferredThisMember: msg.deferredThisMember!,
             byMethod: msg.byMethod!,
           });
+        } else if (msg.type === 'synth-result' && msg.id !== undefined) {
+          pw.busy--;
+          const waiter = this.synthWaiters.get(msg.id);
+          this.synthWaiters.delete(msg.id);
+          waiter?.resolve({ edges: msg.edges ?? [], ms: msg.ms ?? 0 });
         } else if (msg.type === 'error') {
           pw.busy--;
           const err = new Error(`resolver worker: ${msg.message}`);
@@ -106,6 +118,10 @@ export class ResolverPool {
             const waiter = this.waiters.get(msg.id)!;
             this.waiters.delete(msg.id);
             waiter.reject(err);
+          } else if (msg.id !== undefined && this.synthWaiters.has(msg.id)) {
+            const waiter = this.synthWaiters.get(msg.id)!;
+            this.synthWaiters.delete(msg.id);
+            waiter.reject(err);
           } else {
             this.fail(err);
           }
@@ -130,6 +146,8 @@ export class ResolverPool {
     if (!this.failed) this.failed = err;
     for (const [, waiter] of this.waiters) waiter.reject(this.failed);
     this.waiters.clear();
+    for (const [, waiter] of this.synthWaiters) waiter.reject(this.failed);
+    this.synthWaiters.clear();
   }
 
   /** Whether this batch is worth fanning out. */
@@ -175,6 +193,23 @@ export class ResolverPool {
     return out;
   }
 
+  /**
+   * Run one synthesis pass (by SYNTH_PASSES name) on the least-busy worker.
+   * The worker reads the committed graph on its own connection and returns
+   * the pass's edge list; the caller merges in canonical order. Rejects on
+   * worker failure — the caller retries the pass on the main thread.
+   */
+  async runSynthPass(passName: string): Promise<SynthPassResult> {
+    if (this.failed) throw this.failed;
+    const id = this.nextId++;
+    const pw = this.workers.reduce((a, b) => (b.busy < a.busy ? b : a));
+    pw.busy++;
+    return new Promise<SynthPassResult>((resolve, reject) => {
+      this.synthWaiters.set(id, { resolve, reject });
+      pw.worker.postMessage({ type: 'synth', id, pass: passName });
+    });
+  }
+
   async destroy(): Promise<void> {
     await Promise.all(
       this.workers.map(

+ 31 - 1
src/resolution/resolver-worker.ts

@@ -24,6 +24,8 @@ import { parentPort } from 'worker_threads';
 import { createDatabase, SqliteDatabase } from '../db/sqlite-adapter';
 import { QueryBuilder } from '../db/queries';
 import { ReferenceResolver } from './index';
+import { SYNTH_PASSES } from './callback-synthesizer';
+import { createYielder } from './cooperative-yield';
 import type { UnresolvedReference } from '../types';
 
 if (!parentPort) {
@@ -32,11 +34,13 @@ if (!parentPort) {
 const port = parentPort;
 
 let db: SqliteDatabase | null = null;
+let queries: QueryBuilder | null = null;
 let resolver: ReferenceResolver | null = null;
 
 type InMessage =
   | { type: 'open'; dbPath: string; projectRoot: string }
   | { type: 'resolve'; id: number; refs: UnresolvedReference[] }
+  | { type: 'synth'; id: number; pass: string }
   | { type: 'close' };
 
 port.on('message', (msg: InMessage) => {
@@ -49,7 +53,7 @@ port.on('message', (msg: InMessage) => {
         db.pragma('busy_timeout = 5000');
         db.pragma('cache_size = -32000');
         const tDb = Date.now();
-        const queries = new QueryBuilder(db);
+        queries = new QueryBuilder(db);
         resolver = new ReferenceResolver(msg.projectRoot, queries);
         resolver.initialize();
         if (process.env.CODEGRAPH_SYNTH_TIMINGS) console.error(`[pool-timing] worker open: db=${tDb - tOpen}ms init=${Date.now() - tDb}ms`);
@@ -64,6 +68,32 @@ port.on('message', (msg: InMessage) => {
         port.postMessage({ type: 'result', id: msg.id, ...out });
         break;
       }
+      case 'synth': {
+        // Run one synthesis pass against this worker's read-only connection.
+        // Passes only READ (graph + source via the resolver's context); their
+        // edges are returned for the main thread's ordered merge. Async, with
+        // its own error propagation — a throwing pass reports {type:'error'}
+        // and the main thread retries it sequentially.
+        if (!resolver || !queries) throw new Error('resolver-worker: synth before open');
+        const pass = SYNTH_PASSES.find((p) => p.name === msg.pass);
+        if (!pass) throw new Error(`resolver-worker: unknown synth pass '${msg.pass}'`);
+        const q = queries;
+        const r = resolver;
+        void (async () => {
+          const t0 = Date.now();
+          try {
+            const edges = await pass.run(q, r.getResolutionContext(), createYielder());
+            port.postMessage({ type: 'synth-result', id: msg.id, edges, ms: Date.now() - t0 });
+          } catch (err) {
+            port.postMessage({
+              type: 'error',
+              id: msg.id,
+              message: err instanceof Error ? err.message : String(err),
+            });
+          }
+        })();
+        break;
+      }
       case 'close': {
         try {
           db?.close();