diff --git a/.github/workflows/destroyer.yml b/.github/workflows/destroyer.yml index 627dbf2..e903d79 100644 --- a/.github/workflows/destroyer.yml +++ b/.github/workflows/destroyer.yml @@ -112,6 +112,11 @@ jobs: - schedule-due-storm-isolation - same-shard-family-fairness - actor-supervision-failpoint + - family-actor-partial-failure-isolation + - same-shard-family-failure-isolation + - family-actor-exhaustion-readiness + - family-actor-degradation-observability + - family-actor-inflight-concurrent-failure steps: - name: Check out Fitz Destroyer uses: actions/checkout@v7 diff --git a/README.md b/README.md index 63db432..7769399 100644 --- a/README.md +++ b/README.md @@ -163,9 +163,9 @@ scales raise the due set to at least 2,000 and 5,000 definitions respectively. family-actor shard, continuously fills one family's Notice lane, and requires every sibling-family delivery canary to complete within the request timeout. -`same-shard-family-failure-isolation` pins authenticated families 1 and 9 to -the same one of eight family-actor shards, then panics family 1's Stream and RPC -actors. The failed family must reject, family 9 must keep progressing without +`same-shard-family-failure-isolation` pins authenticated families 1 and 5 to +the same one of four family-actor shards, then panics family 1's Stream and RPC +actors. The failed family must reject, family 5 must keep progressing without cross-family delivery, and broker readiness must remain healthy. `family-actor-exhaustion-readiness` fails the two provisioned Stream families @@ -321,14 +321,20 @@ fails the run; recovered session-cleanup retries are recorded in the artifacts. `domain-pressure` runs a short, continuously bombarding client fleet without injecting faults. Use `--domains` to isolate one domain or an interference pair. It requires every selected domain to make progress on every client in each -ten-second window and fails on definite operation errors. Queue operations with -an unknown durable outcome are accepted only when exact reconciliation proves -that every deterministic sequence resolved at most once. +ten-second window and fails on definite operation errors. A typed Queue 4005 +rejection is retried with bounded exponential backoff because the request was +not accepted; each retry remains explicit in stage evidence. Queue operations +with an unknown durable outcome are accepted only when exact reconciliation +proves that every deterministic sequence resolved at most once. KV and Schedule repeatedly update one logical resource per client, while Stream -appends to one resource per client. This keeps route and current-state cardinality -bounded so RSS diagnostics are not dominated by intentionally abandoned KV keys -or Schedule definitions; Stream history still grows according to its durable -append-only contract. +appends to one resource per client. Durable loops run at a steady maximum of one +cycle per second per client; dedicated overload scenarios own saturation tests. +This keeps route and current-state cardinality bounded so RSS diagnostics are not +dominated by intentionally abandoned KV keys or Schedule definitions; Stream +history still grows according to its durable append-only contract. Live-domain +loops use a shorter cadence, including pacing Notice publication so its +fire-and-forget loop cannot monopolize the shared family lane and starve the +response-bearing domain probes. `pressure-evidence.json` contains per-client/domain/stage totals, latency percentiles, normalized error samples, Queue reconciliation, broker snapshots, and diagnostic warnings. diff --git a/compose.destroyer.yml b/compose.destroyer.yml index 4f5cea0..21462e2 100644 --- a/compose.destroyer.yml +++ b/compose.destroyer.yml @@ -27,6 +27,7 @@ services: - "127.0.0.1::9100" fitz: + cpus: "${FITZ_CPU_LIMIT:-0}" image: "${FITZ_IMAGE:-ghcr.io/cntryl/fitz:latest}" restart: "no" stop_grace_period: 15s @@ -50,6 +51,7 @@ services: FITZ_METRICS_BIND_ADDR: 0.0.0.0 FITZ_METRICS_PORT: "9090" FITZ_DRAIN_GRACE_SECONDS: "1" + FITZ_DESTROYER_FAILPOINTS: "${FITZ_DESTROYER_FAILPOINTS:-disabled}" FITZ_STORAGE_MODE: cloud FITZ_STORAGE_PROVIDER: sqrzl-s3 FITZ_STORAGE_ENDPOINT: http://storage-proxy:9000 diff --git a/package-lock.json b/package-lock.json index df97ed0..1c90589 100644 --- a/package-lock.json +++ b/package-lock.json @@ -8,7 +8,7 @@ "name": "fitz-destroyer", "version": "0.1.0", "dependencies": { - "@cntryl/fitz": "0.0.21" + "@cntryl/fitz": "0.0.22" }, "devDependencies": { "@types/node": "^24.10.0", @@ -20,9 +20,9 @@ } }, "node_modules/@cntryl/fitz": { - "version": "0.0.21", - "resolved": "https://registry.npmjs.org/@cntryl/fitz/-/fitz-0.0.21.tgz", - "integrity": "sha512-CTBcJ3bRlXuZCY003Sxl6Ewe+OAaJafCthghEzRQzLwOrA8C/NSLlQ3i5MMlva2aSlAgrGVNnl9bTSZefEJKWg==", + "version": "0.0.22", + "resolved": "https://registry.npmjs.org/@cntryl/fitz/-/fitz-0.0.22.tgz", + "integrity": "sha512-z8uOuFJBvkbD3nMpG9nFHRU72r3ZDj/iUgcraYS/ZFCUYmoJVPOyYmc45no8iqZRjx/9cazq95pdieEF1lSm5w==", "license": "Apache-2.0", "dependencies": { "ws": "^8.21.3" diff --git a/package.json b/package.json index ce5c120..f8c228c 100644 --- a/package.json +++ b/package.json @@ -16,7 +16,7 @@ "check": "npm run typecheck && npm test" }, "dependencies": { - "@cntryl/fitz": "0.0.21" + "@cntryl/fitz": "0.0.22" }, "devDependencies": { "@types/node": "^24.10.0", diff --git a/src/family-shard-topology.ts b/src/family-shard-topology.ts new file mode 100644 index 0000000..062b6df --- /dev/null +++ b/src/family-shard-topology.ts @@ -0,0 +1,8 @@ +export const DESTROYER_FAMILY_ACTOR_SHARD_COUNT = 4; +export const DESTROYER_PRIMARY_FAMILY = 1; +export const DESTROYER_SAME_SHARD_FAMILY = DESTROYER_FAMILY_ACTOR_SHARD_COUNT + 1; + +export const DESTROYER_FAMILY_ACTOR_FAMILIES = Array.from( + { length: DESTROYER_SAME_SHARD_FAMILY }, + (_, index) => index + 1, +); diff --git a/src/orchestration/actor-supervision-failpoint.ts b/src/orchestration/actor-supervision-failpoint.ts index b001cdf..0f3abfb 100644 --- a/src/orchestration/actor-supervision-failpoint.ts +++ b/src/orchestration/actor-supervision-failpoint.ts @@ -24,6 +24,13 @@ export function assertActorSupervisionEvidence(record: Readonly { const startedAt = performance.now(); let domainsInjected = 0; @@ -97,7 +104,11 @@ export async function runActorSupervisionFailpointScenario(stack: ComposeStack, await stack.waitForAllClientDomains(activeFaultStartedAt, config.clientReplicas); const faultStartedAt = new Date(); await injectAndWaitForReadinessWithdrawal("all-domain", config); - const activeFaultErrors = await stack.waitForAllClientErrors(faultStartedAt, config.clientReplicas); + const activeFaultErrors = await stack.waitForAllClientErrors( + faultStartedAt, + config.clientReplicas, + activeFaultObservationTimeoutMs(config.requestTimeoutMs, config.startupTimeoutMs), + ); await stack.stopBombardClientsAndCapture("actor-supervision-active-fault"); await stack.stopFitz(); await stack.restartFitz(); diff --git a/src/orchestration/compose-jobs.ts b/src/orchestration/compose-jobs.ts index 7efcb32..ecc78f2 100644 --- a/src/orchestration/compose-jobs.ts +++ b/src/orchestration/compose-jobs.ts @@ -12,7 +12,6 @@ type ComposeRunner = ( interface StorageFaultRecoveryOperations { stopFitz: () => Promise; - restartStorage: () => Promise; startFitz: () => Promise; } @@ -20,7 +19,6 @@ export async function executeStorageFaultRecovery( operations: StorageFaultRecoveryOperations, ): Promise { await operations.stopFitz(); - await operations.restartStorage(); await operations.startFitz(); } diff --git a/src/orchestration/compose.ts b/src/orchestration/compose.ts index 5b88bf0..69ab0d2 100644 --- a/src/orchestration/compose.ts +++ b/src/orchestration/compose.ts @@ -1,6 +1,12 @@ import { access } from "node:fs/promises"; import { join } from "node:path"; import type { RunConfig } from "../config.js"; +import { + DESTROYER_FAMILY_ACTOR_FAMILIES, + DESTROYER_FAMILY_ACTOR_SHARD_COUNT, + DESTROYER_PRIMARY_FAMILY, + DESTROYER_SAME_SHARD_FAMILY, +} from "../family-shard-topology.js"; import { type Domain, type WorkloadShape } from "../workloads/model.js"; import { Artifacts } from "./artifacts.js"; import { runCommand, type CommandResult } from "./command.js"; @@ -78,6 +84,8 @@ export class ComposeStack { #jobSequence = 0; #roleSequence = 0; #metricsUrl: string | undefined; + #pressureRssBytes: number | undefined; + #pressureRssSampledAt = 0; constructor( config: RunConfig, @@ -121,10 +129,7 @@ export class ComposeStack { FITZ_ASSUME_EXTERNAL_TLS: "true", FITZ_JWT_HMAC_SECRET: "fitz-destroyer-local-auth-only", FITZ_JWT_AUDIENCES: "fitz-destroyer", - FITZ_ROUTE_FAMILIES: config.scenario === "same-shard-family-fairness" || config.scenario === "same-shard-family-failure-isolation" || config.scenario === "family-actor-inflight-concurrent-failure" ? "1,2,3,4,5,6,7,8,9" : "1,2", - FITZ_ROUTE_FAMILY_MAP: config.scenario === "family-actor-inflight-concurrent-failure" ? "identity-a=1,identity-b=2,identity-c=9" : config.scenario === "same-shard-family-fairness" || config.scenario === "same-shard-family-failure-isolation" ? "identity-a=1,identity-b=9" : "identity-a=1,identity-b=2", - FITZ_ROUTE_FAMILY_CLAIM: "tid", - ...(config.scenario === "same-shard-family-fairness" || config.scenario === "same-shard-family-failure-isolation" || config.scenario === "family-actor-inflight-concurrent-failure" ? { FITZ_CPU_LIMIT: "8" } : {}), + ...authenticatedRouteFamilyEnvironment(config.scenario), } : {}), }; @@ -366,20 +371,12 @@ export class ComposeStack { if (containers.length !== 1 || container === undefined) { throw new Error(`Expected one running Fitz container for pressure snapshot, found ${containers.length}`); } - const [queue, rpc, prometheus, stats] = await Promise.all([ + const [queue, rpc, prometheus, rssBytes] = await Promise.all([ this.fetchJson("/api/v1/all/queue/stats"), this.fetchJson("/api/v1/all/rpc/stats"), this.fetchTextAt(`${metricsUrl}/metrics`, "Prometheus /metrics"), - runCommand( - "docker", - ["stats", "--no-stream", "--format", "{{json .}}", container], - { cwd: this.#config.rootDir }, - ), + this.pressureRss(container), ]); - const record = JSON.parse(stats.stdout.trim()) as { MemUsage?: unknown }; - if (typeof record.MemUsage !== "string") { - throw new Error(`Docker stats omitted Fitz MemUsage: ${stats.stdout.trim()}`); - } return { timestamp: new Date().toISOString(), queue, @@ -418,10 +415,29 @@ export class ComposeStack { "fitz_router_high_lane_backpressure_total", ), }, - rssBytes: parseDockerMemoryUsage(record.MemUsage), + rssBytes, }; } + private async pressureRss(container: string): Promise { + const now = Date.now(); + if (this.#pressureRssBytes !== undefined && now - this.#pressureRssSampledAt < 10_000) { + return this.#pressureRssBytes; + } + const stats = await runCommand( + "docker", + ["stats", "--no-stream", "--format", "{{json .}}", container], + { cwd: this.#config.rootDir }, + ); + const record = JSON.parse(stats.stdout.trim()) as { MemUsage?: unknown }; + if (typeof record.MemUsage !== "string") { + throw new Error(`Docker stats omitted Fitz MemUsage: ${stats.stdout.trim()}`); + } + this.#pressureRssBytes = parseDockerMemoryUsage(record.MemUsage); + this.#pressureRssSampledAt = now; + return this.#pressureRssBytes; + } + async prometheusMetricValue(name: string): Promise { const metricsUrl = await this.metricsUrl(); const prometheus = await this.fetchTextAt(`${metricsUrl}/metrics`, "Prometheus /metrics"); @@ -488,10 +504,6 @@ export class ComposeStack { "fitz-after-storage-exhaustion", true, ), - restartStorage: async () => { - await this.killAndRemoveService("sqrzl", "sqrzl-after-exhaustion", true); - await this.compose(["up", "-d", "--no-deps", "sqrzl"], { stream: true }); - }, startFitz: async () => { await this.compose(["up", "-d", "--no-deps", "--no-build", "fitz"], { stream: true }); await this.waitReady(); @@ -646,8 +658,12 @@ export class ComposeStack { throw new Error(`Timed out waiting for fresh client success in every domain: ${lastStatus}`); } - async waitForAllClientErrors(since: Date, replicas: number): Promise { - const deadline = Date.now() + this.#config.requestTimeoutMs; + async waitForAllClientErrors( + since: Date, + replicas: number, + timeoutMs = this.#config.requestTimeoutMs, + ): Promise { + const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { const containers = await this.serviceContainers("client", true); if (containers.length === replicas) { @@ -1069,3 +1085,29 @@ export class ComposeStack { ); } } + +function authenticatedRouteFamilyEnvironment( + scenario: RunConfig["scenario"], +): Readonly> { + if (scenario === "family-actor-inflight-concurrent-failure") { + return { + FITZ_ROUTE_FAMILIES: DESTROYER_FAMILY_ACTOR_FAMILIES.join(","), + FITZ_ROUTE_FAMILY_MAP: `identity-a=${DESTROYER_PRIMARY_FAMILY},identity-b=2,identity-c=${DESTROYER_SAME_SHARD_FAMILY}`, + FITZ_ROUTE_FAMILY_CLAIM: "tid", + FITZ_CPU_LIMIT: String(DESTROYER_FAMILY_ACTOR_SHARD_COUNT), + }; + } + if (scenario === "same-shard-family-fairness" || scenario === "same-shard-family-failure-isolation") { + return { + FITZ_ROUTE_FAMILIES: DESTROYER_FAMILY_ACTOR_FAMILIES.join(","), + FITZ_ROUTE_FAMILY_MAP: `identity-a=${DESTROYER_PRIMARY_FAMILY},identity-b=${DESTROYER_SAME_SHARD_FAMILY}`, + FITZ_ROUTE_FAMILY_CLAIM: "tid", + FITZ_CPU_LIMIT: String(DESTROYER_FAMILY_ACTOR_SHARD_COUNT), + }; + } + return { + FITZ_ROUTE_FAMILIES: "1,2", + FITZ_ROUTE_FAMILY_MAP: "identity-a=1,identity-b=2", + FITZ_ROUTE_FAMILY_CLAIM: "tid", + }; +} diff --git a/src/orchestration/pressure.ts b/src/orchestration/pressure.ts index 76f9542..2880fa2 100644 --- a/src/orchestration/pressure.ts +++ b/src/orchestration/pressure.ts @@ -125,11 +125,16 @@ export async function runPressureScenario( } try { + const verificationCompletedAt = pressureVerificationCompletedAtMs( + pressureStartedAt.getTime(), + pressureCompletedAt.getTime(), + requestedDurationMs, + ); assertProgressWindows( clientLogs, config.bombardDomains, pressureStartedAt.getTime(), - pressureCompletedAt.getTime(), + verificationCompletedAt, ); } catch (error) { assertionFailures.push(errorMessage(error)); @@ -284,6 +289,14 @@ export function assertProgressWindows( } } +export function pressureVerificationCompletedAtMs( + startedAtMs: number, + completedAtMs: number, + requestedDurationMs: number, +): number { + return Math.min(completedAtMs, startedAtMs + requestedDurationMs); +} + async function samplePressure( stack: ComposeStack, artifacts: Artifacts, @@ -313,10 +326,10 @@ export function pressureUnexpectedErrors( return clients.flatMap((client) => domains.flatMap((domain) => { const evidence = client.domains[domain]; - // Queue enqueue/complete timeouts have an explicitly unknown durable - // outcome. They are correctness failures only when exact reconciliation - // fails; definite stage failures still fail the pressure run here. - const count = domain === "queue" + // Queue outcomes are reconciled after the run. Stream outcomes are + // reconciled inline against the latest committed offset. Definite stage + // failures still fail the pressure run here. + const count = domain === "queue" || domain === "stream" ? Object.values(evidence?.stages ?? {}).reduce( (total, stage) => total + stage.failed, 0, @@ -432,6 +445,9 @@ function parseStage(value: unknown, label: string): EvidenceStage { succeeded: numericValue(record.succeeded, `${label}.succeeded`), failed: numericValue(record.failed, `${label}.failed`), ambiguous: numericValue(record.ambiguous, `${label}.ambiguous`), + retryableBackpressure: record.retryableBackpressure === undefined + ? 0 + : numericValue(record.retryableBackpressure, `${label}.retryableBackpressure`), expectedShutdownCancellations: { failed: numericValue(shutdownCancellations.failed, `${label} shutdown cancellation failures`), ambiguous: numericValue(shutdownCancellations.ambiguous, `${label} ambiguous shutdown cancellations`), @@ -486,6 +502,7 @@ function emptyEvidenceStage(): EvidenceStage { succeeded: 0, failed: 0, ambiguous: 0, + retryableBackpressure: 0, expectedShutdownCancellations: { failed: 0, ambiguous: 0 }, latencyHistogram, latency: latencySummary(latencyHistogram), @@ -499,6 +516,7 @@ function mergeStage(target: EvidenceStage, source: EvidenceStage): void { target.succeeded += source.succeeded; target.failed += source.failed; target.ambiguous += source.ambiguous; + target.retryableBackpressure += source.retryableBackpressure; target.expectedShutdownCancellations.failed += source.expectedShutdownCancellations.failed; target.expectedShutdownCancellations.ambiguous += source.expectedShutdownCancellations.ambiguous; target.latencyHistogram = mergeLatencyHistograms([ diff --git a/src/orchestration/storage-faults.ts b/src/orchestration/storage-faults.ts index 7871bc2..010f715 100644 --- a/src/orchestration/storage-faults.ts +++ b/src/orchestration/storage-faults.ts @@ -132,21 +132,18 @@ export function classifyStorageFailure(error: string): StorageFailureLayer { return "admission"; } -async function runStorageFaultIteration( +export async function runStorageFaultIteration( stack: ComposeStack, shape: WorkloadShape, environment: Readonly>, iteration: StorageFaultIteration, + delay: (milliseconds: number) => Promise = sleep, ): Promise> { await stack.setFaultProxy("storage-proxy", { mode: "healthy" }); const waitForGate = iteration.fault !== "bounded-latency" && iteration.fault !== "provider-recovery"; if (iteration.fault === "bounded-latency") { await stack.setFaultProxy("storage-proxy", { mode: "latency", latencyMs: 250 }); - } else if (iteration.fault === "connection-reset") { - await stack.setFaultProxy("storage-proxy", { mode: "reset" }); - } else if (iteration.fault === "provider-partition" || iteration.fault === "fitz-crash-inflight") { - await stack.setFaultProxy("storage-proxy", { mode: "partition" }); } const writer = await stack.startRoleContainers("durability-writer", 1, shape, { @@ -157,15 +154,20 @@ async function runStorageFaultIteration( }); if (waitForGate) { await stack.waitForRoleEvent(writer, "live_producer_ready"); + if (iteration.fault === "connection-reset") { + await stack.setFaultProxy("storage-proxy", { mode: "reset" }); + } else if (iteration.fault === "provider-partition" || iteration.fault === "fitz-crash-inflight") { + await stack.setFaultProxy("storage-proxy", { mode: "partition" }); + } await stack.signalRoleContainers(writer, "SIGUSR1"); await stack.waitForRoleEvent(writer, "durability_operations_dispatched"); } if (iteration.fault === "connection-reset") { - await sleep(250); + await delay(250); await stack.setFaultProxy("storage-proxy", { mode: "healthy" }); } else if (iteration.fault === "provider-partition") { - await sleep(5_000); + await delay(5_000); await stack.setFaultProxy("storage-proxy", { mode: "healthy" }); } else if (iteration.fault === "fitz-crash-inflight") { await stack.killFitz(); diff --git a/src/pressure.ts b/src/pressure.ts index 2afa398..423e183 100644 --- a/src/pressure.ts +++ b/src/pressure.ts @@ -26,6 +26,7 @@ export type StageMetrics = { succeeded: number; failed: number; ambiguous: number; + retryableBackpressure: number; expectedShutdownCancellations: { failed: number; ambiguous: number; @@ -105,6 +106,7 @@ export function createStageMetrics(): MutableStageMetrics { succeeded: 0, failed: 0, ambiguous: 0, + retryableBackpressure: 0, expectedShutdownCancellations: { failed: 0, ambiguous: 0 }, latency: { count: 0, diff --git a/src/worker.ts b/src/worker.ts index 9081445..6a78c18 100644 --- a/src/worker.ts +++ b/src/worker.ts @@ -45,10 +45,14 @@ import { } from "./pressure.js"; import { decodePressureQueueSequence, + isPressureStreamCleanupPending, pressureKvWrite, + pressureLoopDelayMs, pressureScheduleRoute, pressureStreamWrite, pressureValue, + replaceAndReconcilePressureStreamClient, + retryPressureQueueBackpressure, runPressureQueueReconciler, } from "./workloads/pressure.js"; import { @@ -720,8 +724,12 @@ async function bombard(client: Client): Promise { ambiguousCompletions: [] as number[], }; const lastErrors: Partial> = {}; + let nextStreamOffset = 0n; const payload = (domain: Domain, counter: number): Uint8Array => pressureValue(namespace, worker, domain, counter); + let streamClient = selectedDomains.includes("stream") + ? await connectedPressureClient() + : undefined; const rpcRoute = `rpc://destroyer/${namespace}/${worker}`; const rpcWorker = selectedDomains.includes("rpc") @@ -748,17 +756,45 @@ async function bombard(client: Client): Promise { const operations: Record Promise> = { queue: async (i, signal) => { const route = `queue://destroyer/${namespace}/${worker}`; + const enqueueMetrics = stageMetrics(stages, "queue", "enqueue"); try { await observeStage(stages, "queue", "enqueue", () => - client.queue.enqueue(route, { body: payload("queue", i), signal }), true); + retryPressureQueueBackpressure( + () => client.queue.enqueue(route, { body: payload("queue", i), signal }), + (attempt, delayMs, error) => { + enqueueMetrics.retryableBackpressure += 1; + log("pressure_backpressure_retry", { + worker, + domain: "queue", + stage: "enqueue", + attempt, + delayMs, + error: errorMessage(error), + }); + }, + ), true); queueOutcome.acknowledged.push(i); } catch (error) { if (isAmbiguousDurableError(error)) queueOutcome.ambiguousEnqueues.push(i); else queueOutcome.failedEnqueues.push(i); throw error; } + const reserveMetrics = stageMetrics(stages, "queue", "reserve"); const items = await observeStage(stages, "queue", "reserve", () => - client.queue.reserve(route, { leaseSeconds: 2, batchSize: 1, signal }), true); + retryPressureQueueBackpressure( + () => client.queue.reserve(route, { leaseSeconds: 2, batchSize: 1, signal }), + (attempt, delayMs, error) => { + reserveMetrics.retryableBackpressure += 1; + log("pressure_backpressure_retry", { + worker, + domain: "queue", + stage: "reserve", + attempt, + delayMs, + error: errorMessage(error), + }); + }, + ), true); const item = items[0]; if (item !== undefined) { const sequence = decodePressureQueueSequence(namespace, worker, item.body); @@ -780,21 +816,50 @@ async function bombard(client: Client): Promise { await tx.put({ key, value, signal }); await tx.commit({ signal }); } catch (error) { - await tx.rollback({ signal }).catch(() => undefined); + await tx.rollback({ signal: pressureCleanupSignal() }).catch(() => undefined); throw error; } }, true); }, stream: async (i, signal) => { - const { route, expectedOffset } = pressureStreamWrite(namespace, worker, i); + if (streamClient === undefined) throw new Error("stream pressure client is unavailable"); + const { route, expectedOffset } = pressureStreamWrite(namespace, worker, nextStreamOffset); await observeStage(stages, "stream", "append", async () => { - const session = await client.stream.begin(route, { signal }); - try { - await session.append({ expectedOffset, body: payload("stream", i), signal }); - await session.commit({ mode: "Sync", signal }); - } catch (error) { - await session.rollback({ signal }).catch(() => undefined); - throw error; + while (true) { + const currentStreamClient = streamClient; + if (currentStreamClient === undefined) { + throw new Error("stream pressure client is unavailable"); + } + let session: Awaited> | undefined; + try { + session = await currentStreamClient.stream.begin(route, { signal }); + await session.append({ expectedOffset, body: payload("stream", i), signal }); + await session.commit({ mode: "Sync", signal }); + nextStreamOffset += 1n; + return; + } catch (error) { + await session?.rollback({ signal: pressureCleanupSignal() }).catch(() => undefined); + if (shutdown.signal.aborted) throw error; + try { + const recovery = await replaceAndReconcilePressureStreamClient( + currentStreamClient, + connectedPressureClient, + async (replacement) => + (await replacement.stream.peek(route, { + signal: pressureReconciliationSignal(), + }))?.offset, + ); + streamClient = recovery.client; + nextStreamOffset = recovery.nextOffset; + } catch { + throw pressureStreamReconciliationError(); + } + if (isPressureStreamCleanupPending(error)) { + await new Promise((resolve) => setTimeout(resolve, 50)); + continue; + } + throw error; + } } }, true); }, @@ -847,6 +912,7 @@ async function bombard(client: Client): Promise { (next) => (window = next), lastErrors, operations[domain], + pressureLoopDelayMs(domain), ), ); @@ -855,10 +921,31 @@ async function bombard(client: Client): Promise { } finally { clearInterval(progressTimer); if (rpcWorker !== undefined) await rpcWorker.unsubscribe().catch(() => undefined); + await streamClient?.close().catch(() => undefined); log("stopped", { worker, totals: counters, stages, queueOutcome }); } } +async function connectedPressureClient(): Promise { + const client = makeClient(true); + await client.connectWhenReady({ timeoutMs: Infinity, signal: shutdown.signal }); + return client; +} + +function pressureCleanupSignal(): AbortSignal { + return AbortSignal.timeout(Math.min(requestTimeoutMs, 2_000)); +} + +function pressureReconciliationSignal(): AbortSignal { + return AbortSignal.timeout(requestTimeoutMs); +} + +function pressureStreamReconciliationError(): Error { + const error = new Error("stream durable outcome could not be reconciled"); + error.name = "PressureStreamReconciliationError"; + return error; +} + async function observeStage( stages: PressureStages, domain: Domain, @@ -893,6 +980,7 @@ async function domainLoop( setWindow: (next: Counters) => void, lastErrors: Partial>, operation: (counter: number, signal: AbortSignal) => Promise, + delayMs: number, ): Promise { let counter = 0; while (!shutdown.signal.aborted) { @@ -910,7 +998,7 @@ async function domainLoop( } counter += 1; setWindow(getWindow()); - await new Promise((resolve) => setTimeout(resolve, 1)); + await new Promise((resolve) => setTimeout(resolve, delayMs)); } } diff --git a/src/workloads/pressure.ts b/src/workloads/pressure.ts index 310ebe2..7f4ff05 100644 --- a/src/workloads/pressure.ts +++ b/src/workloads/pressure.ts @@ -2,6 +2,15 @@ import type { Client } from "@cntryl/fitz"; import type { LiveLog } from "./live.js"; import type { Domain } from "./model.js"; +const STREAM_SESSION_ALREADY_ACTIVE = 2_002; +const QUEUE_FULL = 4_005; +const NOTICE_PRESSURE_DELAY_MS = 100; +const LIVE_PRESSURE_DELAY_MS = 250; +const DURABLE_PRESSURE_DELAY_MS = 1_000; +const RECONCILIATION_BATCH_SIZE = 32; +const QUEUE_BACKPRESSURE_INITIAL_DELAY_MS = 25; +const QUEUE_BACKPRESSURE_MAX_DELAY_MS = 250; + export function pressureValue( namespace: string, worker: string, @@ -27,18 +36,81 @@ export function pressureKvWrite( export function pressureStreamWrite( namespace: string, worker: string, - counter: number, + expectedOffset: bigint, ): { route: string; expectedOffset: bigint } { return { route: `stream://destroyer/${namespace}/${worker}-stream`, - expectedOffset: BigInt(counter), + expectedOffset, }; } +export function nextPressureStreamOffset(latestOffset: bigint | undefined): bigint { + return latestOffset === undefined ? 0n : latestOffset + 1n; +} + +export function isPressureStreamCleanupPending(error: unknown): boolean { + return ( + typeof error === "object" && + error !== null && + "domainCode" in error && + error.domainCode === STREAM_SESSION_ALREADY_ACTIVE + ); +} + +export async function replaceAndReconcilePressureStreamClient }>( + current: T, + connect: () => Promise, + peekOffset: (client: T) => Promise, +): Promise<{ client: T; nextOffset: bigint }> { + await current.close().catch(() => undefined); + const client = await connect(); + try { + return { client, nextOffset: nextPressureStreamOffset(await peekOffset(client)) }; + } catch (error) { + await client.close().catch(() => undefined); + throw error; + } +} + export function pressureScheduleRoute(namespace: string, worker: string): string { return `schedule://destroyer/${namespace}/${worker}/job`; } +export function pressureLoopDelayMs(domain: Domain): number { + if (domain === "notice") return NOTICE_PRESSURE_DELAY_MS; + if (domain === "lease" || domain === "rpc") return LIVE_PRESSURE_DELAY_MS; + return DURABLE_PRESSURE_DELAY_MS; +} + +export function isPressureQueueBackpressure(error: unknown): boolean { + return typeof error === "object" && + error !== null && + "domainCode" in error && + error.domainCode === QUEUE_FULL; +} + +export async function retryPressureQueueBackpressure( + operation: () => Promise, + onRetry: (attempt: number, delayMs: number, error: unknown) => void, + wait: (delayMs: number) => Promise = pressureDelay, +): Promise { + let retries = 0; + while (true) { + try { + return await operation(); + } catch (error) { + if (!isPressureQueueBackpressure(error)) throw error; + retries += 1; + const delayMs = Math.min( + QUEUE_BACKPRESSURE_INITIAL_DELAY_MS * (2 ** (retries - 1)), + QUEUE_BACKPRESSURE_MAX_DELAY_MS, + ); + onRetry(retries, delayMs, error); + await wait(delayMs); + } + } +} + export type PressureReconcileOptions = { namespace: string; workers: readonly string[]; @@ -60,7 +132,7 @@ export async function runPressureQueueReconciler( while (emptyPolls < 3) { const items = await client.queue.reserve(route, { leaseSeconds: 30, - batchSize: 1_024, + batchSize: RECONCILIATION_BATCH_SIZE, waitSeconds: 1, signal: operationSignal(options), }); @@ -104,3 +176,7 @@ export function decodePressureQueueSequence( function operationSignal(options: PressureReconcileOptions): AbortSignal { return AbortSignal.any([options.signal, AbortSignal.timeout(options.requestTimeoutMs)]); } + +function pressureDelay(milliseconds: number): Promise { + return new Promise((resolve) => setTimeout(resolve, milliseconds)); +} diff --git a/src/workloads/recovery.ts b/src/workloads/recovery.ts index fa735dc..515286d 100644 --- a/src/workloads/recovery.ts +++ b/src/workloads/recovery.ts @@ -18,10 +18,18 @@ export async function loadRecoveryWorkload( shape: WorkloadShape, ): Promise { await Promise.all([ - mapResources(shape.resources, (resource) => loadQueue(client, shape, resource)), - mapResources(shape.resources, (resource) => loadKv(client, shape, resource)), - mapResources(shape.resources, (resource) => loadStream(client, shape, resource)), - mapResources(shape.resources, (resource) => loadSchedules(client, shape, resource)), + runResourceOperationsSequentially(shape.resources, (resource) => + loadQueue(client, shape, resource), + ), + runResourceOperationsSequentially(shape.resources, (resource) => + loadKv(client, shape, resource), + ), + runResourceOperationsSequentially(shape.resources, (resource) => + loadStream(client, shape, resource), + ), + runResourceOperationsSequentially(shape.resources, (resource) => + loadSchedules(client, shape, resource), + ), ]); return totalDurableEntries(shape); } @@ -31,9 +39,15 @@ export async function verifyRecoveryWorkload( shape: WorkloadShape, ): Promise { await Promise.all([ - mapResources(shape.resources, (resource) => verifyQueue(client, shape, resource)), - mapResources(shape.resources, (resource) => verifyKv(client, shape, resource)), - mapResources(shape.resources, (resource) => verifyStream(client, shape, resource)), + runResourceOperationsSequentially(shape.resources, (resource) => + verifyQueue(client, shape, resource), + ), + runResourceOperationsSequentially(shape.resources, (resource) => + verifyKv(client, shape, resource), + ), + runResourceOperationsSequentially(shape.resources, (resource) => + verifyStream(client, shape, resource), + ), verifySchedules(client, shape), ]); return totalDurableEntries(shape); @@ -203,9 +217,15 @@ async function verifySchedules(client: Client, shape: WorkloadShape): Promise Promise, ): Promise { - await Promise.all(Array.from({ length: resources }, (_, resource) => operation(resource))); + // The Fitz wire has no correlation ID for these domain operations. Its + // contract makes concurrent requests of the same message type on one + // connection undefined, so preserve one in-flight operation per domain + // while the four durable domains still run concurrently above. + for (let resource = 0; resource < resources; resource += 1) { + await operation(resource); + } } diff --git a/src/workloads/same-shard-family-failure-isolation.ts b/src/workloads/same-shard-family-failure-isolation.ts index e492697..e55017d 100644 --- a/src/workloads/same-shard-family-failure-isolation.ts +++ b/src/workloads/same-shard-family-failure-isolation.ts @@ -1,5 +1,10 @@ import { createClient, type Client } from "@cntryl/fitz"; import { createDestroyerToken } from "../auth-token.js"; +import { + DESTROYER_FAMILY_ACTOR_SHARD_COUNT, + DESTROYER_PRIMARY_FAMILY, + DESTROYER_SAME_SHARD_FAMILY, +} from "../family-shard-topology.js"; import type { LiveCommonOptions, LiveLog } from "./live.js"; export type SameShardFamilyFailureOptions = LiveCommonOptions & { @@ -7,9 +12,9 @@ export type SameShardFamilyFailureOptions = LiveCommonOptions & { failpointUrl: string; }; -const SHARD_COUNT = 8; -const FAILED_FAMILY = 1; -const SIBLING_FAMILY = 9; +const SHARD_COUNT = DESTROYER_FAMILY_ACTOR_SHARD_COUNT; +const FAILED_FAMILY = DESTROYER_PRIMARY_FAMILY; +const SIBLING_FAMILY = DESTROYER_SAME_SHARD_FAMILY; const encoder = new TextEncoder(); const decoder = new TextDecoder(); diff --git a/tests/actor-supervision-failpoint.test.ts b/tests/actor-supervision-failpoint.test.ts index 95a72f9..dcebd2a 100644 --- a/tests/actor-supervision-failpoint.test.ts +++ b/tests/actor-supervision-failpoint.test.ts @@ -1,6 +1,10 @@ import assert from "node:assert/strict"; import test from "node:test"; -import { assertActorSupervisionEvidence, recoveryCanaryNamespace } from "../src/orchestration/actor-supervision-failpoint.js"; +import { + activeFaultObservationTimeoutMs, + assertActorSupervisionEvidence, + recoveryCanaryNamespace, +} from "../src/orchestration/actor-supervision-failpoint.js"; import { parseWindowErrors } from "../src/orchestration/compose-evidence.js"; test("should_require_fail_closed_actor_supervision_and_restart_recovery", () => { @@ -62,3 +66,16 @@ test("should_isolate_each_post_restart_canary_route_namespace", () => { assert.notEqual(schedule, stream); assert.notEqual(stream, rpc); }); + +test("should_allow_active clients to finish an in-flight request after a correlated fault", () => { + // Arrange + const requestTimeoutMs = 10_000; + + // Act + const hostedBudget = activeFaultObservationTimeoutMs(requestTimeoutMs, 180_000); + const boundedBudget = activeFaultObservationTimeoutMs(requestTimeoutMs, 20_000); + + // Assert + assert.equal(hostedBudget, 30_000); + assert.equal(boundedBudget, 20_000); +}); diff --git a/tests/automation-contract.test.ts b/tests/automation-contract.test.ts new file mode 100644 index 0000000..c012fcc --- /dev/null +++ b/tests/automation-contract.test.ts @@ -0,0 +1,28 @@ +import assert from "node:assert/strict"; +import { readFile } from "node:fs/promises"; +import test from "node:test"; + +const ROOT = new URL("../", import.meta.url); + +test("should_forward_actor_controls_through_the_published_image_compose_file", async () => { + const compose = await readFile(new URL("compose.destroyer.yml", ROOT), "utf8"); + + assert.match(compose, /cpus: "\$\{FITZ_CPU_LIMIT:-0\}"/u); + assert.match( + compose, + /FITZ_DESTROYER_FAILPOINTS: "\$\{FITZ_DESTROYER_FAILPOINTS:-disabled\}"/u, + ); +}); + +test("should_run_every_family_actor_scenario_in_the_hosted_matrix", async () => { + const workflow = await readFile(new URL(".github/workflows/destroyer.yml", ROOT), "utf8"); + const scenarios = [ + "family-actor-partial-failure-isolation", + "same-shard-family-failure-isolation", + "family-actor-exhaustion-readiness", + "family-actor-degradation-observability", + "family-actor-inflight-concurrent-failure", + ]; + + for (const scenario of scenarios) assert.match(workflow, new RegExp(`- ${scenario}\\n`, "u")); +}); diff --git a/tests/cache-and-disk-exhaustion.test.ts b/tests/cache-and-disk-exhaustion.test.ts index 552b3ee..638c73e 100644 --- a/tests/cache-and-disk-exhaustion.test.ts +++ b/tests/cache-and-disk-exhaustion.test.ts @@ -23,14 +23,13 @@ test("should_isolate_each_exhaustion_phase_with_its_own_durable_baseline", () => }); }); -test("should_stop_fitz_before_restarting_exhausted_storage", async () => { +test("should_keep_tmpfs_storage_mounted_while_restarting_fitz", async () => { const calls: string[] = []; await executeStorageFaultRecovery({ stopFitz: async () => { calls.push("stop-fitz"); }, - restartStorage: async () => { calls.push("restart-storage"); }, startFitz: async () => { calls.push("start-fitz"); }, }); - assert.deepEqual(calls, ["stop-fitz", "restart-storage", "start-fitz"]); + assert.deepEqual(calls, ["stop-fitz", "start-fitz"]); }); diff --git a/tests/family-shard-topology.test.ts b/tests/family-shard-topology.test.ts new file mode 100644 index 0000000..d702c0d --- /dev/null +++ b/tests/family-shard-topology.test.ts @@ -0,0 +1,17 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { + DESTROYER_FAMILY_ACTOR_FAMILIES, + DESTROYER_FAMILY_ACTOR_SHARD_COUNT, + DESTROYER_PRIMARY_FAMILY, + DESTROYER_SAME_SHARD_FAMILY, +} from "../src/family-shard-topology.js"; + +test("should_fit_the_same_shard_topology_on_hosted_runners", () => { + assert.deepEqual(DESTROYER_FAMILY_ACTOR_FAMILIES, [1, 2, 3, 4, 5]); + assert.equal(DESTROYER_FAMILY_ACTOR_SHARD_COUNT, 4); + assert.equal( + (DESTROYER_PRIMARY_FAMILY - 1) % DESTROYER_FAMILY_ACTOR_SHARD_COUNT, + (DESTROYER_SAME_SHARD_FAMILY - 1) % DESTROYER_FAMILY_ACTOR_SHARD_COUNT, + ); +}); diff --git a/tests/operational-guidance.test.ts b/tests/operational-guidance.test.ts index e6faad9..5ccff00 100644 --- a/tests/operational-guidance.test.ts +++ b/tests/operational-guidance.test.ts @@ -374,7 +374,7 @@ function fixtureEvents(scenario: ConcreteScenario): object[] { if (scenario === "same-shard-family-fairness") return [{ event: "same_shard_family_fairness_complete", noisyCompleted: 2_048, canariesAttempted: 32, canariesCompleted: 32, canaryErrors: 0, longestCanaryMs: 10, requestTimeoutMs: 1_000, elapsedMs: 1_000 }]; if (scenario === "actor-supervision-failpoint") return [{ event: "actor_supervision_failpoint_complete", domainsInjected: 7, correlatedDomainsInjected: 7, activeFaultClients: 4, activeFaultErrors: 4, readinessWithdrawals: 8, restartsRecovered: 8, canaryDeliveries: 4, queueRecovered: 1, kvRecovered: 1, leaseRecovered: 1, scheduleRecovered: 1, streamRecovered: 1, rpcRecovered: 1, correlatedRecoveryOperations: 7, elapsedMs: 1_000 }]; if (scenario === "family-actor-partial-failure-isolation") return [{ event: "family_actor_partial_failure_isolation_complete", targetedFamilies: 2, failedFamilyRejections: 2, siblingOperations: 2, readinessChecks: 2, crossFamilyDeliveries: 0, elapsedMs: 1_000 }]; - if (scenario === "same-shard-family-failure-isolation") return [{ event: "same_shard_family_failure_isolation_complete", shardCount: 8, failedFamily: 1, siblingFamily: 9, failedFamilyRejections: 2, siblingOperations: 2, readinessChecks: 2, crossFamilyDeliveries: 0, elapsedMs: 1_000 }]; + if (scenario === "same-shard-family-failure-isolation") return [{ event: "same_shard_family_failure_isolation_complete", shardCount: 4, failedFamily: 1, siblingFamily: 5, failedFamilyRejections: 2, siblingOperations: 2, readinessChecks: 2, crossFamilyDeliveries: 0, elapsedMs: 1_000 }]; if (scenario === "family-actor-exhaustion-readiness") return [{ event: "family_actor_exhaustion_readiness_complete", domainsExhausted: 2, partialReadinessChecks: 2, totalReadinessWithdrawals: 2, failedFamilyRejections: 4, restartRecoveries: 2, elapsedMs: 1_000 }]; if (scenario === "family-actor-degradation-observability") return [{ event: "family_actor_degradation_observability_complete", domainsObserved: 2, metricIncrements: 2, duplicateMetricIncrements: 0, siblingCanaries: 2, readinessChecks: 2, restartMetricResets: 2, elapsedMs: 1_000 }]; if (scenario === "family-actor-inflight-concurrent-failure") return [{ event: "family_actor_inflight_concurrent_failure_complete", cycles: 3, concurrentFailures: 12, streamInflightRejections: 6, rpcInflightTerminations: 6, siblingOperations: 6, readinessChecks: 6, metricIncrements: 12, restartRecoveries: 3, crossFamilyDeliveries: 0, elapsedMs: 1_000 }]; diff --git a/tests/pressure.test.ts b/tests/pressure.test.ts index c609697..30dfbf5 100644 --- a/tests/pressure.test.ts +++ b/tests/pressure.test.ts @@ -17,19 +17,26 @@ import { prometheusMetric, } from "../src/orchestration/compose-evidence.js"; import { + isPressureStreamCleanupPending, + isPressureQueueBackpressure, + nextPressureStreamOffset, pressureKvWrite, + pressureLoopDelayMs, pressureScheduleRoute, pressureStreamWrite, + replaceAndReconcilePressureStreamClient, + retryPressureQueueBackpressure, } from "../src/workloads/pressure.js"; import { analyzePressureLogs, assertProgressWindows, + pressureVerificationCompletedAtMs, pressureUnexpectedErrors, } from "../src/orchestration/pressure.js"; test("should_keep_sustained_pressure_durable_identities_bounded", () => { - const firstStream = pressureStreamWrite("run", "worker", 0); - const laterStream = pressureStreamWrite("run", "worker", 42); + const firstStream = pressureStreamWrite("run", "worker", 0n); + const laterStream = pressureStreamWrite("run", "worker", 42n); const firstKv = pressureKvWrite("run", "worker", 0); const laterKv = pressureKvWrite("run", "worker", 42); @@ -41,6 +48,86 @@ test("should_keep_sustained_pressure_durable_identities_bounded", () => { assert.equal(pressureScheduleRoute("run", "worker"), pressureScheduleRoute("run", "worker")); }); +test("should_apply_steady_pressure_cadence_by_domain_semantics", () => { + assert.equal(pressureLoopDelayMs("notice"), 100); + for (const domain of ["lease", "rpc"] as const) { + assert.equal(pressureLoopDelayMs(domain), 250); + } + for (const domain of ["queue", "kv", "stream", "schedule"] as const) { + assert.equal(pressureLoopDelayMs(domain), 1_000); + } +}); + +test("should_retry_only_typed_queue_backpressure_with_bounded_exponential_delay", async () => { + const delays: number[] = []; + let attempts = 0; + const result = await retryPressureQueueBackpressure( + async () => { + attempts += 1; + if (attempts < 6) throw { domainCode: 4_005 }; + return "accepted"; + }, + (_attempt, delayMs) => delays.push(delayMs), + async () => undefined, + ); + + assert.equal(result, "accepted"); + assert.equal(attempts, 6); + assert.deepEqual(delays, [25, 50, 100, 200, 250]); + assert.equal(isPressureQueueBackpressure({ domainCode: 4_005 }), true); + assert.equal(isPressureQueueBackpressure({ domainCode: 5_007 }), false); + await assert.rejects( + retryPressureQueueBackpressure( + async () => { + throw new Error("definite Queue failure"); + }, + () => assert.fail("non-backpressure error must not retry"), + async () => undefined, + ), + /definite Queue failure/u, + ); +}); + +test("should_resume_stream_pressure_after_the_latest_committed_offset", () => { + assert.equal(nextPressureStreamOffset(undefined), 0n); + assert.equal(nextPressureStreamOffset(41n), 42n); +}); + +test("should_retry_stream_pressure_while_the_previous_session_is_cleaning_up", () => { + assert.equal(isPressureStreamCleanupPending({ domainCode: 2_002 }), true); + assert.equal(isPressureStreamCleanupPending({ domainCode: 2_003 }), false); +}); + +test("should_replace_the_stream_connection_before_reconciling_its_offset", async () => { + const calls: string[] = []; + const current = { + close: async () => { + calls.push("close-current"); + }, + }; + const replacement = { + close: async () => { + calls.push("close-replacement"); + }, + }; + + const result = await replaceAndReconcilePressureStreamClient( + current, + async () => { + calls.push("connect-replacement"); + return replacement; + }, + async () => { + calls.push("peek-replacement"); + return 41n; + }, + ); + + assert.equal(result.client, replacement); + assert.equal(result.nextOffset, 42n); + assert.deepEqual(calls, ["close-current", "connect-replacement", "peek-replacement"]); +}); + test("should_summarize_latency_percentiles_from_bounded_histograms", () => { const metrics = createStageMetrics(); for (const value of [1, 2, 3, 10, 100, 1_000, 6_000, 70_000]) { @@ -208,6 +295,33 @@ test("should_defer_ambiguous_queue_outcomes_to_exact_reconciliation", () => { assert.deepEqual(pressureUnexpectedErrors(clients, ["queue"]), ["worker-a/queue=1"]); }); +test("should_accept_only_reconciled_ambiguous_stream_outcomes", () => { + const metrics = createStageMetrics(); + metrics.ambiguous = 1; + metrics.errorClasses.timeout = 1; + const clients = [{ + container: "container-a", + worker: "worker-a", + domains: { + stream: { + succeeded: 4, + failed: 1, + stages: { + append: { + ...metrics, + latencyHistogram: metrics.latency, + latency: latencySummary(metrics.latency), + }, + }, + }, + }, + }]; + + assert.deepEqual(pressureUnexpectedErrors(clients, ["stream"]), []); + clients[0]!.domains.stream.stages.append.failed = 1; + assert.deepEqual(pressureUnexpectedErrors(clients, ["stream"]), ["worker-a/stream=1"]); +}); + test("should_classify_expected_worker_shutdown_signals_as_cancelled", () => { assert.equal(normalizeErrorClass(new Error("received SIGTERM")), "cancelled"); assert.equal(normalizeErrorClass(new Error("received SIGINT")), "cancelled"); @@ -254,6 +368,25 @@ test("should_not_require_progress_in_a_trailing_partial_window", () => { ); }); +test("should_not_extend_progress_windows_when_broker_sampling_overruns", () => { + const records = [ + progress("2026-08-25T12:00:05.000Z", 1, 1), + progress("2026-08-25T12:00:15.000Z", 1, 1), + ].join("\n"); + const start = Date.parse("2026-08-25T12:00:00.000Z"); + const requestedDurationMs = 20_000; + const samplingCompletedAt = Date.parse("2026-08-25T12:00:50.000Z"); + + assert.doesNotThrow(() => + assertProgressWindows( + new Map([["client", records]]), + ["queue", "rpc"], + start, + pressureVerificationCompletedAtMs(start, samplingCompletedAt, requestedDurationMs), + ), + ); +}); + test("should_parse_binary_docker_memory_units", () => { assert.equal(parseDockerMemoryUsage("128MiB / 1GiB"), 128 * 1_024 * 1_024); assert.equal(parseDockerMemoryUsage("1.5GiB / 8GiB"), 1.5 * 1_024 ** 3); diff --git a/tests/recovery.test.ts b/tests/recovery.test.ts new file mode 100644 index 0000000..01d4950 --- /dev/null +++ b/tests/recovery.test.ts @@ -0,0 +1,20 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { runResourceOperationsSequentially } from "../src/workloads/recovery.js"; + +test("should_keep_same_domain_recovery_operations_sequential", async () => { + let active = 0; + let peak = 0; + const completed: number[] = []; + + await runResourceOperationsSequentially(4, async (resource) => { + active += 1; + peak = Math.max(peak, active); + await new Promise((resolve) => setImmediate(resolve)); + completed.push(resource); + active -= 1; + }); + + assert.equal(peak, 1); + assert.deepEqual(completed, [0, 1, 2, 3]); +}); diff --git a/tests/same-shard-family-failure-isolation.test.ts b/tests/same-shard-family-failure-isolation.test.ts index 5bb6f14..55468cf 100644 --- a/tests/same-shard-family-failure-isolation.test.ts +++ b/tests/same-shard-family-failure-isolation.test.ts @@ -1,14 +1,19 @@ import assert from "node:assert/strict"; import { describe, it } from "node:test"; +import { + DESTROYER_FAMILY_ACTOR_SHARD_COUNT, + DESTROYER_PRIMARY_FAMILY, + DESTROYER_SAME_SHARD_FAMILY, +} from "../src/family-shard-topology.js"; import { assertSameShardFamilyFailureEvidence } from "../src/workloads/same-shard-family-failure-isolation.js"; describe("same-shard family failure isolation evidence", () => { it("should_accept_isolated_stream_and_rpc_failures_on_one_shared_shard", () => { // Arrange const evidence = { - shardCount: 8, - failedFamily: 1, - siblingFamily: 9, + shardCount: DESTROYER_FAMILY_ACTOR_SHARD_COUNT, + failedFamily: DESTROYER_PRIMARY_FAMILY, + siblingFamily: DESTROYER_SAME_SHARD_FAMILY, failedFamilyRejections: 2, siblingOperations: 2, readinessChecks: 2, @@ -22,8 +27,8 @@ describe("same-shard family failure isolation evidence", () => { it("should_reject_a_sibling_that_does_not_share_the_failed_family_shard", () => { // Arrange const evidence = { - shardCount: 8, - failedFamily: 1, + shardCount: DESTROYER_FAMILY_ACTOR_SHARD_COUNT, + failedFamily: DESTROYER_PRIMARY_FAMILY, siblingFamily: 2, failedFamilyRejections: 2, siblingOperations: 2, diff --git a/tests/storage-faults.test.ts b/tests/storage-faults.test.ts index 516c8b0..c52fc85 100644 --- a/tests/storage-faults.test.ts +++ b/tests/storage-faults.test.ts @@ -3,8 +3,11 @@ import test from "node:test"; import { STORAGE_FAULT_IDENTITIES, classifyStorageFailure, + runStorageFaultIteration, storageFaultPlan, } from "../src/orchestration/storage-faults.js"; +import type { ComposeStack } from "../src/orchestration/compose.js"; +import type { WorkloadShape } from "../src/workloads/model.js"; test("should_generate_seeded_storage_fault_ledger_iterations", () => { const plan = storageFaultPlan(7, 1234); @@ -20,3 +23,57 @@ test("should_classify_storage_fault_layers_from_observed_errors", () => { assert.equal(classifyStorageFailure("S3 provider persistence failed"), "persistence"); assert.equal(classifyStorageFailure("hydrate recovery failed"), "recovery"); }); + +test("should_connect_the_writer_before_cutting_its_storage_path", async () => { + const calls: string[] = []; + const stack = { + setFaultProxy: async (_service: string, fault: { mode: string }) => { + calls.push(`proxy-${fault.mode}`); + }, + startRoleContainers: async () => { + calls.push("writer-started"); + return []; + }, + waitForRoleEvent: async (_role: unknown, event: string) => { + calls.push(event); + }, + signalRoleContainers: async () => { + calls.push("writer-signalled"); + }, + finishRoleContainers: async () => { + calls.push("writer-finished"); + return new Map(); + }, + ensureReady: async () => { + calls.push("broker-ready"); + }, + } as unknown as ComposeStack; + const shape: WorkloadShape = { + namespace: "storage-fault", + resources: 1, + entriesPerResource: 1, + payloadBytes: 1, + seed: 42, + }; + + await runStorageFaultIteration( + stack, + shape, + {}, + { iteration: 1, sequence: 1, seed: 42, fault: "connection-reset" }, + async () => undefined, + ); + + assert.deepEqual(calls, [ + "proxy-healthy", + "writer-started", + "live_producer_ready", + "proxy-reset", + "writer-signalled", + "durability_operations_dispatched", + "proxy-healthy", + "writer-finished", + "proxy-healthy", + "broker-ready", + ]); +});