feat: parallelize dev ci cd artifact flow

This commit is contained in:
Codex
2026-05-24 04:36:38 +00:00
parent 65583b79d7
commit 86e7df92a0
15 changed files with 1190 additions and 150 deletions
+321 -34
View File
@@ -1,6 +1,6 @@
#!/usr/bin/env node
import { constants as fsConstants } from "node:fs";
import { access, readFile, writeFile } from "node:fs/promises";
import { access, mkdir, readFile, writeFile } from "node:fs/promises";
import { request as httpRequest } from "node:http";
import path from "node:path";
import { fileURLToPath } from "node:url";
@@ -88,6 +88,8 @@ const forbiddenActions = [
"force-push"
];
const runtimeIdentityEnvNames = ["HWLAB_COMMIT_ID", "HWLAB_IMAGE", "HWLAB_IMAGE_TAG"];
const digestPattern = /^sha256:[a-f0-9]{64}$/u;
const defaultCdConcurrency = parseBoundedInteger(process.env.HWLAB_DEV_CD_CONCURRENCY, 4, 1, 8);
export function parseArgs(argv) {
const flags = new Set();
@@ -95,6 +97,7 @@ export function parseArgs(argv) {
let kubeconfig = null;
let kubeconfigSpecified = false;
let reportPath = defaultReportPath;
let rolloutConcurrency = defaultCdConcurrency;
for (let index = 0; index < argv.length; index += 1) {
const arg = argv[index];
@@ -141,6 +144,22 @@ export function parseArgs(argv) {
}
continue;
}
if (arg === "--rollout-concurrency") {
flags.add("--rollout-concurrency");
const next = argv[index + 1];
if (typeof next !== "string" || next.startsWith("--")) {
errors.push("--rollout-concurrency requires a numeric value");
} else {
rolloutConcurrency = parseBoundedInteger(next, defaultCdConcurrency, 1, 8);
index += 1;
}
continue;
}
if (arg.startsWith("--rollout-concurrency=")) {
flags.add("--rollout-concurrency");
rolloutConcurrency = parseBoundedInteger(arg.slice("--rollout-concurrency=".length), defaultCdConcurrency, 1, 8);
continue;
}
flags.add(arg);
}
@@ -151,6 +170,7 @@ export function parseArgs(argv) {
skipLiveProbe: flags.has("--skip-live-probe"),
writeReport: flags.has("--report-output"),
reportPath,
rolloutConcurrency,
kubeconfig,
kubeconfigSpecified,
errors,
@@ -158,6 +178,12 @@ export function parseArgs(argv) {
};
}
function parseBoundedInteger(value, fallback, min, max) {
const parsed = Number.parseInt(String(value ?? ""), 10);
if (!Number.isInteger(parsed)) return fallback;
return Math.min(Math.max(parsed, min), max);
}
function addBlocker(blockers, type, scope, summary) {
const key = `${type}::${scope}`;
if (!blockers.some((blocker) => `${blocker.type}::${blocker.scope}` === key)) {
@@ -169,6 +195,21 @@ function oneLine(value) {
return String(value).replace(/\s+/g, " ").trim();
}
async function mapWithConcurrency(items, concurrency, worker) {
if (items.length === 0) return [];
const results = new Array(items.length);
let nextIndex = 0;
const workerCount = Math.min(Math.max(concurrency, 1), items.length);
await Promise.all(Array.from({ length: workerCount }, async () => {
while (nextIndex < items.length) {
const index = nextIndex;
nextIndex += 1;
results[index] = await worker(items[index], index);
}
}));
return results;
}
async function readJson(relativePath, blockers) {
try {
return JSON.parse(await readFile(path.join(repoRoot, relativePath), "utf8"));
@@ -675,6 +716,146 @@ function parseImageTag(image) {
return image.slice(colonIndex + 1);
}
function parseTaggedImageReference(image) {
if (typeof image !== "string" || image.includes("@")) return null;
const firstSlash = image.indexOf("/");
const lastSlash = image.lastIndexOf("/");
const colonIndex = image.lastIndexOf(":");
if (firstSlash <= 0 || colonIndex <= lastSlash) return null;
return {
image,
registry: image.slice(0, firstSlash),
repository: image.slice(firstSlash + 1, colonIndex),
tag: image.slice(colonIndex + 1)
};
}
function registryManifestUrlForImage(image) {
const parsed = parseTaggedImageReference(image);
if (!parsed) return null;
const repositoryPath = parsed.repository.split("/").map((part) => encodeURIComponent(part)).join("/");
return {
...parsed,
url: `http://${parsed.registry}/v2/${repositoryPath}/manifests/${encodeURIComponent(parsed.tag)}`
};
}
function registryAcceptHeader() {
return [
"application/vnd.docker.distribution.manifest.v2+json",
"application/vnd.oci.image.manifest.v1+json",
"application/vnd.docker.distribution.manifest.list.v2+json",
"application/vnd.oci.image.index.v1+json",
"application/vnd.docker.distribution.manifest.v1+json"
].join(", ");
}
function httpRequestText(url, { method = "GET", headers = {}, timeoutMs = 15000 } = {}) {
return new Promise((resolve, reject) => {
const request = httpRequest(url, { method, headers, timeout: timeoutMs }, (response) => {
response.setEncoding("utf8");
let body = "";
response.on("data", (chunk) => {
body += chunk;
if (body.length > 4096) {
request.destroy(new Error("registry manifest response exceeded 4096 bytes during CD verification"));
}
});
response.on("end", () => resolve({
statusCode: response.statusCode ?? 0,
headers: response.headers,
body
}));
});
request.on("timeout", () => {
request.destroy(new Error(`timeout after ${timeoutMs}ms`));
});
request.on("error", reject);
request.end();
});
}
async function verifyCatalogRegistryManifest(service, blockers) {
const image = service.image;
const expectedDigest = service.digest;
const parsed = registryManifestUrlForImage(image);
const base = {
serviceId: service.serviceId,
image,
imageTag: service.imageTag ?? parseImageTag(image),
expectedDigest,
url: parsed?.url ?? null
};
if (!parsed) {
const reason = `${service.serviceId} image is not a tagged registry image`;
addBlocker(blockers, "contract_blocker", `registry-manifest-${service.serviceId}`, reason);
return { ...base, status: "blocked", reason };
}
if (!digestPattern.test(expectedDigest ?? "")) {
const reason = `${service.serviceId} catalog digest is not an immutable sha256 digest`;
addBlocker(blockers, "contract_blocker", `registry-manifest-${service.serviceId}`, reason);
return { ...base, status: "blocked", reason };
}
try {
const response = await httpRequestText(parsed.url, {
method: "HEAD",
headers: {
Accept: registryAcceptHeader()
},
timeoutMs: 15000
});
const observedDigest = String(response.headers["docker-content-digest"] ?? "").trim();
const okStatus = response.statusCode >= 200 && response.statusCode < 300;
const digestMatches = observedDigest === expectedDigest;
if (!okStatus || !digestMatches) {
const reason = !okStatus
? `${service.serviceId} registry manifest returned HTTP ${response.statusCode}`
: `${service.serviceId} registry digest mismatch: expected ${expectedDigest}, observed ${observedDigest || "missing"}`;
addBlocker(blockers, "environment_blocker", `registry-manifest-${service.serviceId}`, reason);
return {
...base,
status: "blocked",
reason,
httpStatus: response.statusCode,
observedDigest,
digestMatches
};
}
return {
...base,
status: "pass",
httpStatus: response.statusCode,
observedDigest,
digestMatches: true
};
} catch (error) {
const reason = `${service.serviceId} registry manifest read failed: ${oneLine(error.message)}`;
addBlocker(blockers, "environment_blocker", `registry-manifest-${service.serviceId}`, reason);
return {
...base,
status: "blocked",
reason
};
}
}
async function verifyCatalogRegistryManifests(catalog, blockers, concurrency = defaultCdConcurrency) {
const services = (catalog?.services ?? []).filter((service) => service.artifactRequired !== false);
const results = await mapWithConcurrency(services, concurrency, async (service) =>
verifyCatalogRegistryManifest(service, blockers)
);
return {
status: results.every((result) => result.status === "pass") ? "pass" : "blocked",
source: "registry-manifest",
releaseGate: true,
serviceCount: results.length,
concurrency,
results
};
}
function expectedRuntimeIdentityEnv(artifact) {
return {
HWLAB_COMMIT_ID: artifact.sourceCommitId,
@@ -1675,9 +1856,9 @@ async function buildDevLegacySimulatorDeploymentCleanupPlan(kubectl, workloads,
return cleanups;
}
async function executeDevTemplateJobReplacements(kubectl, replacements, blockers) {
async function executeDevTemplateJobReplacements(kubectl, replacements, blockers, concurrency = defaultCdConcurrency) {
const planned = replacements.filter((replacement) => replacement.replace && replacement.result === "planned");
for (const replacement of planned) {
await mapWithConcurrency(planned, concurrency, async (replacement) => {
const commandArgs = [
"-n",
replacement.namespace,
@@ -1694,17 +1875,17 @@ async function executeDevTemplateJobReplacements(kubectl, replacements, blockers
replacement.result = "delete_failed";
replacement.reason = `Failed to delete live suspended template Job before apply: ${oneLine(commandOutput(result))}`;
addBlocker(blockers, "environment_blocker", `template-job-replace-${replacement.jobName}`, replacement.reason);
continue;
return;
}
replacement.result = "deleted_pending_recreate";
replacement.reason = "Live suspended template Job was deleted; kubectl apply must recreate it from the desired manifest.";
}
});
return planned.length;
}
async function executeDevLegacySimulatorDeploymentCleanups(kubectl, cleanups, blockers) {
async function executeDevLegacySimulatorDeploymentCleanups(kubectl, cleanups, blockers, concurrency = defaultCdConcurrency) {
const planned = cleanups.filter((cleanup) => cleanup.cleanup && cleanup.result === "planned");
for (const cleanup of planned) {
await mapWithConcurrency(planned, concurrency, async (cleanup) => {
const commandArgs = [
"-n",
cleanup.namespace,
@@ -1721,14 +1902,107 @@ async function executeDevLegacySimulatorDeploymentCleanups(kubectl, cleanups, bl
cleanup.result = "delete_failed";
cleanup.reason = `Failed to delete stale legacy simulator Deployment before apply: ${oneLine(commandOutput(result))}`;
addBlocker(blockers, "environment_blocker", `legacy-simulator-deployment-cleanup-${cleanup.deploymentName}`, cleanup.reason);
continue;
return;
}
cleanup.result = "deleted_pending_apply";
cleanup.reason = "Live stale legacy simulator Deployment was deleted; kubectl apply must leave the desired indexed StatefulSet in place.";
}
});
return planned.length;
}
function rolloutApiName(kind) {
if (kind === "Deployment") return "deployment";
if (kind === "StatefulSet") return "statefulset";
return null;
}
function rolloutTargetsFromWorkloads(workloads) {
return listItems(workloads)
.map((item) => {
const apiName = rolloutApiName(item?.kind);
if (!apiName) return null;
return {
kind: item.kind,
apiName,
name: item?.metadata?.name ?? "unknown",
namespace: item?.metadata?.namespace ?? namespace,
replicas: replicaPlan(item),
serviceIds: uniqueStrings(containersFor(item).map((container) => serviceIdFor(item, container)))
};
})
.filter(Boolean);
}
async function verifyRolloutTarget(kubectl, target, blockers) {
const commandArgs = [
"-n",
target.namespace,
"rollout",
"status",
`${target.apiName}/${target.name}`,
"--timeout=180s"
];
const base = {
...target,
command: kubectlCommand(kubectl, commandArgs)
};
if (kubectl.status !== "ready") {
addBlocker(blockers, "environment_blocker", `rollout-status-${target.name}`, kubectl.reason);
return {
...base,
status: "blocked",
reason: kubectl.reason
};
}
const result = await kubectlResult(kubectl, commandArgs, 190000);
if (!result.ok) {
const reason = oneLine(commandOutput(result));
addBlocker(blockers, "runtime_blocker", `rollout-status-${target.name}`, reason);
return {
...base,
status: "blocked",
reason,
stdout: result.redactedStdout ?? result.stdout,
stderr: result.redactedStderr ?? result.stderr
};
}
return {
...base,
status: "pass",
stdout: result.redactedStdout ?? result.stdout,
stderr: result.redactedStderr ?? result.stderr
};
}
async function verifyRolloutsAfterApply(args, kubectl, workloads, blockers, applyStep) {
const targets = rolloutTargetsFromWorkloads(workloads);
if (!args.apply) {
return {
status: "not_run",
reason: "dry-run mode",
concurrency: args.rolloutConcurrency,
targets
};
}
if (applyStep.status !== "pass") {
return {
status: "not_run",
reason: "apply did not complete",
concurrency: args.rolloutConcurrency,
targets
};
}
const results = await mapWithConcurrency(targets, args.rolloutConcurrency, async (target) =>
verifyRolloutTarget(kubectl, target, blockers)
);
return {
status: results.every((result) => result.status === "pass") ? "pass" : "blocked",
concurrency: args.rolloutConcurrency,
targetCount: targets.length,
results
};
}
function isExpectedTemplateJobImmutableFailure(result, replacementsNeeded) {
const plannedNames = replacementsNeeded.map((replacement) => replacement.jobName);
if (result.ok || plannedNames.length === 0) return false;
@@ -1758,8 +2032,8 @@ async function runApplyStep(args, kubectl, blockers, templateJobReplacements, le
const replacementsNeeded = templateJobReplacements.filter((replacement) => replacement.replace);
if (args.apply) {
const replaceCount = await executeDevTemplateJobReplacements(kubectl, templateJobReplacements, blockers);
const cleanupCount = await executeDevLegacySimulatorDeploymentCleanups(kubectl, legacySimulatorDeploymentCleanups, blockers);
const replaceCount = await executeDevTemplateJobReplacements(kubectl, templateJobReplacements, blockers, args.rolloutConcurrency);
const cleanupCount = await executeDevLegacySimulatorDeploymentCleanups(kubectl, legacySimulatorDeploymentCleanups, blockers, args.rolloutConcurrency);
mutationAttempted = replaceCount > 0 || cleanupCount > 0;
if (blockers.length > 0) {
return {
@@ -1858,6 +2132,7 @@ export async function runDevDeployApply(argv, io = {}) {
const commitId = resolveApplySourceCommit(deploy, catalog, gitHeadCommitId);
const artifactEvidence = validateDeployAndCatalog(deploy, catalog, commitId, blockers);
const registryManifests = await verifyCatalogRegistryManifests(catalog, blockers, args.rolloutConcurrency);
const k8sManifest = validateK8s(devKustomization, namespaceDoc, workloads, services, healthContract, deploy, catalog, blockers);
const cloudApiDb = await checkCloudApiDb(deploy, workloads, services, blockers);
const codeAgentProviderDesiredState = inspectCodeAgentProviderDesiredState(deploy, workloads);
@@ -1886,33 +2161,41 @@ export async function runDevDeployApply(argv, io = {}) {
templateJobReplacements,
legacySimulatorDeploymentCleanups
);
const cloudWebRolloutAfterApply = args.apply && applyStep.status === "pass"
? await observeDeploymentRollout(kubectl, deploy, catalog, "hwlab-cloud-web", blockers, {
let rolloutStatusAfterApply;
let cloudWebRolloutAfterApply;
let codeAgentProviderLiveAfterApply;
if (args.apply && applyStep.status === "pass") {
[rolloutStatusAfterApply, cloudWebRolloutAfterApply, codeAgentProviderLiveAfterApply] = await Promise.all([
verifyRolloutsAfterApply(args, kubectl, workloads, blockers, applyStep),
observeDeploymentRollout(kubectl, deploy, catalog, "hwlab-cloud-web", blockers, {
observationPhase: "after_apply",
blockRuntimeEnvDrift: true
})
: skippedDeploymentRolloutObservation(
kubectl,
deploy,
catalog,
"hwlab-cloud-web",
"after_apply",
args.apply ? "apply did not complete successfully; post-apply rollout env verification was not run" : "dry-run mode; no post-apply rollout env verification"
);
const cloudWebRollout = args.apply ? cloudWebRolloutAfterApply : cloudWebRolloutBeforeApply;
const codeAgentProviderLiveAfterApply = args.apply && applyStep.status === "pass"
? await observeCodeAgentProviderLiveDeployment(kubectl, blockers, {
}),
observeCodeAgentProviderLiveDeployment(kubectl, blockers, {
phase: "after_apply",
blockOnMismatch: true
})
: {
phase: "after_apply",
status: "not_run",
reason: args.apply ? "apply did not complete" : "dry-run mode",
secretValuesRead: false,
kubernetesSecretDataRead: false,
valuesRedacted: true
};
]);
} else {
rolloutStatusAfterApply = await verifyRolloutsAfterApply(args, kubectl, workloads, blockers, applyStep);
cloudWebRolloutAfterApply = skippedDeploymentRolloutObservation(
kubectl,
deploy,
catalog,
"hwlab-cloud-web",
"after_apply",
args.apply ? "apply did not complete successfully; post-apply rollout env verification was not run" : "dry-run mode; no post-apply rollout env verification"
);
codeAgentProviderLiveAfterApply = {
phase: "after_apply",
status: "not_run",
reason: args.apply ? "apply did not complete" : "dry-run mode",
secretValuesRead: false,
kubernetesSecretDataRead: false,
valuesRedacted: true
};
}
const cloudWebRollout = args.apply ? cloudWebRolloutAfterApply : cloudWebRolloutBeforeApply;
const status = blockers.length > 0 ? "blocked" : "pass";
const artifactPlan = buildArtifactPlan(deploy, catalog, commitId, artifactEvidence);
const workloadPlan = buildWorkloadPlan(workloads);
@@ -1964,6 +2247,7 @@ export async function runDevDeployApply(argv, io = {}) {
evidence: [
`artifact services checked: ${artifactEvidence.length}`,
`expected artifact commit: ${artifactPlan.expectedArtifactCommit}`,
`registry manifests: ${registryManifests.status} (${registryManifests.serviceCount} services)`,
`namespace: ${namespace}`,
`workloads planned: ${workloadPlan.length}`,
`kubectl executor: ${kubectl.executor ?? "missing"}`,
@@ -1972,6 +2256,7 @@ export async function runDevDeployApply(argv, io = {}) {
`live health: ${liveProbe.status}`,
`code agent provider desired-state: ${codeAgentProviderDesiredState.status}`,
`code agent provider live env before apply: ${codeAgentProviderLiveBeforeApply.status}`,
`rollout status concurrency: ${args.rolloutConcurrency}`,
`template job replacements: ${templateJobReplacements.filter((replacement) => replacement.replace).length}`,
`legacy simulator deployment cleanups: ${legacySimulatorDeploymentCleanups.filter((cleanup) => cleanup.cleanup).length}`
],
@@ -1980,7 +2265,7 @@ export async function runDevDeployApply(argv, io = {}) {
devPreconditions: {
status,
requirements: [
"DEV artifact catalog must contain CI publish evidence and registry digests",
"DEV artifact catalog must contain CI publish expectations and registry digests; CD must verify those tag/digest pairs directly against registry manifests",
"kubectl must be available for D601 hwlab-dev and must not target PROD",
"hwlab-cloud-api /health/live must report serviceId, dev environment, and DB env readiness without exposing secret values",
"hwlab-cloud-api must declare and preserve OPENAI_API_KEY from hwlab-code-agent-provider/openai-api-key plus the DEV Code Agent egress proxy env without reading Secret values",
@@ -2009,6 +2294,7 @@ export async function runDevDeployApply(argv, io = {}) {
cloudWebRollout,
cloudWebRolloutBeforeApply,
cloudWebRolloutAfterApply,
rolloutStatusAfterApply,
codeAgentProvider: {
desiredState: codeAgentProviderDesiredState,
liveBeforeApply: codeAgentProviderLiveBeforeApply,
@@ -2024,6 +2310,7 @@ export async function runDevDeployApply(argv, io = {}) {
rollbackHint: buildRollbackHint(kubectl, workloads),
remainingBlockers,
artifactEvidence,
registryManifests,
cloudApiDb,
k8sManifest,
clusterObservation,