mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 17:53:39 +00:00
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
This commit is contained in:
parent
0da6157185
commit
1210fd93e9
7 changed files with 86 additions and 24 deletions
|
|
@ -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 });
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -102,6 +102,16 @@ function expectWorkStartError(
|
|||
}
|
||||
}
|
||||
|
||||
function deferCapture() {
|
||||
const started = createDeferredCore();
|
||||
const capture = createDeferredCore<SessionDiffBaseline>();
|
||||
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<SessionDiffBaseline>();
|
||||
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<SessionDiffBaseline>();
|
||||
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<SessionDiffBaseline>();
|
||||
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<SessionDiffBaseline>();
|
||||
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 },
|
||||
|
|
|
|||
|
|
@ -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<void>((resolve, reject) => {
|
||||
completionServer.close((error) => (error ? reject(error) : resolve()));
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
expect(listRunningSessions().filter((session) => session.scopeKey === scopeKey)).toHaveLength(
|
||||
|
|
|
|||
|
|
@ -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([
|
||||
|
|
|
|||
|
|
@ -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 });
|
||||
}
|
||||
},
|
||||
);
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue