From 1210fd93e90c243fbd2fa1da0fb2ca3b63afa4cc Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sat, 26 Sep 2026 10:58:09 -0700 Subject: [PATCH] test: stabilize flaky fixtures (batch f037) (#159078) * test(gateway): stabilize abort cascade cleanup * test(doctor): stabilize config write fixtures * test(scripts): stabilize benchmark drain setup Load the canonical active-work owner before the timed case and verify the admitted root. Retain the real default drain call, pending check, completion check, and cleanup. * test(scripts): stabilize Telegram diagnostic socket paths Use the existing short tracked socket directory pattern to stay below the macOS Unix socket path limit and clean up startup failures. * test(sessions): stabilize deferred baseline capture * test(worker): stabilize background process completion * test(doctor): preserve nested fixture isolation after replay --- ...doctor-config-flow.include-refusal.test.ts | 2 + ...tor-config-flow.validation-refusal.test.ts | 1 + .../server-methods/chat.abort-cascade.test.ts | 20 +++++++- src/sessions/session-diff-baseline.test.ts | 48 ++++++++++++------- src/worker/worker.runtime.test.ts | 28 +++++++++-- test/scripts/bench-agent-concurrency.test.ts | 2 + test/scripts/mantis-telegram-failure.test.ts | 9 ++-- 7 files changed, 86 insertions(+), 24 deletions(-) diff --git a/src/commands/doctor-config-flow.include-refusal.test.ts b/src/commands/doctor-config-flow.include-refusal.test.ts index 487baa19c6cc..ad5a44ac34b0 100644 --- a/src/commands/doctor-config-flow.include-refusal.test.ts +++ b/src/commands/doctor-config-flow.include-refusal.test.ts @@ -268,6 +268,7 @@ describe("doctor --fix include write ownership", () => { const configPath = await writeOpenClawConfig(home, { agents: { entries: { main: { $include: "./config/main-parent.json5" } } }, gateway: { mode: "local" }, + plugins: { enabled: false }, }); const fragmentDir = path.join(path.dirname(configPath), "config"); await fs.mkdir(fragmentDir); @@ -525,6 +526,7 @@ describe("doctor --fix include write ownership", () => { agents: { list: [{ id: "ops" }] }, browser: { $include: "./browser.json" }, gateway: { mode: "local" }, + plugins: { enabled: false }, }); const includePath = path.join(path.dirname(configPath), "browser.json"); const includeRaw = JSON.stringify({ enabled: true, actionTimeoutMs: 5000 }); diff --git a/src/commands/doctor-config-flow.validation-refusal.test.ts b/src/commands/doctor-config-flow.validation-refusal.test.ts index fe1be14d63eb..88e2f7dc6697 100644 --- a/src/commands/doctor-config-flow.validation-refusal.test.ts +++ b/src/commands/doctor-config-flow.validation-refusal.test.ts @@ -29,6 +29,7 @@ describe("doctor --fix with a validation-blocked candidate", () => { const configPath = await writeOpenClawConfig(home, { gatway: { port: 12345 }, agents: { defaults: { heartbeat: { every: 5 } } }, + plugins: { enabled: false }, }); const rawBefore = await fs.readFile(configPath, "utf-8"); const ctx = await prepareDoctorContext(configPath); diff --git a/src/gateway/server-methods/chat.abort-cascade.test.ts b/src/gateway/server-methods/chat.abort-cascade.test.ts index f370fb55c9c1..c027b57f5526 100644 --- a/src/gateway/server-methods/chat.abort-cascade.test.ts +++ b/src/gateway/server-methods/chat.abort-cascade.test.ts @@ -1,5 +1,5 @@ /** Real handler and registry proof for session-wide descendant cancellation ownership. */ -import { afterEach, describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi, type MockInstance } from "vitest"; import { createDeferred } from "../../../test/helpers/promise.js"; import { registerSubagentRun } from "../../agents/subagents/registry/subagent-registry.js"; import { settleSubagentRegistryPersistenceWork } from "../../agents/subagents/registry/subagent-registry.persistence.test-support.js"; @@ -12,6 +12,7 @@ import { enqueueSwarmRun, releaseSwarmRun } from "../../agents/subagents/swarm/s import { testing as swarmSchedulerTesting } from "../../agents/subagents/swarm/swarm-scheduler.test-support.js"; import { resolveSessionStorePathCore } from "../../config/sessions/paths.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; +import * as gatewayWorkAdmission from "../../process/gateway-work-admission.js"; import { beginSessionWorkAdmission } from "../../sessions/session-lifecycle-admission.js"; import { createWorkerInferenceCancellationService } from "../worker-environments/inference-control.test-helpers.js"; import { handleChatAbortRequestWithLifecycle } from "./chat-abort-handler.js"; @@ -43,7 +44,24 @@ vi.mock("../../agents/subagents/registry/subagent-registry-state.js", async (imp })); describe("descendant cascade ownership", () => { + let rootWork: MockInstance< + typeof gatewayWorkAdmission.runWithGatewayIndependentRootWorkAdmission + >; + + beforeEach(() => { + rootWork = vi.spyOn(gatewayWorkAdmission, "runWithGatewayIndependentRootWorkAdmission"); + }); afterEach(async () => { + await vi.dynamicImportSettled(); + // Join real detached finalization before checking for leaked registry roots. + let settled = 0; + while (settled < rootWork.mock.results.length) { + const pending = rootWork.mock.results.slice(settled); + settled += pending.length; + await Promise.allSettled( + pending.flatMap((result) => (result.type === "return" ? [result.value] : [])), + ); + } await settleSubagentRegistryPersistenceWork(); resetSubagentRegistryForTests({ persist: false }); swarmSchedulerTesting.reset(); diff --git a/src/sessions/session-diff-baseline.test.ts b/src/sessions/session-diff-baseline.test.ts index 4123c0eb8d2f..a78ed2c76249 100644 --- a/src/sessions/session-diff-baseline.test.ts +++ b/src/sessions/session-diff-baseline.test.ts @@ -102,6 +102,16 @@ function expectWorkStartError( } } +function deferCapture() { + const started = createDeferredCore(); + const capture = createDeferredCore(); + captureMocks.capture.mockImplementation(() => { + started.resolve(); + return capture.promise; + }); + return { started: started.promise, resolve: capture.resolve }; +} + describe("ensureSessionDiffBaseline", () => { beforeEach(() => { captureMocks.capture.mockReset(); @@ -159,17 +169,23 @@ describe("ensureSessionDiffBaseline", () => { const sessionId = "concurrent-session"; const entry = makeEntry(sessionId); const target = await seedEntry({ entry }); - const capture = createDeferredCore(); - captureMocks.capture.mockReturnValue(capture.promise); + const capture = deferCapture(); const first = ensure(target, true); const second = ensure(target, true); - await vi.waitFor(() => expect(captureMocks.capture).toHaveBeenCalledTimes(1)); - capture.resolve(baseline(sessionId)); + try { + await capture.started; + expect(captureMocks.capture).toHaveBeenCalledTimes(1); + capture.resolve(baseline(sessionId)); - const [firstResult, secondResult] = await Promise.all([first, second]); - expect(firstResult.sessionDiffBaseline).toEqual(baseline(sessionId)); - expect(secondResult.sessionDiffBaseline).toEqual(baseline(sessionId)); + const [firstResult, secondResult] = await Promise.all([first, second]); + expect(captureMocks.capture).toHaveBeenCalledTimes(1); + expect(firstResult.sessionDiffBaseline).toEqual(baseline(sessionId)); + expect(secondResult.sessionDiffBaseline).toEqual(baseline(sessionId)); + } finally { + capture.resolve(baseline(sessionId)); + await Promise.allSettled([first, second]); + } }); it("rejects a stale cached baseline after the authoritative generation rotates", async () => { @@ -384,11 +400,11 @@ describe("ensureSessionDiffBaseline", () => { sessionDiffBaselineCapture: createSessionDiffBaselineCaptureClaim(), }); const target = await seedEntry({ entry }); - const capture = createDeferredCore(); - captureMocks.capture.mockReturnValue(capture.promise); + const capture = deferCapture(); const completion = ensure(target); const outcome = Promise.allSettled([completion]); - await vi.waitFor(() => expect(captureMocks.capture).toHaveBeenCalledOnce()); + await capture.started; + expect(captureMocks.capture).toHaveBeenCalledOnce(); await deleteSessionEntryLifecycle({ archiveTranscript: false, storePath: target.storePath, @@ -411,11 +427,11 @@ describe("ensureSessionDiffBaseline", () => { sessionDiffBaselineCapture: oldClaim, }); const target = await seedEntry({ entry }); - const capture = createDeferredCore(); - captureMocks.capture.mockReturnValue(capture.promise); + const capture = deferCapture(); const oldCompletions = [ensure(target), ensure(target)]; const outcomes = Promise.allSettled(oldCompletions); - await vi.waitFor(() => expect(captureMocks.capture).toHaveBeenCalledTimes(1)); + await capture.started; + expect(captureMocks.capture).toHaveBeenCalledTimes(1); const freshClaim = createSessionDiffBaselineCaptureClaim(); await replaceSessionEntry( @@ -442,11 +458,11 @@ describe("ensureSessionDiffBaseline", () => { sessionDiffBaselineCapture: claim, }); const target = await seedEntry({ entry }); - const capture = createDeferredCore(); - captureMocks.capture.mockReturnValue(capture.promise); + const capture = deferCapture(); const completion = ensure(target); const outcome = Promise.allSettled([completion]); - await vi.waitFor(() => expect(captureMocks.capture).toHaveBeenCalledOnce()); + await capture.started; + expect(captureMocks.capture).toHaveBeenCalledOnce(); await replaceSessionEntry( { sessionKey: target.sessionKey, storePath: target.storePath }, diff --git a/src/worker/worker.runtime.test.ts b/src/worker/worker.runtime.test.ts index f589d7d87ee0..8b2c2c5eb9c5 100644 --- a/src/worker/worker.runtime.test.ts +++ b/src/worker/worker.runtime.test.ts @@ -1770,13 +1770,27 @@ describe("worker runtime", () => { "text", ], ...(processState === "completed" - ? { backgroundCommand: `${JSON.stringify(process.execPath)} finish-on-file.cjs` } + ? { backgroundCommand: `${JSON.stringify(process.execPath)} finish-on-release.cjs` } : {}), }); + const releaseBackground = createDeferred(); + let completionServer: Server | undefined; if (processState === "completed") { + completionServer = createServer((_request, response) => { + void releaseBackground.promise.then(() => response.end("background-finished")); + }); + const listening = once(completionServer, "listening"); + completionServer.listen(0, "127.0.0.1"); + await listening; + const address = completionServer.address(); + if (!address || typeof address === "string") { + throw new Error("background completion server did not allocate a TCP port"); + } + // An explicit response also releases a child that starts after the turn finishes. + // Filesystem watch notifications can be lost while this shared machine is busy. await writeFile( - path.join(workspaceDir, "finish-on-file.cjs"), - "const fs = require('node:fs'); const finish = () => { if (fs.existsSync('finish-marker')) { process.stdout.write('background-finished'); watcher.close(); } }; const watcher = fs.watch('.', finish); finish();", + path.join(workspaceDir, "finish-on-release.cjs"), + `require('node:http').get('http://127.0.0.1:${address.port}', response => response.pipe(process.stdout));`, ); } const scopeKey = `worker:${SESSION_ID}`; @@ -1810,7 +1824,7 @@ describe("worker runtime", () => { const sessionId = running[0]!.id; expect(settled).not.toHaveBeenCalled(); if (processState === "completed") { - await writeFile(path.join(workspaceDir, "finish-marker"), "finish"); + releaseBackground.resolve(); await waitForExecScope(scopeKey); await waitForFast(() => expect( @@ -1871,12 +1885,18 @@ describe("worker runtime", () => { expect(results[1]?.retainWorker).toBe(false); } } finally { + releaseBackground.resolve(); input.end(); try { await command; } finally { supervisor.cancelScope(scopeKey, "manual-cancel"); await waitForExecScope(scopeKey); + if (completionServer) { + await new Promise((resolve, reject) => { + completionServer.close((error) => (error ? reject(error) : resolve())); + }); + } } } expect(listRunningSessions().filter((session) => session.scopeKey === scopeKey)).toHaveLength( diff --git a/test/scripts/bench-agent-concurrency.test.ts b/test/scripts/bench-agent-concurrency.test.ts index e7d328dc18f4..87a39345f52d 100644 --- a/test/scripts/bench-agent-concurrency.test.ts +++ b/test/scripts/bench-agent-concurrency.test.ts @@ -7,6 +7,7 @@ import { type WorkerResult, type WorkerScenario, } from "../../scripts/bench-agent-concurrency.ts"; +import { createGatewayActiveWorkSnapshot } from "../../src/infra/gateway-active-work.js"; import { resetGatewayWorkAdmission, runWithGatewayIndependentRootWorkAdmission, @@ -120,6 +121,7 @@ describe("agent concurrency benchmark", () => { const deferred = createDeferred(); const rootWork = runWithGatewayIndependentRootWorkAdmission(() => deferred.promise); try { + expect(createGatewayActiveWorkSnapshot().counts.rootRequests).toBe(1); const drain = workerTesting.drainSpawnSampleActiveWork(); await expect( Promise.race([ diff --git a/test/scripts/mantis-telegram-failure.test.ts b/test/scripts/mantis-telegram-failure.test.ts index 30f9fd2682f1..5ea6a222d315 100644 --- a/test/scripts/mantis-telegram-failure.test.ts +++ b/test/scripts/mantis-telegram-failure.test.ts @@ -5,13 +5,16 @@ import http from "node:http"; import os from "node:os"; import path from "node:path"; import JSZip from "jszip"; -import { describe, expect, it } from "vitest"; +import { afterEach, describe, expect, it } from "vitest"; import { createRequestReceipt, createTelegramFailureDiagnostic, requestIdentitySchema, } from "../../scripts/mantis/request-proof.ts"; import { startTelegramProofIngress } from "../../scripts/mantis/telegram-proof-ingress.mts"; +import { useAutoCleanupTempDirTracker } from "../helpers/temp-dir.js"; + +const tempDirs = useAutoCleanupTempDirTracker(afterEach); const identity = requestIdentitySchema.parse({ request_id: "a".repeat(64), @@ -134,7 +137,8 @@ describe("Telegram failure diagnostics", () => { ] as const)( "records %s through the real ingress without raw error or request data", async (expected) => { - const root = await mkdtemp(path.join(os.tmpdir(), "tg-diagnostic-ingress-")); + // Keep the socket below Darwin's 104-byte limit inside the worker temp namespace. + const root = tempDirs.make("sock-"); const socket = process.platform === "win32" ? `\\\\.\\pipe\\mantis-${randomUUID()}` @@ -206,7 +210,6 @@ describe("Telegram failure diagnostics", () => { expect(ingress.getDiagnostics()).toHaveLength(16); } finally { await ingress.close(); - await rm(root, { recursive: true, force: true }); } }, );