mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-07 02:58:02 +00:00
fix(community): run Leiden in a worker so its timeout can fire (#3478)
This commit is contained in:
parent
df49f90889
commit
047354c3f0
4 changed files with 225 additions and 29 deletions
|
|
@ -289,6 +289,7 @@ const SPAWN_CLI = [
|
|||
// Worker threads tests — exercise real worker_threads which have
|
||||
// platform-specific behavior (thread spawning, IPC, exit handling)
|
||||
const WORKER_THREADS = [
|
||||
'test/unit/community-processor.test.ts',
|
||||
'test/integration/worker-pool.test.ts',
|
||||
'test/integration/parse-impl-quarantine-cache-skip.test.ts',
|
||||
];
|
||||
|
|
|
|||
|
|
@ -29,12 +29,9 @@ const _require = createRequire(import.meta.url);
|
|||
/** Graphology Graph instance type (AbstractGraph from graphology-types avoids CJS/ESM interop namespace issue) */
|
||||
type GraphInstance = AbstractGraph<Attributes, Attributes, Attributes>;
|
||||
|
||||
const leiden: LeidenModule = _require(leidenPath);
|
||||
|
||||
/** Vendored Leiden algorithm module shape */
|
||||
interface LeidenModule {
|
||||
detailed: (graph: GraphInstance, options: Record<string, unknown>) => LeidenDetailedResult;
|
||||
}
|
||||
// The Leiden worker loads its own copy; this load keeps a missing vendor/ asset
|
||||
// failing at import time and keeps it visible to dockerfile-runtime-asset-parity.
|
||||
_require(leidenPath);
|
||||
|
||||
/** Result returned by leiden.detailed() */
|
||||
interface LeidenDetailedResult {
|
||||
|
|
@ -110,15 +107,6 @@ interface IcebugWorkerFailure {
|
|||
* incremental-indexing equivalence test (incremental ≡ full rebuild).
|
||||
*/
|
||||
const LEIDEN_SEED = 0xc0de;
|
||||
function createSeededRng(seed: number): () => number {
|
||||
let s = seed >>> 0;
|
||||
return () => {
|
||||
s = (s + 0x6d2b79f5) >>> 0;
|
||||
let t = Math.imul(s ^ (s >>> 15), 1 | s);
|
||||
t = (t + Math.imul(t ^ (t >>> 7), 61 | t)) ^ t;
|
||||
return ((t ^ (t >>> 14)) >>> 0) / 4294967296;
|
||||
};
|
||||
}
|
||||
|
||||
const COMMUNITY_ENGINE_ENV = 'GITNEXUS_COMMUNITY_ENGINE';
|
||||
/**
|
||||
|
|
@ -474,28 +462,101 @@ const runCommunityEngine = async (
|
|||
35,
|
||||
);
|
||||
const fallback = await runGraphologyLeiden(graph, projection.isLarge, engineRequested);
|
||||
return { ...fallback, fallbackReason };
|
||||
return {
|
||||
...fallback,
|
||||
fallbackReason: fallback.fallbackReason
|
||||
? `${fallbackReason}; ${fallback.fallbackReason}`
|
||||
: fallbackReason,
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
const runGraphologyLeiden = async (
|
||||
/**
|
||||
* Runs the vendored Leiden in a worker so LEIDEN_TIMEOUT_MS can actually fire.
|
||||
* leiden.detailed is synchronous and, on some graphs, never returns
|
||||
* (graphology/graphology#557); a timer on the main thread cannot interrupt it.
|
||||
* The worker is pure JS, so terminate() is safe (unlike the icebug N-API worker).
|
||||
*/
|
||||
const GRAPHOLOGY_WORKER_SOURCE = `
|
||||
const { parentPort, workerData } = require('node:worker_threads');
|
||||
try {
|
||||
const Graph = require(workerData.graphologyPath);
|
||||
const leiden = require(workerData.leidenPath);
|
||||
const graph = Graph.from(workerData.graph);
|
||||
const rng = (() => {
|
||||
let s = workerData.seed >>> 0;
|
||||
return () => {
|
||||
s = (s + 0x6d2b79f5) >>> 0;
|
||||
let t = Math.imul(s ^ (s >>> 15), 1 | s);
|
||||
t = (t + Math.imul(t ^ (t >>> 7), 61 | t)) ^ t;
|
||||
return ((t ^ (t >>> 14)) >>> 0) / 4294967296;
|
||||
};
|
||||
})();
|
||||
const d = leiden.detailed(graph, {
|
||||
resolution: workerData.resolution,
|
||||
maxIterations: workerData.maxIterations,
|
||||
rng,
|
||||
});
|
||||
parentPort.postMessage({
|
||||
ok: true,
|
||||
communities: d.communities,
|
||||
count: d.count,
|
||||
modularity: d.modularity,
|
||||
});
|
||||
} catch (error) {
|
||||
parentPort.postMessage({ ok: false, error: error instanceof Error ? error.message : String(error) });
|
||||
}
|
||||
`;
|
||||
|
||||
export const runGraphologyLeiden = async (
|
||||
graph: GraphInstance,
|
||||
isLarge: boolean,
|
||||
engineRequested: CommunityDetectionEngine,
|
||||
timeoutMs: number = LEIDEN_TIMEOUT_MS,
|
||||
): Promise<CommunityEngineResult> => {
|
||||
try {
|
||||
const details = await Promise.race([
|
||||
Promise.resolve(
|
||||
leiden.detailed(graph, {
|
||||
const details = await new Promise<LeidenDetailedResult>((resolve, reject) => {
|
||||
const worker = new Worker(GRAPHOLOGY_WORKER_SOURCE, {
|
||||
eval: true,
|
||||
workerData: {
|
||||
graphologyPath: _require.resolve('graphology'),
|
||||
leidenPath,
|
||||
graph: graph.export(),
|
||||
seed: LEIDEN_SEED,
|
||||
resolution: isLarge ? 2.0 : 1.0,
|
||||
maxIterations: isLarge ? 3 : 0,
|
||||
rng: createSeededRng(LEIDEN_SEED),
|
||||
}),
|
||||
),
|
||||
new Promise<never>((_, reject) =>
|
||||
setTimeout(() => reject(new Error('Leiden timeout')), LEIDEN_TIMEOUT_MS),
|
||||
),
|
||||
]);
|
||||
},
|
||||
});
|
||||
let settled = false;
|
||||
const timeout = setTimeout(() => {
|
||||
settled = true;
|
||||
void worker.terminate();
|
||||
reject(new Error('Leiden timeout'));
|
||||
}, timeoutMs);
|
||||
worker.once('message', (message) => {
|
||||
settled = true;
|
||||
clearTimeout(timeout);
|
||||
if (message.ok === true) {
|
||||
resolve({
|
||||
communities: message.communities,
|
||||
count: message.count,
|
||||
modularity: message.modularity,
|
||||
});
|
||||
} else {
|
||||
reject(new Error(message.error));
|
||||
}
|
||||
});
|
||||
worker.once('error', (error) => {
|
||||
settled = true;
|
||||
clearTimeout(timeout);
|
||||
reject(error);
|
||||
});
|
||||
worker.once('exit', (code) => {
|
||||
if (settled) return;
|
||||
clearTimeout(timeout);
|
||||
reject(new Error(`Graphology Leiden worker exited with code ${code}`));
|
||||
});
|
||||
});
|
||||
return { ...details, engine: 'graphology', engineRequested };
|
||||
} catch (e: any) {
|
||||
if (e.message !== 'Leiden timeout') {
|
||||
|
|
|
|||
38
gitnexus/test/helpers/leiden-timeout-child.ts
Normal file
38
gitnexus/test/helpers/leiden-timeout-child.ts
Normal file
|
|
@ -0,0 +1,38 @@
|
|||
import {
|
||||
buildGraphologyGraph,
|
||||
runGraphologyLeiden,
|
||||
} from '../../src/core/ingestion/community-processor.js';
|
||||
|
||||
// 105-node / 128-edge graph from #3476 on which the vendored Leiden never
|
||||
// terminates (graphology/graphology#557). Preserve the original insertion order:
|
||||
// sorting either the nodes or edges changes the path through the algorithm.
|
||||
const EDGES =
|
||||
'39-26 65-66 72-52 10-38 98-68 37-10 26-81 37-38 30-66 18-66 30-97 84-81 50-96 36-76 ' +
|
||||
'33-29 54-87 4-75 44-31 29-31 60-44 91-10 48-41 68-15 83-90 82-57 50-71 28-55 100-11 ' +
|
||||
'42-19 28-42 74-92 55-5 7-65 96-41 61-27 60-85 89-67 97-34 94-64 89-6 61-51 12-13 ' +
|
||||
'35-9 91-20 94-23 91-7 43-99 80-85 80-101 75-50 38-57 4-104 35-6 65-103 83-45 62-91 ' +
|
||||
'16-37 94-33 103-17 67-66 39-73 68-2 26-31 48-19 43-88 88-47 54-42 26-66 100-65 7-30 ' +
|
||||
'32-13 56-2 22-24 78-24 84-53 90-71 60-12 76-49 86-35 55-68 96-46 33-30 4-1 24-68 ' +
|
||||
'23-14 93-3 58-56 64-27 85-87 61-101 36-37 37-85 60-63 37-26 41-74 103-87 78-36 69-97 ' +
|
||||
'77-42 97-52 57-25 27-34 68-28 95-104 18-96 4-15 60-48 17-14 51-31 83-40 48-75 16-99 ' +
|
||||
'79-65 70-14 3-27 0-62 21-101 44-59 4-24 34-52 65-8 102-66 67-24 11-26 39-78 6-103 ' +
|
||||
'4-71 70-4';
|
||||
|
||||
const graph = buildGraphologyGraph({
|
||||
nodes: Array.from({ length: 105 }, (_, i) => ({
|
||||
id: String(i),
|
||||
name: String(i),
|
||||
filePath: `/src/${i}.ts`,
|
||||
type: 'Function',
|
||||
})),
|
||||
edges: EDGES.split(/\s+/).map((edge) => {
|
||||
const [source, target] = edge.split('-').map(Number);
|
||||
return [source, target];
|
||||
}),
|
||||
symbolCount: 105,
|
||||
isLarge: false,
|
||||
});
|
||||
|
||||
const result = await runGraphologyLeiden(graph, false, 'graphology', 500);
|
||||
console.log(JSON.stringify(result));
|
||||
// No explicit exit: a leaked worker must keep this process alive for the parent to detect.
|
||||
|
|
@ -1,9 +1,12 @@
|
|||
import { spawnSync } from 'node:child_process';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { Worker } from 'node:worker_threads';
|
||||
import { describe, it, expect, vi } from 'vitest';
|
||||
import { tsxLoaderUrl } from '../helpers/cli-entry.js';
|
||||
import { createKnowledgeGraph } from '../../src/core/graph/graph.js';
|
||||
import type { GraphNode, GraphRelationship } from 'gitnexus-shared';
|
||||
import {
|
||||
|
|
@ -197,9 +200,19 @@ describe('community-processor', () => {
|
|||
vi.resetModules();
|
||||
vi.doMock('node:worker_threads', () => {
|
||||
class MockWorker extends EventEmitter {
|
||||
constructor() {
|
||||
constructor(_source?: string, options?: { workerData?: { graphologyPath?: string } }) {
|
||||
super();
|
||||
queueMicrotask(() => {
|
||||
if (options?.workerData?.graphologyPath) {
|
||||
// The Graphology Leiden worker: answer as the real one would.
|
||||
this.emit('message', {
|
||||
ok: true,
|
||||
communities: { 'fn:a': 0, 'fn:b': 0 },
|
||||
count: 1,
|
||||
modularity: 0,
|
||||
});
|
||||
return;
|
||||
}
|
||||
this.emit('message', { ok: true, partition: [0, 0], modularity: Number.NaN });
|
||||
});
|
||||
}
|
||||
|
|
@ -232,7 +245,7 @@ describe('community-processor', () => {
|
|||
|
||||
expect(result.stats.engineRequested).toBe('icebug');
|
||||
expect(result.stats.engine).toBe('graphology');
|
||||
expect(result.stats.fallbackReason).toContain('modularity');
|
||||
expect(result.stats.fallbackReason).toBe('optional icebug modularity was not finite');
|
||||
expect(progress.some((message) => message.includes('falling back to Graphology'))).toBe(
|
||||
true,
|
||||
);
|
||||
|
|
@ -246,6 +259,64 @@ describe('community-processor', () => {
|
|||
}
|
||||
});
|
||||
|
||||
it.each(['icebug', 'auto'] as const)(
|
||||
'preserves both fallback reasons when %s fails and Graphology times out',
|
||||
async (engine) => {
|
||||
vi.resetModules();
|
||||
vi.useFakeTimers();
|
||||
const icebugError = 'optional icebug could not load';
|
||||
vi.doMock('node:worker_threads', () => {
|
||||
class MockWorker extends EventEmitter {
|
||||
constructor(_source?: string, options?: { workerData?: { graphologyPath?: string } }) {
|
||||
super();
|
||||
if (!options?.workerData?.graphologyPath) {
|
||||
queueMicrotask(() => this.emit('message', { ok: false, error: icebugError }));
|
||||
}
|
||||
// The Graphology worker stays silent until its real deadline fires.
|
||||
}
|
||||
|
||||
terminate(): Promise<number> {
|
||||
return Promise.resolve(0);
|
||||
}
|
||||
|
||||
unref(): void {}
|
||||
}
|
||||
|
||||
return { Worker: MockWorker };
|
||||
});
|
||||
|
||||
try {
|
||||
const { processCommunities: processCommunitiesWithMockWorker } =
|
||||
await import('../../src/core/ingestion/community-processor.js');
|
||||
const graph = createKnowledgeGraph();
|
||||
graph.addNode(makeNode('fn:a', 'a'));
|
||||
graph.addNode(makeNode('fn:b', 'b'));
|
||||
graph.addRelationship(makeRel('rel:ab', 'fn:a', 'fn:b'));
|
||||
|
||||
const pending = processCommunitiesWithMockWorker(graph, undefined, { engine });
|
||||
await vi.advanceTimersByTimeAsync(60_000);
|
||||
const result = await pending;
|
||||
|
||||
expect(result.stats).toMatchObject({
|
||||
engineRequested: engine,
|
||||
engine: 'graphology',
|
||||
fallbackReason: `${icebugError}; Graphology Leiden timeout`,
|
||||
totalCommunities: 1,
|
||||
modularity: 0,
|
||||
nodesProcessed: 2,
|
||||
});
|
||||
expect(result.memberships).toEqual([
|
||||
{ nodeId: 'fn:a', communityId: 'comm_0' },
|
||||
{ nodeId: 'fn:b', communityId: 'comm_0' },
|
||||
]);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
vi.doUnmock('node:worker_threads');
|
||||
vi.resetModules();
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it('falls back before icebug worker launch for nondeterministic options', async () => {
|
||||
const graph = createKnowledgeGraph();
|
||||
graph.addNode(makeNode('fn:a', 'a', 'Function', '/src/group/a.ts'));
|
||||
|
|
@ -372,6 +443,31 @@ module.exports = {
|
|||
});
|
||||
});
|
||||
|
||||
describe('graphology Leiden timeout', () => {
|
||||
it('returns the fallback and exits naturally after terminating the hanging worker', () => {
|
||||
const childPath = fileURLToPath(
|
||||
new URL('../helpers/leiden-timeout-child.ts', import.meta.url),
|
||||
);
|
||||
// A fallback alone is not enough: the child must exit without process.exit().
|
||||
// SIGKILL bounds cleanup if a regression leaves its worker running forever.
|
||||
const child = spawnSync(process.execPath, ['--import', tsxLoaderUrl(), childPath], {
|
||||
encoding: 'utf8',
|
||||
timeout: 10_000,
|
||||
killSignal: 'SIGKILL',
|
||||
});
|
||||
|
||||
expect(child.error).toBeUndefined();
|
||||
expect(child.signal).toBeNull();
|
||||
expect(child.status, child.stderr).toBe(0);
|
||||
const result = JSON.parse(child.stdout);
|
||||
|
||||
expect(result.fallbackReason).toBe('Graphology Leiden timeout');
|
||||
expect(result.count).toBe(1);
|
||||
expect(new Set(Object.values(result.communities))).toEqual(new Set([0]));
|
||||
expect(Object.keys(result.communities)).toHaveLength(105);
|
||||
}, 15_000);
|
||||
});
|
||||
|
||||
describe('vendored Leiden partitioning', () => {
|
||||
// Golden values for the seeded graph below, captured from the vendored
|
||||
// implementation. They pin the partition, not just its shape.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue