tools/core/orchestrators/pipeline.orchestrator.ts
tools/core/orchestrators/pipeline.orchestrator.ts is a file in Coordination Surface. 150 lines of code and 32 definitions.
import type { RuleRun, RunOptions, RunResult } from "../types/pipeline.types.ts";
import type { StepOptions, StepOutcome } from "../types/rule.types.ts";
import { existsSync, statSync } from "node:fs";
import { selectedRules, walkRule } from "../coordinators/rule.coordinator.ts";
import { cleanStage } from "../steps/source.step.ts";
import { discoverRules } from "../registries/rule.registry.ts";
import { gateStage } from "../steps/gate.step.ts";
import { loadTaxonomy } from "../resolvers/taxonomy.resolver.ts";
import { qualityStage } from "../steps/quality.step.ts";
import { readSource } from "../iterators/file.iterator.ts";
import { resolve } from "node:path";
import { resolveScope } from "../resolvers/scope.resolver.ts";
import { snapshotStage } from "../steps/snapshot.step.ts";
import { typecheckStage } from "../steps/typecheck.step.ts";
const WHOLE_SCOPE = "whole";
const withinDeclaredScope = function withinDeclaredScope(
path: string,
scope: string,
handed: readonly string[],
): boolean {
if (scope === WHOLE_SCOPE) {
return true;
}
return handed.includes(path);
};
const narrowingsOf = function narrowingsOf(options: RunOptions): string[] {
return [
options.scope === null ? null : `path=${options.scope}`,
options.ruleId === null ? null : `rule=${options.ruleId}`,
options.stage === null ? null : `stage=${options.stage}`,
options.bypass.length > 0 ? `bypassed=${options.bypass.join("+")}` : null,
].filter((entry): entry is string => entry !== null);
};
const stampOf = function stampOf(absolute: string): number {
return existsSync(absolute) ? statSync(absolute).mtimeMs : 0;
};
const createReader = function createReader(
repoRoot: string,
paths: readonly string[],
): { read: (path: string) => string; stamps: Map<string, number> } {
const cache = new Map<string, string>();
const stamps = new Map<string, number>();
const read = (path: string): string => {
const hit = cache.get(path);
if (hit !== undefined) {
return hit;
}
const absolute = resolve(repoRoot, path);
const source = readSource(absolute);
cache.set(path, source);
stamps.set(path, stampOf(absolute));
return source;
};
for (const path of paths) {
stamps.set(path, stampOf(resolve(repoRoot, path)));
}
return { read, stamps };
};
const runSteps = function runSteps(
steps: readonly ((options: StepOptions) => StepOutcome)[],
stepOptions: StepOptions,
paths: readonly string[],
): StepOutcome[] {
return [...steps.map((step) => step(stepOptions)), snapshotStage(stepOptions, paths)];
};
const movements = function movements(
repoRoot: string,
stamps: ReadonlyMap<string, number>,
written: readonly string[],
): { moved: string[]; movedByThisRun: string[] } {
const changed = [...stamps]
.filter(([path, stamp]) => stampOf(resolve(repoRoot, path)) !== stamp)
.map(([path]) => path);
return {
moved: changed.filter((path) => !written.includes(path)),
movedByThisRun: changed.filter((path) => written.includes(path)),
};
};
const unresolvedRun = function unresolvedRun(
authoritative: boolean,
registered: number,
scope: string,
unresolvedScope: string,
): RunResult {
return {
authoritative,
bypassedAny: false,
escaped: [],
findings: [],
incomparable: [],
moved: [],
movedByThisRun: [],
registered,
scanned: 0,
scope,
stages: [],
unfulfilled: [],
unresolvedScope,
written: [],
};
};
export const runPipeline = async function runPipeline(options: RunOptions): Promise<RunResult> {
const narrowings = narrowingsOf(options);
const scope = narrowings.length === 0 ? WHOLE_SCOPE : narrowings.join(" · ");
const authoritative = narrowings.length === 0;
const taxonomy = loadTaxonomy();
const registry = await discoverRules(options.repoRoot);
const { byJurisdiction } = resolveScope(options.repoRoot, taxonomy, options.scope);
const paths = byJurisdiction.all;
if (options.scope !== null && paths.length === 0) {
return unresolvedRun(authoritative, registry.rules.length, scope, options.scope);
}
const { read, stamps } = createReader(options.repoRoot, paths);
const stepOptions = {
authoritative,
bypass: options.bypass,
fix: options.fix,
repoRoot: options.repoRoot,
scanned: paths.length,
scope,
};
const whole = options.ruleId === null && options.stage === null;
const steps = whole ? [cleanStage, typecheckStage, qualityStage, gateStage] : [cleanStage];
const outcomes = runSteps(steps, stepOptions, paths);
const run: RuleRun = { authoritative, byJurisdiction, options, read, scope, taxonomy };
const walked = selectedRules(registry.rules, options).map(({ registered, stage }) =>
walkRule(registered, stage, run),
);
const written = walked.flatMap((outcome) => outcome.written);
const { moved, movedByThisRun } = movements(options.repoRoot, stamps, written);
return {
authoritative,
bypassedAny: outcomes.some((outcome) => outcome.stage.bypassed) || walked.some((outcome) => outcome.bypassed),
escaped: written.filter((path) => !withinDeclaredScope(path, scope, paths)),
findings: [
...registry.findings,
...outcomes.flatMap((outcome) => outcome.findings),
...walked.flatMap((outcome) => outcome.findings),
],
incomparable: walked.filter((outcome) => outcome.incomparable).map((outcome) => outcome.stage.rule),
moved,
movedByThisRun,
registered: registry.rules.length,
scanned: paths.length,
scope,
stages: [...outcomes.map((outcome) => outcome.stage), ...walked.map((outcome) => outcome.stage)],
unfulfilled: written.filter((path) => stamps.has(path) && !movedByThisRun.includes(path)),
written,
};
};