feat: improve resource reconciliation and observability

This commit is contained in:
2026-09-11 10:51:57 +00:00 Unverified
parent 0067a2cdcd
commit e3df2f4d76
28 changed files with 2407 additions and 217 deletions
+265 -52
View File
@@ -1,4 +1,4 @@
import { createHash } from "node:crypto";
import { createHash, randomUUID } from "node:crypto";
import { rm } from "node:fs/promises";
import { join } from "node:path";
import {
@@ -14,9 +14,11 @@ import {
BUILD_RECORD_API_VERSION,
BuildStoreConflictError,
type BuildRecord,
type BuildReconciliationLease,
type BuildStore,
type UploadRecord,
} from "./build-store";
import type { WorkspaceLease, WorkspaceLeaseProvider } from "./operation-store";
import {
materializeWorkspace,
parseWorkspaceManifest,
@@ -28,6 +30,8 @@ export const DEFAULT_MAX_BLOB_BYTES = 1024 * 1024 * 1024;
export const DEFAULT_MAX_UPLOAD_CHUNK_BYTES = 8 * 1024 * 1024;
const SNAPSHOT_CHECK_CONCURRENCY = 20;
export const DEFAULT_MAX_LOG_BYTES = 1024 * 1024;
const BUILD_RECONCILE_LEASE_TTL_MS = 30_000;
const BUILD_RECONCILE_HEARTBEAT_MS = BUILD_RECONCILE_LEASE_TTL_MS / 3;
export interface BuildCas extends MaterializeCas {
has(digest: Sha256Digest): Promise<boolean>;
@@ -53,6 +57,13 @@ export interface BuildKubernetesOperations {
deleteJob(namespace: string, name: string): Promise<void>;
}
type ReconcileOptions = { signal?: AbortSignal };
function throwIfAborted(signal?: AbortSignal): void {
if (signal?.aborted)
throw new Error("Build reconciliation scan was cancelled");
}
export interface BuildControllerOptions {
cas: BuildCas;
store: BuildStore;
@@ -76,6 +87,12 @@ export interface BuildControllerOptions {
now?: () => Date;
materialize?: typeof materializeWorkspace;
resolveDigest?: typeof resolveRegistryDigest;
onReconcileFailure?: (id: string, error: unknown) => void;
reconcileLeases?: WorkspaceLeaseProvider;
reconcileHolder?: string;
reconcileHeartbeatMs?: number;
reconcileSetTimeout?: typeof setTimeout;
reconcileClearTimeout?: typeof clearTimeout;
}
export class BuildValidationError extends Error {
@@ -204,6 +221,8 @@ export class BuildController {
private readonly materializer: typeof materializeWorkspace;
private readonly digestResolver: typeof resolveRegistryDigest;
private readonly imageSubmissionLocks = new Map<string, Promise<unknown>>();
private readonly reconcileLocks = new Map<string, Promise<BuildStatus>>();
private readonly reconcileHolder: string;
constructor(private readonly options: BuildControllerOptions) {
this.now = options.now ?? (() => new Date());
@@ -213,6 +232,7 @@ export class BuildController {
this.maxLogBytes = options.maxLogBytes ?? DEFAULT_MAX_LOG_BYTES;
this.materializer = options.materialize ?? materializeWorkspace;
this.digestResolver = options.resolveDigest ?? resolveRegistryDigest;
this.reconcileHolder = options.reconcileHolder ?? randomUUID();
}
async negotiateSnapshot(
@@ -534,47 +554,215 @@ export class BuildController {
);
}
async reconcileBuild(id: string): Promise<BuildStatus> {
async reconcileBuild(
id: string,
options: ReconcileOptions = {},
): Promise<BuildStatus> {
const current = this.reconcileLocks.get(id);
if (current) return current;
const reconciliation = this.reconcileBuildWithLease(id, options.signal);
this.reconcileLocks.set(id, reconciliation);
try {
return await reconciliation;
} finally {
if (this.reconcileLocks.get(id) === reconciliation)
this.reconcileLocks.delete(id);
}
}
private async reconcileBuildWithLease(
id: string,
signal?: AbortSignal,
): Promise<BuildStatus> {
throwIfAborted(signal);
let lease: WorkspaceLease | undefined;
let reconcileLease: BuildReconciliationLease | undefined;
let heartbeat: ReturnType<typeof setTimeout> | undefined;
let heartbeatActive = false;
let leaseOwnershipLost = false;
try {
lease = this.options.reconcileLeases
? await this.options.reconcileLeases.acquire(
`build-reconcile:${id}`,
this.reconcileHolder,
BUILD_RECONCILE_LEASE_TTL_MS,
)
: undefined;
throwIfAborted(signal);
if (this.options.reconcileLeases && !lease)
return this.getBuildStatus(id);
reconcileLease = await this.options.store.acquireBuildReconciliationLease(
id,
this.reconcileHolder,
BUILD_RECONCILE_LEASE_TTL_MS,
);
throwIfAborted(signal);
if (!reconcileLease) return this.getBuildStatus(id);
const reconcileLeaseToken = reconcileLease.token;
const requireLeaseOwnership = async () => {
throwIfAborted(signal);
if (leaseOwnershipLost)
throw new BuildConflictError(
"Build reconciliation lease ownership was lost",
);
const ownsLease = lease
? await lease.renew(BUILD_RECONCILE_LEASE_TTL_MS)
: true;
const ownsReconcileLease =
await this.options.store.renewBuildReconciliationLease(
id,
reconcileLeaseToken,
BUILD_RECONCILE_LEASE_TTL_MS,
);
throwIfAborted(signal);
if (!ownsLease || !ownsReconcileLease) {
leaseOwnershipLost = true;
throw new BuildConflictError(
"Build reconciliation lease ownership was lost",
);
}
};
const heartbeatMs = Math.min(
BUILD_RECONCILE_HEARTBEAT_MS,
this.options.reconcileHeartbeatMs &&
Number.isFinite(this.options.reconcileHeartbeatMs) &&
this.options.reconcileHeartbeatMs > 0
? this.options.reconcileHeartbeatMs
: BUILD_RECONCILE_HEARTBEAT_MS,
);
const scheduleHeartbeat = () => {
if (!heartbeatActive || leaseOwnershipLost) return;
heartbeat = (this.options.reconcileSetTimeout ?? setTimeout)(() => {
void requireLeaseOwnership()
.catch(() => {
leaseOwnershipLost = true;
})
.finally(scheduleHeartbeat);
}, heartbeatMs);
};
heartbeatActive = true;
scheduleHeartbeat();
await requireLeaseOwnership();
return await this.reconcileBuildInternal(
id,
requireLeaseOwnership,
reconcileLease,
);
} finally {
heartbeatActive = false;
if (heartbeat !== undefined)
(this.options.reconcileClearTimeout ?? clearTimeout)(heartbeat);
try {
if (reconcileLease)
await this.options.store.releaseBuildReconciliationLease(
id,
reconcileLease.token,
);
} finally {
await lease?.release();
}
}
}
private async reconcileBuildInternal(
id: string,
requireLeaseOwnership: () => Promise<void>,
reconcileLease: BuildReconciliationLease,
): Promise<BuildStatus> {
let record = await this.requireBuild(id);
if (record.status.state === "succeeded" || record.status.state === "failed")
return recordStatus(record);
if (record.status.jobCreated) record = await this.captureLogs(record);
if (record.status.jobCreated)
record = await this.captureLogs(
record,
requireLeaseOwnership,
reconcileLease,
);
const observation = await this.options.kubernetes.getJob(
this.options.namespace,
record.spec.jobName,
);
if (!observation) return recordStatus(record);
if (observation.phase === "running" && record.status.state === "queued") {
record = await this.setState(record, "running", {
startedAt: observation.startedAt ?? this.now().toISOString(),
});
await requireLeaseOwnership();
record = await this.setState(
record,
"running",
{
startedAt: observation.startedAt ?? this.now().toISOString(),
},
reconcileLease,
);
} else if (observation.phase === "failed") {
record = await this.setState(record, "failed", {
startedAt: record.status.startedAt ?? observation.startedAt,
finishedAt: observation.finishedAt ?? this.now().toISOString(),
error: observation.error ?? "BuildKit Job failed",
});
await requireLeaseOwnership();
record = await this.setState(
record,
"failed",
{
startedAt: record.status.startedAt ?? observation.startedAt,
finishedAt: observation.finishedAt ?? this.now().toISOString(),
error: observation.error ?? "BuildKit Job failed",
},
reconcileLease,
);
} else if (observation.phase === "succeeded") {
let digest: Sha256Digest;
try {
const digest = await this.digestResolver(
record.spec.request.spec.image,
);
digest = await this.digestResolver(record.spec.request.spec.image);
assertSha256Digest(digest);
record = await this.setState(record, "succeeded", {
} catch (error) {
await requireLeaseOwnership();
record = await this.setState(
record,
"failed",
{
finishedAt: this.now().toISOString(),
error: `Unable to resolve pushed image digest: ${error instanceof Error ? error.message : String(error)}`,
},
reconcileLease,
);
return recordStatus(record);
}
await requireLeaseOwnership();
record = await this.setState(
record,
"succeeded",
{
startedAt: record.status.startedAt ?? observation.startedAt,
finishedAt: observation.finishedAt ?? this.now().toISOString(),
digest,
});
} catch (error) {
record = await this.setState(record, "failed", {
finishedAt: this.now().toISOString(),
error: `Unable to resolve pushed image digest: ${error instanceof Error ? error.message : String(error)}`,
});
}
},
reconcileLease,
);
}
return recordStatus(record);
}
async reconcilePendingBuilds(options: ReconcileOptions = {}): Promise<void> {
throwIfAborted(options.signal);
const records = await this.options.store.listBuilds();
throwIfAborted(options.signal);
for (const record of records) {
throwIfAborted(options.signal);
if (
!record.status.jobCreated ||
record.status.state === "succeeded" ||
record.status.state === "failed"
)
continue;
try {
await this.reconcileBuild(record.metadata.name, options);
} catch (error) {
if (options.signal?.aborted) throw error;
try {
this.options.onReconcileFailure?.(record.metadata.name, error);
} catch {
// Failure reporting must not interrupt reconciliation of other builds.
}
}
}
}
async cancelBuild(id: string): Promise<BuildStatus> {
const record = await this.requireBuild(id);
if (record.status.state === "succeeded" || record.status.state === "failed")
@@ -637,15 +825,26 @@ export class BuildController {
private async updateBuild(
record: BuildRecord,
change: (next: BuildRecord) => void,
reconcileLease?: BuildReconciliationLease,
): Promise<BuildRecord> {
// A heartbeat may have advanced this record's resourceVersion while an
// unabortable Kubernetes observation was pending.
if (reconcileLease) record = await this.requireBuild(record.metadata.name);
const next = clone(record);
next.metadata.resourceVersion = String(
Number(record.metadata.resourceVersion) + 1,
);
change(next);
if (reconcileLease) {
if (next.status.reconcileLease?.token !== reconcileLease.token)
throw new BuildConflictError(
"Build reconciliation lease ownership was lost",
);
}
await this.options.store.replaceBuild(
next,
record.metadata.resourceVersion,
reconcileLease?.token,
);
return this.requireBuild(record.metadata.name);
}
@@ -654,12 +853,17 @@ export class BuildController {
record: BuildRecord,
state: BuildStatus["state"],
values: Partial<BuildRecord["status"]>,
reconcileLease?: BuildReconciliationLease,
): Promise<BuildRecord> {
if (record.status.state === state && state === "running") return record;
return this.updateBuild(record, (next) => {
Object.assign(next.status, values, { state });
next.status.events.push({ type: "status", status: recordStatus(next) });
});
return this.updateBuild(
record,
(next) => {
Object.assign(next.status, values, { state });
next.status.events.push({ type: "status", status: recordStatus(next) });
},
reconcileLease,
);
}
private async failBuild(id: string, message: string): Promise<void> {
@@ -675,7 +879,11 @@ export class BuildController {
});
}
private async captureLogs(record: BuildRecord): Promise<BuildRecord> {
private async captureLogs(
record: BuildRecord,
requireLeaseOwnership: () => Promise<void>,
reconcileLease: BuildReconciliationLease,
): Promise<BuildRecord> {
let raw: string | Uint8Array;
try {
raw = await this.options.kubernetes.getJobLogs(
@@ -689,31 +897,36 @@ export class BuildController {
const offset =
bytes.byteLength < record.status.logOffset ? 0 : record.status.logOffset;
if (bytes.byteLength === offset) return record;
await requireLeaseOwnership();
const delta = bytes.subarray(offset);
return this.updateBuild(record, (next) => {
next.status.logOffset = bytes.byteLength;
for (const message of delta
.toString("utf8")
.split(/(?<=\n)/)
.filter(Boolean)) {
const event: BuildEvent = {
type: "log",
id: record.metadata.name,
sequence: next.status.nextSequence++,
message,
};
next.status.events.push(event);
next.status.logBytes += Buffer.byteLength(message);
}
while (next.status.logBytes > this.maxLogBytes) {
const index = next.status.events.findIndex(
(event) => event.type === "log",
);
if (index === -1) break;
const [removed] = next.status.events.splice(index, 1);
if (removed?.type === "log")
next.status.logBytes -= Buffer.byteLength(removed.message);
}
});
return this.updateBuild(
record,
(next) => {
next.status.logOffset = bytes.byteLength;
for (const message of delta
.toString("utf8")
.split(/(?<=\n)/)
.filter(Boolean)) {
const event: BuildEvent = {
type: "log",
id: record.metadata.name,
sequence: next.status.nextSequence++,
message,
};
next.status.events.push(event);
next.status.logBytes += Buffer.byteLength(message);
}
while (next.status.logBytes > this.maxLogBytes) {
const index = next.status.events.findIndex(
(event) => event.type === "log",
);
if (index === -1) break;
const [removed] = next.status.events.splice(index, 1);
if (removed?.type === "log")
next.status.logBytes -= Buffer.byteLength(removed.message);
}
},
reconcileLease,
);
}
}