feat: harden self-managed reconciliation

This commit is contained in:
2026-09-05 10:17:55 +00:00 Unverified
parent 8e9d207915
commit 321f4e807a
40 changed files with 4977 additions and 404 deletions
+239 -34
View File
@@ -12,7 +12,7 @@ import {
type Sha256Digest,
} from "../shared/build-protocol";
import { resolveComposeArch } from "./arch";
import { apiRequest, type ApiRequestInit } from "./api";
import { apiRequest, type ApiRequestInit, type ApiRequestOptions } from "./api";
import { DEFAULT_REGISTRY } from "./config";
import {
enumerateWorkspace,
@@ -23,10 +23,132 @@ import {
const execFileAsync = promisify(execFile);
const UPLOAD_CHUNK_BYTES = 8 * 1024 * 1024;
const DEFAULT_POLL_INTERVAL_MS = 1_000;
const BUILD_POLL_REQUEST_TIMEOUT_MS = 300_000;
const MAX_BUILD_POLL_ATTEMPTS = 3;
export const MAX_CONCURRENT_REQUESTS = 20;
export const MAX_REQUESTS_PER_SECOND = 40;
export type SchedulerClock = {
now(): number;
};
export type SchedulerSleep = (ms: number) => Promise<void>;
export class TaskScheduler {
private inflight = 0;
private maxInflight: number;
private maxPerSecond: number;
private requestStarts: number[] = [];
private waiting: Array<() => void> = [];
private rateGate: Promise<void> = Promise.resolve();
private clock: SchedulerClock;
private sleep: SchedulerSleep;
constructor(options?: {
maxInflight?: number;
maxPerSecond?: number;
clock?: SchedulerClock;
sleep?: SchedulerSleep;
}) {
this.maxInflight = options?.maxInflight ?? MAX_CONCURRENT_REQUESTS;
this.maxPerSecond = options?.maxPerSecond ?? MAX_REQUESTS_PER_SECOND;
if (!Number.isSafeInteger(this.maxInflight) || this.maxInflight < 1)
throw new RangeError("maxInflight must be a positive integer");
if (!Number.isSafeInteger(this.maxPerSecond) || this.maxPerSecond < 1)
throw new RangeError("maxPerSecond must be a positive integer");
this.clock = options?.clock ?? { now: () => Date.now() };
this.sleep = options?.sleep ?? ((ms) => Bun.sleep(ms));
}
get currentInflight(): number {
return this.inflight;
}
get maxConcurrent(): number {
return this.maxInflight;
}
get currentRequestStarts(): number {
this.pruneOldStarts();
return this.requestStarts.length;
}
private pruneOldStarts(): void {
const cutoff = this.clock.now() - 1000;
while (this.requestStarts.length > 0 && this.requestStarts[0]! <= cutoff) {
this.requestStarts.shift();
}
}
private async acquire(): Promise<void> {
if (this.inflight < this.maxInflight) {
this.inflight++;
return;
}
await new Promise<void>((resolve) => {
this.waiting.push(resolve);
});
}
private release(): void {
if (this.waiting.length > 0) {
const next = this.waiting.shift()!;
next();
} else {
this.inflight--;
}
}
private waitForRateLimit(): Promise<void> {
const reservation = this.rateGate.then(async () => {
for (;;) {
this.pruneOldStarts();
if (this.requestStarts.length < this.maxPerSecond) {
this.requestStarts.push(this.clock.now());
return;
}
const oldest = this.requestStarts[0]!;
await this.sleep(Math.max(1, oldest + 1000 - this.clock.now()));
}
});
this.rateGate = reservation.catch(() => {});
return reservation;
}
async run<T>(fn: () => Promise<T>): Promise<T> {
await this.acquire();
try {
await this.waitForRateLimit();
return await fn();
} finally {
this.release();
}
}
}
async function runConcurrent<T>(
values: T[],
concurrency: number,
run: (value: T) => Promise<void>,
): Promise<void> {
let index = 0;
const worker = async () => {
for (;;) {
const current = index++;
if (current >= values.length) return;
await run(values[current]!);
}
};
await Promise.all(
Array.from({ length: Math.min(concurrency, values.length) }, worker),
);
}
export type ApiRequester = <T>(
path: string,
init?: ApiRequestInit,
options?: ApiRequestOptions,
) => Promise<T>;
export type BuildOptions = {
@@ -35,6 +157,7 @@ export type BuildOptions = {
pollIntervalMs?: number;
sleep?: (milliseconds: number) => Promise<void>;
snapshot?: WorkspaceSnapshot;
scheduler?: TaskScheduler;
};
type BuildPlan = {
@@ -191,29 +314,42 @@ async function uploadBlob(
digest: Sha256Digest,
data: Uint8Array,
request: ApiRequester,
scheduler: TaskScheduler,
project?: string,
): Promise<void> {
const path = `/blobs/${encodeURIComponent(digest)}/uploads`;
const progress = await request<{ offset: number; complete: boolean }>(path, {
method: "POST",
json: { size: data.byteLength },
});
const uploadPath = `/blobs/${encodeURIComponent(digest)}/uploads`;
const projectQuery = project ? `?project=${encodeURIComponent(project)}` : "";
const path = `${uploadPath}${projectQuery}`;
const progress = await scheduler.run(() =>
request<{ offset: number; complete: boolean }>(path, {
method: "POST",
json: { size: data.byteLength },
}),
);
let offset = progress.offset;
while (!progress.complete && offset < data.byteLength) {
const chunk = data.subarray(offset, offset + UPLOAD_CHUNK_BYTES);
const uploaded = await request<{ offset: number }>(path, {
method: "PATCH",
headers: {
"content-type": "application/octet-stream",
"upload-offset": String(offset),
},
body: chunk,
});
const uploaded = await scheduler.run(() =>
request<{ offset: number }>(path, {
method: "PATCH",
headers: {
"content-type": "application/octet-stream",
"upload-offset": String(offset),
},
body: chunk,
}),
);
if (uploaded.offset <= offset)
throw new Error(`Blob upload for ${digest} made no progress`);
offset = uploaded.offset;
}
if (!progress.complete) {
await request(`${path}/complete`, { method: "POST", json: {} });
await scheduler.run(() =>
request(`${uploadPath}/complete${projectQuery}`, {
method: "POST",
json: {},
}),
);
}
}
@@ -221,30 +357,36 @@ export async function uploadWorkspaceSnapshot(
snapshot: WorkspaceSnapshot,
request: ApiRequester = apiRequest,
reporter?: BuildReporter,
scheduler?: TaskScheduler,
project?: string,
): Promise<void> {
const blobs = new Map(snapshot.blobs.map((blob) => [blob.digest, blob.data]));
blobs.set(snapshot.digest, serializeWorkspaceManifest(snapshot.manifest));
const requestScheduler = scheduler ?? new TaskScheduler();
for (;;) {
const negotiation = await request<SnapshotNegotiation>(
"/snapshots/negotiate",
{
const negotiation = await requestScheduler.run(() =>
request<SnapshotNegotiation>("/snapshots/negotiate", {
method: "POST",
json: { workspace: snapshot.digest },
},
json: { workspace: snapshot.digest, ...(project && { project }) },
}),
);
if (negotiation.ready) return;
if (negotiation.missing.length === 0)
throw new Error(
"Snapshot negotiation is incomplete but reported no missing blobs",
);
for (const digest of negotiation.missing) {
const data = blobs.get(digest);
if (!data)
throw new Error(`Server requested unknown workspace blob ${digest}`);
await reporter?.progress?.(`Uploading ${digest}`);
await uploadBlob(digest, data, request);
}
await runConcurrent(
negotiation.missing,
requestScheduler.maxConcurrent,
async (digest) => {
const data = blobs.get(digest);
if (!data)
throw new Error(`Server requested unknown workspace blob ${digest}`);
await reporter?.progress?.(`Uploading ${digest}`);
await uploadBlob(digest, data, request, requestScheduler, project);
},
);
}
}
@@ -265,6 +407,52 @@ async function reportBuildEvent(
return event.sequence;
}
function isTransientBuildPollError(error: unknown): boolean {
if (!(error instanceof Error) || error.name === "AbortError") return false;
if (error.name === "TimeoutError" || error instanceof TypeError) return true;
const code =
"code" in error && typeof error.code === "string"
? error.code
: error.cause &&
typeof error.cause === "object" &&
"code" in error.cause &&
typeof error.cause.code === "string"
? error.cause.code
: undefined;
return (
code === "ECONNABORTED" ||
code === "ECONNRESET" ||
code === "ECONNREFUSED" ||
code === "EAI_AGAIN" ||
code === "ETIMEDOUT"
);
}
async function requestBuildPoll<T>(
request: ApiRequester,
path: string,
init: ApiRequestInit | undefined,
pollIntervalMs: number,
sleep: (milliseconds: number) => Promise<void>,
): Promise<T> {
for (let attempt = 1; attempt <= MAX_BUILD_POLL_ATTEMPTS; attempt++) {
try {
return await request<T>(path, init, {
timeoutMs: BUILD_POLL_REQUEST_TIMEOUT_MS,
});
} catch (error) {
if (
!isTransientBuildPollError(error) ||
attempt === MAX_BUILD_POLL_ATTEMPTS
) {
throw error;
}
await sleep(pollIntervalMs);
}
}
throw new Error("Build poll retries exhausted");
}
async function waitForBuild(
id: string,
request: ApiRequester,
@@ -277,8 +465,12 @@ async function waitForBuild(
let sequence = 0;
const reportedStates = new Set<BuildStatus["state"]>();
for (;;) {
const events = await request<BuildEvent[]>(
const events = await requestBuildPoll<BuildEvent[]>(
request,
`/builds/${encodeURIComponent(id)}/events?after=${sequence}`,
undefined,
pollIntervalMs,
sleep,
);
for (const event of events)
sequence = Math.max(
@@ -287,12 +479,15 @@ async function waitForBuild(
);
if (status.state === "succeeded" || status.state === "failed")
return status;
status = await request<BuildStatus>(
status = await requestBuildPoll<BuildStatus>(
request,
`/builds/${encodeURIComponent(id)}/reconcile`,
{
method: "POST",
json: {},
},
pollIntervalMs,
sleep,
);
if (status.state !== "succeeded" && status.state !== "failed")
await sleep(pollIntervalMs);
@@ -351,7 +546,13 @@ export async function buildServices(
return plan ? [plan] : [];
},
);
await uploadWorkspaceSnapshot(snapshot, request, reporter);
await uploadWorkspaceSnapshot(
snapshot,
request,
reporter,
options.scheduler,
project,
);
const images: Record<string, string> = {};
for (const plan of plans) {
@@ -372,10 +573,14 @@ export async function buildServices(
workspace: snapshot.digest,
},
};
const initial = await request<BuildStatus>("/builds", {
method: "POST",
json: buildRequest,
});
const initial = await request<BuildStatus>(
"/builds",
{
method: "POST",
json: buildRequest,
},
{ timeoutMs: 300_000 },
);
const status = await waitForBuild(
id,
request,
+20 -5
View File
@@ -236,7 +236,15 @@ function toManagedRole(claim: PostgresClaim): ManagedRole {
};
}
async function reconcileManagedRoles(claims: PostgresClaim[]): Promise<void> {
function throwIfAborted(signal?: AbortSignal): void {
if (!signal?.aborted) return;
throw new Error("Workspace operation execution was cancelled");
}
async function reconcileManagedRoles(
claims: PostgresClaim[],
signal?: AbortSignal,
): Promise<void> {
if (claims.length === 0) return;
const cluster = await readObject<
@@ -264,6 +272,7 @@ async function reconcileManagedRoles(claims: PostgresClaim[]): Promise<void> {
}
try {
throwIfAborted(signal);
await applyResource({
apiVersion: "postgresql.cnpg.io/v1",
kind: "Cluster",
@@ -287,6 +296,7 @@ async function reconcileManagedRoles(claims: PostgresClaim[]): Promise<void> {
async function reconcileDatabases(
project: string,
claims: PostgresClaim[],
signal?: AbortSignal,
): Promise<void> {
const uniqueDatabases = new Map<string, PostgresClaim>();
for (const claim of claims) {
@@ -295,6 +305,7 @@ async function reconcileDatabases(
for (const claim of uniqueDatabases.values()) {
try {
throwIfAborted(signal);
await applyResource({
apiVersion: "postgresql.cnpg.io/v1",
kind: "Database",
@@ -328,16 +339,19 @@ async function reconcileDatabases(
export async function reconcilePostgresClaim(
project: string,
claim: PostgresClaim,
signal?: AbortSignal,
): Promise<RoleCredentials> {
throwIfAborted(signal);
const credentials = await ensureRoleSecret(claim);
await reconcileManagedRoles([claim]);
await reconcileDatabases(project, [claim]);
await reconcileManagedRoles([claim], signal);
await reconcileDatabases(project, [claim], signal);
return credentials;
}
export async function reconcilePostgresClaims(
project: string,
compose: ComposeSpecification,
signal?: AbortSignal,
): Promise<Record<string, Record<string, string>>> {
const claims = getComposePostgresClaims(compose);
if (claims.length === 0) return {};
@@ -345,11 +359,12 @@ export async function reconcilePostgresClaims(
const credentialsBySecret = new Map<string, RoleCredentials>();
for (const claim of claims) {
if (credentialsBySecret.has(claim.secretName)) continue;
throwIfAborted(signal);
credentialsBySecret.set(claim.secretName, await ensureRoleSecret(claim));
}
await reconcileManagedRoles(claims);
await reconcileDatabases(project, claims);
await reconcileManagedRoles(claims, signal);
await reconcileDatabases(project, claims, signal);
return Object.fromEntries(
claims.map((claim) => {
+2 -1
View File
@@ -161,7 +161,8 @@ export async function openExecSession(
try {
if (!live) {
live = true;
for (const frame of pending.splice(0)) socket.send(JSON.stringify(frame));
for (const frame of pending.splice(0))
socket.send(JSON.stringify(frame));
}
const frame = parseFrame(event.data);
const reader = readers.shift();
+36 -6
View File
@@ -138,15 +138,41 @@ export function getComposeS3Claims(compose: ComposeSpecification): S3Claim[] {
return claims;
}
function throwIfAborted(signal?: AbortSignal): void {
if (!signal?.aborted) return;
throw new Error("Workspace operation execution was cancelled");
}
async function waitForInterval(
delayMs: number,
signal?: AbortSignal,
): Promise<void> {
if (!signal) return new Promise((resolve) => setTimeout(resolve, delayMs));
throwIfAborted(signal);
await new Promise<void>((resolve, reject) => {
const timer = setTimeout(() => {
signal.removeEventListener("abort", cancelWait);
resolve();
}, delayMs);
const cancelWait = () => {
clearTimeout(timer);
reject(new Error("Workspace operation execution was cancelled"));
};
signal.addEventListener("abort", cancelWait, { once: true });
});
}
async function reconcileBuckets(
project: string,
claims: S3Claim[],
signal?: AbortSignal,
): Promise<void> {
const buckets = new Map<string, S3Claim>();
for (const claim of claims) buckets.set(claim.bucket, claim);
for (const claim of buckets.values()) {
try {
throwIfAborted(signal);
await applyResource({
apiVersion: GARAGE_API_VERSION,
kind: "GarageBucket",
@@ -177,12 +203,14 @@ async function reconcileBuckets(
async function reconcileKeys(
project: string,
claims: S3Claim[],
signal?: AbortSignal,
): Promise<void> {
const keys = new Map<string, S3Claim>();
for (const claim of claims) keys.set(claim.key, claim);
for (const claim of keys.values()) {
try {
throwIfAborted(signal);
await applyResource({
apiVersion: GARAGE_API_VERSION,
kind: "GarageKey",
@@ -227,11 +255,13 @@ async function reconcileKeys(
export async function getS3Credentials(
claim: S3Claim,
signal?: AbortSignal,
): Promise<Record<string, string>> {
const deadline = Date.now() + SECRET_WAIT_TIMEOUT_MS;
let lastPhase: string | undefined;
while (Date.now() < deadline) {
throwIfAborted(signal);
const key = await readObject<GarageKeyStatus>({
apiVersion: GARAGE_API_VERSION,
kind: "GarageKey",
@@ -274,9 +304,7 @@ export async function getS3Credentials(
}
}
await new Promise((resolve) =>
setTimeout(resolve, SECRET_WAIT_INTERVAL_MS),
);
await waitForInterval(SECRET_WAIT_INTERVAL_MS, signal);
}
throw new Error(
@@ -287,17 +315,19 @@ export async function getS3Credentials(
export async function reconcileS3Claims(
project: string,
compose: ComposeSpecification,
signal?: AbortSignal,
): Promise<Record<string, Record<string, string>>> {
const claims = getComposeS3Claims(compose);
if (claims.length === 0) return {};
await reconcileBuckets(project, claims);
await reconcileKeys(project, claims);
await reconcileBuckets(project, claims, signal);
await reconcileKeys(project, claims, signal);
const environmentsByKey = new Map<string, Record<string, string>>();
for (const claim of claims) {
if (environmentsByKey.has(claim.key)) continue;
environmentsByKey.set(claim.key, await getS3Credentials(claim));
throwIfAborted(signal);
environmentsByKey.set(claim.key, await getS3Credentials(claim, signal));
}
return Object.fromEntries(
+105
View File
@@ -0,0 +1,105 @@
import { createHash, randomUUID } from "node:crypto";
import {
chmod,
mkdir,
readFile,
realpath,
rename,
writeFile,
} from "node:fs/promises";
import { join } from "node:path";
export const TRUST_SCHEMA_VERSION = 1;
export type TrustIdentity = { project: string; fingerprint: string };
export type LocalTrustRecord = TrustIdentity;
type TrustFile = { version: number; records: LocalTrustRecord[] };
export function getTrustPath(): string {
return join(
process.env.XDG_CONFIG_HOME ?? join(process.env.HOME ?? "/tmp", ".config"),
"kuber",
"trust.json",
);
}
export async function resolveTrustIdentity(
project: string,
cwd = process.cwd(),
): Promise<TrustIdentity> {
const path = await realpath(cwd);
return {
project,
fingerprint: createHash("sha256").update(path).digest("hex"),
};
}
function validRecord(value: unknown): value is LocalTrustRecord {
return (
!!value &&
typeof value === "object" &&
typeof (value as Record<string, unknown>).project === "string" &&
typeof (value as Record<string, unknown>).fingerprint === "string" &&
/^[a-f0-9]{64}$/.test(
(value as Record<string, unknown>).fingerprint as string,
)
);
}
export async function readTrust(): Promise<LocalTrustRecord[]> {
try {
const value: unknown = JSON.parse(await readFile(getTrustPath(), "utf8"));
if (
!value ||
typeof value !== "object" ||
(value as TrustFile).version !== TRUST_SCHEMA_VERSION ||
!Array.isArray((value as TrustFile).records)
)
return [];
return (value as TrustFile).records.filter(validRecord);
} catch {
return [];
}
}
export async function writeTrust(records: LocalTrustRecord[]): Promise<void> {
const path = getTrustPath();
const directory = path.slice(0, path.lastIndexOf("/"));
await mkdir(directory, { recursive: true, mode: 0o700 });
await chmod(directory, 0o700);
const temporary = `${path}.${randomUUID()}.tmp`;
await writeFile(
temporary,
JSON.stringify({ version: TRUST_SCHEMA_VERSION, records }, null, 2) + "\n",
{ flag: "wx", mode: 0o600 },
);
await rename(temporary, path);
await chmod(path, 0o600);
}
export async function updateTrust(
update: (records: LocalTrustRecord[]) => LocalTrustRecord[],
): Promise<void> {
await writeTrust(update(await readTrust()));
}
export function trustHeaders(identity: TrustIdentity): Record<string, string> {
return {
"x-kuber-trust-project": identity.project,
"x-kuber-trust-fingerprint": identity.fingerprint,
};
}
export async function requireLocalTrust(
identity: TrustIdentity,
): Promise<LocalTrustRecord> {
const record = (await readTrust()).find(
(candidate) =>
candidate.project === identity.project &&
candidate.fingerprint === identity.fingerprint,
);
if (record) return record;
throw new Error(
"TRUST_REQUIRED: this directory is not trusted for this namespace. Run kuber trust before kuber up.",
);
}