mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-09-22 00:31:17 +00:00
* Initial plan * refactor: move language-specific container node logic into LanguageProvider - Add resolveEnclosingOwner hook to LanguageProviderConfig - Add staticOwnerTypes to MethodExtractionConfig - Implement Ruby resolveEnclosingOwner (singleton_class → class/module) - Replace hardcoded STATIC_OWNER_TYPES with config.staticOwnerTypes - Move Ruby static types to rubyMethodConfig - Move Kotlin static types to kotlinMethodConfig - Remove Ruby singleton_class branch from findEnclosingClassInfo - Collapse seqFindEnclosingClassNode/seqFindRawEnclosingContainerNode into single provider-aware seqFindEnclosingOwnerNode - Update worker path to pass provider.resolveEnclosingOwner Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/bc9f9d4d-f749-4872-9ff2-17fc86e08787 Co-authored-by: magyargergo <11230420+magyargergo@users.noreply.github.com> * test: add regression tests for config-driven staticOwnerTypes and resolveEnclosingOwner hook Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/bc9f9d4d-f749-4872-9ff2-17fc86e08787 Co-authored-by: magyargergo <11230420+magyargergo@users.noreply.github.com> * refactor: implement DAG-based pipeline architecture with phase extraction Restructure the ingestion pipeline from a ~1800-line monolithic orchestrator into a DAG (Directed Acyclic Graph) of named phases with explicit dependencies. New files under pipeline-phases/: - types.ts: PipelinePhase, PipelineContext, PhaseResult contracts - runner.ts: DAG runner with topological sort validation - scan.ts, structure.ts, markdown.ts, cobol.ts: early phases - parse.ts + parse-impl.ts: chunked parse + resolve (the core) - routes.ts, tools.ts, orm.ts: post-parse enrichment phases - cross-file.ts + cross-file-impl.ts: cross-file binding propagation - mro.ts, communities.ts, processes.ts: graph analysis phases - index.ts: barrel export pipeline.ts reduced from ~1960 lines to ~184 lines: - DAG phase array declaration - runPipelineFromRepo as thin orchestrator - topologicalLevelSort retained for backward compat Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/136bf9c3-2f4f-449b-9fff-001332c8371c * test: add DAG runner unit tests, update ARCHITECTURE.md with phase DAG docs Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/136bf9c3-2f4f-449b-9fff-001332c8371c * fix: address code review - pass resolutionContext through parse output, fix worker URL path Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/136bf9c3-2f4f-449b-9fff-001332c8371c * fix: declare transitive parse dependency explicitly in mro/communities/processes phases Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/136bf9c3-2f4f-449b-9fff-001332c8371c * refactor: improve pipeline-phases clean code and folder structure - Extract synthesizeWildcardImportBindings to wildcard-synthesis.ts - Extract extractORMQueriesInline to orm-extraction.ts - Create shared constants.ts for AST_CACHE_CAP - Fix inline type import in orm.ts (use proper top-level import) - Add comprehensive JSDoc to getPhaseOutput explaining type safety - Move isDev to module level in cross-file.ts (consistency) - Improve module-level documentation across files - Organize barrel exports in index.ts with section comments Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/2bd6d4aa-6271-4009-8dd2-332ea8ec73ab Co-authored-by: magyargergo <11230420+magyargergo@users.noreply.github.com> * address review feedback: fix circular dep, allFetchCalls mutation, progress bugs, remove DAG naming, extract isDev, fix _item naming, fix O(n²) line calc Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/6cf53c9b-d55d-4c6f-bf3d-7bfb82d512b6 Co-authored-by: magyargergo <11230420+magyargergo@users.noreply.github.com> * improve JSDoc on lineNumberAtOffset binary search Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/6cf53c9b-d55d-4c6f-bf3d-7bfb82d512b6 Co-authored-by: magyargergo <11230420+magyargergo@users.noreply.github.com> * address review: filter deps in runner, move totalFiles to ctx, fix cycle JSDoc, centralize isDev, remove DAG naming Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/b388424f-b939-4a94-97de-3855f9465564 Co-authored-by: magyargergo <11230420+magyargergo@users.noreply.github.com> * fix doc consistency in graph-sort.ts module-level and function-level JSDoc Agent-Logs-Url: https://github.com/abhigyanpatwari/GitNexus/sessions/b388424f-b939-4a94-97de-3855f9465564 Co-authored-by: magyargergo <11230420+magyargergo@users.noreply.github.com> * fix(pipeline): wrap phase errors with phase name and emit terminal error progress event Restores phase diagnostics at CLI/MCP boundary. runPipeline now wraps phase.execute() in try/catch and rethrows with 'Phase <name> failed: ...' preserving the original via { cause }. Also emits a terminal { phase: 'error' } progress event so subscribers see the failure before the rejection propagates. Handler errors during error reporting are swallowed to keep the original cause authoritative. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U1) * fix(pipeline): move bindingAccumulator dispose into crossFile try/finally; make single-use crossFile.execute() now wraps its body in try/finally so the accumulator is released on both the happy path and when runCrossFileBindingPropagation throws. Dev-mode telemetry stays inside the try block before dispose (all three counters return 0 after dispose clears internal maps). BindingAccumulator becomes single-use: appendFile after dispose now throws 'BindingAccumulator: use after dispose' instead of silently re-animating via the old _disposed auto-clear. Docs updated; the only production construction site (parse-impl) always creates a fresh instance per run, so no caller relied on the re-use contract. Residual risk documented in crossFile module JSDoc: a future phase inserted between parse and crossFile that throws would still leak the accumulator. Any such phase must manage accumulator lifetime explicitly. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U2) * docs(pipeline): explain why importCtx teardown is safe before crossFile Investigation (plan U3) confirms: `importCtx` (ImportResolutionContext) is a scratch workspace with no downstream consumer after parse. `resolutionContext` (returned to crossFile) is a distinct object that owns importMap / namedImportMap / packageMap / moduleAliasMap / model, and never closes over importCtx. cross-file-impl consumes only that ctx via processCalls. The two confusingly-similar "context" names were the root of the adversarial reviewer's concern — comment locks in the invariant so the next reader sees it. No behavioral change. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U3) * refactor(pipeline): remove ctx.totalFiles side-channel; promote to ParseOutput totalFiles was a hidden mutable field on PipelineContext written by parse and read by mro/communities/processes — five reviewers flagged this as a violation of the immutable-context invariant. Removed from PipelineContext, which is now fully readonly, and made the implicit temporal dep explicit: mro/communities/processes now declare 'parse' as a dep and read totalFiles via getPhaseOutput<ParseOutput>(...). No behavior change. Topo-sort unchanged because parse was already a transitive dep through crossFile. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U4) * feat(method-extractor): runtime staticOwnerTypes guard at factory chokepoint createMethodExtractor now rejects MethodExtractionConfigs that list companion_object / singleton_class / object_declaration in typeDeclarationNodes but omit the matching entry from staticOwnerTypes. Fails loudly at provider construction time instead of producing silent isStatic=false on the 50000th file analyzed. Opt-out convention preserved: an explicit `new Set()` (empty Set) signals intentional exclusion and passes the guard (memory obs #30588). All 13 existing language configs pass the guard; the new negative test fails without it. Test-first. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U5) * fix(pipeline): wrap sequential-fallback in try/finally so cleanup survives throws The sequential-fallback block in runChunkedParseAndResolve now runs inside a try/finally that guarantees astCache.clear(), accumulator finalize, and enrichExportedTypeMap execute even if readFileContents or processCalls throws mid-fallback. Cleanup failures are caught inside the finally so they can't mask the original error. Accumulator disposal ownership remains with crossFile (U2) — U6 only adds astCache cleanup and preserves finalize ordering on the error path. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U6) * test(pipeline): direct unit coverage for wildcard-synthesis and cross-file-impl Both modules previously had zero direct unit coverage — branches were exercised only through integration tests' happy paths. wildcard-synthesis.test.ts covers: Go graph-IMPORTS fallback, Python moduleAliasMap build, MAX_SYNTHETIC_BINDINGS_PER_FILE cap, dedup against existing namedImportMap entries, and empty-exportedSymbols early return. cross-file-impl.test.ts covers: gapRatio below threshold no-op, MAX_CROSS_FILE_REPROCESS cap, graph-only exportedTypeMap fallback, and empty namedImportMap short-circuit. Tests assert current behavior — any future regression flips them. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U7) * test(pipeline): golden-file graph-parity regression guard on mini-repo fixture Pins the current post-P1/P2 graph output (57 symbols, 92 relationships, 4 processes, deterministic edge digest) so future silent refactors cannot drift behavior unnoticed. If any count changes or any edge rewires, the test fails with a readable diff listing what changed and a copy-pasteable UPDATE_GOLDEN=1 regen command. Edge digest keyed by symbolic (label, name, filePath) triples rather than raw generateId output — stays meaningful across id-encoding refactors while still catching real semantic rewiring. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U8) * fix(pipeline): minimal cycle reporting + resolveEnclosingOwner loop safeguards U9: runner cycle detection now reports only the SCC members via DFS back-edge trace ('Cycle detected: A -> B -> C -> A') rather than everything with inDegree > 0 (which mixed cycle members with blocked dependents). Also emits the 'error' progress event for graph- validation failures, symmetric with U1's runtime-error path. U16: findEnclosingClassInfo now defends against language-provider hooks that return non-container nodes — visitedContainers Set breaks repeat-visit loops, MAX_ENCLOSING_WALK_ITERATIONS is belt-and-braces. Documented the hook contract invariant so future provider authors know the walk-continues-upward expectation. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U9, U16) * refactor(pipeline): type hygiene, dead code cleanup, shared allPathSet, graph-sort naming Bundles plan units U10, U11, U12, U14, U15: U10 — Type hygiene: readonly ParseOutput arrays (allExtractedRoutes, allDecoratorRoutes, allToolDefs, allORMQueries, allPaths); removed redundant 'as string[] | undefined' cast in routes.ts and 'as URL' in parse-impl.ts; WorkerPool is now 'import type'. Readonly contract propagated into processORMQueries (only iterates). U11 — Dead code & shims: deleted constants.ts shim (AST_CACHE_CAP inlined into its sole real consumer cross-file-impl.ts; isDev consumers now import directly from ../utils/env.js). Removed internal utility re-exports from pipeline-phases/index.ts (no external consumers). Removed topologicalLevelSort re-export from pipeline.ts; updated topological-sort.test.ts to import from the canonical utils/graph-sort.js. Stripped 'Phase 3+4:' stale JSDoc from parse-impl.ts. U12 — Perf: StructureOutput now carries allPathSet (ReadonlySet<string>) built once; cobol, markdown, and cross-file-impl consume the shared set instead of allocating their own. Parse forwards it via ParseOutput.allPathSet; processCobol/processMarkdown widened to ReadonlySet<string>. U14 — graph-sort.ts: renamed local 'inDegree' to 'pendingImportsPerFile' with expanded JSDoc explaining the reverse- graph Kahn's formulation and warning future maintainers not to 'correct' it to standard in-degree semantics. Added self-edge test. U15 — Unconditional worker-fallback logging: removed isDev guard on the worker-pool-creation-failure console.warn so operators can diagnose perf degradations in production. No behavior change. U8 golden-file test confirms pipeline output is byte-identical. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U10, U11, U12, U14, U15) * docs: fix ARCHITECTURE.md table integrity; bump AGENTS.md/CLAUDE.md to 1.3.0 U13 — documentation fixes: ARCHITECTURE.md: the prior insertion of the 'Pipeline Phase DAG' section orphaned 7 rows from the 'Where to change what' header. Moved those 7 rows back up under their header so the table reads contiguously; DAG section now follows the completed table. AGENTS.md + CLAUDE.md: bumped version 1.2.0 -> 1.3.0, updated Last reviewed to 2026-04-13, added matching Changelog row documenting the GitNexus index stats refresh after the DAG refactor. Stat bumps (symbols/relationships/execution flows) that were sitting uncommitted in the working tree are now landed under a proper changelog entry per each file's own documented schema. Plan: docs/plans/2026-04-13-001-fix-pipeline-dag-refactor-review-findings-plan.md (U13) * refactor(pipeline): drop spurious parse deps, true-readonly ParseOutput.exportedTypeMap, skip redundant wildcard synth - mro/communities/processes: switch redundant `parse` dep to `structure` — totalFiles originates in structure, so depending on parse for it was a spurious data dep that obscured the real DAG. - ParseOutput.exportedTypeMap: typed as truly ReadonlyMap<...,ReadonlyMap>>; graph→exports enrichment moved into parse-impl so the snapshot is fully populated at parse return. crossFile builds its own local mutable working copy for per-file re-resolution writes — no cast at the boundary. - parse-impl: hasSynthesized flag guards the unconditional final synthesizeWildcardImportBindings call when per-chunk/fallback synthesis already ran (graph-global + idempotent across chunks). - cross-file-impl: documented the intentional `phase: 'parsing'` progress label so telemetry bucketing stays consistent with the parse phase. - cross-file-impl test: replaced the now-moved fallback-enrichment assertion with a stronger one — crossFile must not mutate the parse-supplied map. Addresses PR #809 review pass 5 carry-overs. --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: magyargergo <11230420+magyargergo@users.noreply.github.com> Co-authored-by: Gergo Magyar <gergomagyar@icloud.com>
417 lines
12 KiB
TypeScript
417 lines
12 KiB
TypeScript
import { describe, it, expect } from 'vitest';
|
|
import { runPipeline } from '../../src/core/ingestion/pipeline-phases/runner.js';
|
|
import type {
|
|
PipelinePhase,
|
|
PipelineContext,
|
|
PhaseResult,
|
|
} from '../../src/core/ingestion/pipeline-phases/types.js';
|
|
import { getPhaseOutput } from '../../src/core/ingestion/pipeline-phases/types.js';
|
|
import { createKnowledgeGraph } from '../../src/core/graph/graph.js';
|
|
|
|
function makeCtx(): PipelineContext {
|
|
return {
|
|
repoPath: '/tmp/test',
|
|
graph: createKnowledgeGraph(),
|
|
onProgress: () => {},
|
|
pipelineStart: Date.now(),
|
|
};
|
|
}
|
|
|
|
describe('runPipeline', () => {
|
|
it('executes phases in dependency order', async () => {
|
|
const order: string[] = [];
|
|
|
|
const phaseA: PipelinePhase<string> = {
|
|
name: 'a',
|
|
deps: [],
|
|
async execute() {
|
|
order.push('a');
|
|
return 'resultA';
|
|
},
|
|
};
|
|
|
|
const phaseB: PipelinePhase<string> = {
|
|
name: 'b',
|
|
deps: ['a'],
|
|
async execute(_ctx, deps) {
|
|
const a = getPhaseOutput<string>(deps, 'a');
|
|
order.push('b');
|
|
return `${a}+B`;
|
|
},
|
|
};
|
|
|
|
const phaseC: PipelinePhase<string> = {
|
|
name: 'c',
|
|
deps: ['a'],
|
|
async execute(_ctx, deps) {
|
|
const a = getPhaseOutput<string>(deps, 'a');
|
|
order.push('c');
|
|
return `${a}+C`;
|
|
},
|
|
};
|
|
|
|
const phaseD: PipelinePhase<string> = {
|
|
name: 'd',
|
|
deps: ['b', 'c'],
|
|
async execute(_ctx, deps) {
|
|
const b = getPhaseOutput<string>(deps, 'b');
|
|
const c = getPhaseOutput<string>(deps, 'c');
|
|
order.push('d');
|
|
return `${b}|${c}`;
|
|
},
|
|
};
|
|
|
|
const results = await runPipeline([phaseD, phaseA, phaseC, phaseB], makeCtx());
|
|
|
|
// A must run before B and C; B and C must run before D
|
|
expect(order.indexOf('a')).toBeLessThan(order.indexOf('b'));
|
|
expect(order.indexOf('a')).toBeLessThan(order.indexOf('c'));
|
|
expect(order.indexOf('b')).toBeLessThan(order.indexOf('d'));
|
|
expect(order.indexOf('c')).toBeLessThan(order.indexOf('d'));
|
|
|
|
// Check outputs are correctly threaded
|
|
expect(results.get('d')?.output).toBe('resultA+B|resultA+C');
|
|
});
|
|
|
|
it('passes shared PipelineContext to every phase', async () => {
|
|
const ctx = makeCtx();
|
|
const seenContexts: PipelineContext[] = [];
|
|
|
|
const phase: PipelinePhase<void> = {
|
|
name: 'test',
|
|
deps: [],
|
|
async execute(c) {
|
|
seenContexts.push(c);
|
|
},
|
|
};
|
|
|
|
await runPipeline([phase], ctx);
|
|
expect(seenContexts).toHaveLength(1);
|
|
expect(seenContexts[0]).toBe(ctx);
|
|
});
|
|
|
|
it('records timing metadata in PhaseResult', async () => {
|
|
const phase: PipelinePhase<number> = {
|
|
name: 'slow',
|
|
deps: [],
|
|
async execute() {
|
|
await new Promise((r) => setTimeout(r, 10));
|
|
return 42;
|
|
},
|
|
};
|
|
|
|
const results = await runPipeline([phase], makeCtx());
|
|
const result = results.get('slow')!;
|
|
expect(result.phaseName).toBe('slow');
|
|
expect(result.output).toBe(42);
|
|
expect(result.durationMs).toBeGreaterThanOrEqual(0);
|
|
});
|
|
|
|
it('rejects duplicate phase names', async () => {
|
|
const phaseA: PipelinePhase = {
|
|
name: 'dup',
|
|
deps: [],
|
|
async execute() {},
|
|
};
|
|
const phaseB: PipelinePhase = {
|
|
name: 'dup',
|
|
deps: [],
|
|
async execute() {},
|
|
};
|
|
|
|
await expect(runPipeline([phaseA, phaseB], makeCtx())).rejects.toThrow(/Duplicate phase name/);
|
|
});
|
|
|
|
it('rejects missing dependencies', async () => {
|
|
const phase: PipelinePhase = {
|
|
name: 'orphan',
|
|
deps: ['nonexistent'],
|
|
async execute() {},
|
|
};
|
|
|
|
await expect(runPipeline([phase], makeCtx())).rejects.toThrow(/depends on 'nonexistent'/);
|
|
});
|
|
|
|
it('rejects cyclic dependencies', async () => {
|
|
const phaseA: PipelinePhase = {
|
|
name: 'x',
|
|
deps: ['y'],
|
|
async execute() {},
|
|
};
|
|
const phaseB: PipelinePhase = {
|
|
name: 'y',
|
|
deps: ['x'],
|
|
async execute() {},
|
|
};
|
|
|
|
await expect(runPipeline([phaseA, phaseB], makeCtx())).rejects.toThrow(/Cycle detected/);
|
|
});
|
|
|
|
it('reports only cycle members (not transitive dependents) in cycle error', async () => {
|
|
// A <-> B is the actual cycle. C, D, E are downstream and would also have
|
|
// inDegree > 0 after Kahn's drains, but they are NOT cycle members.
|
|
const phases: PipelinePhase[] = [
|
|
{ name: 'a', deps: ['b'], async execute() {} },
|
|
{ name: 'b', deps: ['a'], async execute() {} },
|
|
{ name: 'c', deps: ['a'], async execute() {} },
|
|
{ name: 'd', deps: ['c'], async execute() {} },
|
|
{ name: 'e', deps: ['c'], async execute() {} },
|
|
];
|
|
|
|
let caught: Error | undefined;
|
|
try {
|
|
await runPipeline(phases, makeCtx());
|
|
} catch (err) {
|
|
caught = err as Error;
|
|
}
|
|
|
|
expect(caught).toBeDefined();
|
|
const msg = caught!.message;
|
|
// Cycle path must include both A and B
|
|
expect(msg).toMatch(/Cycle detected in pipeline phases: /);
|
|
expect(msg).toMatch(/\ba\b/);
|
|
expect(msg).toMatch(/\bb\b/);
|
|
// Transitive dependents must NOT appear in the cycle path itself,
|
|
// they should be summarized in the parenthetical.
|
|
const pathSection = msg.split('(')[0];
|
|
expect(pathSection).not.toMatch(/\bc\b/);
|
|
expect(pathSection).not.toMatch(/\bd\b/);
|
|
expect(pathSection).not.toMatch(/\be\b/);
|
|
expect(msg).toMatch(/3 transitive dependents blocked/);
|
|
});
|
|
|
|
it('reports the full path for a 3-phase cycle', async () => {
|
|
// A -> B -> C -> A
|
|
const phases: PipelinePhase[] = [
|
|
{ name: 'a', deps: ['c'], async execute() {} },
|
|
{ name: 'b', deps: ['a'], async execute() {} },
|
|
{ name: 'c', deps: ['b'], async execute() {} },
|
|
];
|
|
|
|
let caught: Error | undefined;
|
|
try {
|
|
await runPipeline(phases, makeCtx());
|
|
} catch (err) {
|
|
caught = err as Error;
|
|
}
|
|
|
|
expect(caught).toBeDefined();
|
|
const msg = caught!.message;
|
|
expect(msg).toMatch(/Cycle detected in pipeline phases: /);
|
|
// All three names must appear
|
|
expect(msg).toMatch(/\ba\b/);
|
|
expect(msg).toMatch(/\bb\b/);
|
|
expect(msg).toMatch(/\bc\b/);
|
|
// Path uses the " -> " arrow separator
|
|
expect(msg).toMatch(/ -> /);
|
|
// No transitive-dependent suffix when every leftover IS a cycle member
|
|
expect(msg).not.toMatch(/transitive dependent/);
|
|
});
|
|
|
|
it("emits a terminal 'error' progress event on cycle detection", async () => {
|
|
const events: { phase: string; message: string; detail?: string }[] = [];
|
|
const ctx: PipelineContext = {
|
|
...makeCtx(),
|
|
onProgress: (p) => {
|
|
events.push({ phase: p.phase, message: p.message, detail: p.detail });
|
|
},
|
|
};
|
|
|
|
const phases: PipelinePhase[] = [
|
|
{ name: 'x', deps: ['y'], async execute() {} },
|
|
{ name: 'y', deps: ['x'], async execute() {} },
|
|
];
|
|
|
|
await expect(runPipeline(phases, ctx)).rejects.toThrow(/Cycle detected/);
|
|
|
|
const errorEvents = events.filter((e) => e.phase === 'error');
|
|
expect(errorEvents).toHaveLength(1);
|
|
expect(errorEvents[0].detail).toMatch(/Cycle detected/);
|
|
});
|
|
|
|
it('executes a single root phase with no deps', async () => {
|
|
const phase: PipelinePhase<string> = {
|
|
name: 'root',
|
|
deps: [],
|
|
async execute() {
|
|
return 'hello';
|
|
},
|
|
};
|
|
|
|
const results = await runPipeline([phase], makeCtx());
|
|
expect(results.get('root')?.output).toBe('hello');
|
|
});
|
|
|
|
it('handles a linear chain correctly', async () => {
|
|
const order: string[] = [];
|
|
|
|
const phases: PipelinePhase<number>[] = [];
|
|
for (let i = 0; i < 5; i++) {
|
|
const idx = i;
|
|
phases.push({
|
|
name: `step${i}`,
|
|
deps: i > 0 ? [`step${i - 1}`] : [],
|
|
async execute(_ctx, deps) {
|
|
if (idx > 0) {
|
|
const prev = getPhaseOutput<number>(deps, `step${idx - 1}`);
|
|
order.push(`step${idx}`);
|
|
return prev + 1;
|
|
}
|
|
order.push(`step${idx}`);
|
|
return 0;
|
|
},
|
|
});
|
|
}
|
|
|
|
const results = await runPipeline(phases, makeCtx());
|
|
expect(results.get('step4')?.output).toBe(4);
|
|
expect(order).toEqual(['step0', 'step1', 'step2', 'step3', 'step4']);
|
|
});
|
|
|
|
it('wraps phase Error with phase name and preserves cause', async () => {
|
|
const original = new Error('boom');
|
|
const phase: PipelinePhase = {
|
|
name: 'failing',
|
|
deps: [],
|
|
async execute() {
|
|
throw original;
|
|
},
|
|
};
|
|
|
|
await expect(runPipeline([phase], makeCtx())).rejects.toThrow(/Phase 'failing' failed: boom/);
|
|
|
|
try {
|
|
await runPipeline([phase], makeCtx());
|
|
throw new Error('expected runPipeline to reject');
|
|
} catch (err) {
|
|
expect(err).toBeInstanceOf(Error);
|
|
expect((err as Error).message).toMatch(/Phase 'failing' failed: boom/);
|
|
expect((err as Error & { cause?: unknown }).cause).toBe(original);
|
|
}
|
|
});
|
|
|
|
it('surfaces phase name when phase throws a non-Error value', async () => {
|
|
const phaseString: PipelinePhase = {
|
|
name: 'string-thrower',
|
|
deps: [],
|
|
async execute() {
|
|
throw 'oops';
|
|
},
|
|
};
|
|
|
|
await expect(runPipeline([phaseString], makeCtx())).rejects.toThrow(
|
|
/Phase 'string-thrower' failed: oops/,
|
|
);
|
|
|
|
const phaseNumber: PipelinePhase = {
|
|
name: 'number-thrower',
|
|
deps: [],
|
|
async execute() {
|
|
throw 42;
|
|
},
|
|
};
|
|
|
|
await expect(runPipeline([phaseNumber], makeCtx())).rejects.toThrow(
|
|
/Phase 'number-thrower' failed: 42/,
|
|
);
|
|
});
|
|
|
|
it("emits a terminal 'error' progress event exactly once on phase failure", async () => {
|
|
const events: { phase: string; message: string; detail?: string }[] = [];
|
|
const ctx: PipelineContext = {
|
|
...makeCtx(),
|
|
onProgress: (p) => {
|
|
events.push({ phase: p.phase, message: p.message, detail: p.detail });
|
|
},
|
|
};
|
|
|
|
const phase: PipelinePhase = {
|
|
name: 'failing',
|
|
deps: [],
|
|
async execute() {
|
|
throw new Error('kaboom');
|
|
},
|
|
};
|
|
|
|
await expect(runPipeline([phase], ctx)).rejects.toThrow(/Phase 'failing' failed/);
|
|
|
|
const errorEvents = events.filter((e) => e.phase === 'error');
|
|
expect(errorEvents).toHaveLength(1);
|
|
expect(errorEvents[0].message).toMatch(/failing/);
|
|
expect(errorEvents[0].detail).toBe('kaboom');
|
|
});
|
|
|
|
it('still rejects when onProgress handler throws during error reporting', async () => {
|
|
const original = new Error('underlying');
|
|
const ctx: PipelineContext = {
|
|
...makeCtx(),
|
|
onProgress: () => {
|
|
throw new Error('handler exploded');
|
|
},
|
|
};
|
|
|
|
const phase: PipelinePhase = {
|
|
name: 'failing',
|
|
deps: [],
|
|
async execute() {
|
|
throw original;
|
|
},
|
|
};
|
|
|
|
try {
|
|
await runPipeline([phase], ctx);
|
|
throw new Error('expected runPipeline to reject');
|
|
} catch (err) {
|
|
// The original phase error must win, not the handler's error.
|
|
expect((err as Error).message).toMatch(/Phase 'failing' failed: underlying/);
|
|
expect((err as Error & { cause?: unknown }).cause).toBe(original);
|
|
}
|
|
});
|
|
|
|
it('only exposes declared deps to each phase', async () => {
|
|
const phaseA: PipelinePhase<string> = {
|
|
name: 'a',
|
|
deps: [],
|
|
async execute() {
|
|
return 'resultA';
|
|
},
|
|
};
|
|
|
|
const phaseB: PipelinePhase<string> = {
|
|
name: 'b',
|
|
deps: ['a'],
|
|
async execute() {
|
|
return 'resultB';
|
|
},
|
|
};
|
|
|
|
// C depends on B but not A — should not see A's result
|
|
const phaseC: PipelinePhase<string> = {
|
|
name: 'c',
|
|
deps: ['b'],
|
|
async execute(_ctx, deps) {
|
|
expect(deps.has('b')).toBe(true);
|
|
expect(deps.has('a')).toBe(false);
|
|
return 'resultC';
|
|
},
|
|
};
|
|
|
|
await runPipeline([phaseA, phaseB, phaseC], makeCtx());
|
|
});
|
|
});
|
|
|
|
describe('getPhaseOutput', () => {
|
|
it('retrieves typed output from dependency map', () => {
|
|
const deps = new Map<string, PhaseResult<unknown>>();
|
|
deps.set('test', { phaseName: 'test', output: { value: 42 }, durationMs: 0 });
|
|
|
|
const result = getPhaseOutput<{ value: number }>(deps, 'test');
|
|
expect(result.value).toBe(42);
|
|
});
|
|
|
|
it('throws for missing phase', () => {
|
|
const deps = new Map<string, PhaseResult<unknown>>();
|
|
|
|
expect(() => getPhaseOutput(deps, 'missing')).toThrow(/Phase 'missing' not found/);
|
|
});
|
|
});
|