refactor(evidence): extract TypeScript lifecycle (#32)

This commit is contained in:
2026-08-24 02:50:56 +02:00
parent 840848df94
commit f375515dc0
8 changed files with 432 additions and 135 deletions
@@ -9,7 +9,7 @@ import {
writeFileSync,
} from "node:fs";
import { dirname, isAbsolute, join } from "node:path";
import { GitWorkspaceRepository } from "./git-repository.js";
import type { GitWorkspaceRepository } from "../git-repository.js";
export interface EvidenceMaterializationLimits {
maxEntries: number;
@@ -81,7 +81,10 @@ function writeExclusiveNoFollow(path: string, contents: Buffer, mode: number): v
}
export interface MaterializeEvidenceTreeOptions {
repository: GitWorkspaceRepository;
repository: Pick<
GitWorkspaceRepository,
"evidenceTreeObjects" | "evidenceTreeId" | "gitObjectSize" | "evidenceBlobBytes"
>;
revision: string;
id: string;
/** The workspace directory (e.g. `<staging>/<id>`) that will receive `evidence/` and the manifest. */
@@ -0,0 +1,177 @@
import { isIP } from "node:net";
import type { WorkspaceDescriptor } from "../schema.js";
type EvidenceConfig = WorkspaceDescriptor["evidence"];
type SemanticFailureCode = "workspace_not_activatable" | "semantic_index_incompatible";
export interface EvidenceJobState {
runId: string;
completedStages: string[];
childRuns: Record<string, string>;
}
export interface EvidencePreprocessingDependencies {
runStage(argv: string[]): Promise<Record<string, unknown>>;
persistJob(): void;
semanticPreflight(): Promise<{ ok: true } | { ok: false; code: SemanticFailureCode }>;
requireRunId(value: unknown): string;
numberRecord(value: unknown): Record<string, number> | undefined;
}
export interface EvidencePreprocessingRequest {
evidence: EvidenceConfig;
job: EvidenceJobState;
dryRun?: boolean;
httpPrivateHostAllowlist?: readonly string[];
}
export interface EvidencePreprocessingOutcome {
status: "succeeded" | "unchanged" | "dry_run" | "failed";
code: "ok" | "egress_policy_refused" | SemanticFailureCode;
runId?: string;
childRuns?: Record<string, string>;
completedStages?: string[];
counts?: Record<string, number>;
warnings?: string[];
}
function isPrivateHost(hostname: string): boolean {
if (hostname === "localhost" || hostname === "metadata.google.internal") return true;
const address = isIP(hostname);
if (address === 4) {
if (/^127\./.test(hostname) || /^10\./.test(hostname) || /^192\.168\./.test(hostname)) {
return true;
}
if (/^169\.254\./.test(hostname) || /^0\./.test(hostname)) return true;
const match = /^172\.(\d+)\./.exec(hostname);
return Boolean(match && Number(match[1]) >= 16 && Number(match[1]) <= 31);
}
if (address === 6) {
const normalized = hostname.toLowerCase();
return normalized === "::1"
|| normalized.startsWith("fe80:")
|| normalized.startsWith("fd")
|| normalized.startsWith("fc");
}
return hostname.endsWith(".internal");
}
function evidencePolicy(
evidence: EvidenceConfig,
httpPrivateHostAllowlist?: readonly string[],
): EvidencePreprocessingOutcome | undefined {
if (!evidence || evidence.source.type === "filesystem") return undefined;
if (evidence.source.type === "http") {
for (const value of evidence.source.uris) {
const host = new URL(value).hostname;
if (
isPrivateHost(host)
&& !(evidence.source.allow_private_hosts && httpPrivateHostAllowlist?.includes(host))
) {
return { status: "failed", code: "egress_policy_refused" };
}
}
return undefined;
}
if (
evidence.source.endpoint_url !== undefined
|| evidence.source.credentials === "ambient"
|| evidence.source.allow_private_endpoint
|| evidence.source.allow_insecure_endpoint
) {
return { status: "failed", code: "egress_policy_refused" };
}
return undefined;
}
function jobResult(job: EvidenceJobState): Pick<
EvidencePreprocessingOutcome,
"runId" | "childRuns" | "completedStages"
> {
return {
runId: job.runId,
childRuns: { ...job.childRuns },
completedStages: [...job.completedStages],
};
}
async function runEvidenceStage(
request: EvidencePreprocessingRequest,
deps: EvidencePreprocessingDependencies,
): Promise<EvidencePreprocessingOutcome> {
const payload = await deps.runStage([
"preprocess",
"evidence",
...(request.dryRun ? ["--dry-run"] : []),
...(request.job.childRuns.evidence
? ["--resume", request.job.childRuns.evidence]
: []),
"--json",
"-c",
"/dev/fd/3",
]);
if (typeof payload.run_id === "string") {
request.job.childRuns.evidence = deps.requireRunId(payload.run_id);
}
if (!request.dryRun && !request.job.completedStages.includes("evidence")) {
request.job.completedStages.push("evidence");
}
deps.persistJob();
return {
status: request.dryRun ? "dry_run" : "succeeded",
code: "ok",
...jobResult(request.job),
counts: deps.numberRecord(payload.counts),
};
}
export async function preprocessEvidence(
request: EvidencePreprocessingRequest,
deps: EvidencePreprocessingDependencies,
): Promise<EvidencePreprocessingOutcome> {
if (!request.evidence) {
return {
status: "unchanged",
code: "ok",
warnings: ["workspace has no Evidence source"],
};
}
const policy = evidencePolicy(request.evidence, request.httpPrivateHostAllowlist);
if (policy) return policy;
const semantic = await deps.semanticPreflight();
if (!semantic.ok) {
return { status: "failed", code: semantic.code, runId: request.job.runId };
}
if (request.job.completedStages.includes("evidence") && !request.dryRun) {
return {
status: "unchanged",
code: "ok",
runId: request.job.runId,
completedStages: [...request.job.completedStages],
};
}
return await runEvidenceStage(request, deps);
}
export async function continueEvidencePreprocessing(
request: Omit<EvidencePreprocessingRequest, "dryRun"> & {
priorCounts?: Record<string, number>;
},
deps: EvidencePreprocessingDependencies,
): Promise<EvidencePreprocessingOutcome> {
if (!request.evidence) {
return {
status: "succeeded",
code: "ok",
...jobResult(request.job),
warnings: ["workspace has no Evidence source"],
...(request.priorCounts ? { counts: request.priorCounts } : {}),
};
}
const policy = evidencePolicy(request.evidence, request.httpPrivateHostAllowlist);
if (policy) return { ...policy, ...jobResult(request.job) };
if (!request.job.completedStages.includes("evidence")) {
return await runEvidenceStage(request, deps);
}
return { status: "unchanged", code: "ok", ...jobResult(request.job) };
}
+47 -128
View File
@@ -1,7 +1,12 @@
import { createHash, randomBytes } from "node:crypto";
import { readdirSync, readFileSync, rmSync, writeFileSync, mkdirSync } from "node:fs";
import { isIP } from "node:net";
import { join } from "node:path";
import {
continueEvidencePreprocessing,
preprocessEvidence as runEvidencePreprocessing,
type EvidencePreprocessingDependencies,
type EvidencePreprocessingOutcome,
} from "./evidence/preprocessing.js";
import type { WorkspaceDescriptor } from "./schema.js";
import {
PreprocessingStateStore,
@@ -104,26 +109,6 @@ function baseResult(
};
}
function isPrivateHost(hostname: string): boolean {
if (hostname === "localhost" || hostname === "metadata.google.internal") return true;
const address = isIP(hostname);
if (address === 4) {
if (/^127\./.test(hostname) || /^10\./.test(hostname) || /^192\.168\./.test(hostname)) return true;
if (/^169\.254\./.test(hostname) || /^0\./.test(hostname)) return true;
const match = /^172\.(\d+)\./.exec(hostname);
return Boolean(match && Number(match[1]) >= 16 && Number(match[1]) <= 31);
}
if (address === 6) {
const normalized = hostname.toLowerCase();
return normalized === "::1" || normalized.startsWith("fe80:") || normalized.startsWith("fd") || normalized.startsWith("fc");
}
return hostname.endsWith(".internal");
}
function noEvidenceWarning(workspace: WorkspaceDescriptor): string[] {
return workspace.evidence === undefined ? ["workspace has no Evidence source"] : [];
}
export class WorkspacePreprocessingService {
constructor(private readonly deps: WorkspacePreprocessingServiceDeps) {}
@@ -332,36 +317,16 @@ export class WorkspacePreprocessingService {
async preprocessEvidence(options: { workspaceId: string; dryRun?: boolean; resumeRunId?: string }): Promise<WorkspaceOperationResult> {
const scope = await this.startRun(options.workspaceId, "preprocess evidence", options.resumeRunId);
if (scope.runtime.workspace.evidence === undefined) {
return baseResult(scope.runtime, "preprocess evidence", "unchanged", "ok", {
warnings: noEvidenceWarning(scope.runtime.workspace),
});
}
const policy = this.evidencePolicy(scope.runtime.workspace);
if (policy !== undefined) return baseResult(scope.runtime, "preprocess evidence", policy.status, policy.code, { warnings: policy.warnings });
const semantic = await this.deps.semanticPreflight(scope.runtime.workspace);
if (!semantic.ok) return baseResult(scope.runtime, "preprocess evidence", "failed", semantic.code, { runId: scope.job.runId });
if (scope.job.completedStages.includes("evidence") && !options.dryRun) {
return baseResult(scope.runtime, "preprocess evidence", "unchanged", "ok", {
runId: scope.job.runId,
completedStages: [...scope.job.completedStages],
});
}
const payload = await this.runJsonStage(scope.runtime, [
"preprocess", "evidence",
...(options.dryRun ? ["--dry-run"] : []),
...(scope.job.childRuns.evidence ? ["--resume", scope.job.childRuns.evidence] : []),
"--json", "-c", "/dev/fd/3",
]);
if (typeof payload.run_id === "string") scope.job.childRuns.evidence = this.requireRunId(payload.run_id);
if (!options.dryRun && !scope.job.completedStages.includes("evidence")) scope.job.completedStages.push("evidence");
this.state(scope.runtime.workspaceId).writeJob(scope.job);
return baseResult(scope.runtime, "preprocess evidence", options.dryRun ? "dry_run" : "succeeded", "ok", {
runId: scope.job.runId,
childRuns: { ...scope.job.childRuns },
completedStages: [...scope.job.completedStages],
counts: this.numberRecord(payload.counts),
});
const outcome = await runEvidencePreprocessing(
{
evidence: scope.runtime.workspace.evidence,
job: scope.job,
dryRun: options.dryRun,
httpPrivateHostAllowlist: this.deps.httpPrivateHostAllowlist,
},
this.evidenceDependencies(scope),
);
return this.evidenceResult(scope, "preprocess evidence", outcome);
}
async run(options: { workspaceId: string; resumeRunId?: string }): Promise<WorkspaceOperationResult> {
@@ -401,59 +366,23 @@ export class WorkspacePreprocessingService {
}
const semantic = await this.deps.semanticPreflight(scope.runtime.workspace);
if (!semantic.ok) return baseResult(scope.runtime, "preprocess run", "failed", semantic.code, { runId: scope.job.runId });
let schemaCounts: Record<string, number> | undefined;
if (!scope.job.completedStages.includes("schema_index")) {
const payload = await this.runJsonStage(scope.runtime, ["vector", "index-schema", "--json", "-c", "/dev/fd/3"]);
scope.job.completedStages.push("schema_index");
this.state(scope.runtime.workspaceId).writeJob(scope.job);
const warnings = noEvidenceWarning(scope.runtime.workspace);
if (scope.runtime.workspace.evidence === undefined) {
return baseResult(scope.runtime, "preprocess run", "succeeded", "ok", {
runId: scope.job.runId,
childRuns: { ...scope.job.childRuns },
completedStages: [...scope.job.completedStages],
counts: this.numberRecord(payload.counts),
warnings,
});
}
schemaCounts = this.numberRecord(payload.counts);
}
if (scope.runtime.workspace.evidence === undefined) {
return baseResult(scope.runtime, "preprocess run", "succeeded", "ok", {
runId: scope.job.runId,
childRuns: { ...scope.job.childRuns },
completedStages: [...scope.job.completedStages],
warnings: noEvidenceWarning(scope.runtime.workspace),
});
}
const policy = this.evidencePolicy(scope.runtime.workspace);
if (policy !== undefined) {
return baseResult(scope.runtime, "preprocess run", policy.status, policy.code, {
runId: scope.job.runId,
childRuns: { ...scope.job.childRuns },
completedStages: [...scope.job.completedStages],
warnings: policy.warnings,
});
}
if (!scope.job.completedStages.includes("evidence")) {
const payload = await this.runJsonStage(scope.runtime, [
"preprocess", "evidence",
...(scope.job.childRuns.evidence ? ["--resume", scope.job.childRuns.evidence] : []),
"--json", "-c", "/dev/fd/3",
]);
if (typeof payload.run_id === "string") scope.job.childRuns.evidence = this.requireRunId(payload.run_id);
scope.job.completedStages.push("evidence");
this.state(scope.runtime.workspaceId).writeJob(scope.job);
return baseResult(scope.runtime, "preprocess run", "succeeded", "ok", {
runId: scope.job.runId,
childRuns: { ...scope.job.childRuns },
completedStages: [...scope.job.completedStages],
counts: this.numberRecord(payload.counts),
});
}
return baseResult(scope.runtime, "preprocess run", "unchanged", "ok", {
runId: scope.job.runId,
childRuns: { ...scope.job.childRuns },
completedStages: [...scope.job.completedStages],
});
const outcome = await continueEvidencePreprocessing(
{
evidence: scope.runtime.workspace.evidence,
job: scope.job,
httpPrivateHostAllowlist: this.deps.httpPrivateHostAllowlist,
priorCounts: schemaCounts,
},
this.evidenceDependencies(scope),
);
return this.evidenceResult(scope, "preprocess run", outcome);
}
private async startRun(workspaceId: string, operation: string, resumeRunId?: string): Promise<RunScope> {
@@ -476,6 +405,25 @@ export class WorkspacePreprocessingService {
return new PreprocessingStateStore({ dataRoot: this.deps.dataRoot, workspaceId });
}
private evidenceDependencies(scope: RunScope): EvidencePreprocessingDependencies {
return {
runStage: async (argv) => await this.runJsonStage(scope.runtime, argv),
persistJob: () => this.state(scope.runtime.workspaceId).writeJob(scope.job),
semanticPreflight: async () => await this.deps.semanticPreflight(scope.runtime.workspace),
requireRunId: (value) => this.requireRunId(value),
numberRecord: (value) => this.numberRecord(value),
};
}
private evidenceResult(
scope: RunScope,
operation: "preprocess evidence" | "preprocess run",
outcome: EvidencePreprocessingOutcome,
): WorkspaceOperationResult {
const { status, code, ...extra } = outcome;
return baseResult(scope.runtime, operation, status, code, extra);
}
private async runSuggestStage(
scope: RunScope,
fromSql: ReadonlyArray<{ name: string; sql: string }>,
@@ -581,33 +529,4 @@ export class WorkspacePreprocessingService {
return undefined;
}
private evidencePolicy(workspace: WorkspaceDescriptor): {
status: WorkspaceOperationResult["status"];
code: WorkspaceOperationResult["code"];
warnings?: string[];
} | undefined {
const evidence = workspace.evidence;
if (!evidence) return undefined;
// P6: filesystem Evidence is materialized from the pinned commit at activation, so the
// engine may proceed directly against the immutable revision content root.
if (evidence.source.type === "filesystem") return undefined;
if (evidence.source.type === "http") {
for (const value of evidence.source.uris) {
const host = new URL(value).hostname;
if (isPrivateHost(host) && !(evidence.source.allow_private_hosts && this.deps.httpPrivateHostAllowlist?.includes(host))) {
return { status: "failed", code: "egress_policy_refused" };
}
}
return undefined;
}
if (
evidence.source.endpoint_url !== undefined
|| evidence.source.credentials === "ambient"
|| evidence.source.allow_private_endpoint
|| evidence.source.allow_insecure_endpoint
) {
return { status: "failed", code: "egress_policy_refused" };
}
return undefined;
}
}
+1 -1
View File
@@ -5,7 +5,7 @@ import { isAbsolute, join } from "node:path";
import { buildInstallationContract, renderWorkspaceDocs } from "./contracts.js";
import { parseAnnotationsYaml } from "./annotations.js";
import { syncAnnotations } from "./annotations-sync.js";
import { materializeEvidenceTree } from "./evidence-materialization.js";
import { materializeEvidenceTree } from "./evidence/materialization.js";
import { assertCatalogMatchesDescriptor, parseWorkspaceCatalogYaml, type WorkspaceCatalog, type WorkspaceCatalogEntry } from "./catalog.js";
import {
GitWorkspaceRepository,