diff --git a/command/main.ts b/command/main.ts index 7651131..4740112 100644 --- a/command/main.ts +++ b/command/main.ts @@ -22,7 +22,7 @@ import { trust } from "./trust"; export const main = defineCommand({ meta: { name: "kuber", - version: "2.3.0", + version: "2.3.1", description: "Docker Compose -> K8s translation layer", }, args: { diff --git a/command/up.ts b/command/up.ts index ef8dae4..ad8b657 100644 --- a/command/up.ts +++ b/command/up.ts @@ -73,6 +73,7 @@ export type OperationResumeOptions = { now?: () => number; sleep?: (milliseconds: number) => Promise; onEvent?: (event: ResourceOperationEvent) => void; + signal?: AbortSignal; }; type OperationEventsResponse = { @@ -335,6 +336,41 @@ function operationResumeDeadline( return now + Math.max(rolloutTimeoutMs, 0) + OPERATION_RESUME_GRACE_MS; } +function abortReason(signal: AbortSignal): unknown { + return ( + signal.reason ?? new DOMException("The operation was aborted", "AbortError") + ); +} + +function throwIfAborted(signal: AbortSignal | undefined): void { + if (signal?.aborted) throw abortReason(signal); +} + +function abortableDelay( + milliseconds: number, + signal: AbortSignal | undefined, +): Promise { + throwIfAborted(signal); + return new Promise((resolve, reject) => { + let timeout: ReturnType | undefined; + const cleanup = () => { + if (timeout !== undefined) clearTimeout(timeout); + signal?.removeEventListener("abort", abort); + }; + const complete = () => { + cleanup(); + resolve(); + }; + const abort = () => { + cleanup(); + reject(abortReason(signal!)); + }; + timeout = setTimeout(complete, milliseconds); + signal?.addEventListener("abort", abort, { once: true }); + if (signal?.aborted) abort(); + }); +} + async function resumeManagedOperation( project: string, request: ApiRequester, @@ -344,18 +380,22 @@ async function resumeManagedOperation( options: OperationResumeOptions, ): Promise { const now = options.now ?? Date.now; - const sleep = options.sleep ?? ((milliseconds) => Bun.sleep(milliseconds)); const deadline = operationResumeDeadline(rolloutTimeoutMs, now()); const headers = new Headers(init.headers); headers.set("idempotency-key", randomUUID()); headers.set("prefer", "respond-async"); - let operationInit = { ...init, headers }; + let operationInit = { + ...init, + headers, + signal: options.signal ?? init.signal, + }; let operationId: string | undefined; let backoffMs = OPERATION_RESUME_INITIAL_BACKOFF_MS; let restartRetryPending = false; let eventCursor = 0; for (;;) { + throwIfAborted(options.signal); if (restartRetryPending && now() >= deadline) throw new Error("Timed out while reconnecting to resume the operation"); try { @@ -365,7 +405,7 @@ async function resumeManagedOperation( project, request, `/operations/${encodeURIComponent(operationId)}`, - {}, + { signal: options.signal }, ); status = operation.status; } else { @@ -386,7 +426,7 @@ async function resumeManagedOperation( project, request, `/operations/${encodeURIComponent(operationId)}/events?after=${eventCursor}`, - {}, + { signal: options.signal }, ); if (hasOperationEventCursorGap(events, eventCursor)) throw new OperationProgressCursorGapError(); @@ -402,6 +442,7 @@ async function resumeManagedOperation( } } catch (error) { if (error instanceof OperationProgressCursorGapError) throw error; + throwIfAborted(options.signal); // Event polling is additive; never delay resumable operation status polling. } } @@ -413,10 +454,15 @@ async function resumeManagedOperation( eventCursor = 0; restartRetryPending = true; headers.set("idempotency-key", randomUUID()); - operationInit = { ...init, headers }; + operationInit = { + ...init, + headers, + signal: options.signal ?? init.signal, + }; } else if (status.state === "failed" || status.state === "cancelled") throw operationFailure(operationId, status); } catch (error) { + throwIfAborted(options.signal); if ( error instanceof KuberApiError && error.code === "OPERATION_INTERRUPTED" @@ -426,7 +472,11 @@ async function resumeManagedOperation( eventCursor = 0; restartRetryPending = true; headers.set("idempotency-key", randomUUID()); - operationInit = { ...init, headers }; + operationInit = { + ...init, + headers, + signal: options.signal ?? init.signal, + }; } else { if (!isRecoverableConnectionInterruption(error)) throw adoptionHint(project, error); @@ -437,7 +487,12 @@ async function resumeManagedOperation( const remainingMs = deadline - now(); if (remainingMs <= 0) throw new Error("Timed out while reconnecting to resume the operation"); - await sleep(Math.min(backoffMs, remainingMs)); + const delayMs = Math.min(backoffMs, remainingMs); + if (options.sleep) { + throwIfAborted(options.signal); + await options.sleep(delayMs); + throwIfAborted(options.signal); + } else await abortableDelay(delayMs, options.signal); backoffMs = Math.min(backoffMs * 2, OPERATION_RESUME_MAX_BACKOFF_MS); } } @@ -561,52 +616,59 @@ function resourceOperationTaskTitle( return `${action} ${resource.kind}/${resource.name}`; } -export async function runLiveResourceOperation( +export function runLiveResourceOperation( task: ListrTaskWrapper, phase: ResourceOperationEvent["data"]["phase"], targets: ResourceOperationTarget[], operation: ( onEvent: (event: ResourceOperationEvent) => void, + signal: AbortSignal, ) => Promise, options: LiveResourceOperationOptions = {}, -): Promise { - if (targets.length === 0) return operation(() => {}); +): Listr | Promise { + if (targets.length === 0) return operation(() => {}, task.signal); const active = new Map< string, ListrTaskWrapper >(); - const complete: Array<() => void> = []; + const complete = new Set<() => void>(); let started = 0; - let markStarted!: () => void; - const ready = new Promise((resolve) => { - markStarted = resolve; - }); - const children = task.newListr( + return task.newListr( targets.map((target) => ({ title: resourceOperationTaskTitle(phase, target), - task: (_target, child) => { + task: async (_target, child) => { active.set(resourceOperationKey(target), child); options.onTaskStarted?.(target, child); started += 1; - if (started === targets.length) markStarted(); - return new Promise((resolve) => complete.push(resolve)); + if (started !== targets.length) + return new Promise((resolve) => { + const completeTask = () => { + task.signal.removeEventListener("abort", completeTask); + complete.delete(completeTask); + resolve(); + }; + complete.add(completeTask); + if (task.signal.aborted) completeTask(); + else + task.signal.addEventListener("abort", completeTask, { + once: true, + }); + }); + try { + await operation((event) => { + const resource = event.data.resource; + const activeTask = active.get(resourceOperationKey(resource)); + if (!activeTask) return; + activeTask.output = event.data.state; + task.output = resourceOperationEventTitle(event); + }, task.signal); + } finally { + for (const resolve of [...complete]) resolve(); + } }, })), + { concurrent: true }, ); - const childRun = children.run(); - await ready; - try { - await operation((event) => { - const resource = event.data.resource; - const child = active.get(resourceOperationKey(resource)); - if (!child) return; - child.output = event.data.state; - task.output = resourceOperationEventTitle(event); - }); - } finally { - for (const resolve of complete) resolve(); - await childRun; - } } export async function runUp( @@ -801,14 +863,14 @@ export async function runUp( phaseTask, "apply", desired, - (onEvent) => + (onEvent, signal) => resumeManagedOperation( project, request, `${resourcePath}/apply`, { method: "POST", json: { resources: plan.desired } }, config.rolloutTimeoutMs, - { onEvent }, + { onEvent, signal }, ), ), }, @@ -829,7 +891,7 @@ export async function runUp( phaseTask, "wait", deployments, - (onEvent) => + (onEvent, signal) => resumeManagedOperation( project, request, @@ -844,7 +906,7 @@ export async function runUp( }, }, config.rolloutTimeoutMs, - { onEvent }, + { onEvent, signal }, ), ), }, @@ -856,14 +918,14 @@ export async function runUp( phaseTask, "delete", stale, - (onEvent) => + (onEvent, signal) => resumeManagedOperation( project, request, `${resourcePath}/delete`, { method: "POST", json: { resources: plan.stale } }, config.rolloutTimeoutMs, - { onEvent }, + { onEvent, signal }, ), ), }, diff --git a/package.json b/package.json index 99f3df4..7689d85 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@dmgnr/kuber", - "version": "2.3.0", + "version": "2.3.1", "description": "Docker Compose to Kubernetes translation layer", "bin": { "kuber": "dist/index.js" diff --git a/tests/command/ci.test.ts b/tests/command/ci.test.ts index 5a49668..c14758e 100644 --- a/tests/command/ci.test.ts +++ b/tests/command/ci.test.ts @@ -1,5 +1,11 @@ import { expect, spyOn, test } from "bun:test"; -import { apiKeyRequest } from "../../command/ci"; +import { mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { ci, apiKeyRequest, runCi } from "../../command/ci"; +import { provideContext } from "../../lib/context"; +import { writeSession } from "../../lib/session"; +import { resolveTrustIdentity } from "../../lib/trust"; test("CI requester supplies an API key without session authentication", async () => { const fetch = spyOn(globalThis, "fetch").mockResolvedValue( @@ -17,3 +23,104 @@ test("CI requester supplies an API key without session authentication", async () fetch.mockRestore(); } }); + +test("CI command requires an API key before loading project context", async () => { + const previous = process.env.KUBER_API_KEY; + delete process.env.KUBER_API_KEY; + try { + await expect( + ci.run!({ + args: { apiKey: undefined, build: false, trust: false }, + } as never), + ).rejects.toThrow("KUBER_API_KEY or --api-key is required"); + } finally { + if (previous === undefined) delete process.env.KUBER_API_KEY; + else process.env.KUBER_API_KEY = previous; + } +}); + +test("CI trust grant and deployment pipeline use the API key requester", async () => { + const root = await mkdtemp(join(tmpdir(), "kuber-ci-")); + const previousCwd = process.cwd(); + const previousConfigHome = process.env.XDG_CONFIG_HOME; + const previousRuntimeDirectory = process.env.XDG_RUNTIME_DIR; + const calls: Array<{ + path: string; + method: string | undefined; + body: string | undefined; + headers: Headers; + }> = []; + const fetch = spyOn(globalThis, "fetch").mockImplementation((async ( + input, + init, + ) => { + const url = new URL(input.toString()); + const path = url.pathname.replace("/api/v2", ""); + calls.push({ + path, + method: init?.method, + body: typeof init?.body === "string" ? init.body : undefined, + headers: new Headers(init?.headers), + }); + if (path === "/workspaces/shop") + return new Response( + JSON.stringify({ + metadata: { name: "shop", uid: "workspace", resourceVersion: "1" }, + }), + ); + if (path.endsWith("/plan")) + return new Response(JSON.stringify({ desired: [], stale: [] })); + if (path === "/snapshots/negotiate") + return new Response( + JSON.stringify({ workspace: "sha256:abc", missing: [], ready: true }), + ); + return new Response(JSON.stringify({ resourcesAdopted: 0 })); + }) as typeof globalThis.fetch); + + try { + process.env.XDG_CONFIG_HOME = join(root, "config"); + process.env.XDG_RUNTIME_DIR = join(root, "runtime"); + await writeSession( + { + token: "session-secret", + expiresAt: "2030-01-01T00:00:00Z", + user: { username: "ci", roles: [] }, + }, + true, + ); + await writeFile(join(root, "compose.yml"), "services: {}\n"); + await writeFile( + join(root, ".kuberrc.ts"), + 'export default { project: "shop" };\n', + ); + const git = Bun.spawn(["git", "init", "-q", root]); + expect(await git.exited).toBe(0); + const fingerprint = (await resolveTrustIdentity("shop", root)).fingerprint; + process.chdir(root); + await provideContext(() => runCi(false, true, "ci-key")); + + expect(calls[0]?.path).toBe("/workspaces/shop/trust"); + expect(calls[0]?.method).toBe("POST"); + expect(JSON.parse(calls[0]?.body ?? "")).toEqual({ fingerprint }); + expect(calls.some(({ path }) => path.endsWith("/resources/plan"))).toBe( + true, + ); + expect(calls.some(({ path }) => path.endsWith("/resources/apply"))).toBe( + true, + ); + expect( + calls.every( + ({ headers }) => headers.get("authorization") === "Bearer ci-key", + ), + ).toBe(true); + } finally { + fetch.mockRestore(); + process.chdir(previousCwd); + if (previousConfigHome === undefined) delete process.env.XDG_CONFIG_HOME; + else process.env.XDG_CONFIG_HOME = previousConfigHome; + if (previousRuntimeDirectory === undefined) + delete process.env.XDG_RUNTIME_DIR; + else process.env.XDG_RUNTIME_DIR = previousRuntimeDirectory; + await rm(root, { recursive: true, force: true }); + } +}); diff --git a/tests/command/up-api.test.ts b/tests/command/up-api.test.ts index cceb955..7f0d7cc 100644 --- a/tests/command/up-api.test.ts +++ b/tests/command/up-api.test.ts @@ -237,6 +237,57 @@ describe("up API pipeline", () => { 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; + let timerCallback: (() => void) | undefined; + let timerDelay: number | undefined; + let timerCleared = false; + let applyAttempts = 0; + let startedTimer!: () => void; + const timerStarted = new Promise((resolve) => { + startedTimer = resolve; + }); + globalThis.setTimeout = ((callback: () => void, milliseconds?: number) => { + timerCallback = callback; + timerDelay = milliseconds; + startedTimer(); + return timer; + }) as typeof setTimeout; + globalThis.clearTimeout = ((handle: ReturnType) => { + if (handle === timer) timerCleared = true; + }) as typeof clearTimeout; + + try { + const request: ApiRequester = async (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 ( @@ -385,14 +436,13 @@ describe("up API pipeline", () => { expect(operationPolls).toBe(0); }); - test("starts live resource subtasks before the operation and updates them from progress", async () => { - let finish!: () => void; - let child: { title: string; output: string } | undefined; - const completed = new Promise((resolve) => (finish = resolve)); + test("attaches concurrent live resource subtasks before driving multi-target operations", async () => { + const children = new Map(); + let operationCalls = 0; const listr = new Listr([ { title: "Apply resources", - task: async (_ctx, task) => + task: (_ctx, task) => runLiveResourceOperation( task, "apply", @@ -402,9 +452,15 @@ describe("up API pipeline", () => { kind: "Service", name: "web", }, + { + apiVersion: "apps/v1", + kind: "Deployment", + name: "api", + }, ], async (onEvent) => { - expect(child?.title).toBe("Apply Service/web"); + operationCalls += 1; + expect(children.size).toBe(2); onEvent({ sequence: 1, data: { @@ -417,20 +473,146 @@ describe("up API pipeline", () => { state: "started", }, }); - expect(child?.output).toBe("started"); - await completed; + 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), }, - { onTaskStarted: (_target, activeTask) => (child = activeTask) }, ), }, ]); - const run = listr.run(); - await Bun.sleep(0); - finish(); - await run; + 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((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((_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 = []; + 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; diff --git a/tests/server/api-keys.test.ts b/tests/server/api-keys.test.ts index 0c55004..c41e8be 100644 --- a/tests/server/api-keys.test.ts +++ b/tests/server/api-keys.test.ts @@ -3,6 +3,9 @@ import { cleanupExpiredSessions, createApp } from "../../server/app"; import { hashToken, MemoryAuthStore } from "../../server/auth"; import { MemoryAuditStore } from "../../server/audit-store"; import { MemoryOperationStore } from "../../server/operation-store"; +import type { ManagementService } from "../../server/management"; +import { MemoryTrustStore } from "../../server/trust-store"; +import { MemoryWorkspaceStore } from "../../server/workspace-store"; const now = Date.parse("2026-09-05T00:00:00.000Z"); @@ -255,6 +258,67 @@ describe("API keys", () => { ).toBe(403); }); + test("allows scoped keys to apply only in their workspace and rejects deleted owners", async () => { + const { store } = await setup(); + const workspaceStore = new MemoryWorkspaceStore({ + uid: () => "workspace-uid", + }); + for (const id of ["shop", "other"]) + await workspaceStore.create({ + id, + source: { uri: `oci://example/${id}`, digest: "sha256:abc" }, + }); + const trustStore = new MemoryTrustStore(); + const fingerprint = "a".repeat(64); + await trustStore.grant("shop", fingerprint); + let applies = 0; + const app = createApp({ + store, + workspaceStore, + trustStore, + operationStore: new MemoryOperationStore(() => new Date(now)), + management: { + applyResources: async () => { + applies += 1; + return []; + }, + } as unknown as ManagementService, + now: () => now, + }); + await store.createApiKey({ + id: "key_apply_scope_1", + tokenHash: hashToken("scoped-apply-key"), + username: "ci", + capabilities: ["kubernetes:write"], + workspace: "shop", + expiresAt: "2026-10-05T00:00:00.000Z", + }); + const apply = (workspace: string) => + app( + request( + `/api/v2/workspaces/${workspace}/resources/apply`, + { + method: "POST", + headers: { + "idempotency-key": `apply-${workspace}`, + "x-kuber-trust-project": "shop", + "x-kuber-trust-fingerprint": fingerprint, + }, + body: JSON.stringify({ resources: [] }), + }, + "scoped-apply-key", + ), + ); + + expect((await apply("shop")).status).toBe(200); + expect((await apply("other")).status).toBe(403); + expect(applies).toBe(1); + await store.deleteUser("ci"); + expect( + (await app(request("/api/v2/me", {}, "scoped-apply-key"))).status, + ).toBe(401); + }); + test("limits workspace-scoped keys to their own audit records", async () => { const { app, store, auditStore } = await setup(); await store.createApiKey({ diff --git a/tests/server/app.test.ts b/tests/server/app.test.ts index 302dbb0..5622687 100644 --- a/tests/server/app.test.ts +++ b/tests/server/app.test.ts @@ -391,6 +391,51 @@ describe("kuber v2 HTTP routes", () => { expect((await apply(app)).status).toBe(403); }); + test("returns OPERATION_CONFLICT when an idempotency key is reused for another body", async () => { + const workspaceStore = new MemoryWorkspaceStore({ + uid: () => "workspace-uid", + }); + await workspaceStore.create({ + id: "demo", + source: { uri: "oci://example/demo", digest: "sha256:abc" }, + }); + const trustStore = new MemoryTrustStore(); + const fingerprint = "c".repeat(64); + await trustStore.grant("demo", fingerprint); + const app = createApp({ + store: await authenticatedStore("operator"), + workspaceStore, + trustStore, + operationStore: new MemoryOperationStore(), + management: { + applyResources: async () => [], + } as unknown as ManagementService, + }); + const apply = (resources: unknown[]) => + app( + request( + "/api/v2/workspaces/demo/resources/apply", + { + method: "POST", + headers: { + "idempotency-key": "same-apply", + "x-kuber-trust-project": "demo", + "x-kuber-trust-fingerprint": fingerprint, + }, + body: JSON.stringify({ resources }), + }, + "token", + ), + ); + + expect((await apply([])).status).toBe(200); + const conflict = await apply([ + { apiVersion: "v1", kind: "Service", metadata: { name: "web" } }, + ]); + expect(conflict.status).toBe(409); + expect(await conflict.json()).toMatchObject({ code: "OPERATION_CONFLICT" }); + }); + test("starts preferred resource operations before returning so progress can be polled", async () => { const workspaceStore = new MemoryWorkspaceStore({ uid: () => "workspace-uid", diff --git a/tests/server/management.test.ts b/tests/server/management.test.ts index fdae5db..f05c38d 100644 --- a/tests/server/management.test.ts +++ b/tests/server/management.test.ts @@ -236,6 +236,35 @@ describe("server management service", () => { }); }); + test("emits aborted resource progress when an execution signal aborts", async () => { + const controller = new AbortController(); + const events: Array<{ state: string; resource: { name: string } }> = []; + const service = createManagementService( + dependencies({ + applyResource: async (_resource, execution) => { + controller.abort(); + if (execution?.signal?.aborted) + throw new Error("Workspace operation execution was cancelled"); + return _resource; + }, + }), + ); + + await expect( + service.applyResources( + workspace, + [object("Service", "web", "service-uid")], + { + signal: controller.signal, + emit: async (event) => { + events.push(event); + }, + }, + ), + ).rejects.toThrow("cancelled"); + expect(events.map(({ state }) => state)).toEqual(["started", "aborted"]); + }); + test("rejects namespace and resource ownership mismatches", async () => { const wrongNamespace = createManagementService( dependencies({