diff --git a/src/agents/embedded-agent-runner/run.compaction-loop-guard.test-support.ts b/src/agents/embedded-agent-runner/run.compaction-loop-guard.test-support.ts index 8f70b2f96739..a417253410a7 100644 --- a/src/agents/embedded-agent-runner/run.compaction-loop-guard.test-support.ts +++ b/src/agents/embedded-agent-runner/run.compaction-loop-guard.test-support.ts @@ -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" }], }); }); diff --git a/src/agents/embedded-agent-runner/run.midturn-precheck-retry.test-support.ts b/src/agents/embedded-agent-runner/run.midturn-precheck-retry.test-support.ts index c3fab166c239..fc132447e63c 100644 --- a/src/agents/embedded-agent-runner/run.midturn-precheck-retry.test-support.ts +++ b/src/agents/embedded-agent-runner/run.midturn-precheck-retry.test-support.ts @@ -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`, }), ); diff --git a/src/agents/embedded-agent-runner/run.recovery-cancellation.test-support.ts b/src/agents/embedded-agent-runner/run.recovery-cancellation.test-support.ts index 4c7c1f286b34..e8881a711fa3 100644 --- a/src/agents/embedded-agent-runner/run.recovery-cancellation.test-support.ts +++ b/src/agents/embedded-agent-runner/run.recovery-cancellation.test-support.ts @@ -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 | 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(); } diff --git a/src/agents/embedded-agent-runner/run.retry-lifecycle.test-support.ts b/src/agents/embedded-agent-runner/run.retry-lifecycle.test-support.ts index a973347c7452..094e67e3fa4f 100644 --- a/src/agents/embedded-agent-runner/run.retry-lifecycle.test-support.ts +++ b/src/agents/embedded-agent-runner/run.retry-lifecycle.test-support.ts @@ -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 | undefined; let waitSettled = false; let pending: ReturnType | 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, diff --git a/src/agents/embedded-agent-runner/run.shared-integration-harness.test-support.ts b/src/agents/embedded-agent-runner/run.shared-integration-harness.test-support.ts index 7784195374ce..d40ce42b9a50 100644 --- a/src/agents/embedded-agent-runner/run.shared-integration-harness.test-support.ts +++ b/src/agents/embedded-agent-runner/run.shared-integration-harness.test-support.ts @@ -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 | undefined; +let sharedSessionState: Promise | undefined; +let sessionSequence = 0; /** * These scenarios intentionally cross several runner owners. Load the mocked @@ -31,48 +36,69 @@ export function loadSharedRunIntegrationHarness(): Promise 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 { + 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 | 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[0]) => + makeAttemptResult({ sessionIdUsed: sessionId, ...overrides }), + cleanup, }; } catch (error) { - await state.cleanup(); + await cleanup(); throw error; } } diff --git a/src/agents/embedded-agent-runner/run.shared-integration.test.ts b/src/agents/embedded-agent-runner/run.shared-integration.test.ts index 4ec2df8d59cb..13bad85a1808 100644 --- a/src/agents/embedded-agent-runner/run.shared-integration.test.ts +++ b/src/agents/embedded-agent-runner/run.shared-integration.test.ts @@ -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); diff --git a/src/agents/embedded-agent-runner/run.timeout-triggered-compaction.test-support.ts b/src/agents/embedded-agent-runner/run.timeout-triggered-compaction.test-support.ts index e0c1990e845b..fee80f5ca35f 100644 --- a/src/agents/embedded-agent-runner/run.timeout-triggered-compaction.test-support.ts +++ b/src/agents/embedded-agent-runner/run.timeout-triggered-compaction.test-support.ts @@ -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, diff --git a/src/gateway/server.sessions.create-worktree-spawn.test-support.ts b/src/gateway/server.sessions.create-worktree-spawn.test-support.ts new file mode 100644 index 000000000000..a49372396528 --- /dev/null +++ b/src/gateway/server.sessions.create-worktree-spawn.test-support.ts @@ -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 { + 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(); + return async (caseRoot: string, name: string): Promise => { + 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); + }; +} diff --git a/src/gateway/server.sessions.create-worktree-spawn.test.ts b/src/gateway/server.sessions.create-worktree-spawn.test.ts index 3d909fbc281e..53d8258659ac 100644 --- a/src/gateway/server.sessions.create-worktree-spawn.test.ts +++ b/src/gateway/server.sessions.create-worktree-spawn.test.ts @@ -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; +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 { - 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; 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(); + 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; diff --git a/src/gateway/test-helpers.acquisition.test.ts b/src/gateway/test-helpers.acquisition.test.ts index f4612b450b05..c68f29a26da3 100644 --- a/src/gateway/test-helpers.acquisition.test.ts +++ b/src/gateway/test-helpers.acquisition.test.ts @@ -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) { + 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()), + resetPreparedGatewayModelCatalogForTest: vi.fn(), +})); +vi.mock("../test-utils/ports.js", async (importOriginal) => ({ + ...(await importOriginal()), + 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) { - 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(); const requests: ReturnType[] = []; 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()), - resetPreparedGatewayModelCatalogForTest: vi.fn(), - })); - vi.doMock("../test-utils/ports.js", async (importOriginal) => ({ - ...(await importOriginal()), - 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()), - 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 }) diff --git a/src/gateway/worker-environments/provider-crabbox-runtime-preflight.test.ts b/src/gateway/worker-environments/provider-crabbox-runtime-preflight.test.ts index 1d9049322015..4eb3d9688c3c 100644 --- a/src/gateway/worker-environments/provider-crabbox-runtime-preflight.test.ts +++ b/src/gateway/worker-environments/provider-crabbox-runtime-preflight.test.ts @@ -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 { } describe("Crabbox runtime preflight cleanup", () => { - support.setupWorkerEnvironmentServiceSuite(); + support.setupWorkerEnvironmentServiceSuite({ reuseReadWorkers: true }); const pluginServices: OpenClawPluginService[] = []; async function registerProvider(): Promise { 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)), diff --git a/src/gateway/worker-environments/service.test-support.ts b/src/gateway/worker-environments/service.test-support.ts index acf1e4da168d..8a156bf72d47 100644 --- a/src/gateway/worker-environments/service.test-support.ts +++ b/src/gateway/worker-environments/service.test-support.ts @@ -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 }, }); diff --git a/src/gateway/worker-environments/worker-turn-launcher.test-support.ts b/src/gateway/worker-environments/worker-turn-launcher.test-support.ts index 113d525ed7f8..1bdd0e77eb73 100644 --- a/src/gateway/worker-environments/worker-turn-launcher.test-support.ts +++ b/src/gateway/worker-environments/worker-turn-launcher.test-support.ts @@ -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 { sessionFile = SESSION_KEY; } -export async function cleanupWorkerTurnLauncherTest(): Promise { +export function cleanupWorkerTurnLauncherTest(): Promise; +export function cleanupWorkerTurnLauncherTest(options: { + reuseReadWorkers: boolean; +}): Promise; +export async function cleanupWorkerTurnLauncherTest( + options: { reuseReadWorkers?: boolean } = {}, +): Promise { 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(); diff --git a/src/gateway/worker-environments/workspace-result-repository.test.ts b/src/gateway/worker-environments/workspace-result-repository.test.ts index a93318502fd5..ac0d285b6a8e 100644 --- a/src/gateway/worker-environments/workspace-result-repository.test.ts +++ b/src/gateway/worker-environments/workspace-result-repository.test.ts @@ -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(); let closeNode: (() => Promise) | 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, diff --git a/src/test-utils/session-state-cleanup.ts b/src/test-utils/session-state-cleanup.ts index 88b7562f6581..06d2f795a87d 100644 --- a/src/test-utils/session-state-cleanup.ts +++ b/src/test-utils/session-state-cleanup.ts @@ -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 { await drainSessionStoreWriterQueuesForTest(); @@ -20,6 +21,12 @@ export async function cleanupSessionStateForTest( } await drainFileLockStateForTest(); clearSessionStoreCacheForTest(); +} + +export async function cleanupSessionStateForTest( + options: { stateDir?: string; rootPath?: string } = {}, +): Promise { + await drainSessionStateForTest(options); if (!options.stateDir) { return; }