improve(tests): reduce agents and Gateway fixture setup time (#155562)

* perf(agents): reuse durable integration session fixtures

* perf(gateway): reuse repository seeds and acquisition fixtures

* perf(workers): reuse reader processes and repository fixtures

* perf(workers): retain prepared readers across workspace result cases
This commit is contained in:
Peter Steinberger 2026-09-22 01:18:25 -07:00 • committed by GitHub
parent 1282d73363
commit 24150a2c44
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
15 changed files with 287 additions and 174 deletions

View file

@ -16,7 +16,7 @@ import type {
} from "../tool-loop-detection.js";
import type { PostCompactionLoopPersistedError as PostCompactionLoopPersistedErrorType } from "./post-compaction-loop-guard.js";
import {
makeAttemptResult,
type makeAttemptResult,
makeCompactionSuccess,
makeOverflowError,
} from "./run.overflow-compaction.fixture.js";
@ -213,7 +213,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
// Attempt 1: overflow triggers compaction.
mockedRunEmbeddedAttempt.mockImplementationOnce(async () =>
makeAttemptResult({
session.makeAttemptResult({
terminal: { kind: "failed", source: "prompt", error: overflowError },
}),
);
@ -239,7 +239,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
if (revoked) {
admission.close();
}
return makeAttemptResult({
return session.makeAttemptResult({
toolMetas: [{ toolName: "gateway" }, { toolName: "gateway" }, { toolName: "gateway" }],
acceptedSessionSpawns: [
{
@ -290,7 +290,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
const overflowError = makeOverflowError();
let attemptAborted = false;
mockedRunEmbeddedAttempt.mockResolvedValueOnce(
makeAttemptResult({
session.makeAttemptResult({
terminal: { kind: "failed", source: "prompt", error: overflowError },
}),
);
@ -339,7 +339,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
expect(pendingTasks.size, "backend still runs after the outer timeout").toBeGreaterThan(0);
await expect(run).rejects.toMatchObject({ name: "CommandLaneTaskTimeoutError" });
} finally {
ignoredAttempt.resolve(makeAttemptResult());
ignoredAttempt.resolve(session.makeAttemptResult());
await run?.catch(() => undefined);
vi.useRealTimers();
}
@ -367,7 +367,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
},
);
if (!stop) {
mockedRunEmbeddedAttempt.mockResolvedValueOnce(makeAttemptResult());
mockedRunEmbeddedAttempt.mockResolvedValueOnce(session.makeAttemptResult());
}
const run = runEmbeddedAgent({
...baseParams,
@ -451,7 +451,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
expect(attemptSignal?.reason).toMatchObject({ name: "CommandLaneTaskTimeoutError" });
await expect(run).rejects.toMatchObject({ name: "CommandLaneTaskTimeoutError" });
} finally {
heldAttempt.resolve(makeAttemptResult());
heldAttempt.resolve(session.makeAttemptResult());
await run?.catch(() => undefined);
vi.useRealTimers();
}
@ -498,7 +498,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
expect(pendingTasks.size, "backend still runs after the outer timeout").toBeGreaterThan(0);
await expect(run).rejects.toMatchObject({ name: "CommandLaneTaskTimeoutError" });
} finally {
heldAttempt.resolve(makeAttemptResult());
heldAttempt.resolve(session.makeAttemptResult());
await run?.catch(() => undefined);
vi.useRealTimers();
}
@ -545,7 +545,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
expect(pendingTasks.size, "backend still runs after the outer timeout").toBeGreaterThan(0);
await expect(run).rejects.toMatchObject({ name: "CommandLaneTaskTimeoutError" });
} finally {
heldAttempt.resolve(makeAttemptResult());
heldAttempt.resolve(session.makeAttemptResult());
await run?.catch(() => undefined);
vi.useRealTimers();
}
@ -555,7 +555,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
const overflowError = makeOverflowError();
mockedRunEmbeddedAttempt.mockImplementationOnce(async () =>
makeAttemptResult({
session.makeAttemptResult({
terminal: { kind: "failed", source: "prompt", error: overflowError },
}),
);
@ -570,7 +570,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
onToolOutcome,
);
}
return makeAttemptResult({
return session.makeAttemptResult({
toolMetas: [{ toolName: "gateway" }, { toolName: "gateway" }, { toolName: "gateway" }],
});
});
@ -619,7 +619,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
// Attempt 1: overflow -> triggers compaction.
mockedRunEmbeddedAttempt.mockImplementationOnce(async () =>
makeAttemptResult({
session.makeAttemptResult({
terminal: { kind: "failed", source: "prompt", error: overflowError },
}),
);
@ -639,7 +639,7 @@ describe("post-compaction loop guard wired into runEmbeddedAgent", () => {
}
// History is still capped at HISTORY_TRIM_CAP after the trim.
expect(sessionState.toolCallHistory?.length).toBe(HISTORY_TRIM_CAP);
return makeAttemptResult({
return session.makeAttemptResult({
toolMetas: [{ toolName: "gateway" }, { toolName: "gateway" }, { toolName: "gateway" }],
});
});

View file

@ -2,11 +2,7 @@
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
import { makeTextToolResult } from "../../../test/helpers/text-tool-result.js";
import { buildEmbeddedRunnerAssistant } from "../test-helpers/embedded-agent-runner-e2e-fixtures.js";
import {
makeAttemptResult,
makeCompactionSuccess,
makeOverflowError,
} from "./run.overflow-compaction.fixture.js";
import { makeCompactionSuccess, makeOverflowError } from "./run.overflow-compaction.fixture.js";
import {
mockedCompactDirect,
mockedRunEmbeddedAttempt,
@ -61,7 +57,7 @@ function makeReplayUnsafeMidTurnOverflow(params?: {
}) {
const activeCount = params?.activeCount ?? 0;
const resultRecorded = params?.resultRecorded ?? true;
return makeAttemptResult({
return session.makeAttemptResult({
...(params?.codeModeEngaged ? { codeModeEngaged: true } : {}),
promptError: makeOverflowError("Context overflow: prompt too large (mid-turn precheck)."),
promptErrorSource: "precheck",
@ -110,7 +106,7 @@ describe("runEmbeddedAgent mid-turn precheck retry", () => {
it("continues once when persisted truncation is already a no-op", async () => {
mockedRunEmbeddedAttempt
.mockResolvedValueOnce(
makeAttemptResult({
session.makeAttemptResult({
preflightRecovery: {
route: "truncate_tool_results_only",
source: "mid-turn",
@ -121,7 +117,7 @@ describe("runEmbeddedAgent mid-turn precheck retry", () => {
latestMcpAppChannelView: { viewId: "view-before-retry" },
}),
)
.mockResolvedValueOnce(makeAttemptResult());
.mockResolvedValueOnce(session.makeAttemptResult());
const result = await runEmbeddedAgent({
...session.runParams,
@ -144,7 +140,7 @@ describe("runEmbeddedAgent mid-turn precheck retry", () => {
it("still compacts after a real provider overflow follows the no-op", async () => {
mockedRunEmbeddedAttempt
.mockResolvedValueOnce(
makeAttemptResult({
session.makeAttemptResult({
preflightRecovery: {
route: "truncate_tool_results_only",
source: "mid-turn",
@ -153,8 +149,8 @@ describe("runEmbeddedAgent mid-turn precheck retry", () => {
},
}),
)
.mockResolvedValueOnce(makeAttemptResult({ promptError: makeOverflowError() }))
.mockResolvedValueOnce(makeAttemptResult());
.mockResolvedValueOnce(session.makeAttemptResult({ promptError: makeOverflowError() }))
.mockResolvedValueOnce(session.makeAttemptResult());
mockedCompactDirect.mockResolvedValueOnce(
makeCompactionSuccess({
summary: "Compacted after provider rejection",
@ -177,7 +173,7 @@ describe("runEmbeddedAgent mid-turn precheck retry", () => {
it("compacts settled replay-unsafe tools and continues from their recorded result", async () => {
mockedRunEmbeddedAttempt
.mockResolvedValueOnce(makeReplayUnsafeMidTurnOverflow())
.mockResolvedValueOnce(makeAttemptResult());
.mockResolvedValueOnce(session.makeAttemptResult());
mockedCompactDirect.mockResolvedValueOnce(
makeCompactionSuccess({
summary: "Compacted after settled exec",
@ -209,7 +205,7 @@ describe("runEmbeddedAgent mid-turn precheck retry", () => {
vi.mocked(truncateOversizedToolResultsInSessionManager).mockImplementation(
actualTruncation.truncateOversizedToolResultsInSessionManager,
);
const successorId = "tool-projection-successor";
const successorId = `${session.runParams.sessionId}-tool-projection-successor`;
const toolResult = makeTextToolResult(
"call-exec",
"exec",
@ -270,14 +266,14 @@ describe("runEmbeddedAgent mid-turn precheck retry", () => {
.mockImplementationOnce(async (attempt) => {
expect(attempt.sessionId).toBe(successorId);
successor = prepareAttemptProjection(attempt, 1_000);
return makeAttemptResult({
return session.makeAttemptResult({
...makeReplayUnsafeMidTurnOverflow(),
sessionIdUsed: successorId,
messagesSnapshot: successor.messages,
preflightRecovery: { route: "compact_then_truncate", source: "mid-turn" },
});
})
.mockResolvedValueOnce(makeAttemptResult({ sessionIdUsed: successorId }));
.mockResolvedValueOnce(session.makeAttemptResult({ sessionIdUsed: successorId }));
mockedCompactDirect
.mockResolvedValueOnce(
makeCompactionSuccess({ summary: "Adopt successor", sessionId: successorId }),
@ -317,7 +313,7 @@ describe("runEmbeddedAgent mid-turn precheck retry", () => {
codeModeSuspended: true,
}),
)
.mockResolvedValueOnce(makeAttemptResult());
.mockResolvedValueOnce(session.makeAttemptResult());
mockedCompactDirect.mockResolvedValueOnce(
makeCompactionSuccess({
summary: "Compacted while exec waits",
@ -356,7 +352,7 @@ describe("runEmbeddedAgent mid-turn precheck retry", () => {
summary: "Compacted into a successor session",
firstKeptEntryId: "entry-rotated-exec",
tokensBefore: 201_000,
sessionId: "rotated-session",
sessionId: `${session.runParams.sessionId}-rotated`,
}),
);

View file

@ -27,10 +27,11 @@ import type {
EmbeddedRunAttemptInternalParams,
} from "./run/internal-params.js";
function timeoutAttempt() {
function timeoutAttempt(sessionIdUsed = "test-session") {
const assistant = makeAssistantMessageFixture();
assistant.usage = { ...assistant.usage, input: 180_000, totalTokens: 180_000 };
return makeAttemptResult({
sessionIdUsed,
terminal: { kind: "timeout", phase: "prompt", source: "idle", aborted: true },
assistantTexts: [],
lastAssistant: assistant,
@ -209,7 +210,7 @@ describe("recovery cancellation through the public run owner", () => {
let replacement: ReturnType<typeof prepareSystemAgentRunAdmission> | undefined;
const abort = new AbortController();
const callerError = new Error("caller stopped after owner changed");
mockedRunEmbeddedAttempt.mockResolvedValueOnce(timeoutAttempt());
mockedRunEmbeddedAttempt.mockResolvedValueOnce(timeoutAttempt(runParams.sessionId));
mockedCompactDirect.mockResolvedValueOnce(
makeCompactionSuccess({ summary: "Committed before owner change", tokensAfter: 40 }),
);
@ -309,7 +310,7 @@ describe("recovery cancellation through the public run owner", () => {
throw attemptError;
}
return {
...timeoutAttempt(),
...timeoutAttempt(attempt.sessionId),
compactionCount: harness.subscription.getCompactionCount(),
compactionTokensAfter: 80,
};
@ -332,7 +333,7 @@ describe("recovery cancellation through the public run owner", () => {
harness.emit({ type: "message_start", message: assistant });
harness.emit({ type: "message_end", message: assistant });
await harness.subscription.waitForPendingEvents();
return makeAttemptResult();
return makeAttemptResult({ sessionIdUsed: attempt.sessionId });
} finally {
harness.subscription.unsubscribe();
}

View file

@ -1,6 +1,6 @@
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../../test/helpers/promise.js";
import { makeAssistantMessageFixture } from "../test-helpers/assistant-message-fixtures.js";
import { makeAttemptResult } from "./run.overflow-compaction.fixture.js";
import {
mockedBuildEmbeddedRunPayloads,
mockedClassifyAssistantFailoverReason,
@ -61,7 +61,7 @@ describe("direct embedded retry lifecycle", () => {
content: failed ? [] : [{ type: "text", text: "Recovered reply" }],
errorMessage: failed ? "An error occurred while processing the request." : undefined,
});
return makeAttemptResult({
return session.makeAttemptResult({
providerRetryMaxRetries: budget,
hasSuccessfulModelResponse: attempts === 2 ? progress : false,
assistantTexts: failed ? [] : ["Recovered reply"],
@ -109,6 +109,7 @@ describe("direct embedded retry lifecycle", () => {
let wait: Promise<void> | undefined;
let waitSettled = false;
let pending: ReturnType<typeof run> | undefined;
const sleepStarted = createDeferred();
vi.useFakeTimers();
try {
mockedSleep.mockImplementation((delayMs, signal) => {
@ -116,6 +117,7 @@ describe("direct embedded retry lifecycle", () => {
wait = sleep(delayMs, signal).finally(() => {
waitSettled = true;
});
sleepStarted.resolve();
return wait;
});
const assistant = makeAssistantMessageFixture({
@ -124,7 +126,7 @@ describe("direct embedded retry lifecycle", () => {
errorMessage: "429 rate limit exceeded; Retry-After: 3600",
});
mockedRunEmbeddedAttempt.mockResolvedValue(
makeAttemptResult({
session.makeAttemptResult({
lastAssistant: assistant,
currentAttemptAssistant: assistant,
}),
@ -138,7 +140,7 @@ describe("direct embedded retry lifecycle", () => {
abortSignal: caller.signal,
});
const outcome = pending.catch((error: unknown) => error);
await vi.waitFor(() => expect(mockedSleep).toHaveBeenCalled(), { timeout: 10_000 });
await Promise.race([sleepStarted.promise, pending]);
expect(mockedSleep).toHaveBeenCalledWith(3_600_000, expect.any(AbortSignal));
expect(waitSettled).toBe(false);
await vi.advanceTimersByTimeAsync(60_001);
@ -177,7 +179,7 @@ describe("direct embedded retry lifecycle", () => {
assistantTranscriptIdempotencyKey: "saved-A",
},
});
return makeAttemptResult({
return session.makeAttemptResult({
assistantTexts: [],
lastAssistant: assistant,
currentAttemptAssistant: assistant,
@ -230,7 +232,7 @@ describe("direct embedded retry lifecycle", () => {
assistantTranscriptIdempotencyKey: `saved-${attempts}`,
},
});
return makeAttemptResult({
return session.makeAttemptResult({
assistantTexts: failed ? [] : ["Recovered reply"],
lastAssistant: assistant,
currentAttemptAssistant: assistant,

View file

@ -1,4 +1,7 @@
import fs from "node:fs/promises";
import path from "node:path";
import type { OpenClawTestState } from "../../test-utils/openclaw-test-state.js";
import { makeAttemptResult } from "./run.overflow-compaction.fixture.js";
import {
loadRunOverflowCompactionHarness,
createOverflowRunParams,
@ -8,6 +11,8 @@ import {
import { guardRunWorkspaceOwnership } from "./run.workspace-ownership.test-support.js";
let sharedRunEmbeddedAgent: Promise<TestRunEmbeddedAgent> | undefined;
let sharedSessionState: Promise<OpenClawTestState> | undefined;
let sessionSequence = 0;
/**
* These scenarios intentionally cross several runner owners. Load the mocked
@ -31,48 +36,69 @@ export function loadSharedRunIntegrationHarness(): Promise<TestRunEmbeddedAgent>
return sharedRunEmbeddedAgent;
}
/** Durable recovery needs a real row; the public lane owner installs its writer claim. */
export async function createSharedRunIntegrationSession(identity?: {
sessionId: string;
sessionKey: string;
}) {
/** Close only after every case has drained its work and released its selectors. */
export async function cleanupSharedRunIntegrationSessions(): Promise<void> {
await (await sharedSessionState)?.cleanup();
}
/** Reuse real databases; each case owns distinct session rows and workspace files. */
export async function createSharedRunIntegrationSession() {
const { createOpenClawTestState } = await import("../../test-utils/openclaw-test-state.js");
const { captureEnv } = await import("../../test-utils/env.js");
const { resetConfigRuntimeState } = await import("../../config/runtime-snapshot.js");
const { drainSessionStateForTest } = await import("../../test-utils/session-state-cleanup.js");
const { loadSessionEntry, replaceSessionEntry } =
await import("../../config/sessions/session-accessor.js");
const { forgetActiveSessionForShutdown } =
await import("../../gateway/active-sessions-shutdown-tracker.js");
const state = await createOpenClawTestState({ label: "run-integration-session" });
const baseRunParams = createOverflowRunParams(state);
const { sessionId, sessionKey } = identity ?? baseRunParams;
sharedSessionState ??= createOpenClawTestState({
label: "run-integration-sessions",
applyEnv: false,
});
const state = await sharedSessionState;
const env = captureEnv(Object.keys(state.envVars));
state.applyEnv();
const caseId = ++sessionSequence;
const sessionId = `test-session-${caseId}`;
const sessionKey = `agent:main:test-key-${caseId}`;
const workspaceDir = path.join(state.workspaceDir, `case-${caseId}`);
const baseRunParams = createOverflowRunParams({ workspaceDir });
const sessionTarget = {
agentId: "main",
sessionId,
sessionKey,
storePath: path.join(state.agentDir(), "openclaw-agent.sqlite"),
};
let cleanupPromise: Promise<void> | undefined;
const cleanup = () =>
(cleanupPromise ??= (async () => {
await drainSessionStateForTest({ stateDir: state.stateDir, rootPath: state.root });
forgetActiveSessionForShutdown(sessionId);
const current = loadSessionEntry({ ...sessionTarget, readConsistency: "latest" });
if (current) {
forgetActiveSessionForShutdown(current.sessionId);
}
env.restore();
resetConfigRuntimeState();
})());
try {
await fs.mkdir(workspaceDir, { recursive: true });
// The public lane owner still installs the real durable writer claim.
await replaceSessionEntry(sessionTarget, { sessionId, updatedAt: 1 });
return {
runParams: {
...baseRunParams,
sessionId,
sessionKey,
runId: `run-integration-${caseId}`,
sessionTarget,
},
cleanup: async () => {
try {
forgetActiveSessionForShutdown(sessionId);
const current = loadSessionEntry({ ...sessionTarget, readConsistency: "latest" });
if (current) {
forgetActiveSessionForShutdown(current.sessionId);
}
} finally {
await state.cleanup();
}
},
makeAttemptResult: (overrides?: Parameters<typeof makeAttemptResult>[0]) =>
makeAttemptResult({ sessionIdUsed: sessionId, ...overrides }),
cleanup,
};
} catch (error) {
await state.cleanup();
await cleanup();
throw error;
}
}

View file

@ -1,4 +1,6 @@
// The imported scenario modules share one mocked runEmbeddedAgent module graph.
import { afterAll } from "vitest";
import { cleanupSharedRunIntegrationSessions } from "./run.shared-integration-harness.test-support.js";
import "./run.before-agent-finalize.test-support.js";
import "./run.before-agent-reply-cron.test-support.js";
import "./run.codex-app-server-recovery.test-support.js";
@ -16,3 +18,5 @@ import "./run.retry-lifecycle.test-support.js";
import "./run.timeout-triggered-compaction.test-support.js";
import "./sessions-yield.orchestration.test-support.js";
import "./usage-reporting.test-support.js";
afterAll(cleanupSharedRunIntegrationSessions);

View file

@ -67,19 +67,19 @@ describe("runEmbeddedAgent timeout recovery composition", () => {
fixture = session;
const successor = {
...session.runParams.sessionTarget,
sessionId: "timeout-rotated-session",
sessionId: `${session.runParams.sessionId}-timeout-rotated`,
};
mockedBuildEmbeddedRunPayloads.mockReturnValue([{ text: "timeout recovery complete" }]);
mockedRunEmbeddedAttempt
.mockImplementationOnce(async (params) => {
params.onUserMessagePersisted?.(makeUserMessage("hello", 1));
return makeAttemptResult({
return session.makeAttemptResult({
timedOut: true,
lastAssistant: { usage: { input: 160_000 } } as never,
});
})
.mockResolvedValueOnce(
makeAttemptResult({
session.makeAttemptResult({
promptError: null,
sessionIdUsed: successor.sessionId,
sessionFileUsed: successor.sessionKey,
@ -112,7 +112,7 @@ describe("runEmbeddedAgent timeout recovery composition", () => {
expect(mockedRunEmbeddedAttempt.mock.calls[1]?.[0]?.prompt).not.toBe(session.runParams.prompt);
const compactParams = mockedCompactDirect.mock.calls[0]?.[0] as CompactParams | undefined;
expect(compactParams).toMatchObject({
sessionId: "test-session",
sessionId: session.runParams.sessionId,
tokenBudget: 200_000,
force: true,
compactionTarget: "budget",
@ -145,7 +145,7 @@ describe("runEmbeddedAgent timeout recovery composition", () => {
}),
).toBe(true);
clearActiveEmbeddedRun(params.sessionId, handle, params.sessionKey, params.sessionFile);
return makeAttemptResult({
return session.makeAttemptResult({
timedOut: true,
lastAssistant: { usage: { input: 160_000 } } as never,
});
@ -232,7 +232,7 @@ describe("runEmbeddedAgent timeout recovery composition", () => {
mode: "api-key",
}));
mockedRunEmbeddedAttempt.mockResolvedValue(
makeAttemptResult({
session.makeAttemptResult({
timedOut: true,
aborted: true,
lastAssistant: { usage: { input: 150_000 } } as never,

View file

@ -0,0 +1,46 @@
import { execFile } from "node:child_process";
import fs from "node:fs/promises";
import path from "node:path";
import { promisify } from "node:util";
const execFileAsync = promisify(execFile);
async function initializeRepositorySeed(root: string, name: string): Promise<void> {
await fs.mkdir(path.join(root, ".openclaw"), { recursive: true });
await fs.writeFile(path.join(root, "README.md"), `${name}\n`);
await fs.writeFile(
path.join(root, ".openclaw", "worktree-setup.sh"),
"#!/bin/sh\ntouch setup-marker.txt\n",
{ mode: 0o755 },
);
await execFileAsync("git", ["init", "-b", "main", root]);
await execFileAsync("git", ["-C", root, "add", "."]);
await execFileAsync("git", [
"-C",
root,
"-c",
"user.name=Test",
"-c",
"user.email=test@example.invalid",
"-c",
"commit.gpgsign=false",
"commit",
"-m",
"Initialize fixture",
]);
}
export function createWorktreeSpawnRepositoryFixture(seedRoot: string) {
const repositorySeeds = new Set<string>();
return async (caseRoot: string, name: string): Promise<string> => {
const seed = path.join(seedRoot, name);
if (!repositorySeeds.has(name)) {
await initializeRepositorySeed(seed, name);
repositorySeeds.add(name);
}
const root = path.join(caseRoot, name);
// Preserve each source's distinct committed README, with private Git metadata per case.
await fs.cp(seed, root, { recursive: true });
return await fs.realpath(root);
};
}

View file

@ -35,6 +35,7 @@ import { GATEWAY_CLIENT_MODES, GATEWAY_CLIENT_NAMES } from "../utils/message-cha
import { waitForChatAbortControllerRemoval } from "./chat-abort-lifecycle-internal.js";
import type { ChatAbortControllerEntry } from "./chat-abort.js";
import { createDirectChatContext } from "./server-chat.agent-events.test-helpers.js";
import { createWorktreeSpawnRepositoryFixture } from "./server.sessions.create-worktree-spawn.test-support.js";
import { settleWorkspaceRuns } from "./server.sessions.create.projects.test-support.js";
import { agentDiscoveryMock, dispatchInboundMessageMock, testState } from "./test-helpers.js";
import {
@ -49,7 +50,12 @@ const projectCloneMocks = vi.hoisted(() => ({
}));
vi.mock("../projects/project-clone.js", () => projectCloneMocks);
const { createSessionStoreDir } = setupGatewaySessionsHandlerTestHarness();
let createRepository: ReturnType<typeof createWorktreeSpawnRepositoryFixture>;
const { createSessionStoreDir } = setupGatewaySessionsHandlerTestHarness(async (makeTempDir) => {
createRepository = createWorktreeSpawnRepositoryFixture(
makeTempDir("openclaw-spawn-repo-seeds-"),
);
});
const execFileAsync = promisify(execFile);
const parentKey = "agent:main:dashboard:project-parent";
const parentCreateParams = {
@ -84,33 +90,6 @@ type CreatedWorktreeSession = {
worktree: { id: string; path: string };
};
async function createRepository(name: string): Promise<string> {
const root = path.join(state.root, name);
await fs.mkdir(path.join(root, ".openclaw"), { recursive: true });
await fs.writeFile(path.join(root, "README.md"), `${name}\n`);
await fs.writeFile(
path.join(root, ".openclaw", "worktree-setup.sh"),
"#!/bin/sh\ntouch setup-marker.txt\n",
{ mode: 0o755 },
);
await execFileAsync("git", ["init", "-b", "main", root]);
await execFileAsync("git", ["-C", root, "add", "."]);
await execFileAsync("git", [
"-C",
root,
"-c",
"user.name=Test",
"-c",
"user.email=test@example.invalid",
"-c",
"commit.gpgsign=false",
"commit",
"-m",
"Initialize fixture",
]);
return await fs.realpath(root);
}
function spawnClient(admin = false, requesterSessionKey = parentKey) {
return {
connect: { scopes: [admin ? "operator.admin" : "operator.write"] },
@ -176,7 +155,7 @@ beforeEach(async () => {
state = await createOpenClawTestState({ layout: "state-only", prefix: "openclaw-spawn-repo-" });
const defaultWorkspace = path.join(state.root, "non-git-workspace");
await fs.mkdir(defaultWorkspace);
repository = await createRepository("selected-project");
repository = await createRepository(state.root, "selected-project");
await closeOpenClawStateDatabaseAsync();
closeOpenClawStateDatabaseForTest();
testState.agentConfig = { workspace: defaultWorkspace };
@ -224,7 +203,7 @@ test.each([
}
const projectName =
required && !worktree ? "non-git-workspace/tool-selected-project" : "tool-selected-project";
const otherRepository = await createRepository(projectName);
const otherRepository = await createRepository(state.root, projectName);
const project = await registerProjectRegistry({ path: otherRepository });
projectCloneMocks.materializeProjectClone.mockResolvedValue(project);
const { getRuntimeConfig } = await getGatewayConfigModule();
@ -315,7 +294,7 @@ test.each([
{ agentId: "main", sessionKey: parentKey, storePath },
{ ...parent, sandbox: "required" },
);
const otherRepository = await createRepository("sandbox-external-project");
const otherRepository = await createRepository(state.root, "sandbox-external-project");
const project = await registerProjectRegistry({ path: otherRepository });
projectCloneMocks.materializeProjectClone.mockResolvedValue(project);
const { getRuntimeConfig } = await getGatewayConfigModule();
@ -594,7 +573,7 @@ test("keyed worktree creation reuses its recorded base after reopening the regis
error: { code: "INVALID_REQUEST", message: expect.stringContaining("already bound") },
});
}
const otherRepository = await createRepository("other-replay-project");
const otherRepository = await createRepository(state.root, "other-replay-project");
const wrongRepository = await directSessionReq(
"sessions.create",
{ ...params, cwd: otherRepository },
@ -640,7 +619,7 @@ test.each([
source === "direct"
? (await createDirectProjectParent()).key
: (await createManagedProjectParent()).key;
const otherRepository = await createRepository("other-project");
const otherRepository = await createRepository(state.root, "other-project");
let params: Record<string, unknown>;
if (selection === "project") {
const project = await registerProjectRegistry({ path: otherRepository });
@ -714,9 +693,13 @@ test("publishes a failed worktree spawn only after its durable session failure",
const key = "agent:main:dashboard:failed-worktree-child";
const target = { agentId: "main", sessionKey: key, storePath };
const preparation = createDeferredCore<never>();
const preparationStarted = createDeferredCore();
const createWorktree = vi
.spyOn(managedWorktrees, "createWithOutcome")
.mockReturnValueOnce(preparation.promise);
.mockImplementationOnce(() => {
preparationStarted.resolve();
return preparation.promise;
});
const failure = new Error(
"git ls-tree -r --format=%(objectsize) c79ad267ba623c1a323f1f6e8b60228bd5a30ce5 -- failed (timed out after 120 seconds; signal SIGTERM):\n4514\n4168\nCheck repository access and disk space.",
);
@ -803,7 +786,8 @@ test("publishes a failed worktree spawn only after its durable session failure",
const released = getSessionWorkAdmissionRelease({ scope: storePath, identities: [key] });
expect(released).toBeDefined();
expect(loadSessionEntry(target)?.pendingWorktree?.workspace).toBe(repository);
await vi.waitFor(() => expect(createWorktree).toHaveBeenCalledOnce());
await preparationStarted.promise;
expect(createWorktree).toHaveBeenCalledOnce();
preparation.reject(failure);
await withTimeout(released!, SESSION_WORK_ADMISSION_DRAIN_TIMEOUT_MS, "failed workspace proof");
await registryProjection;

View file

@ -24,14 +24,41 @@ const { WebSocket, WebSocketServer }: typeof import("ws") = require(
);
type WebSocket = WebSocketClient;
const acquisitionFixture = vi.hoisted(() => ({
start: vi.fn(),
port: 0,
observeClient: undefined as ((client: WebSocket) => void) | undefined,
}));
async function observeWebSocket() {
const actual = await vi.importActual<
typeof import("../../packages/gateway-client/src/websocket.js")
>("../../packages/gateway-client/src/websocket.js");
class ObservedWebSocket extends actual.WebSocket {
constructor(...args: ConstructorParameters<typeof WebSocket>) {
super(...args);
acquisitionFixture.observeClient?.(this);
}
}
return { ...actual, default: ObservedWebSocket, WebSocket: ObservedWebSocket };
}
vi.mock("ws", () => observeWebSocket());
vi.mock("../../packages/gateway-client/src/websocket.js", () => observeWebSocket());
vi.mock("./server.js", () => ({ startGatewayServer: acquisitionFixture.start }));
vi.mock("../agents/prepared-model-runtime.test-support.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../agents/prepared-model-runtime.test-support.js")>()),
resetPreparedGatewayModelCatalogForTest: vi.fn(),
}));
vi.mock("../test-utils/ports.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../test-utils/ports.js")>()),
getDeterministicFreePortBlock: async () => acquisitionFixture.port,
}));
afterEach(() => {
vi.doUnmock("ws");
vi.doUnmock("../../packages/gateway-client/src/websocket.js");
vi.doUnmock("./server.js");
vi.doUnmock("../agents/prepared-model-runtime.test-support.js");
vi.doUnmock("../test-utils/ports.js");
vi.doUnmock("../infra/device-pairing.js");
vi.resetModules();
acquisitionFixture.start.mockReset();
acquisitionFixture.observeClient = undefined;
vi.restoreAllMocks();
});
type PeerBehavior =
@ -74,33 +101,17 @@ async function withAcquisitionPeer(
const rejectAuth = behavior === "reject auth" || behavior === "reject auth without close";
// Observe the real dependency; keep otherwise-unhandled errors local to this case.
// Counting the remaining listeners makes a removed owner handler observable.
const observeWebSocket = async () => {
const actual = await vi.importActual<
typeof import("../../packages/gateway-client/src/websocket.js")
>("../../packages/gateway-client/src/websocket.js");
class ObservedWebSocket extends actual.WebSocket {
constructor(...args: ConstructorParameters<typeof WebSocket>) {
super(...args);
clients.push(this);
this.once("close", () => closed.add(this));
this.on("error", (error) => {
errors.push(error);
transportFailure.resolve(error);
if (this.listenerCount("error") === 1) {
unownedErrors.push(error);
}
});
acquisitionFixture.observeClient = (client) => {
clients.push(client);
client.once("close", () => closed.add(client));
client.on("error", (error) => {
errors.push(error);
transportFailure.resolve(error);
if (client.listenerCount("error") === 1) {
unownedErrors.push(error);
}
}
return {
...actual,
default: ObservedWebSocket,
WebSocket: ObservedWebSocket,
WebSocketServer,
};
});
};
vi.doMock("ws", observeWebSocket);
vi.doMock("../../packages/gateway-client/src/websocket.js", observeWebSocket);
const sockets = new Set<Socket>();
const requests: ReturnType<typeof parseMinimalGatewayRequestFrame>[] = [];
const connectReceived = createDeferred();
@ -253,22 +264,12 @@ function mockPeerGateway(peer: AcquisitionPeer, close = peer.close) {
startupSettled: Promise.resolve(),
getTailscaleIngressEndpoint: () => undefined,
} satisfies import("./server.js").GatewayServer;
const start = vi.fn(async () => {
acquisitionFixture.port = peer.port;
const start = acquisitionFixture.start.mockImplementation(async () => {
// Mirror the real startup-owned selector for close-order assertions.
process.env.OPENCLAW_GATEWAY_PORT = String(peer.port);
return server;
});
vi.doMock("./server.js", () => ({
startGatewayServer: start,
}));
vi.doMock("../agents/prepared-model-runtime.test-support.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../agents/prepared-model-runtime.test-support.js")>()),
resetPreparedGatewayModelCatalogForTest: vi.fn(),
}));
vi.doMock("../test-utils/ports.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../test-utils/ports.js")>()),
getDeterministicFreePortBlock: async () => peer.port,
}));
return start;
}
@ -570,15 +571,13 @@ describe("raw Gateway helper acquisition ownership", () => {
const release = createDeferred();
const preparationError = new Error("synthetic preparation failure");
let preparationFinished = false;
vi.doMock("../infra/device-pairing.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../infra/device-pairing.js")>()),
getPairedDevice: async () => {
preparing.resolve();
await release.promise;
preparationFinished = true;
throw preparationError;
},
}));
const devicePairing = await import("../infra/device-pairing.js");
vi.spyOn(devicePairing, "getPairedDevice").mockImplementation(async () => {
preparing.resolve();
await release.promise;
preparationFinished = true;
throw preparationError;
});
const { connectWebchatClient } = await import("./test-helpers.server.js");
let settled = false;
const acquisition = connectWebchatClient({ port: peer.port })

View file

@ -1,13 +1,13 @@
import { fileURLToPath } from "node:url";
import { expectDefined } from "@openclaw/normalization-core";
import type { OpenClawPluginService, WorkerProvider } from "openclaw/plugin-sdk/plugin-entry";
import { createPluginStateKeyedStoreForTests } from "openclaw/plugin-sdk/plugin-state-test-runtime";
import { createTestPluginApi } from "openclaw/plugin-sdk/plugin-test-api";
import { createPluginRuntimeMock } from "openclaw/plugin-sdk/plugin-test-runtime";
import * as processRuntime from "openclaw/plugin-sdk/process-runtime";
import type { SpawnResult } from "openclaw/plugin-sdk/process-runtime";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { importFreshModule } from "../../plugin-sdk/test-helpers/import-fresh.js";
import { createPluginRuntimeMock } from "../../plugin-sdk/test-helpers/plugin-runtime-mock.js";
import { createPluginStateKeyedStore } from "../../plugin-state/plugin-state-store.js";
import { resolvePluginModuleExport } from "../../plugins/module-export.js";
import * as support from "./service.test-support.js";
@ -35,7 +35,7 @@ function commandResult(overrides: Partial<SpawnResult> = {}): SpawnResult {
}
describe("Crabbox runtime preflight cleanup", () => {
support.setupWorkerEnvironmentServiceSuite();
support.setupWorkerEnvironmentServiceSuite({ reuseReadWorkers: true });
const pluginServices: OpenClawPluginService[] = [];
async function registerProvider(): Promise<WorkerProvider> {
let registered: WorkerProvider | undefined;
@ -50,7 +50,7 @@ describe("Crabbox runtime preflight cleanup", () => {
id: "crabbox",
runtime: createPluginRuntimeMock({
state: {
openKeyedStore: (options) => createPluginStateKeyedStoreForTests("crabbox", options),
openKeyedStore: (options) => createPluginStateKeyedStore("crabbox", options),
},
}),
rootDir: fileURLToPath(new URL("../../../extensions/crabbox/", import.meta.url)),

View file

@ -2,7 +2,7 @@ import fs from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import { expectDefined } from "@openclaw/normalization-core";
import { afterEach, beforeEach, vi } from "vitest";
import { afterAll, afterEach, beforeEach, vi } from "vitest";
import type { OpenClawConfig } from "../../config/types.js";
import type {
WorkerDesktopEndpoint,
@ -12,6 +12,7 @@ import type {
} from "../../plugins/types.js";
import {
closeOpenClawStateDatabaseAsync,
closeOpenClawStateDatabaseByPathAsync,
closeOpenClawStateDatabaseForTest,
openOpenClawStateDatabase,
type OpenClawStateDatabase,
@ -112,12 +113,14 @@ export const testState = {} as {
config: OpenClawConfig;
nowMs: number;
providersEnabled: boolean;
reuseReadWorkers: boolean;
prepareInstallation: WorkerEnvironmentServiceOptions["prepareInstallation"];
bootstrapWorker: WorkerEnvironmentServiceOptions["bootstrapWorker"];
};
export function setupWorkerEnvironmentServiceSuite() {
export function setupWorkerEnvironmentServiceSuite(options: { reuseReadWorkers?: boolean } = {}) {
beforeEach(async () => {
testState.reuseReadWorkers = options.reuseReadWorkers === true;
testState.root = await fs.mkdtemp(
path.join(await fs.realpath(os.tmpdir()), "openclaw-worker-service-"),
);
@ -155,10 +158,26 @@ export function setupWorkerEnvironmentServiceSuite() {
// Shutdown may schedule cleanup after a test leaves fake timers installed.
vi.useRealTimers();
await testState.service?.stop();
await closeOpenClawStateDatabaseAsync();
closeOpenClawStateDatabaseForTest();
await closeWorkerEnvironmentDatabase();
await fs.rm(testState.root, { recursive: true, force: true });
});
if (options.reuseReadWorkers) {
afterAll(async () => {
await closeOpenClawStateDatabaseAsync();
closeOpenClawStateDatabaseForTest();
});
}
}
async function closeWorkerEnvironmentDatabase() {
if (testState.reuseReadWorkers) {
// Close native handles and admission for this case; retain only the reader worker code.
await closeOpenClawStateDatabaseByPathAsync(testState.stateDb.path);
} else {
await closeOpenClawStateDatabaseAsync();
}
closeOpenClawStateDatabaseForTest();
}
export function getDevelopmentProfile() {
@ -171,8 +190,7 @@ export function getDevelopmentProfile() {
export async function reopenWorkerEnvironmentStore() {
await testState.service?.stop();
testState.service = undefined;
await closeOpenClawStateDatabaseAsync();
closeOpenClawStateDatabaseForTest();
await closeWorkerEnvironmentDatabase();
testState.stateDb = openOpenClawStateDatabase({
env: { OPENCLAW_STATE_DIR: testState.root },
});

View file

@ -15,6 +15,7 @@ import { upsertSessionEntryCore } from "../../config/sessions/session-accessor.j
import { resetAgentEventsForTest } from "../../infra/agent-events.js";
import {
closeOpenClawStateDatabaseAsync,
closeOpenClawStateDatabaseByPathAsync,
closeOpenClawStateDatabaseForTest,
openOpenClawStateDatabase,
type OpenClawStateDatabase,
@ -99,11 +100,22 @@ export async function setupWorkerTurnLauncherTest(): Promise<void> {
sessionFile = SESSION_KEY;
}
export async function cleanupWorkerTurnLauncherTest(): Promise<void> {
export function cleanupWorkerTurnLauncherTest(): Promise<void>;
export function cleanupWorkerTurnLauncherTest(options: {
reuseReadWorkers: boolean;
}): Promise<void>;
export async function cleanupWorkerTurnLauncherTest(
options: { reuseReadWorkers?: boolean } = {},
): Promise<void> {
cleanupAdmissionSink?.();
cleanupAdmissionSink = undefined;
clearRuntimeConfigSnapshot();
await closeOpenClawStateDatabaseAsync();
if (options.reuseReadWorkers) {
// Retain reader execution only; this case's native handles and admission still close.
await closeOpenClawStateDatabaseByPathAsync(database.path);
} else {
await closeOpenClawStateDatabaseAsync();
}
closeOpenClawStateDatabaseForTest();
resetAgentEventsForTest();
await testState.cleanup();

View file

@ -2,7 +2,8 @@ import fs from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import { pathToFileURL } from "node:url";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js";
import { setRuntimeConfigSnapshot } from "../../config/io.js";
import {
loadSessionEntry,
@ -14,6 +15,7 @@ import { NodeWorkerWorkspaceRuntime } from "../../node-host/node-worker-workspac
import { runCommandWithTimeout } from "../../process/exec.js";
import type { DB } from "../../state/openclaw-state-db.generated.js";
import {
closeOpenClawStateDatabaseAsync,
closeOpenClawStateDatabaseForTest,
openOpenClawStateDatabase,
} from "../../state/openclaw-state-db.js";
@ -64,6 +66,8 @@ vi.mock("./worker-github-binding.js", () => ({
}));
describe("repository workspace result ownership", () => {
const seedDirs = useAutoCleanupTempDirTracker(afterAll);
const originSeeds = new Map<boolean, string>();
let closeNode: (() => Promise<void>) | undefined;
beforeEach(setupWorkerTurnLauncherTest);
afterEach(async () => {
@ -71,14 +75,15 @@ describe("repository workspace result ownership", () => {
await closeNode?.();
} finally {
closeNode = undefined;
await cleanupWorkerTurnLauncherTest();
await cleanupWorkerTurnLauncherTest({ reuseReadWorkers: true });
}
});
afterAll(async () => {
await closeOpenClawStateDatabaseAsync();
closeOpenClawStateDatabaseForTest();
});
async function fixture(executionMode: "worker-turn" | "remote-exec", runSetupScript = false) {
setRuntimeConfigSnapshot({ session: { store: sessionTarget.storePath } });
const origin = path.join(root, "origin");
await fs.mkdir(origin);
async function initializeOriginSeed(origin: string, runSetupScript: boolean) {
if (runSetupScript) {
await fs.mkdir(path.join(origin, ".openclaw"));
await fs.writeFile(
@ -92,7 +97,7 @@ describe("repository workspace result ownership", () => {
timeoutMs: 10_000,
baseEnv: {
PATH: process.env.PATH,
HOME: root,
HOME: origin,
GIT_CONFIG_GLOBAL: os.devNull,
GIT_CONFIG_NOSYSTEM: "1",
},
@ -112,6 +117,19 @@ describe("repository workspace result ownership", () => {
"-m",
"base",
);
}
async function fixture(executionMode: "worker-turn" | "remote-exec", runSetupScript = false) {
setRuntimeConfigSnapshot({ session: { store: sessionTarget.storePath } });
let seed = originSeeds.get(runSetupScript);
if (!seed) {
seed = seedDirs.make("openclaw-repository-result-seed-");
await initializeOriginSeed(seed, runSetupScript);
originSeeds.set(runSetupScript, seed);
}
const origin = path.join(root, "origin");
// Only pristine source bytes are shared; checkpoints and Git refs stay case-owned.
await fs.cp(seed, origin, { recursive: true });
const store = getSessionRepositoryWorkspaceStore();
const repository = store.create({
agentId: sessionTarget.agentId,

View file

@ -8,7 +8,8 @@ import { closeOpenClawAgentDatabasesAsync } from "../state/openclaw-agent-db-lif
import { closeOpenClawStateDatabaseByPathAsync } from "../state/openclaw-state-db.js";
import { resolveOpenClawStateSqlitePath } from "../state/openclaw-state-db.paths.js";
export async function cleanupSessionStateForTest(
/** Settle case-owned work while a suite fixture retains its database workers. */
export async function drainSessionStateForTest(
options: { stateDir?: string; rootPath?: string } = {},
): Promise<void> {
await drainSessionStoreWriterQueuesForTest();
@ -20,6 +21,12 @@ export async function cleanupSessionStateForTest(
}
await drainFileLockStateForTest();
clearSessionStoreCacheForTest();
}
export async function cleanupSessionStateForTest(
options: { stateDir?: string; rootPath?: string } = {},
): Promise<void> {
await drainSessionStateForTest(options);
if (!options.stateDir) {
return;
}