// Imported by agent.test.ts to keep its mocked suite in one Vitest module graph. import fs from "node:fs/promises"; import { afterEach, describe, expect, it, vi } from "vitest"; import { ErrorCodes } from "../../../packages/gateway-protocol/src/index.js"; import type { CronCreatorAuthorityCapability } from "../../agents/cron-creator-authority-context.js"; import { createAgentRunDirectAbortError, createAgentRunRestartAbortError, isAgentRunDirectAbortReason, isAgentRunRestartAbortReason, } from "../../agents/run-termination.js"; import { beginSessionWorkAdmission, cancelSessionWorkAdmissionHandoff, interruptSessionWorkAdmissions, runExclusiveSessionLifecycleMutation, } from "../../sessions/session-lifecycle-admission.js"; import { createDeferredCore } from "../../shared/deferred.js"; import { withTestDir } from "../../test-helpers/temp-dir.js"; import { normalizeSessionDeliveryState } from "../../utils/delivery-context.shared.js"; import { getAgentTestMocks, makeContext, type AgentHandlerArgs, type AgentParams, type AgentCommandCall, setDateOnlyFakeClockActive, waitForAssertion, requireValue, expectRecordFields, expectSqliteSessionFileMarkerForEntry, mockCallArg, expectRespondError, flushScheduledDispatchStep, mockMainSessionEntry, buildExistingMainStoreEntry, useTestStateDir, primeMainAgentRun, runMainAgent, runMainAgentAndCaptureEntry, backendGatewayClient, operatorWriteCliClient, waitForAgentCommandCall, waitForAgentCommandCallAfter, invokeAgent, describe0AfterEach0, } from "./agent.test-harness.js"; import { handleChatAbortRequest } from "./chat-abort-handler.js"; import type { GatewayRequestContext } from "./types.js"; const mocks = getAgentTestMocks(); describe("gateway agent handler", () => { afterEach(describe0AfterEach0); it.each(["cleared", "no-op", "discarded result"])( "returns the replacement projection result with a %s store fixture", async (mode) => { mocks.updateSessionStore.mockReset(); if (mode === "no-op") { mocks.updateSessionStore.mockResolvedValue(undefined); } else if (mode === "discarded result") { mocks.updateSessionStore.mockImplementation(async (_path, updater) => { await updater({}); }); } const result = { transition: "empty-fixture" }; const update = vi.fn(async () => ({ result })); await expect( mocks.applySessionEntryReplacements({ storePath: "/tmp/sessions.json", update }), ).resolves.toBe(result); expect(update).toHaveBeenCalledOnce(); expect(update).toHaveBeenCalledWith([]); }, ); it("keeps backing store replacements when its delegate discards the result", async () => { const sessionKey = "agent:main:main"; const entry = { sessionId: "session-1", updatedAt: 1 }; const store = { [sessionKey]: entry }; const replacement = { ...entry, updatedAt: 2 }; const result = { transition: "replaced" }; mocks.updateSessionStore.mockImplementation(async (_path, updater) => { await updater(store); }); const update = vi.fn(async () => ({ replacements: [{ sessionKey, entry: replacement }], result, })); await expect( mocks.applySessionEntryReplacements({ storePath: "/tmp/sessions.json", sessionKeys: [sessionKey], update, }), ).resolves.toBe(result); expect(update).toHaveBeenCalledWith([{ sessionKey, entry }]); expect(store[sessionKey]).toEqual(replacement); expect(store[sessionKey]).not.toBe(replacement); }); it("does not repeat a replacement projection whose result is undefined", async () => { mocks.updateSessionStore.mockImplementation(async (_path, updater) => { await updater({}); }); const update = vi.fn(async () => ({ result: undefined })); await expect( mocks.applySessionEntryReplacements({ storePath: "/tmp/sessions.json", update }), ).resolves.toBeUndefined(); expect(update).toHaveBeenCalledOnce(); }); it("does not run a replacement projection after its store fixture rejects", async () => { const failure = new Error("fixture write rejected"); mocks.updateSessionStore.mockRejectedValueOnce(failure); const update = vi.fn(async () => ({ result: "unexpected projection" })); await expect( mocks.applySessionEntryReplacements({ storePath: "/tmp/sessions.json", update }), ).rejects.toBe(failure); expect(update).not.toHaveBeenCalled(); }); it("passes resolved maintenance config to the gateway admission store write", async () => { primeMainAgentRun({ cfg: { session: { maintenance: { mode: "enforce", maxEntries: 42, }, }, }, }); await runMainAgent("hi", "idem-maintenance-config"); const updateOptions = mocks.updateSessionStore.mock.calls.at(-1)?.[2]; expect(updateOptions).toMatchObject({ takeCacheOwnership: true, maintenanceConfig: { mode: "enforce", maxEntries: 42, }, }); }); it("carries exact cron creator authority through direct local agent RPC", async () => { const runId = "direct-agent-cron-authority"; let capability: CronCreatorAuthorityCapability | undefined; primeMainAgentRun(); mocks.agentCommand.mockImplementation(async (opts: AgentCommandCall) => { capability = opts.cronCreatorAuthorityCapability as | CronCreatorAuthorityCapability | undefined; expect(capability).toMatchObject({ active: true, runId }); return { payloads: [{ text: "ok" }], meta: { durationMs: 100 } }; }); await invokeAgent( { message: "create an automation", agentId: "main", sessionKey: "agent:main:main", idempotencyKey: runId, }, { client: { ...operatorWriteCliClient(["operator.admin"]), internal: { isLocalClient: true }, } as AgentHandlerArgs["client"], }, ); await waitForAssertion(() => expect(capability?.active).toBe(false)); }); it("resolves explicit recipient sessions before Gateway admission", async () => { const sessionKey = "agent:ops:whatsapp:work:direct:+15551234567"; mocks.listAgentIds.mockReturnValue(["main", "ops"]); mocks.loadConfigReturn = { session: { dmScope: "per-account-channel-peer" } }; mocks.resolveAgentExplicitRecipientSession.mockResolvedValue({ sessionKey, channel: "whatsapp", to: "user:+15551234567", accountId: "work", threadId: "topic-42", }); mocks.loadSessionEntry.mockReturnValue({ cfg: mocks.loadConfigReturn, storePath: "/tmp/sessions.json", entry: { sessionId: "recipient-session", updatedAt: Date.now() }, canonicalKey: sessionKey, }); let persistedEntry: Record | undefined; mocks.updateSessionStore.mockImplementation(async (_path, updater) => { const store = { [sessionKey]: { sessionId: "recipient-session", updatedAt: Date.now() }, }; const result = await updater(store); persistedEntry = store[sessionKey]; return result; }); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); await invokeAgent({ message: "hi", agentId: "ops", channel: "whatsapp", to: "+15551234567", threadId: "topic-42", idempotencyKey: "recipient-session-route", }); expect(mocks.resolveAgentExplicitRecipientSession).toHaveBeenCalledWith({ cfg: mocks.loadConfigReturn, agentId: "ops", channel: "whatsapp", to: "+15551234567", accountId: undefined, threadId: "topic-42", }); const call = await waitForAgentCommandCall<{ sessionKey?: string; channel?: string; to?: string; accountId?: string; threadId?: string; }>(); expect(call.sessionKey).toBe(sessionKey); expect(call).toMatchObject({ channel: "whatsapp", to: "user:+15551234567", accountId: "work", threadId: "topic-42", }); expect(persistedEntry?.delivery).toEqual( normalizeSessionDeliveryState({ context: { channel: "whatsapp", to: "user:+15551234567", accountId: "work", threadId: "topic-42", }, }), ); }); it.each(["webchat", "LAST"])( "keeps the agent main session for non-deliverable channel hint %s", async (channel) => { const sessionKey = "agent:ops:main"; mocks.listAgentIds.mockReturnValue(["main", "ops"]); mocks.resolveExplicitAgentSessionKey.mockReturnValue(sessionKey); mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: "/tmp/sessions.json", entry: { sessionId: "ops-main", updatedAt: Date.now() }, canonicalKey: sessionKey, }); mocks.updateSessionStore.mockImplementation( async (_path, updater) => await updater({ [sessionKey]: { sessionId: "ops-main", updatedAt: Date.now() }, }), ); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); await invokeAgent({ message: "hi", agentId: "ops", channel, to: "+15551234567", idempotencyKey: `non-deliverable-${channel}`, }); expect(mocks.resolveAgentExplicitRecipientSession).not.toHaveBeenCalled(); const call = await waitForAgentCommandCall<{ sessionKey?: string }>(); expect(call.sessionKey).toBe(sessionKey); }, ); it("dedupes retries while explicit recipient session routing is pending", async () => { const sessionKey = "agent:ops:whatsapp:work:direct:+15551234567"; const runId = "recipient-session-route-pending"; const { promise: routePending, resolve: finishRoute } = createDeferredCore<{ sessionKey: string; }>(); mocks.listAgentIds.mockReturnValue(["main", "ops"]); mocks.loadConfigReturn = { session: { dmScope: "per-account-channel-peer" } }; mocks.resolveAgentExplicitRecipientSession.mockReturnValue(routePending); mocks.loadSessionEntry.mockReturnValue({ cfg: mocks.loadConfigReturn, storePath: "/tmp/sessions.json", entry: { sessionId: "recipient-session", updatedAt: Date.now() }, canonicalKey: sessionKey, }); mocks.updateSessionStore.mockImplementation( async (_path, updater) => await updater({ [sessionKey]: { sessionId: "recipient-session", updatedAt: Date.now() }, }), ); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); const context = makeContext(); const request = { message: "hi", agentId: "ops", channel: "whatsapp", accountId: "work", to: "+15551234567", idempotencyKey: runId, } satisfies AgentParams; const first = invokeAgent(request, { context, reqId: runId }); await waitForAssertion(() => { expectRecordFields(context.dedupe.get(`agent:${runId}`)?.payload, { runId, status: "accepted", }); }); expect(context.dedupe.get(`agent:${runId}`)?.payload).not.toHaveProperty("sessionKey"); const duplicateRespond = vi.fn(); await invokeAgent(request, { context, reqId: `${runId}-duplicate`, respond: duplicateRespond, flushDispatch: false, }); expect(duplicateRespond).toHaveBeenCalledWith( true, { runId, status: "in_flight", agentId: "ops", admissionPending: true }, undefined, { cached: true, runId, }, ); expect(mocks.resolveAgentExplicitRecipientSession).toHaveBeenCalledTimes(1); finishRoute({ sessionKey }); await first; }); it("honors owner cancellation while explicit recipient session routing is pending", async () => { const sessionKey = "agent:ops:whatsapp:work:direct:+15551234567"; const runId = "recipient-session-route-abort"; const { promise: routePending, resolve: finishRoute } = createDeferredCore<{ sessionKey: string; }>(); mocks.listAgentIds.mockReturnValue(["main", "ops"]); mocks.loadConfigReturn = { session: { dmScope: "per-account-channel-peer" } }; mocks.resolveAgentExplicitRecipientSession.mockReturnValue(routePending); mocks.loadSessionEntry.mockReturnValue({ cfg: mocks.loadConfigReturn, storePath: "/tmp/sessions.json", entry: { sessionId: "recipient-session", updatedAt: Date.now() }, canonicalKey: sessionKey, }); mocks.updateSessionStore.mockImplementation( async (_path, updater) => await updater({ [sessionKey]: { sessionId: "recipient-session", updatedAt: Date.now() }, }), ); mocks.agentCommand.mockClear(); const context = makeContext(); const ownerClient = { connId: "owner-conn" } as AgentHandlerArgs["client"]; const pending = invokeAgent( { message: "hi", agentId: "ops", channel: "whatsapp", to: "+15551234567", idempotencyKey: runId, }, { context, client: ownerClient, reqId: runId, flushDispatch: false }, ); await waitForAssertion(() => { expectRecordFields(context.dedupe.get(`agent:${runId}`)?.payload, { runId, agentId: "ops", status: "accepted", }); }); expect(context.dedupe.get(`agent:${runId}`)?.payload).not.toHaveProperty("sessionKey"); const abortRespond = vi.fn(); await handleChatAbortRequest({ params: { sessionKey: "agent:ops:main", runId }, respond: abortRespond as never, context, req: { type: "req", id: "abort-recipient-route", method: "chat.abort" }, client: ownerClient, isWebchatConnect: () => false, }); expectRecordFields(mockCallArg(abortRespond, 0, 1), { aborted: true, runIds: [runId], }); expectRecordFields(context.dedupe.get(`agent:${runId}`)?.payload, { runId, status: "timeout", stopReason: "rpc", }); expect(context.dedupe.get(`agent:${runId}`)?.payload).not.toHaveProperty("sessionKey"); finishRoute({ sessionKey }); await pending; await flushScheduledDispatchStep(); expect(mocks.agentCommand).not.toHaveBeenCalled(); }); it("clears pending dedupe when explicit recipient session routing fails", async () => { const runId = "recipient-session-route-error"; mocks.listAgentIds.mockReturnValue(["main", "ops"]); mocks.resolveAgentExplicitRecipientSession.mockResolvedValue({ error: new Error("ambiguous recipient"), }); mocks.agentCommand.mockClear(); const context = makeContext(); const respond = vi.fn(); await invokeAgent( { message: "hi", agentId: "ops", channel: "whatsapp", to: "team", idempotencyKey: runId, }, { context, respond, reqId: runId, flushDispatch: false }, ); expectRespondError(respond, { code: ErrorCodes.INVALID_REQUEST, message: "ambiguous recipient", }); expect(context.dedupe.has(`agent:${runId}`)).toBe(false); expect(mocks.agentCommand).not.toHaveBeenCalled(); }); it("clears a sessionless reservation when content validation fails", async () => { const runId = "sessionless-invalid-reply-channel"; const context = makeContext(); const respond = vi.fn(); await invokeAgent( { message: "hi", replyChannel: "unknown-channel", idempotencyKey: runId, }, { context, respond, reqId: runId, flushDispatch: false }, ); expectRespondError(respond, { message: "invalid agent params: unknown channel: unknown-channel", }); expect(context.dedupe.has(`agent:${runId}`)).toBe(false); expect(mocks.agentCommand).not.toHaveBeenCalled(); }); it.each([ { state: "archived", entry: { archivedAt: 1 }, reason: "is archived. Restore it before starting new work.", }, { state: "project preparation", entry: { pendingProjectGitUrl: "https://github.com/openclaw/openclaw.git" }, reason: "workspace is not ready. Wait for setup to finish or retry in chat.", }, { state: "worktree preparation", entry: { pendingWorktree: { workspace: "/tmp/project", titleSource: "Prepare workspace", }, }, reason: "workspace is not ready. Wait for setup to finish or retry in chat.", }, ])( "clears pending dedupe when the routed recipient session awaits $state", async ({ entry, reason }) => { const sessionKey = "agent:ops:whatsapp:direct:+15551234567"; const runId = "recipient-session-route-archived"; mocks.listAgentIds.mockReturnValue(["main", "ops"]); mocks.resolveAgentExplicitRecipientSession.mockResolvedValue({ sessionKey }); mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: "/tmp/sessions.json", entry: { sessionId: "recipient-session", updatedAt: 1, ...entry }, canonicalKey: sessionKey, }); mocks.agentCommand.mockClear(); const context = makeContext(); const respond = vi.fn(); await invokeAgent( { message: "hi", agentId: "ops", channel: "whatsapp", to: "+15551234567", idempotencyKey: runId, }, { context, respond, reqId: runId, flushDispatch: false }, ); expectRespondError(respond, { message: `Session "${sessionKey}" ${reason}`, }); expect(context.dedupe.has(`agent:${runId}`)).toBe(false); expect(mocks.agentCommand).not.toHaveBeenCalled(); }, ); it("rejects agent RPC creation in an agent harness-owned namespace", async () => { const sessionKey = "agent:main:harness:codex:supervision:native-thread"; const runId = "agent-harness-reserved"; mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: "/tmp/sessions.json", entry: undefined, canonicalKey: sessionKey, }); mocks.agentCommand.mockClear(); const context = makeContext(); const respond = vi.fn(); await invokeAgent( { message: "claim reserved session", agentId: "main", sessionKey, idempotencyKey: runId, }, { context, respond, reqId: runId, flushDispatch: false }, ); expectRespondError(respond, { code: ErrorCodes.INVALID_REQUEST, message: "Session key namespace is reserved for agent harness-owned sessions.", }); expect(context.dedupe.has(`agent:${runId}`)).toBe(false); expect(mocks.agentCommand).not.toHaveBeenCalled(); }); it.each(["agent:main:harness:codex:supervision:native-thread", "agent:main:ordinary-locked"])( "rejects agent RPC session-id rotation for locked session %s", async (sessionKey) => { const runId = "agent-harness-session-id-rotation"; mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: "/tmp/sessions.json", entry: { agentHarnessId: "codex", modelSelectionLocked: true, sessionId: "native-session", }, canonicalKey: sessionKey, }); mocks.agentCommand.mockClear(); const updateSessionStoreCallsBefore = mocks.updateSessionStore.mock.calls.length; const context = makeContext(); const respond = vi.fn(); await invokeAgent( { message: "replace native transcript identity", agentId: "main", sessionKey, sessionId: "replacement-session", idempotencyKey: runId, }, { context, respond, reqId: runId, flushDispatch: false }, ); expectRespondError(respond, { code: ErrorCodes.INVALID_REQUEST, message: "Agent harness-owned session identity is locked and cannot be replaced or shared.", }); expect(context.dedupe.has(`agent:${runId}`)).toBe(false); expect(mocks.updateSessionStore).toHaveBeenCalledTimes(updateSessionStoreCallsBefore); expect(mocks.agentCommand).not.toHaveBeenCalled(); }, ); it("rejects one-shot model runs against harness-owned sessions", async () => { const sessionKey = "agent:main:harness:codex:supervision:native-thread"; const runId = "agent-harness-model-run"; mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: "/tmp/sessions.json", entry: { agentHarnessId: "codex", modelSelectionLocked: true, sessionId: "native-session", }, canonicalKey: sessionKey, }); mocks.agentCommand.mockClear(); const context = makeContext(); const respond = vi.fn(); await invokeAgent( { message: "run through another model", agentId: "main", sessionKey, modelRun: true, idempotencyKey: runId, }, { context, respond, reqId: runId, flushDispatch: false }, ); expectRespondError(respond, { code: ErrorCodes.INVALID_REQUEST, message: "Agent harness-owned sessions cannot be used for one-shot model runs.", }); expect(context.dedupe.has(`agent:${runId}`)).toBe(false); expect(mocks.agentCommand).not.toHaveBeenCalled(); }); it("allows raw model runs against grandfathered unlocked harness-prefixed sessions", async () => { const sessionKey = "agent:main:harness:notes"; const runId = "legacy-harness-model-run"; mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: "/tmp/sessions.json", entry: { agentHarnessId: "codex", modelSelectionLocked: false, sessionId: "legacy-session", updatedAt: Date.now(), }, canonicalKey: sessionKey, }); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "pong" }], meta: { durationMs: 100 }, }); await invokeAgent( { message: "Reply exactly: pong", agentId: "main", sessionKey, modelRun: true, promptMode: "none", idempotencyKey: runId, }, { reqId: runId, client: operatorWriteCliClient(), }, ); expectRecordFields(await waitForAgentCommandCall(), { modelRun: true, promptMode: "none", sessionEffects: "internal", sessionId: "legacy-session", sessionKey, }); }); it.each([false, true])( "responds with the admission error after model cleanup (cleanup fails: %s)", async (cleanupFails) => { primeMainAgentRun(); const context = makeContext(); const runId = `admission-cleanup-${cleanupFails}`; const inputError = new Error("input persistence failed"); const cleanupError = new Error("model cleanup failed"); const cleanupEntered = createDeferredCore(); const releaseCleanup = createDeferredCore(); const order: string[] = []; const dispose = vi.fn(async () => { order.push("cleanup started"); cleanupEntered.resolve(); await releaseCleanup.promise; order.push("cleanup finished"); if (cleanupFails) { throw cleanupError; } }); const runtime = await import("../../agents/prepared-model-runtime.js"); const acquire = vi.mocked(runtime.acquireAgentRunPreparedModelRuntime); const createLease = requireValue( acquire.getMockImplementation(), "model lease fixture missing", ); acquire.mockImplementationOnce(async (...args) => ({ ...(await createLease(...args)), [Symbol.asyncDispose]: dispose, })); mocks.stageSessionPendingInput.mockRejectedValueOnce(inputError); const respond = vi.fn(() => order.push("response")); const outcome = invokeAgent( { message: "persist this input", agentId: "main", sessionKey: "agent:main:main", idempotencyKey: runId, }, { context, respond, flushDispatch: false }, ).then( () => undefined, (error: unknown) => error, ); try { await Promise.race([ cleanupEntered.promise, outcome.then(() => { throw new Error("Admission completed without joining model cleanup"); }), ]); expect(mocks.stageSessionPendingInput).toHaveBeenCalledOnce(); expect(order).toEqual(["cleanup started"]); expect(respond).not.toHaveBeenCalled(); expect(mocks.agentCommand).not.toHaveBeenCalled(); releaseCleanup.resolve(); await expect(outcome).resolves.toBe(cleanupFails ? cleanupError : undefined); expect(dispose).toHaveBeenCalledOnce(); expect(order).toEqual(["cleanup started", "cleanup finished", "response"]); expect(respond.mock.calls).toEqual([ [false, undefined, { code: ErrorCodes.UNAVAILABLE, message: inputError.message }], ]); expect(mocks.agentCommand).not.toHaveBeenCalled(); expect(context.chatAbortControllers.has(runId)).toBe(false); expect(context.dedupe.has(`agent:${runId}`)).toBe(false); } finally { releaseCleanup.resolve(); await outcome; acquire.mockReset().mockImplementation(createLease); } }, ); it("passes a canonical user-turn recorder to gateway agent runs", async () => { primeMainAgentRun(); await runMainAgent("persist me", "idem-user-turn-recorder"); const call = await waitForAgentCommandCall< AgentCommandCall & { userTurnTranscriptRecorder?: { persistApproved: () => Promise; }; } >(); expect(call.userTurnTranscriptRecorder).toEqual( expect.objectContaining({ persistApproved: expect.any(Function), persistFallback: expect.any(Function), }), ); }); it("attributes gateway agent prompts to the authenticated connection user", async () => { primeMainAgentRun(); await invokeAgent( { message: "persist me", agentId: "main", sessionKey: "agent:main:main", idempotencyKey: "idem-attributed-user-turn-recorder", }, { reqId: "idem-attributed-user-turn-recorder", client: { ...requireValue(operatorWriteCliClient(), "expected operator client"), authenticatedUserId: "alice@example.com", }, }, ); const call = await waitForAgentCommandCall< AgentCommandCall & { userTurnTranscriptRecorder?: { message?: unknown } } >(); expect(call.userTurnTranscriptRecorder?.message).toMatchObject({ role: "user", content: "persist me", __openclaw: { senderId: "alice@example.com" }, }); }); it("dispatches with the session id reloaded during lifecycle admission", async () => { const sessionKey = "agent:main:main"; const initialSessionId = "session-before-reset"; const admittedSessionId = "session-after-reset"; let currentSessionId = initialSessionId; mocks.loadSessionEntry.mockImplementation(() => ({ cfg: {}, storePath: "/tmp/sessions.json", entry: { sessionId: currentSessionId, updatedAt: Date.now() }, canonicalKey: sessionKey, })); mocks.updateSessionStore.mockImplementation(async (_path, updater) => { const result = await updater({ [sessionKey]: { sessionId: initialSessionId, updatedAt: Date.now() }, }); currentSessionId = admittedSessionId; return result; }); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); await runMainAgent("hi", "idem-reset-before-admission"); const call = await waitForAgentCommandCall<{ sessionId?: string }>(); expect(call.sessionId).toBe(admittedSessionId); }); it("does not recreate a session deleted before lifecycle admission", async () => { const sessionKey = "agent:main:main"; let deleted = false; mocks.loadSessionEntry.mockImplementation(() => ({ cfg: {}, storePath: "/tmp/sessions.json", entry: deleted ? undefined : { sessionId: "session-before-delete", updatedAt: Date.now() }, canonicalKey: sessionKey, })); mocks.updateSessionStore.mockImplementation(async (_path, updater) => { const result = await updater({ [sessionKey]: { sessionId: "session-before-delete", updatedAt: Date.now() }, }); deleted = true; return result; }); mocks.agentCommand.mockClear(); const respond = vi.fn(); await invokeAgent( { message: "do not recreate the deleted session", agentId: "main", sessionKey, idempotencyKey: "idem-deleted-before-admission", }, { respond, reqId: "idem-deleted-before-admission" }, ); expectRespondError(respond, { message: `Session "${sessionKey}" was deleted while starting work. Retry.`, }); expect(mocks.agentCommand).not.toHaveBeenCalled(); }); it("does not recreate a session deleted before its initial store touch", async () => { const sessionKey = "agent:main:main"; mocks.loadSessionEntry.mockImplementation(() => ({ cfg: {}, storePath: "/tmp/sessions.json", entry: { sessionId: "session-before-delete", updatedAt: Date.now() }, canonicalKey: sessionKey, })); mocks.updateSessionStore.mockImplementation(async (_path, updater) => await updater({})); mocks.agentCommand.mockClear(); const respond = vi.fn(); await invokeAgent( { message: "do not recreate the deleted session", agentId: "main", sessionKey, idempotencyKey: "idem-deleted-before-touch", }, { respond, reqId: "idem-deleted-before-touch" }, ); expectRespondError(respond, { message: `Session "${sessionKey}" was deleted while starting work. Retry.`, }); expect(mocks.agentCommand).not.toHaveBeenCalled(); }); it("does not recreate a newly touched session deleted before lifecycle admission", async () => { const sessionKey = "agent:main:slack:group:new-session"; let storedEntry: Record | undefined; mocks.loadSessionEntry.mockImplementation(() => ({ cfg: {}, storePath: "/tmp/sessions.json", entry: storedEntry, canonicalKey: sessionKey, })); mocks.updateSessionStore.mockImplementation(async (_path, updater) => { const result = await updater({}); storedEntry = undefined; return result; }); mocks.agentCommand.mockClear(); const respond = vi.fn(); await invokeAgent( { message: "do not revive the newly deleted session", agentId: "main", sessionKey, idempotencyKey: "idem-new-session-deleted-before-admission", }, { respond, reqId: "idem-new-session-deleted-before-admission" }, ); expectRespondError(respond, { message: `Session "${sessionKey}" was deleted while starting work. Retry.`, }); expect(mocks.agentCommand).not.toHaveBeenCalled(); }); it("blocks the initial store touch and preserves an abort while admission waits", async () => { primeMainAgentRun(); mocks.agentCommand.mockClear(); mocks.updateSessionStore.mockClear(); const sessionKey = "agent:main:main"; const runId = "idem-abort-during-admission"; let releaseMutation = () => {}; const { promise: mutationStarted, resolve: markMutationStarted } = createDeferredCore(); const mutation = runExclusiveSessionLifecycleMutation({ scope: "/tmp/sessions.json", identities: [sessionKey, "existing-session-id"], run: async () => { markMutationStarted(); await new Promise((release) => { releaseMutation = release; }); }, }); await mutationStarted; const context = makeContext(); const respond = vi.fn(); const request = invokeAgent( { message: "do not run after abort", agentId: "main", sessionKey, idempotencyKey: runId, }, { context, respond, reqId: runId }, ); await waitForAssertion(() => { expect(context.dedupe.has(`agent:${runId}`)).toBe(true); }); expect(mocks.updateSessionStore).not.toHaveBeenCalled(); const abortRespond = vi.fn(); await handleChatAbortRequest({ params: { sessionKey, runId }, respond: abortRespond as never, context, req: { type: "req", id: "abort-req", method: "chat.abort" }, client: null, isWebchatConnect: () => false, }); releaseMutation(); await mutation; await request; expect(mockCallArg(abortRespond)).toBe(true); expect(mocks.agentCommand).not.toHaveBeenCalled(); expect(context.chatAbortControllers.has(runId)).toBe(false); expect( respond.mock.calls.some( ([ok, payload]) => ok === true && (payload as { status?: string })?.status === "timeout", ), ).toBe(true); }); it("does not revive an expired reservation after lifecycle admission waits", async () => { primeMainAgentRun(); mocks.agentCommand.mockClear(); const sessionKey = "agent:main:main"; const runId = "idem-expired-during-admission"; let releaseMutation = () => {}; const { promise: mutationStarted, resolve: markMutationStarted } = createDeferredCore(); const mutation = runExclusiveSessionLifecycleMutation({ scope: "/tmp/sessions.json", identities: [sessionKey, "existing-session-id"], run: async () => { markMutationStarted(); await new Promise((release) => { releaseMutation = release; }); }, }); await mutationStarted; const context = makeContext(); const respond = vi.fn(); const request = invokeAgent( { message: "do not run after queue expiry", agentId: "main", sessionKey, idempotencyKey: runId, }, { context, respond, reqId: runId }, ); const dedupeKey = `agent:${runId}`; await waitForAssertion(() => { expect(context.dedupe.has(dedupeKey)).toBe(true); }); const reserved = context.dedupe.get(dedupeKey); context.dedupe.set(dedupeKey, { ts: reserved?.ts ?? Date.now(), ok: true, payload: { ...(reserved?.payload as Record), expiresAtMs: Date.now() - 1, }, }); releaseMutation(); await mutation; await request; expect(mocks.agentCommand).not.toHaveBeenCalled(); expect(context.chatAbortControllers.has(runId)).toBe(false); expect( respond.mock.calls.some( ([ok, payload]) => ok === true && (payload as { status?: string; stopReason?: string })?.status === "timeout" && (payload as { stopReason?: string }).stopReason === "timeout", ), ).toBe(true); }); it("preserves a newer terminal reservation result after lifecycle admission waits", async () => { primeMainAgentRun(); mocks.agentCommand.mockClear(); const sessionKey = "agent:main:main"; const runId = "idem-terminal-during-admission"; let releaseMutation = () => {}; const { promise: mutationStarted, resolve: markMutationStarted } = createDeferredCore(); const mutation = runExclusiveSessionLifecycleMutation({ scope: "/tmp/sessions.json", identities: [sessionKey, "existing-session-id"], run: async () => { markMutationStarted(); await new Promise((release) => { releaseMutation = release; }); }, }); await mutationStarted; const context = makeContext(); const respond = vi.fn(); const request = invokeAgent( { message: "do not overwrite the replacement result", agentId: "main", sessionKey, idempotencyKey: runId, }, { context, respond, reqId: runId }, ); const dedupeKey = `agent:${runId}`; await waitForAssertion(() => { expect(context.dedupe.has(dedupeKey)).toBe(true); }); const replacement = { ts: Date.now(), ok: true, payload: { runId, status: "ok", summary: "replacement completed" }, }; context.dedupe.set(dedupeKey, replacement); releaseMutation(); await mutation; await request; expect(context.dedupe.get(dedupeKey)).toBe(replacement); expect(mocks.agentCommand).not.toHaveBeenCalled(); expect(respond).toHaveBeenLastCalledWith( true, replacement.payload, undefined, expect.objectContaining({ cached: true, runId }), ); }); it("runs gateway agent work inside its lifecycle admission context", async () => { primeMainAgentRun(); const sessionKey = "agent:main:main"; const sessionId = "existing-session-id"; mocks.agentCommand.mockImplementationOnce(async () => { await expect( interruptSessionWorkAdmissions({ scope: "/tmp/sessions.json", identities: [sessionKey, sessionId], timeoutMs: 5, }), ).resolves.toBe(true); return { payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }; }); await invokeAgent({ message: "run in-band lifecycle work", agentId: "main", sessionKey, idempotencyKey: "idem-agent-admission-context", }); expect(mocks.agentCommand).toHaveBeenCalledOnce(); }); it.each(["restart", "rpc"] as const)( "adopts a recovery admission interrupted by %s before the RPC", async (stopReason) => { const reason = stopReason === "rpc" ? createAgentRunDirectAbortError() : undefined; const sessionKey = "agent:main:main"; const sessionId = "existing-session-id"; const runId = "idem-recovery-admission-handoff"; const scope = "/tmp/sessions.json"; primeMainAgentRun({ sessionId }); mocks.agentCommand.mockClear(); const admission = await beginSessionWorkAdmission({ scope, identities: [sessionKey, sessionId], assertAllowed: () => {}, }); const handoffId = admission.createHandoff(); const { promise: mutationStarted, resolve: markMutationStarted } = createDeferredCore(); let mutationRan = false; const mutation = runExclusiveSessionLifecycleMutation({ scope, identities: [sessionKey, sessionId], prepare: async () => { markMutationStarted(); expect( await interruptSessionWorkAdmissions({ scope, identities: [sessionKey, sessionId], reason, timeoutMs: 1_000, }), ).toBe(true); }, run: async () => { mutationRan = true; }, }); await mutationStarted; const respond = vi.fn(); try { await invokeAgent( { message: "resume the admitted turn", agentId: "main", sessionKey, expectedExistingSessionId: sessionId, internalRuntimeHandoffId: handoffId, idempotencyKey: runId, }, { client: backendGatewayClient(), flushDispatch: false, reqId: runId, respond, }, ); await mutation; expect(cancelSessionWorkAdmissionHandoff(handoffId)).toBe(false); expect(mutationRan).toBe(true); expect(mocks.agentCommand).not.toHaveBeenCalled(); expect(respond).toHaveBeenCalledWith( true, expect.objectContaining({ runId, status: "timeout", stopReason, providerStarted: false }), undefined, expect.objectContaining({ cached: true, runId }), ); } finally { cancelSessionWorkAdmissionHandoff(handoffId); admission.release(); await mutation; } }, ); it.each(["generic", "explicit restart", "terminal Stop", "already stopped"] as const)( "preserves Gateway lifecycle interruption semantics: %s", async (interruption) => { primeMainAgentRun(); const sessionKey = "agent:main:main"; const sessionId = "existing-session-id"; const terminal = interruption === "terminal Stop" || interruption === "already stopped"; const reason = terminal ? createAgentRunDirectAbortError() : interruption === "explicit restart" ? createAgentRunRestartAbortError() : undefined; const settled = createDeferredCore(); let observedAbortReason: unknown; mocks.agentCommand.mockImplementationOnce(async (opts: { abortSignal?: AbortSignal }) => { await new Promise((resolve) => { const finish = () => { observedAbortReason = opts.abortSignal?.reason; resolve(); }; if (opts.abortSignal?.aborted) { finish(); return; } opts.abortSignal?.addEventListener("abort", finish, { once: true }); }); await settled.promise; throw observedAbortReason; }); const context = makeContext(); const runId = `idem-agent-lifecycle-${interruption}`; await invokeAgent( { message: "interrupt this admitted turn", agentId: "main", sessionKey, idempotencyKey: runId, }, { context, reqId: runId }, ); await waitForAgentCommandCall(); const abortEntry = requireValue( context.chatAbortControllers.get(runId), "admitted controller", ); if (interruption === "already stopped") { abortEntry.abortStopReason = "rpc"; abortEntry.controller.abort(reason); } const draining = interruptSessionWorkAdmissions({ scope: "/tmp/sessions.json", identities: [sessionKey, sessionId], reason: interruption === "already stopped" ? createAgentRunRestartAbortError() : reason, }); try { expect(abortEntry.abortStopReason).toBe(terminal ? "rpc" : "restart"); if (terminal) { expect(abortEntry.controller.signal.reason).toBe(reason); } } finally { settled.resolve(); await draining; } await flushScheduledDispatchStep(); expect(isAgentRunRestartAbortReason(observedAbortReason)).toBe(!terminal); expect(isAgentRunDirectAbortReason(observedAbortReason)).toBe(terminal); if (terminal) { expect(observedAbortReason).toBe(reason); } expectRecordFields(context.dedupe.get(`agent:${runId}`)?.payload, { runId, status: "timeout", stopReason: terminal ? "rpc" : "restart", }); }, ); it("does not mutate a session archived before the initial store update", async () => { primeMainAgentRun(); mocks.agentCommand.mockClear(); const sessionKey = "agent:main:main"; const archivedEntry = { sessionId: "existing-session-id", updatedAt: Date.now(), archivedAt: Date.now(), }; const store = { [sessionKey]: structuredClone(archivedEntry) }; mocks.updateSessionStore.mockImplementation(async (_path, updater) => await updater(store)); const respond = vi.fn(); await invokeAgent( { message: "do not touch the archive", agentId: "main", sessionKey, idempotencyKey: "idem-archive-before-store-update", }, { respond, reqId: "idem-archive-before-store-update" }, ); expectRespondError(respond, { message: `Session "${sessionKey}" is archived. Restore it before starting new work.`, }); expect(store).toEqual({ [sessionKey]: archivedEntry }); expect(mocks.agentCommand).not.toHaveBeenCalled(); }); it("uses the freshest alias when checking archive state before migration", async () => { const cfg = { session: { mainKey: "work" }, agents: { list: [{ id: "main", default: true }] }, }; mocks.loadConfigReturn = cfg; mocks.loadSessionEntry.mockReturnValue({ cfg, storePath: "/tmp/sessions.json", entry: { sessionId: "restored-session", updatedAt: 2 }, canonicalKey: "agent:main:work", // Post-flip the accessor target carries the alias set prepared at request // start; freshest-entry resolution happens inside the patch transaction. storeKeys: ["agent:main:work", "agent:main:main"], }); const store = { "agent:main:work": { sessionId: "archived-session", updatedAt: 1, archivedAt: 1, }, "agent:main:main": { sessionId: "restored-session", updatedAt: 2, }, }; mocks.updateSessionStore.mockImplementation(async (_path, updater) => await updater(store)); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); await invokeAgent({ message: "continue restored session", agentId: "main", sessionKey: "agent:main:main", idempotencyKey: "idem-restored-alias", }); expect(mocks.agentCommand).toHaveBeenCalled(); expect(store["agent:main:work"]?.archivedAt).toBeUndefined(); expect(store["agent:main:main"]).toBeUndefined(); }); it("does not re-persist a committed gateway user turn after the session key is rebound", async () => { primeMainAgentRun({ sessionId: "accepted-session-id" }); mocks.agentCommand.mockImplementationOnce( async (opts: { userTurnTranscriptRecorder?: { persistApproved: () => Promise }; }) => { await requireValue( opts.userTurnTranscriptRecorder, "user turn recorder missing", ).persistApproved(); return { payloads: [{ text: "ok" }], meta: { durationMs: 100 } }; }, ); await runMainAgent("stale after reset", "idem-user-turn-rebound"); const call = await waitForAgentCommandCall< AgentCommandCall & { userTurnTranscriptRecorder?: { hasPersisted: () => boolean; persistApproved: () => Promise; }; } >(); expect(call.userTurnTranscriptRecorder?.hasPersisted()).toBe(true); expect(mocks.persistSessionTranscriptTurn).toHaveBeenCalledTimes(1); mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: "/tmp/sessions.json", entry: { sessionId: "new-session-id", updatedAt: Date.now(), }, canonicalKey: "agent:main:main", store: { "agent:main:main": { sessionId: "new-session-id", updatedAt: Date.now(), }, }, }); await expect(call.userTurnTranscriptRecorder?.persistApproved()).resolves.toBeDefined(); expect(mocks.persistSessionTranscriptTurn).toHaveBeenCalledTimes(1); }); it("durably admits managed media for inline image agent runs", async () => { await withTestDir({ prefix: "openclaw-gateway-agent-inline-image-" }, async (root) => { useTestStateDir(root); mockMainSessionEntry({ sessionId: "existing-session-id", model: "vision-model", modelProvider: "test", modelOverride: "vision-model", modelOverrideSource: "user", providerOverride: "test", }); mocks.updateSessionStore.mockResolvedValue(undefined); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); const context = { ...makeContext(), loadGatewayModelCatalog: vi.fn(async () => [ { id: "vision-model", name: "vision-model", provider: "test", input: ["image"], }, ]), } as unknown as GatewayRequestContext; await invokeAgent( { message: "describe this image", agentId: "main", sessionKey: "agent:main:main", idempotencyKey: "idem-image-user-turn-recorder", attachments: [ { type: "file", mimeType: "image/png", fileName: "test.png", content: Buffer.from("fake-png-data").toString("base64"), }, ], }, { context, reqId: "idem-image-user-turn-recorder" }, ); const call = await waitForAgentCommandCall< AgentCommandCall & { userTurnTranscriptRecorder?: { message?: { __openclaw?: Record }; hasPersisted: () => boolean; }; } >(); expect(call.images).toEqual([ expect.objectContaining({ type: "image", mimeType: "image/png", }), ]); expect(call.userTurnTranscriptRecorder?.hasPersisted()).toBe(false); expect(mocks.persistSessionTranscriptTurn).not.toHaveBeenCalled(); expect(call.userTurnTranscriptRecorder?.message?.["__openclaw"]).toMatchObject({ media: [expect.objectContaining({ contentType: "image/png", kind: "image" })], mediaImageLayout: { slots: [{ kind: "inline", factIndex: 0 }] }, }); }); }); it("durably admits managed media for offloaded image agent runs", async () => { await withTestDir({ prefix: "openclaw-gateway-agent-offloaded-image-" }, async (root) => { useTestStateDir(root); mockMainSessionEntry({ sessionId: "existing-session-id", model: "vision-model", modelProvider: "test", modelOverride: "vision-model", modelOverrideSource: "user", providerOverride: "test", }); mocks.updateSessionStore.mockResolvedValue(undefined); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); const context = { ...makeContext(), loadGatewayModelCatalog: vi.fn(async () => [ { id: "vision-model", name: "vision-model", provider: "test", input: ["image"], }, ]), } as unknown as GatewayRequestContext; await invokeAgent( { message: "describe this large image", agentId: "main", sessionKey: "agent:main:main", idempotencyKey: "idem-offloaded-image-user-turn-recorder", attachments: [ { type: "file", mimeType: "image/png", fileName: "large.png", content: Buffer.alloc(2_000_001, 1).toString("base64"), }, ], }, { context, reqId: "idem-offloaded-image-user-turn-recorder" }, ); const call = await waitForAgentCommandCall< AgentCommandCall & { userTurnTranscriptRecorder?: { message?: { __openclaw?: Record }; hasPersisted: () => boolean; }; } >(); expect(call.images).toEqual([]); expect(call.imageOrder).toEqual(["offloaded"]); expect(call.message).toContain("[media attached: media://inbound/"); expect(call.userTurnTranscriptRecorder?.hasPersisted()).toBe(false); expect(mocks.persistSessionTranscriptTurn).not.toHaveBeenCalled(); expect(call.userTurnTranscriptRecorder?.message?.["__openclaw"]).toMatchObject({ media: [expect.objectContaining({ contentType: "image/png", kind: "image" })], mediaImageLayout: { slots: [{ kind: "offloaded", factIndex: 0 }] }, }); }); }); it("preserves ACP metadata from the current stored session entry", async () => { const existingAcpMeta = { backend: "acpx", agent: "codex", runtimeSessionName: "runtime-1", mode: "persistent", state: "idle", lastActivityAt: Date.now(), }; mockMainSessionEntry({ acp: existingAcpMeta, }); let capturedEntry: Record | undefined; mocks.updateSessionStore.mockImplementation(async (_path, updater) => { const store: Record = { "agent:main:main": buildExistingMainStoreEntry({ acp: existingAcpMeta }), }; const result = await updater(store); capturedEntry = store["agent:main:main"] as Record; return result; }); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); await runMainAgent("test", "test-idem-acp-meta"); expect(mocks.updateSessionStore).toHaveBeenCalled(); expect(requireValue(capturedEntry, "updated session entry missing").acp).toEqual( existingAcpMeta, ); }); it("clears automatic recovery quarantine state when a user turn rotates the session id", async () => { vi.useFakeTimers({ toFake: ["Date"] }); setDateOnlyFakeClockActive(true); vi.setSystemTime(new Date("2026-05-07T12:00:00.000Z")); const staleEntry = { sessionId: "quarantined-session-id", updatedAt: 0, sessionStartedAt: 0, lastInteractionAt: 0, abortedLastRun: true, restartRecoveryRuns: [ { runId: "initial-wedged-run", lifecycleGeneration: "gen-1" }, { runId: "recovery-run-1", lifecycleGeneration: "gen-2" }, ], mainRestartRecovery: { automaticAttempts: 2, lastAttemptAt: 3, lastRunId: "recovery-run-1", }, subagentRecovery: { automaticAttempts: 2, lastAttemptAt: 3, wedgedAt: 4, wedgedReason: "automatic_attempt_budget_exceeded", }, }; mockMainSessionEntry(staleEntry); const capturedEntry = await runMainAgentAndCaptureEntry("test-idem-rotated-recovery-clear"); expect(capturedEntry.sessionId).not.toBe("quarantined-session-id"); expect(capturedEntry.abortedLastRun).toBeUndefined(); expect(capturedEntry.restartRecoveryRuns).toBeUndefined(); expect(capturedEntry.mainRestartRecovery).toBeUndefined(); expect(capturedEntry.subagentRecovery).toEqual({ automaticAttempts: 2, lastAttemptAt: 3, wedgedAt: 4, wedgedReason: "automatic_attempt_budget_exceeded", }); }); it("drops a stale transcript path when a stale session rotates ids", async () => { vi.useFakeTimers({ toFake: ["Date"] }); setDateOnlyFakeClockActive(true); vi.setSystemTime(new Date("2026-05-07T12:00:00.000Z")); const staleEntry = { sessionId: "old-session-id", sessionFile: "/tmp/openclaw/agents/main/sessions/old-session-id.jsonl", updatedAt: 0, sessionStartedAt: 0, }; mockMainSessionEntry(staleEntry); let capturedEntry: Record | undefined; mocks.updateSessionStore.mockImplementation(async (_path, updater) => { const store: Record = { "agent:main:main": { ...staleEntry }, }; const result = await updater(store); capturedEntry = result as Record; return result; }); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); await runMainAgent("test", "test-idem-stale-transcript"); expect(capturedEntry?.sessionId).not.toBe("old-session-id"); expectSqliteSessionFileMarkerForEntry(capturedEntry); }); it("rotates a failed session instead of resuming when its transcript is missing", async () => { const now = Date.parse("2026-05-18T09:45:00.000Z"); vi.useFakeTimers({ toFake: ["Date"] }); setDateOnlyFakeClockActive(true); vi.setSystemTime(now); const missingTranscriptEntry = { sessionId: "failed-missing-session-id", sessionFile: "/tmp/openclaw/missing/failed-missing-session-id.jsonl", status: "failed", updatedAt: now, sessionStartedAt: now, lastInteractionAt: now, startedAt: now - 2_000, endedAt: now - 1_000, runtimeMs: 1_000, abortedLastRun: true, }; mockMainSessionEntry(missingTranscriptEntry); const capturedEntry = await runMainAgentAndCaptureEntry("test-idem-failed-missing-transcript"); const call = await waitForAgentCommandCall<{ sessionId?: string }>(); expect(call.sessionId).not.toBe("failed-missing-session-id"); expect(capturedEntry?.sessionId).not.toBe("failed-missing-session-id"); expect(capturedEntry?.status).toBeUndefined(); expect(capturedEntry?.startedAt).toBeUndefined(); expect(capturedEntry?.endedAt).toBeUndefined(); expect(capturedEntry?.runtimeMs).toBeUndefined(); expect(capturedEntry?.abortedLastRun).toBeUndefined(); expectSqliteSessionFileMarkerForEntry(capturedEntry); }); it.each([ { name: "status-done row", status: "done" as const, expectReuse: true }, { name: "status-killed row", status: "killed" as const, expectReuse: false }, { name: "endedAt-only row", status: undefined, expectReuse: false }, ])( "handles a terminal main session from a $name when its transcript is newer", async (scenario) => { const now = Date.parse("2026-05-18T09:47:00.000Z"); vi.useFakeTimers({ toFake: ["Date"] }); setDateOnlyFakeClockActive(true); vi.setSystemTime(now); mocks.readTranscriptMutationStateSync.mockReturnValue({ observedAt: null, updatedAt: now - 1_000, }); await withTestDir({ prefix: "openclaw-gateway-terminal-main-newer-" }, async (root) => { const sessionsDir = `${root}/sessions`; const sessionFile = "terminal-main-session.jsonl"; mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: `${sessionsDir}/sessions.json`, entry: { sessionId: "terminal-main-session", sessionFile, ...(scenario.status ? { status: scenario.status } : {}), updatedAt: now - 10_000, sessionStartedAt: now - 60_000, lastInteractionAt: now - 10_000, startedAt: now - 20_000, endedAt: now - 15_000, runtimeMs: 5_000, cliSessionBindings: { "claude-cli": { sessionId: "old-claude-cli-session" }, "codex-cli": { sessionId: "old-codex-cli-session" }, }, cliSessionIds: { "claude-cli": "old-claude-cli-session", "codex-cli": "old-codex-cli-session", }, claudeCliSessionId: "old-claude-cli-session", }, canonicalKey: "agent:main:main", }); const commandCallCount = mocks.agentCommand.mock.calls.length; const capturedEntry = await runMainAgentAndCaptureEntry( "test-idem-terminal-main-newer-transcript", ); const call = await waitForAgentCommandCallAfter<{ sessionId?: string }>(commandCallCount); if (scenario.expectReuse) { expect(call.sessionId).toBe("terminal-main-session"); expect(capturedEntry?.sessionId).toBe("terminal-main-session"); expect(mocks.readTranscriptMutationStateSync).not.toHaveBeenCalled(); return; } expect(call.sessionId).not.toBe("terminal-main-session"); expect(capturedEntry?.sessionId).not.toBe("terminal-main-session"); expect(capturedEntry?.status).toBeUndefined(); expect(capturedEntry?.startedAt).toBeUndefined(); expect(capturedEntry?.endedAt).toBeUndefined(); expect(capturedEntry?.runtimeMs).toBeUndefined(); expectSqliteSessionFileMarkerForEntry(capturedEntry); expect(capturedEntry?.cliSessionBindings).toBeUndefined(); expect(capturedEntry?.cliSessionIds).toBeUndefined(); expect(capturedEntry?.claudeCliSessionId).toBeUndefined(); }); }, ); it("reuses terminal main sessions when the fresh store row has the transcript marker", async () => { const now = Date.parse("2026-05-18T09:47:30.000Z"); vi.useFakeTimers({ toFake: ["Date"] }); setDateOnlyFakeClockActive(true); vi.setSystemTime(now); await withTestDir({ prefix: "openclaw-gateway-terminal-main-fresh-marker-" }, async (root) => { const sessionsDir = `${root}/sessions`; await fs.mkdir(sessionsDir, { recursive: true }); const sessionFile = "terminal-main-session.jsonl"; const transcriptPath = `${sessionsDir}/${sessionFile}`; await fs.writeFile( transcriptPath, `${JSON.stringify({ type: "session", id: "terminal-main-session" })}\n`, "utf8", ); await fs.utimes(transcriptPath, new Date(now - 1_000), new Date(now - 1_000)); const staleEntry = { sessionId: "terminal-main-session", sessionFile, status: "done", updatedAt: now - 10_000, cliSessionBindings: { "claude-cli": { sessionId: "existing-claude-cli-session" }, }, cliSessionIds: { "claude-cli": "existing-claude-cli-session", }, claudeCliSessionId: "existing-claude-cli-session", }; mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: `${sessionsDir}/sessions.json`, entry: staleEntry, canonicalKey: "agent:main:main", }); let capturedEntry: Record | undefined; mocks.updateSessionStore.mockImplementation(async (_path, updater) => { const store = { "agent:main:main": { ...staleEntry, updatedAt: now, }, }; const result = await updater(store); capturedEntry = result as Record; return result; }); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); await runMainAgent("hi", "test-idem-terminal-main-fresh-marker"); const call = await waitForAgentCommandCall<{ sessionId?: string }>(); expect(call.sessionId).toBe("terminal-main-session"); expect(capturedEntry?.sessionId).toBe("terminal-main-session"); expectSqliteSessionFileMarkerForEntry(capturedEntry); expect(capturedEntry?.cliSessionIds).toEqual({ "claude-cli": "existing-claude-cli-session", }); expect(capturedEntry?.claudeCliSessionId).toBe("existing-claude-cli-session"); }); }); it("honors explicit gateway session-id resumes for terminal main rows", async () => { const now = Date.parse("2026-05-18T09:48:00.000Z"); vi.useFakeTimers({ toFake: ["Date"] }); setDateOnlyFakeClockActive(true); vi.setSystemTime(now); await withTestDir( { prefix: "openclaw-gateway-terminal-main-explicit-resume-" }, async (root) => { const sessionsDir = `${root}/sessions`; await fs.mkdir(sessionsDir, { recursive: true }); const sessionFile = "terminal-main-session.jsonl"; const transcriptPath = `${sessionsDir}/${sessionFile}`; await fs.writeFile( transcriptPath, `${JSON.stringify({ type: "session", id: "terminal-main-session" })}\n`, "utf8", ); await fs.utimes(transcriptPath, new Date(now - 1_000), new Date(now - 1_000)); const existingEntry = { sessionId: "terminal-main-session", sessionFile, status: "done", updatedAt: now - 10_000, sessionStartedAt: now - 60_000, lastInteractionAt: now - 10_000, startedAt: now - 20_000, endedAt: now - 15_000, runtimeMs: 5_000, }; mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: `${sessionsDir}/sessions.json`, entry: existingEntry, canonicalKey: "agent:main:main", }); let capturedEntry: Record | undefined; mocks.updateSessionStore.mockImplementation(async (_path, updater) => { const store: Record = { "agent:main:main": { ...existingEntry }, }; const result = await updater(store); capturedEntry = result as Record; return result; }); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); await invokeAgent({ message: "resume terminal main", agentId: "main", sessionKey: "agent:main:main", sessionId: "terminal-main-session", idempotencyKey: "test-idem-terminal-main-explicit-resume", } as AgentParams); const call = await waitForAgentCommandCall<{ sessionId?: string }>(); expect(call.sessionId).toBe("terminal-main-session"); expect(capturedEntry?.sessionId).toBe("terminal-main-session"); expectSqliteSessionFileMarkerForEntry(capturedEntry); expect(capturedEntry?.status).toBe("done"); expect(capturedEntry?.startedAt).toBe(now - 20_000); expect(capturedEntry?.endedAt).toBe(now - 15_000); expect(capturedEntry?.runtimeMs).toBe(5_000); }, ); }); it.each(["heartbeat", "cron"] as const)( "preserves terminal main session reuse for %s gateway runs", async (runKind) => { const now = Date.parse("2026-05-18T09:49:00.000Z"); vi.useFakeTimers({ toFake: ["Date"] }); setDateOnlyFakeClockActive(true); vi.setSystemTime(now); await withTestDir( { prefix: `openclaw-gateway-terminal-main-${runKind}-reuse-` }, async (root) => { const sessionsDir = `${root}/sessions`; await fs.mkdir(sessionsDir, { recursive: true }); const sessionFile = `terminal-main-${runKind}.jsonl`; const transcriptPath = `${sessionsDir}/${sessionFile}`; await fs.writeFile( transcriptPath, `${JSON.stringify({ type: "session", id: "terminal-main-session" })}\n`, "utf8", ); await fs.utimes(transcriptPath, new Date(now - 1_000), new Date(now - 1_000)); const existingEntry = { sessionId: "terminal-main-session", sessionFile, status: "done", updatedAt: now - 10_000, sessionStartedAt: now - 60_000, lastInteractionAt: now - 10_000, startedAt: now - 20_000, endedAt: now - 15_000, runtimeMs: 5_000, }; mocks.loadSessionEntry.mockReturnValue({ cfg: {}, storePath: `${sessionsDir}/sessions.json`, entry: existingEntry, canonicalKey: "agent:main:main", }); let capturedEntry: Record | undefined; mocks.updateSessionStore.mockImplementation(async (_path, updater) => { const store: Record = { "agent:main:main": { ...existingEntry }, }; const result = await updater(store); capturedEntry = result as Record; return result; }); mocks.agentCommand.mockResolvedValue({ payloads: [{ text: "ok" }], meta: { durationMs: 100 }, }); await invokeAgent({ message: `${runKind} probe`, agentId: "main", sessionKey: "agent:main:main", bootstrapContextRunKind: runKind, idempotencyKey: `test-idem-terminal-main-${runKind}-reuse`, } as AgentParams); const call = await waitForAgentCommandCall<{ sessionId?: string }>(); expect(call.sessionId).toBe("terminal-main-session"); expect(capturedEntry?.sessionId).toBe("terminal-main-session"); expectSqliteSessionFileMarkerForEntry(capturedEntry); }, ); }, ); it("rotates a failed session when its default transcript is missing", async () => { const now = Date.parse("2026-05-18T09:48:00.000Z"); vi.useFakeTimers({ toFake: ["Date"] }); setDateOnlyFakeClockActive(true); vi.setSystemTime(now); const missingDefaultTranscriptEntry = { sessionId: "failed-missing-default-session-id", status: "failed", updatedAt: now, sessionStartedAt: now, lastInteractionAt: now, }; mockMainSessionEntry(missingDefaultTranscriptEntry); const capturedEntry = await runMainAgentAndCaptureEntry( "test-idem-failed-missing-default-transcript", ); const call = await waitForAgentCommandCall<{ sessionId?: string }>(); expect(call.sessionId).not.toBe("failed-missing-default-session-id"); expect(capturedEntry?.sessionId).not.toBe("failed-missing-default-session-id"); expect(capturedEntry?.status).toBeUndefined(); expectSqliteSessionFileMarkerForEntry(capturedEntry); }); }); /* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */