1368 lines
49 KiB
TypeScript
1368 lines
49 KiB
TypeScript
import { describe, expect, spyOn, test } from "bun:test";
|
||
import { DefaultRenderer, Listr } from "listr2";
|
||
import { mkdtemp, rm, writeFile } from "node:fs/promises";
|
||
import { tmpdir } from "node:os";
|
||
import { join } from "node:path";
|
||
import { Writable } from "node:stream";
|
||
import { stripVTControlCharacters } from "node:util";
|
||
import type { ApiRequestInit, ApiRequestOptions } from "../../lib/api";
|
||
import { KuberApiError } from "../../lib/api";
|
||
import type { ApiRequester } from "../../lib/build";
|
||
import {
|
||
ensureWorkspace,
|
||
reconcileResources,
|
||
runLiveResourceOperation,
|
||
runUp,
|
||
workspaceAdoptionRoute,
|
||
} from "../../command/up";
|
||
import { provideContext } from "../../lib/context";
|
||
import { resolveTrustIdentity, updateTrust } from "../../lib/trust";
|
||
import { workspaceManifestDigest } from "../../lib/workspace";
|
||
|
||
const manifest = { version: 1 as const, files: [] };
|
||
const snapshot = {
|
||
manifest,
|
||
digest: workspaceManifestDigest(manifest),
|
||
blobs: [],
|
||
};
|
||
|
||
// DefaultRenderer.create returns a complete TTY frame before the cursor redraw.
|
||
function frameRows(frame: string): string[] {
|
||
return stripVTControlCharacters(frame).split("\n");
|
||
}
|
||
|
||
function taskRow(rows: string[], title: string): string {
|
||
return rows.find((row) => row.includes(title)) ?? "";
|
||
}
|
||
|
||
function depth(row: string): number {
|
||
return row.match(/^\s*/)?.[0].length ?? 0;
|
||
}
|
||
|
||
describe("up API pipeline", () => {
|
||
test("renders build lifecycle in child titles without duplicate phase output", async () => {
|
||
const root = await mkdtemp(join(tmpdir(), "kuber-up-api-"));
|
||
const previousCwd = process.cwd();
|
||
const tty = Object.getOwnPropertyDescriptor(process.stdout, "isTTY");
|
||
const frames: string[] = [];
|
||
const create = DefaultRenderer.prototype.create;
|
||
const renderSpy = spyOn(DefaultRenderer.prototype, "create").mockImplementation(function (this: DefaultRenderer, options) {
|
||
const frame = create.call(this, options);
|
||
frames.push(frame);
|
||
return frame;
|
||
});
|
||
const writes = spyOn(process.stdout, "write").mockImplementation((() => true) as typeof process.stdout.write);
|
||
const polls = new Map<string, number>();
|
||
const builds = new Map<string, string>();
|
||
try {
|
||
Object.defineProperty(process.stdout, "isTTY", { configurable: true, value: true });
|
||
await writeFile(join(root, "compose.yml"), "services:\n app:\n build: .\n client:\n build:\n context: .\n args:\n ROLE: client\n");
|
||
await writeFile(join(root, ".kuberrc.ts"), 'export default { project: "shop" };\n');
|
||
process.chdir(root);
|
||
const trust = await resolveTrustIdentity("shop", root);
|
||
const request: ApiRequester = async <T>(path: string, init?: ApiRequestInit) => {
|
||
if (path === "/snapshots/negotiate") return { ready: true } as T;
|
||
if (path === "/builds") {
|
||
const build = init!.json as { id: string; service: string };
|
||
builds.set(build.id, build.service);
|
||
return { state: "queued" } as T;
|
||
}
|
||
const id = path.split("/")[2]!;
|
||
const service = builds.get(id)!;
|
||
if (path.includes("/events"))
|
||
return [{ type: "status", status: {
|
||
state: "queued",
|
||
phase: ["queued", "creating", "starting", "running", "done"][polls.get(service) ?? 0],
|
||
} }] as T;
|
||
if (path.endsWith("/reconcile")) {
|
||
const count = (polls.get(service) ?? 0) + 1;
|
||
polls.set(service, count);
|
||
await Bun.sleep(110);
|
||
return count < 4
|
||
? { state: "queued" } as T
|
||
: service === "client"
|
||
? { state: "failed", error: "build stopped\nstack detail" } as T
|
||
: { state: "succeeded" } as T;
|
||
}
|
||
if (path.endsWith("/result")) return { reference: "image:app" } as T;
|
||
throw new Error(path);
|
||
};
|
||
await expect(provideContext(() => runUp(true, request, { trust }))).rejects.toThrow("build stopped");
|
||
expect(polls.get("client")).toBe(4);
|
||
const active = frames.map(frameRows).filter((rows) =>
|
||
taskRow(rows, "Build images") && taskRow(rows, "Build app") && taskRow(rows, "Build client"));
|
||
expect(active.length).toBeGreaterThan(0);
|
||
for (const rows of active) {
|
||
const parent = taskRow(rows, "Build images");
|
||
for (const name of ["app", "client"]) {
|
||
const children = rows.filter((row) => row.includes(`Build ${name}`));
|
||
expect(children).toHaveLength(1);
|
||
expect(depth(children[0]!)).toBeGreaterThan(depth(parent));
|
||
}
|
||
expect(rows.filter((row) => /^\s*› Build (?:queued|creating|starting|running|done)/.test(row))).toHaveLength(0);
|
||
}
|
||
for (const phase of ["queued", "creating", "starting", "running", "done"])
|
||
expect(frames.some((frame) => frame.includes(`Build app: ${phase}`))).toBe(true);
|
||
const final = frameRows(frames.at(-1)!);
|
||
expect(taskRow(final, "Build app: done")).toMatch(/✔/);
|
||
expect(taskRow(final, "Build client: failed")).toMatch(/✖/);
|
||
expect(final.join("\n")).toContain("build stopped");
|
||
expect(final.join("\n")).toContain("stack detail");
|
||
expect(final.join("\n").split("build stopped")).toHaveLength(2);
|
||
} finally {
|
||
renderSpy.mockRestore();
|
||
writes.mockRestore();
|
||
if (tty) Object.defineProperty(process.stdout, "isTTY", tty);
|
||
else Reflect.deleteProperty(process.stdout, "isTTY");
|
||
process.chdir(previousCwd);
|
||
await rm(root, { recursive: true, force: true });
|
||
}
|
||
}, 5_000);
|
||
|
||
test("attaches each build log to its started image child", async () => {
|
||
const root = await mkdtemp(join(tmpdir(), "kuber-up-api-"));
|
||
const previousCwd = process.cwd();
|
||
const originalRun = Listr.prototype.run;
|
||
const lines = new Map<string, string[]>();
|
||
const acquired: string[] = [];
|
||
const runSpy = spyOn(Listr.prototype, "run").mockImplementation(function (this: Listr) {
|
||
if (this.tasks[0]?.title === "Build web") {
|
||
for (const entry of this.tasks) {
|
||
if (!entry.title) continue;
|
||
const name = entry.title.replace(/^Build /, "");
|
||
const executable = entry as unknown as { taskFn: typeof entry.task.task };
|
||
const originalTask = executable.taskFn;
|
||
executable.taskFn = async (ctx, child) => {
|
||
const output: string[] = [];
|
||
lines.set(name, output);
|
||
child.stdout = () => {
|
||
acquired.push(name);
|
||
return new Writable({ write(chunk, _encoding, done) {
|
||
output.push(String(chunk));
|
||
done();
|
||
} });
|
||
};
|
||
return originalTask(ctx, child);
|
||
};
|
||
}
|
||
}
|
||
return originalRun.call(this);
|
||
});
|
||
try {
|
||
await writeFile(join(root, "compose.yml"),
|
||
"services:\n web:\n build: .\n worker:\n build: .\n admin:\n build:\n context: .\n args:\n ROLE: admin\n");
|
||
await writeFile(join(root, ".kuberrc.ts"), 'export default { project: "shop" };\n');
|
||
process.chdir(root);
|
||
const trust = await resolveTrustIdentity("shop", root);
|
||
const builds = new Map<string, string>();
|
||
const request: ApiRequester = async <T>(path: string, init?: ApiRequestInit) => {
|
||
if (path === "/snapshots/negotiate") return { ready: true } as T;
|
||
if (path === "/builds") {
|
||
const build = init!.json as { id: string; service: string };
|
||
builds.set(build.id, build.service);
|
||
return { state: "queued" } as T;
|
||
}
|
||
const service = builds.get(path.split("/")[2]!)!;
|
||
if (path.includes("/events"))
|
||
return path.includes("after=0")
|
||
? [{ type: "log", sequence: 1, message: `${service} log\n` }] as T
|
||
: [] as T;
|
||
if (path.endsWith("/reconcile"))
|
||
return { state: service === "admin" ? "failed" : "succeeded", error: "admin failed" } as T;
|
||
if (path.endsWith("/result"))
|
||
return { reference: "image:web", references: { worker: "image:worker" } } as T;
|
||
throw new Error(path);
|
||
};
|
||
await expect(provideContext(() => runUp(true, request, { trust }))).rejects.toThrow("admin failed");
|
||
expect(builds.size).toBe(2);
|
||
expect(acquired.sort()).toEqual(["admin", "web", "worker"]);
|
||
expect(lines.get("web")).toEqual(["web log\n"]);
|
||
expect(lines.get("worker")).toEqual(["web log\n"]);
|
||
expect(lines.get("admin")).toEqual(["admin log\n"]);
|
||
} finally {
|
||
runSpy.mockRestore();
|
||
process.chdir(previousCwd);
|
||
await rm(root, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
test("renders a failing image under its child and sends the stack to the parent bottom bar", async () => {
|
||
const root = await mkdtemp(join(tmpdir(), "kuber-up-api-"));
|
||
const previousCwd = process.cwd();
|
||
const rendered: string[] = [];
|
||
const writes = spyOn(process.stdout, "write").mockImplementation(((chunk: string | Uint8Array) => {
|
||
rendered.push(String(chunk));
|
||
return true;
|
||
}) as typeof process.stdout.write);
|
||
try {
|
||
await writeFile(join(root, "compose.yml"), "services:\n client:\n build: .\n server:\n build: .\n");
|
||
await writeFile(join(root, ".kuberrc.ts"), 'export default { project: "shop" };\n');
|
||
process.chdir(root);
|
||
const trust = await resolveTrustIdentity("shop", root);
|
||
const request: ApiRequester = async <T>(path: string, init?: ApiRequestInit) => {
|
||
if (path === "/snapshots/negotiate") return { ready: true } as T;
|
||
if (path === "/builds") return { id: (init!.json as { id: string }).id, state: "failed", error: "buildkit failed\nstack detail" } as T;
|
||
if (path.includes("/events")) return [] as T;
|
||
throw new Error(path);
|
||
};
|
||
await expect(provideContext(() => runUp(true, request, { trust }))).rejects.toThrow();
|
||
const output = rendered.join("");
|
||
expect(output).toContain("Build client");
|
||
expect(output).toContain("Build server");
|
||
expect(output).toContain("buildkit failed");
|
||
expect(output).toContain("stack detail");
|
||
} finally {
|
||
writes.mockRestore();
|
||
process.chdir(previousCwd);
|
||
await rm(root, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
test("forwards the build timeout through the trusted requester", async () => {
|
||
const root = await mkdtemp(join(tmpdir(), "kuber-up-api-"));
|
||
const previousCwd = process.cwd();
|
||
const previousConfigHome = process.env.XDG_CONFIG_HOME;
|
||
const configHome = join(root, "config");
|
||
const calls: Array<{
|
||
path: string;
|
||
init?: ApiRequestInit;
|
||
options?: ApiRequestOptions;
|
||
}> = [];
|
||
|
||
try {
|
||
await writeFile(
|
||
join(root, "compose.yml"),
|
||
"services:\n web:\n build: .\n",
|
||
);
|
||
await writeFile(
|
||
join(root, ".kuberrc.ts"),
|
||
'export default { project: "shop" };\n',
|
||
);
|
||
process.chdir(root);
|
||
process.env.XDG_CONFIG_HOME = configHome;
|
||
const identity = await resolveTrustIdentity("shop", root);
|
||
await updateTrust((records) => [...records, identity]);
|
||
|
||
const request: ApiRequester = async <T>(
|
||
path: string,
|
||
init?: ApiRequestInit,
|
||
options?: ApiRequestOptions,
|
||
) => {
|
||
calls.push({ path, init, options });
|
||
if (path === "/snapshots/negotiate")
|
||
return { workspace: snapshot.digest, missing: [], ready: true } as T;
|
||
if (path === "/builds") {
|
||
const buildRequest = init?.json as { id?: unknown } | undefined;
|
||
if (typeof buildRequest?.id !== "string")
|
||
throw new Error("Expected build request ID");
|
||
return {
|
||
version: 1,
|
||
id: buildRequest.id,
|
||
state: "succeeded",
|
||
createdAt: "2026-01-01T00:00:00Z",
|
||
} as T;
|
||
}
|
||
if (path.includes("/events")) return [] as T;
|
||
if (path.endsWith("/result"))
|
||
return {
|
||
image: "registry.server/kuber/shop-web",
|
||
digest: `sha256:${"a".repeat(64)}`,
|
||
reference: `registry.server/kuber/shop-web@sha256:${"a".repeat(64)}`,
|
||
} as T;
|
||
if (path === "/workspaces/shop")
|
||
throw new KuberApiError("missing", 404);
|
||
if (path === "/workspaces")
|
||
return {
|
||
metadata: { name: "shop", uid: "workspace", resourceVersion: "1" },
|
||
} as T;
|
||
if (path.endsWith("/adopt")) return { resourcesAdopted: 0 } as T;
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
return {} as T;
|
||
};
|
||
|
||
await provideContext(() => runUp(true, request));
|
||
|
||
const build = calls.find(({ path }) => path === "/builds");
|
||
if (!build) throw new Error("Expected build submission");
|
||
expect(build.options).toEqual({ timeoutMs: 300_000 });
|
||
expect(
|
||
new Headers(build.init?.headers).get("x-kuber-trust-project"),
|
||
).toBe("shop");
|
||
} finally {
|
||
process.chdir(previousCwd);
|
||
if (previousConfigHome === undefined) delete process.env.XDG_CONFIG_HOME;
|
||
else process.env.XDG_CONFIG_HOME = previousConfigHome;
|
||
await rm(root, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
test("adopts before overlapping database and storage reconciliation, then renders merged env", async () => {
|
||
const root = await mkdtemp(join(tmpdir(), "kuber-up-api-"));
|
||
const previousCwd = process.cwd();
|
||
const order: string[] = [];
|
||
const started: string[] = [];
|
||
let rootTasks: Listr["tasks"] | undefined;
|
||
const originalRun = Listr.prototype.run;
|
||
const frames: string[] = [];
|
||
const create = DefaultRenderer.prototype.create;
|
||
const renderSpy = spyOn(DefaultRenderer.prototype, "create").mockImplementation(function (this: DefaultRenderer, options) {
|
||
const frame = create.call(this, options);
|
||
frames.push(frame);
|
||
return frame;
|
||
});
|
||
const tty = Object.getOwnPropertyDescriptor(process.stdout, "isTTY");
|
||
const writes = spyOn(process.stdout, "write").mockImplementation((() => true) as typeof process.stdout.write);
|
||
const runSpy = spyOn(Listr.prototype, "run").mockImplementation(function (this: Listr) {
|
||
if (this.tasks[0]?.title === "Read compose") rootTasks = this.tasks;
|
||
return originalRun.call(this);
|
||
});
|
||
let bothReady!: () => void;
|
||
const bothStarted = new Promise<void>((resolve) => {
|
||
bothReady = resolve;
|
||
});
|
||
let releaseDatabase!: () => void;
|
||
let releaseStorage!: () => void;
|
||
const database = new Promise<void>((resolve) => {
|
||
releaseDatabase = resolve;
|
||
});
|
||
const storage = new Promise<void>((resolve) => {
|
||
releaseStorage = resolve;
|
||
});
|
||
try {
|
||
Object.defineProperty(process.stdout, "isTTY", { configurable: true, value: true });
|
||
await writeFile(
|
||
join(root, "compose.yml"),
|
||
"services:\n app:\n image: nginx\n volumes:\n - postgresql:app\n - s3:assets\n",
|
||
);
|
||
await writeFile(
|
||
join(root, ".kuberrc.ts"),
|
||
'export default { project: "shop" };\n',
|
||
);
|
||
expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0);
|
||
process.chdir(root);
|
||
const trust = await resolveTrustIdentity("shop", root);
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
order.push(path);
|
||
if (path === "/workspaces/shop") throw new KuberApiError("missing", 404);
|
||
if (path === "/workspaces")
|
||
return {
|
||
metadata: { name: "shop", uid: "workspace", resourceVersion: "1" },
|
||
} as T;
|
||
if (path.endsWith("/adopt")) return { resourcesAdopted: 0 } as T;
|
||
if (path.endsWith("/databases")) {
|
||
started.push("database");
|
||
if (started.length === 2) bothReady();
|
||
await database;
|
||
return { app: { DATABASE_URL: "postgresql://app", SHARED: "database" } } as T;
|
||
}
|
||
if (path.endsWith("/storage")) {
|
||
started.push("storage");
|
||
if (started.length === 2) bothReady();
|
||
await storage;
|
||
return { app: { AWS_ACCESS_KEY_ID: "key", SHARED: "storage" } } as T;
|
||
}
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
return {} as T;
|
||
};
|
||
const run = provideContext(() => runUp(false, request, { trust }));
|
||
try {
|
||
await bothStarted;
|
||
const titles = rootTasks!.map(({ title }) => title);
|
||
expect(titles.indexOf("Reconcile backing services")).toBeLessThan(
|
||
titles.indexOf("Render manifests"),
|
||
);
|
||
expect(titles.indexOf("Render manifests")).toBeLessThan(
|
||
titles.indexOf("Reconcile resources"),
|
||
);
|
||
const backing = rootTasks!.find(({ title }) => title === "Reconcile backing services")!;
|
||
expect(backing.subtasks.map(({ title }) => title)).toEqual([
|
||
"Reconcile databases",
|
||
"Reconcile S3 storage",
|
||
]);
|
||
const active = frames.map(frameRows).filter((rows) =>
|
||
taskRow(rows, "Reconcile backing services") &&
|
||
taskRow(rows, "Reconcile databases") &&
|
||
taskRow(rows, "Reconcile S3 storage"));
|
||
expect(active.length).toBeGreaterThan(0);
|
||
for (const rows of active) {
|
||
const parent = taskRow(rows, "Reconcile backing services");
|
||
expect(depth(taskRow(rows, "Reconcile databases"))).toBeGreaterThan(depth(parent));
|
||
expect(depth(taskRow(rows, "Reconcile S3 storage"))).toBeGreaterThan(depth(parent));
|
||
}
|
||
expect(rootTasks!.filter(({ title }) => title === "Reconcile databases")).toEqual([]);
|
||
expect(started).toEqual(["database", "storage"]);
|
||
expect(order.indexOf("/workspaces/shop/adopt")).toBeLessThan(
|
||
order.indexOf("/workspaces/shop/databases"),
|
||
);
|
||
expect(order).not.toContain("/workspaces/shop/resources/plan");
|
||
releaseStorage();
|
||
await Promise.resolve();
|
||
expect(order).not.toContain("/workspaces/shop/resources/plan");
|
||
releaseDatabase();
|
||
const result = await run;
|
||
expect(result.serviceEnv).toEqual({
|
||
app: {
|
||
DATABASE_URL: "postgresql://app",
|
||
SHARED: "storage",
|
||
AWS_ACCESS_KEY_ID: "key",
|
||
},
|
||
});
|
||
expect(order.indexOf("/workspaces/shop/resources/plan")).toBeGreaterThan(
|
||
order.indexOf("/workspaces/shop/storage"),
|
||
);
|
||
} finally {
|
||
releaseStorage();
|
||
releaseDatabase();
|
||
await run.catch(() => {});
|
||
}
|
||
} finally {
|
||
runSpy.mockRestore();
|
||
renderSpy.mockRestore();
|
||
writes.mockRestore();
|
||
if (tty) Object.defineProperty(process.stdout, "isTTY", tty);
|
||
else Reflect.deleteProperty(process.stdout, "isTTY");
|
||
process.chdir(previousCwd);
|
||
await rm(root, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
test("shows database API problem details on the database task without starting resources", async () => {
|
||
const root = await mkdtemp(join(tmpdir(), "kuber-up-api-"));
|
||
const previousCwd = process.cwd();
|
||
const frames: string[] = [];
|
||
const create = DefaultRenderer.prototype.create;
|
||
const renderSpy = spyOn(DefaultRenderer.prototype, "create").mockImplementation(function (this: DefaultRenderer, options) {
|
||
const frame = create.call(this, options);
|
||
frames.push(frame);
|
||
return frame;
|
||
});
|
||
const tty = Object.getOwnPropertyDescriptor(process.stdout, "isTTY");
|
||
const writes = spyOn(process.stdout, "write").mockImplementation((() => true) as typeof process.stdout.write);
|
||
const calls: string[] = [];
|
||
try {
|
||
Object.defineProperty(process.stdout, "isTTY", { configurable: true, value: true });
|
||
await writeFile(join(root, "compose.yml"), "services:\n web:\n image: nginx\n volumes:\n - postgresql:web_db\n");
|
||
await writeFile(join(root, ".kuberrc.ts"), 'export default { project: "shop" };\n');
|
||
process.chdir(root);
|
||
const trust = await resolveTrustIdentity("shop", root);
|
||
const detail = "Database reconciliation failed during database apply for database web_db (service web, role web): Forbidden";
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
calls.push(path);
|
||
if (path === "/workspaces/shop") throw new KuberApiError("missing", 404);
|
||
if (path === "/workspaces") return { metadata: { name: "shop", uid: "workspace", resourceVersion: "1" } } as T;
|
||
if (path.endsWith("/adopt")) return { resourcesAdopted: 0 } as T;
|
||
if (path.endsWith("/databases")) throw new KuberApiError(detail, 500, {
|
||
title: "Operation failed", status: 500, code: "DATABASE_RECONCILE_FAILED", detail,
|
||
});
|
||
throw new Error(`Unexpected request: ${path}`);
|
||
};
|
||
await expect(provideContext(() => runUp(false, request, { trust }))).rejects.toThrow(detail);
|
||
const rows = frameRows(frames.at(-1)!);
|
||
expect(depth(taskRow(rows, "Reconcile databases"))).toBeGreaterThan(depth(taskRow(rows, "Reconcile backing services")));
|
||
expect(taskRow(rows, "Reconcile databases")).toMatch(/✖/);
|
||
const visible = rows.join(" ").replace(/\s+/g, " ");
|
||
expect(visible.split(detail)).toHaveLength(2);
|
||
expect(calls).not.toContain("/workspaces/shop/resources/plan");
|
||
} finally {
|
||
renderSpy.mockRestore();
|
||
writes.mockRestore();
|
||
if (tty) Object.defineProperty(process.stdout, "isTTY", tty);
|
||
else Reflect.deleteProperty(process.stdout, "isTTY");
|
||
process.chdir(previousCwd);
|
||
await rm(root, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
test("creates workspace metadata without embedding source blobs", async () => {
|
||
const calls: Array<{ path: string; init?: ApiRequestInit }> = [];
|
||
const request: ApiRequester = async <T>(
|
||
path: string,
|
||
init?: ApiRequestInit,
|
||
) => {
|
||
calls.push({ path, init });
|
||
if (!init) throw new KuberApiError("missing", 404);
|
||
return {
|
||
metadata: { name: "shop", uid: "uid", resourceVersion: "1" },
|
||
} as T;
|
||
};
|
||
await ensureWorkspace(
|
||
"shop",
|
||
{ services: { web: { image: "nginx" } } },
|
||
snapshot,
|
||
request,
|
||
);
|
||
|
||
expect(calls.map(({ path }) => path)).toEqual([
|
||
"/workspaces/shop",
|
||
"/workspaces",
|
||
]);
|
||
const body = calls[1]!.init?.json as Record<string, unknown>;
|
||
expect(body).toMatchObject({
|
||
id: "shop",
|
||
source: { uri: `cas://${snapshot.digest}`, digest: snapshot.digest },
|
||
});
|
||
expect(JSON.stringify(body)).not.toContain('"blob"');
|
||
expect(JSON.stringify(body)).not.toContain('"files"');
|
||
});
|
||
|
||
test("updates workspace metadata with the current resource-version ETag", async () => {
|
||
const calls: Array<{ path: string; init?: ApiRequestInit }> = [];
|
||
const request: ApiRequester = async <T>(
|
||
path: string,
|
||
init?: ApiRequestInit,
|
||
) => {
|
||
calls.push({ path, init });
|
||
return {
|
||
metadata: {
|
||
name: "shop",
|
||
uid: "uid",
|
||
resourceVersion: init ? "8" : "7",
|
||
},
|
||
} as T;
|
||
};
|
||
await ensureWorkspace("shop", { services: {} }, snapshot, request);
|
||
expect(calls[1]?.init?.method).toBe("PUT");
|
||
expect(new Headers(calls[1]?.init?.headers).get("if-match")).toBe('"7"');
|
||
});
|
||
|
||
test("retries one conflict using a fresh ETag", async () => {
|
||
const calls: Array<{ path: string; init?: ApiRequestInit }> = [];
|
||
let gets = 0;
|
||
const request: ApiRequester = async <T>(
|
||
path: string,
|
||
init?: ApiRequestInit,
|
||
) => {
|
||
calls.push({ path, init });
|
||
if (!init) {
|
||
gets += 1;
|
||
return {
|
||
metadata: {
|
||
name: "shop",
|
||
uid: "uid",
|
||
resourceVersion: gets === 1 ? "7" : "8",
|
||
},
|
||
spec: { source: { uri: "cas://old", digest: "old" }, config: {} },
|
||
} as T;
|
||
}
|
||
if (calls.length === 2)
|
||
throw new KuberApiError("conflict", 409, {
|
||
title: "Conflict",
|
||
status: 409,
|
||
code: "WORKSPACE_CONFLICT",
|
||
});
|
||
return {
|
||
metadata: { name: "shop", uid: "uid", resourceVersion: "9" },
|
||
} as T;
|
||
};
|
||
|
||
await ensureWorkspace("shop", { services: {} }, snapshot, request);
|
||
|
||
expect(calls.map(({ path }) => path)).toEqual([
|
||
"/workspaces/shop",
|
||
"/workspaces/shop",
|
||
"/workspaces/shop",
|
||
"/workspaces/shop",
|
||
]);
|
||
expect(new Headers(calls[3]?.init?.headers).get("if-match")).toBe('"8"');
|
||
});
|
||
|
||
test("accepts a fresh workspace that already matches after a conflict", async () => {
|
||
let gets = 0;
|
||
let puts = 0;
|
||
const desired = {
|
||
source: { uri: `cas://${snapshot.digest}`, digest: snapshot.digest },
|
||
config: { compose: { services: {} } },
|
||
};
|
||
const request: ApiRequester = async <T>(
|
||
_path: string,
|
||
init?: ApiRequestInit,
|
||
) => {
|
||
if (!init) {
|
||
gets += 1;
|
||
return {
|
||
metadata: { name: "shop", uid: "uid", resourceVersion: "8" },
|
||
spec:
|
||
gets === 1
|
||
? { ...desired, source: { uri: "cas://old", digest: "old" } }
|
||
: desired,
|
||
} as T;
|
||
}
|
||
puts += 1;
|
||
throw new KuberApiError("conflict", 409, {
|
||
title: "Conflict",
|
||
status: 409,
|
||
code: "WORKSPACE_CONFLICT",
|
||
});
|
||
};
|
||
|
||
await ensureWorkspace("shop", { services: {} }, snapshot, request);
|
||
|
||
expect({ gets, puts }).toEqual({ gets: 2, puts: 1 });
|
||
});
|
||
|
||
test("recovers an ambiguous transport failure after the update committed", async () => {
|
||
let gets = 0;
|
||
let puts = 0;
|
||
const desired = {
|
||
source: { uri: `cas://${snapshot.digest}`, digest: snapshot.digest },
|
||
config: { compose: { services: {} } },
|
||
};
|
||
const request: ApiRequester = async <T>(
|
||
_path: string,
|
||
init?: ApiRequestInit,
|
||
) => {
|
||
if (!init) {
|
||
gets += 1;
|
||
return {
|
||
metadata: { name: "shop", uid: "uid", resourceVersion: "8" },
|
||
spec:
|
||
gets === 1
|
||
? { ...desired, source: { uri: "cas://old", digest: "old" } }
|
||
: desired,
|
||
} as T;
|
||
}
|
||
puts += 1;
|
||
throw new TypeError("connection reset");
|
||
};
|
||
|
||
await ensureWorkspace("shop", { services: {} }, snapshot, request);
|
||
|
||
expect({ gets, puts }).toEqual({ gets: 2, puts: 1 });
|
||
});
|
||
|
||
test("recovers an ambiguous transport failure with one fresh-ETag retry", async () => {
|
||
const etags: string[] = [];
|
||
let gets = 0;
|
||
const request: ApiRequester = async <T>(
|
||
_path: string,
|
||
init?: ApiRequestInit,
|
||
) => {
|
||
if (!init) {
|
||
gets += 1;
|
||
return {
|
||
metadata: {
|
||
name: "shop",
|
||
uid: "uid",
|
||
resourceVersion: gets === 1 ? "7" : "8",
|
||
},
|
||
spec: { source: { uri: "cas://old", digest: "old" }, config: {} },
|
||
} as T;
|
||
}
|
||
etags.push(new Headers(init.headers).get("if-match")!);
|
||
if (etags.length === 1) throw new TypeError("connection reset");
|
||
return {
|
||
metadata: { name: "shop", uid: "uid", resourceVersion: "9" },
|
||
} as T;
|
||
};
|
||
|
||
await ensureWorkspace("shop", { services: {} }, snapshot, request);
|
||
|
||
expect(etags).toEqual(['"7"', '"8"']);
|
||
});
|
||
|
||
test("does not retry a generic workspace update failure", async () => {
|
||
const failure = new KuberApiError("failed", 500, {
|
||
title: "Internal server error",
|
||
status: 500,
|
||
code: "INTERNAL_ERROR",
|
||
});
|
||
let calls = 0;
|
||
const request: ApiRequester = async <T>() => {
|
||
calls += 1;
|
||
if (calls === 1)
|
||
return {
|
||
metadata: { name: "shop", uid: "uid", resourceVersion: "7" },
|
||
} as T;
|
||
throw failure;
|
||
};
|
||
|
||
await expect(
|
||
ensureWorkspace("shop", { services: {} }, snapshot, request),
|
||
).rejects.toBe(failure);
|
||
expect(calls).toBe(2);
|
||
});
|
||
|
||
test("plans, applies, runs the hook, waits, then deletes stale identities", async () => {
|
||
const order: string[] = [];
|
||
const desired = [
|
||
{ apiVersion: "apps/v1", kind: "Deployment", metadata: { name: "web" } },
|
||
];
|
||
const stale = [
|
||
{
|
||
apiVersion: "v1",
|
||
kind: "Secret",
|
||
name: "old",
|
||
uid: "secret-uid",
|
||
workspaceUid: "workspace-uid",
|
||
},
|
||
];
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
order.push(path.split("/").at(-1)!);
|
||
if (path.endsWith("/plan")) return { desired, stale } as T;
|
||
return {} as T;
|
||
};
|
||
|
||
await reconcileResources(
|
||
"shop",
|
||
desired,
|
||
42_000,
|
||
() => {
|
||
order.push("hook");
|
||
},
|
||
request,
|
||
);
|
||
expect(order).toEqual(["plan", "apply", "hook", "wait", "delete"]);
|
||
});
|
||
|
||
test("resumes an interrupted reconcile operation by its persisted ID", async () => {
|
||
const calls: string[] = [];
|
||
const applyIdempotencyKeys: Array<string | null> = [];
|
||
let applyAttempts = 0;
|
||
let operationPolls = 0;
|
||
const request: ApiRequester = async <T>(
|
||
path: string,
|
||
init?: ApiRequestInit,
|
||
) => {
|
||
calls.push(path);
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
if (path.endsWith("/apply")) {
|
||
applyIdempotencyKeys.push(
|
||
new Headers(init?.headers).get("idempotency-key"),
|
||
);
|
||
applyAttempts++;
|
||
if (applyAttempts === 1) throw new TypeError("fetch failed");
|
||
return {
|
||
operationId: "operation-apply",
|
||
operation: { status: { state: "running" } },
|
||
} as T;
|
||
}
|
||
if (path === "/operations/operation-apply") {
|
||
operationPolls++;
|
||
if (operationPolls === 1) throw new TypeError("connection reset");
|
||
return { status: { state: "succeeded" } } as T;
|
||
}
|
||
throw new Error(`Unexpected request: ${path}`);
|
||
};
|
||
|
||
await reconcileResources("shop", [], 1, undefined, request, {
|
||
sleep: async () => {},
|
||
});
|
||
|
||
expect(calls).toEqual([
|
||
"/workspaces/shop/resources/plan",
|
||
"/workspaces/shop/resources/apply",
|
||
"/workspaces/shop/resources/apply",
|
||
"/operations/operation-apply/events?after=0",
|
||
"/operations/operation-apply",
|
||
"/operations/operation-apply",
|
||
"/operations/operation-apply/events?after=0",
|
||
]);
|
||
expect(applyIdempotencyKeys[0]).toBeTruthy();
|
||
expect(applyIdempotencyKeys[1]).toBe(applyIdempotencyKeys[0]);
|
||
});
|
||
|
||
test("clears the production retry timer and does not retry when cancelled", async () => {
|
||
const controller = new AbortController();
|
||
const originalSetTimeout = globalThis.setTimeout;
|
||
const originalClearTimeout = globalThis.clearTimeout;
|
||
const timer = {} as ReturnType<typeof setTimeout>;
|
||
let timerCallback: (() => void) | undefined;
|
||
let timerDelay: number | undefined;
|
||
let timerCleared = false;
|
||
let applyAttempts = 0;
|
||
let startedTimer!: () => void;
|
||
const timerStarted = new Promise<void>((resolve) => {
|
||
startedTimer = resolve;
|
||
});
|
||
globalThis.setTimeout = ((callback: () => void, milliseconds?: number) => {
|
||
timerCallback = callback;
|
||
timerDelay = milliseconds;
|
||
startedTimer();
|
||
return timer;
|
||
}) as typeof setTimeout;
|
||
globalThis.clearTimeout = ((handle: ReturnType<typeof setTimeout>) => {
|
||
if (handle === timer) timerCleared = true;
|
||
}) as typeof clearTimeout;
|
||
|
||
try {
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
if (path.endsWith("/apply")) {
|
||
applyAttempts += 1;
|
||
throw new TypeError("connection reset");
|
||
}
|
||
throw new Error(`Unexpected request: ${path}`);
|
||
};
|
||
const cancelled = new DOMException("Cancelled", "AbortError");
|
||
const run = reconcileResources("shop", [], 1, undefined, request, {
|
||
signal: controller.signal,
|
||
});
|
||
|
||
await timerStarted;
|
||
controller.abort(cancelled);
|
||
|
||
await expect(run).rejects.toBe(cancelled);
|
||
expect(timerDelay).toBe(250);
|
||
expect(timerCleared).toBe(true);
|
||
expect(timerCallback).toBeDefined();
|
||
expect(applyAttempts).toBe(1);
|
||
} finally {
|
||
globalThis.setTimeout = originalSetTimeout;
|
||
globalThis.clearTimeout = originalClearTimeout;
|
||
}
|
||
});
|
||
|
||
test("restarts after an interrupted apply is persisted before its ID is received", async () => {
|
||
const applyRequests: Array<{ key: string | null; json: unknown }> = [];
|
||
const request: ApiRequester = async <T>(
|
||
path: string,
|
||
init?: ApiRequestInit,
|
||
) => {
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
if (path.endsWith("/apply")) {
|
||
applyRequests.push({
|
||
key: new Headers(init?.headers).get("idempotency-key"),
|
||
json: init?.json,
|
||
});
|
||
if (applyRequests.length === 1) throw new TypeError("fetch failed");
|
||
if (applyRequests.length === 2)
|
||
throw new KuberApiError("operation interrupted", 500, {
|
||
title: "Operation interrupted",
|
||
status: 500,
|
||
code: "OPERATION_INTERRUPTED",
|
||
});
|
||
return {
|
||
operationId: "operation-apply",
|
||
operation: { status: { state: "succeeded" } },
|
||
} as T;
|
||
}
|
||
throw new Error(`Unexpected request: ${path}`);
|
||
};
|
||
|
||
await reconcileResources("shop", [], 1, undefined, request, {
|
||
sleep: async () => {},
|
||
});
|
||
|
||
expect(applyRequests).toHaveLength(3);
|
||
expect(applyRequests[0]?.key).toBeTruthy();
|
||
expect(applyRequests[1]?.key).toBe(applyRequests[0]?.key);
|
||
expect(applyRequests[2]?.key).toBeTruthy();
|
||
expect(applyRequests[2]?.key).not.toBe(applyRequests[0]?.key);
|
||
expect(applyRequests[2]?.json).toEqual(applyRequests[0]?.json);
|
||
});
|
||
|
||
test("polls persisted resource progress without duplicating events after reconnect", async () => {
|
||
const progress: string[] = [];
|
||
let poll = 0;
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
if (path.endsWith("/apply"))
|
||
return {
|
||
operationId: "operation-apply",
|
||
operation: { status: { state: "running" } },
|
||
} as T;
|
||
if (path.includes("/events")) {
|
||
poll++;
|
||
return {
|
||
items:
|
||
poll === 1
|
||
? [
|
||
{
|
||
sequence: 1,
|
||
data: null,
|
||
},
|
||
{
|
||
sequence: 2,
|
||
data: {
|
||
resource: {
|
||
apiVersion: "v1",
|
||
kind: "Service",
|
||
name: "web",
|
||
},
|
||
phase: "apply",
|
||
state: "started",
|
||
},
|
||
},
|
||
]
|
||
: [
|
||
{
|
||
sequence: 2,
|
||
data: {
|
||
resource: {
|
||
apiVersion: "v1",
|
||
kind: "Service",
|
||
name: "web",
|
||
},
|
||
phase: "apply",
|
||
state: "started",
|
||
},
|
||
},
|
||
{
|
||
sequence: 3,
|
||
data: {
|
||
resource: {
|
||
apiVersion: "v1",
|
||
kind: "Service",
|
||
name: "web",
|
||
},
|
||
phase: "apply",
|
||
state: "succeeded",
|
||
},
|
||
},
|
||
],
|
||
} as T;
|
||
}
|
||
if (path === "/operations/operation-apply") {
|
||
return { status: { state: poll > 1 ? "succeeded" : "running" } } as T;
|
||
}
|
||
throw new Error(`Unexpected request: ${path}`);
|
||
};
|
||
|
||
await reconcileResources("shop", [], 1, undefined, request, {
|
||
sleep: async () => {},
|
||
onEvent: (event) =>
|
||
progress.push(`${event.sequence}:${event.data.state}`),
|
||
});
|
||
expect(progress).toEqual(["2:started", "3:succeeded"]);
|
||
});
|
||
|
||
test("fails visibly when retained operation progress has a cursor gap", async () => {
|
||
let applyAttempts = 0;
|
||
let operationPolls = 0;
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
if (path.endsWith("/apply")) {
|
||
applyAttempts += 1;
|
||
return {
|
||
operationId: "operation-apply",
|
||
operation: { status: { state: "running" } },
|
||
} as T;
|
||
}
|
||
if (path.includes("/events"))
|
||
return {
|
||
retainedFirstSequence: 3,
|
||
cursorGap: true,
|
||
items: [],
|
||
} as T;
|
||
if (path === "/operations/operation-apply") {
|
||
operationPolls += 1;
|
||
return { status: { state: "succeeded" } } as T;
|
||
}
|
||
throw new Error(`Unexpected request: ${path}`);
|
||
};
|
||
|
||
await expect(
|
||
reconcileResources("shop", [], 1, undefined, request, {
|
||
sleep: async () => {},
|
||
}),
|
||
).rejects.toThrow("progress history was truncated");
|
||
expect(applyAttempts).toBe(1);
|
||
expect(operationPolls).toBe(0);
|
||
});
|
||
|
||
test("attaches concurrent live resource subtasks before driving multi-target operations", async () => {
|
||
const children = new Map<string, { title: string; output: string }>();
|
||
let operationCalls = 0;
|
||
const listr = new Listr([
|
||
{
|
||
title: "Apply resources",
|
||
task: (_ctx, task) =>
|
||
runLiveResourceOperation(
|
||
task,
|
||
"apply",
|
||
[
|
||
{
|
||
apiVersion: "v1",
|
||
kind: "Service",
|
||
name: "web",
|
||
},
|
||
{
|
||
apiVersion: "apps/v1",
|
||
kind: "Deployment",
|
||
name: "api",
|
||
},
|
||
],
|
||
async (onEvent) => {
|
||
operationCalls += 1;
|
||
expect(children.size).toBe(2);
|
||
onEvent({
|
||
sequence: 1,
|
||
data: {
|
||
resource: {
|
||
apiVersion: "v1",
|
||
kind: "Service",
|
||
name: "web",
|
||
},
|
||
phase: "apply",
|
||
state: "started",
|
||
},
|
||
});
|
||
onEvent({
|
||
sequence: 2,
|
||
data: {
|
||
resource: {
|
||
apiVersion: "apps/v1",
|
||
kind: "Deployment",
|
||
name: "api",
|
||
},
|
||
phase: "apply",
|
||
state: "succeeded",
|
||
},
|
||
});
|
||
expect(children.get("web")?.output).toBe("started");
|
||
expect(children.get("api")?.output).toBe("succeeded");
|
||
},
|
||
{
|
||
onTaskStarted: (target, activeTask) =>
|
||
children.set(target.name, activeTask),
|
||
},
|
||
),
|
||
},
|
||
]);
|
||
|
||
await listr.run();
|
||
|
||
expect(operationCalls).toBe(1);
|
||
expect(listr.tasks[0]?.subtasks).toHaveLength(2);
|
||
expect(listr.tasks[0]?.subtasks.map((task) => task.title)).toEqual([
|
||
"Apply Service/web",
|
||
"Apply Deployment/api",
|
||
]);
|
||
}, 1_000);
|
||
|
||
test("renders a single live resource target and runs zero-target operations directly", async () => {
|
||
const single = new Listr([
|
||
{
|
||
title: "Apply resources",
|
||
task: (_ctx, task) =>
|
||
runLiveResourceOperation(
|
||
task,
|
||
"apply",
|
||
[{ apiVersion: "v1", kind: "Service", name: "web" }],
|
||
async (onEvent) => {
|
||
onEvent({
|
||
sequence: 1,
|
||
data: {
|
||
resource: { apiVersion: "v1", kind: "Service", name: "web" },
|
||
phase: "apply",
|
||
state: "succeeded",
|
||
},
|
||
});
|
||
},
|
||
),
|
||
},
|
||
]);
|
||
await single.run();
|
||
expect(single.tasks[0]?.subtasks.map(({ title }) => title)).toEqual([
|
||
"Apply Service/web",
|
||
]);
|
||
|
||
let called = false;
|
||
await runLiveResourceOperation(
|
||
{ signal: new AbortController().signal } as never,
|
||
"delete",
|
||
[],
|
||
async () => {
|
||
called = true;
|
||
},
|
||
);
|
||
expect(called).toBe(true);
|
||
});
|
||
|
||
test("cancels a multi-resource live operation without leaving child tasks waiting", async () => {
|
||
const children: Array<{ signal: AbortSignal }> = [];
|
||
let phaseTask: { cancel: () => void } | undefined;
|
||
let operationObservedAbort = false;
|
||
let startOperation: (() => void) | undefined;
|
||
const operationStarted = new Promise<void>((resolve) => {
|
||
startOperation = resolve;
|
||
});
|
||
const listr = new Listr([
|
||
{
|
||
title: "Apply resources",
|
||
task: (_ctx, task) => {
|
||
phaseTask = task;
|
||
return runLiveResourceOperation(
|
||
task,
|
||
"apply",
|
||
[
|
||
{ apiVersion: "v1", kind: "Service", name: "web" },
|
||
{ apiVersion: "apps/v1", kind: "Deployment", name: "api" },
|
||
],
|
||
async (_onEvent, signal) => {
|
||
startOperation?.();
|
||
await new Promise<void>((_resolve, reject) => {
|
||
signal.addEventListener(
|
||
"abort",
|
||
() => {
|
||
operationObservedAbort = true;
|
||
reject(signal.reason);
|
||
},
|
||
{ once: true },
|
||
);
|
||
});
|
||
},
|
||
{ onTaskStarted: (_target, child) => children.push(child) },
|
||
);
|
||
},
|
||
},
|
||
]);
|
||
const originalExit = process.exit;
|
||
const exitCodes: Array<number | undefined> = [];
|
||
process.exit = ((code?: number) => {
|
||
exitCodes.push(code);
|
||
return undefined as never;
|
||
}) as typeof process.exit;
|
||
|
||
try {
|
||
const run = listr.run();
|
||
await operationStarted;
|
||
phaseTask?.cancel();
|
||
await Promise.race([
|
||
run,
|
||
Bun.sleep(100).then(() => {
|
||
throw new Error("Cancelled live operation did not settle promptly");
|
||
}),
|
||
]);
|
||
} finally {
|
||
process.exit = originalExit;
|
||
}
|
||
|
||
expect(operationObservedAbort).toBe(true);
|
||
expect(exitCodes).toEqual([127]);
|
||
expect(children).toHaveLength(2);
|
||
expect(children.every(({ signal }) => signal.aborted)).toBe(true);
|
||
expect(
|
||
listr.tasks[0]?.subtasks.every(({ state }) => state === "CANCELLED"),
|
||
).toBe(true);
|
||
}, 1_000);
|
||
|
||
test("resubmits with a fresh key after server restart interruption", async () => {
|
||
const applyRequests: Array<{ key: string | null; json: unknown }> = [];
|
||
let operationPolls = 0;
|
||
const request: ApiRequester = async <T>(
|
||
path: string,
|
||
init?: ApiRequestInit,
|
||
) => {
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
if (path.endsWith("/apply")) {
|
||
applyRequests.push({
|
||
key: new Headers(init?.headers).get("idempotency-key"),
|
||
json: init?.json,
|
||
});
|
||
return {
|
||
operationId: `operation-apply-${applyRequests.length}`,
|
||
operation: { status: { state: "running" } },
|
||
} as T;
|
||
}
|
||
if (path === "/operations/operation-apply-1") {
|
||
operationPolls++;
|
||
if (operationPolls === 1)
|
||
return {
|
||
status: {
|
||
state: "failed",
|
||
error: {
|
||
code: "OPERATION_INTERRUPTED",
|
||
message: "server restarted",
|
||
},
|
||
},
|
||
} as T;
|
||
return { status: { state: "succeeded" } } as T;
|
||
}
|
||
if (path === "/operations/operation-apply-2")
|
||
return { status: { state: "succeeded" } } as T;
|
||
throw new Error(`Unexpected request: ${path}`);
|
||
};
|
||
|
||
await reconcileResources("shop", [], 1, undefined, request, {
|
||
sleep: async () => {},
|
||
});
|
||
|
||
expect(applyRequests).toHaveLength(2);
|
||
expect(applyRequests[0]?.key).toBeTruthy();
|
||
expect(applyRequests[1]?.key).toBeTruthy();
|
||
expect(applyRequests[1]?.key).not.toBe(applyRequests[0]?.key);
|
||
expect(applyRequests[1]?.json).toEqual(applyRequests[0]?.json);
|
||
});
|
||
|
||
test.each([
|
||
["failed", "OPERATION_FAILED"],
|
||
["cancelled", "OPERATION_CANCELLED"],
|
||
])("does not retry arbitrary %s operations", async (state, code) => {
|
||
let applyAttempts = 0;
|
||
let operationPolls = 0;
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
if (path.endsWith("/apply")) {
|
||
applyAttempts++;
|
||
return {
|
||
operationId: "operation-apply",
|
||
operation: { status: { state: "running" } },
|
||
} as T;
|
||
}
|
||
if (path === "/operations/operation-apply") {
|
||
operationPolls++;
|
||
return {
|
||
status: { state, error: { code, message: "terminal" } },
|
||
} as T;
|
||
}
|
||
throw new Error(`Unexpected request: ${path}`);
|
||
};
|
||
|
||
await expect(
|
||
reconcileResources("shop", [], 1, undefined, request, {
|
||
sleep: async () => {},
|
||
}),
|
||
).rejects.toMatchObject({ code });
|
||
expect(applyAttempts).toBe(1);
|
||
expect(operationPolls).toBe(1);
|
||
});
|
||
|
||
test("bounds restart recovery sleeps by the resume deadline", async () => {
|
||
let currentTime = 0;
|
||
const sleeps: Array<{ startedAt: number; milliseconds: number }> = [];
|
||
let operationNumber = 0;
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
if (path.endsWith("/apply")) {
|
||
operationNumber++;
|
||
return {
|
||
operationId: `operation-${operationNumber}`,
|
||
operation: { status: { state: "running" } },
|
||
} as T;
|
||
}
|
||
return {
|
||
status: {
|
||
state: "failed",
|
||
error: { code: "OPERATION_INTERRUPTED", message: "restarted" },
|
||
},
|
||
} as T;
|
||
};
|
||
|
||
await expect(
|
||
reconcileResources("shop", [], 0, undefined, request, {
|
||
now: () => currentTime,
|
||
sleep: async (milliseconds) => {
|
||
sleeps.push({ startedAt: currentTime, milliseconds });
|
||
currentTime += milliseconds;
|
||
},
|
||
}),
|
||
).rejects.toThrow("Timed out while reconnecting to resume the operation");
|
||
|
||
expect(sleeps).not.toHaveLength(0);
|
||
expect(
|
||
sleeps.every(
|
||
({ startedAt, milliseconds }) => startedAt + milliseconds <= 60_000,
|
||
),
|
||
).toBe(true);
|
||
expect(currentTime).toBe(60_000);
|
||
});
|
||
|
||
test("does not restart after an apply interruption reaches the resume deadline", async () => {
|
||
let currentTime = 0;
|
||
let applyAttempts = 0;
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
applyAttempts++;
|
||
if (applyAttempts === 1) throw new TypeError("connection reset");
|
||
throw new KuberApiError("operation interrupted", 500, {
|
||
title: "Operation interrupted",
|
||
status: 500,
|
||
code: "OPERATION_INTERRUPTED",
|
||
});
|
||
};
|
||
|
||
await expect(
|
||
reconcileResources("shop", [], 0, undefined, request, {
|
||
now: () => currentTime,
|
||
sleep: async (milliseconds) => {
|
||
currentTime += milliseconds;
|
||
if (currentTime < 60_000) currentTime = 60_000;
|
||
},
|
||
}),
|
||
).rejects.toMatchObject({
|
||
code: "OPERATION_INTERRUPTED",
|
||
status: 500,
|
||
});
|
||
expect(applyAttempts).toBe(2);
|
||
expect(currentTime).toBe(60_000);
|
||
});
|
||
|
||
test("fails immediately for a typed 503 operation-store error", async () => {
|
||
const unavailable = new KuberApiError(
|
||
"Operation storage is not configured",
|
||
503,
|
||
{
|
||
title: "Service unavailable",
|
||
status: 503,
|
||
code: "OPERATION_STORE_UNAVAILABLE",
|
||
},
|
||
);
|
||
const calls: string[] = [];
|
||
const sleeps: number[] = [];
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
calls.push(path);
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
throw unavailable;
|
||
};
|
||
|
||
await expect(
|
||
reconcileResources("shop", [], 1, undefined, request, {
|
||
sleep: async (milliseconds) => {
|
||
sleeps.push(milliseconds);
|
||
},
|
||
}),
|
||
).rejects.toBe(unavailable);
|
||
|
||
expect(calls).toEqual([
|
||
"/workspaces/shop/resources/plan",
|
||
"/workspaces/shop/resources/apply",
|
||
]);
|
||
expect(sleeps).toEqual([]);
|
||
});
|
||
|
||
test("never sleeps past the operation resume deadline", async () => {
|
||
let currentTime = 0;
|
||
const sleeps: Array<{ startedAt: number; milliseconds: number }> = [];
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
if (path.endsWith("/plan")) return { desired: [], stale: [] } as T;
|
||
throw new TypeError("connection reset");
|
||
};
|
||
|
||
await expect(
|
||
reconcileResources("shop", [], 0, undefined, request, {
|
||
now: () => currentTime,
|
||
sleep: async (milliseconds) => {
|
||
sleeps.push({ startedAt: currentTime, milliseconds });
|
||
currentTime += milliseconds;
|
||
},
|
||
}),
|
||
).rejects.toThrow("connection reset");
|
||
|
||
expect(sleeps).not.toHaveLength(0);
|
||
expect(
|
||
sleeps.every(
|
||
({ startedAt, milliseconds }) => startedAt + milliseconds <= 60_000,
|
||
),
|
||
).toBe(true);
|
||
expect(sleeps.at(-1)).toEqual({
|
||
startedAt: 57_750,
|
||
milliseconds: 2_250,
|
||
});
|
||
expect(currentTime).toBe(60_000);
|
||
});
|
||
|
||
test("fails closed with the precise missing adoption route", async () => {
|
||
const request: ApiRequester = async <T>(path: string) => {
|
||
if (path.endsWith("/plan"))
|
||
throw new Error(
|
||
"Namespace shop is external; refusing workspace mutation",
|
||
);
|
||
return {} as T;
|
||
};
|
||
await expect(
|
||
reconcileResources("shop", [], 1, undefined, request),
|
||
).rejects.toThrow(`POST ${workspaceAdoptionRoute("shop")}`);
|
||
});
|
||
});
|