import { createPluginRuntimeMock } from "openclaw/plugin-sdk/plugin-test-runtime"; import { beforeEach, describe, expect, it, vi } from "vitest"; import type { CodexSessionCatalogControl, CodexSessionCatalogControlFactory, } from "../session-catalog-types.js"; import type { CodexThreadForkParams, CodexTurn } from "./protocol.js"; import type { CodexAppServerBindingStore } from "./session-binding.js"; const boundaryMocks = vi.hoisted(() => ({ listTurns: vi.fn(), })); const linkMocks = vi.hoisted(() => ({ delete: vi.fn(), upsert: vi.fn(), })); const transcriptMocks = vi.hoisted(() => ({ importHistory: vi.fn(), })); const boundary = { beforeTurnId: "turn-2", targetTurnId: "turn-2", retainedMarker: { turnId: "turn-1", userMessageCount: 1 }, } as const; vi.mock("openclaw/plugin-sdk/session-catalog", async (importOriginal) => ({ ...(await importOriginal()), deleteSessionUpstreamLink: linkMocks.delete, upsertSessionUpstreamLink: linkMocks.upsert, })); vi.mock("./transcript-mirror.js", () => ({ importCodexThreadHistoryToTranscript: transcriptMocks.importHistory, })); vi.mock("./upstream-fork-boundary.js", () => ({ resolveCodexUpstreamForkBoundary: vi.fn(async () => ({ ok: true, boundary, editorText: "edit me", })), listCodexUpstreamTurns: boundaryMocks.listTurns, precheckCodexUpstreamForkBoundary: vi.fn(() => ({ ok: true, boundary })), })); import { forkCodexUpstreamSession } from "./upstream-session-fork.js"; function turn(id: string, text: string): CodexTurn { return { id, status: "completed", items: [ { aggregatedOutput: null, changes: [], command: null, cwd: null, id: `${id}-user`, name: null, query: null, server: null, status: null, text: "", title: null, tool: null, content: [{ type: "text", text, textElements: [] }], type: "userMessage", }, ], }; } function forkResponse(threadId = "thread-forked") { return { approvalPolicy: "never", approvalsReviewer: "user", cwd: "/tmp", model: "gpt-5.4", modelProvider: "openai", sandbox: { type: "dangerFullAccess" }, thread: { id: threadId, sessionId: "session-forked", cliVersion: "0.148.0", createdAt: 1715299200, updatedAt: 1715299200, cwd: "/tmp", ephemeral: false, modelProvider: "openai", preview: "forked thread", source: "appServer", status: { type: "notLoaded" }, turns: [], }, }; } function forkParams() { return { targetKey: "agent:main:dashboard:forked", source: { agentId: "main", sessionId: "session-source", sessionKey: "agent:main:source", storePath: "/tmp/sessions.db", entryId: "entry-2", }, upstream: { catalogId: "codex", hostId: "gateway:local", kind: "codex-app-server" as const, threadId: "thread-source", ref: { connectionFingerprint: "fingerprint", threadId: "thread-source" }, }, }; } type ForkThreadStub = (params: CodexThreadForkParams) => Promise; function factoryForControl(control: CodexSessionCatalogControl): CodexSessionCatalogControlFactory { return { forRequest: () => control, homesForAgent: () => [], forUpstream: (_agentId, fingerprint) => fingerprint === control.connectionFingerprint ? control : undefined, }; } function forkControl( forkThread: ForkThreadStub = vi.fn(async () => forkResponse()), connectionFingerprint = "fingerprint", ) { const archiveThread = vi.fn(async () => undefined); const control = { archiveThread, clientId: "client-pinned", connectionFingerprint, forkThread, } as unknown as CodexSessionCatalogControl; control.withPinnedConnection = async (run) => await run(control); return { archiveThread, control, controlFactory: factoryForControl(control), forkThread }; } beforeEach(() => { boundaryMocks.listTurns.mockReset(); linkMocks.delete.mockReset(); linkMocks.upsert.mockReset().mockReturnValue(true); transcriptMocks.importHistory.mockReset().mockResolvedValue({ importedMessages: 1, omittedMessages: 0, }); }); describe("forkCodexUpstreamSession", () => { it("verifies the cut, imports the fork history, then links before binding", async () => { const retainedTurn = turn("turn-1", "one"); boundaryMocks.listTurns .mockResolvedValueOnce([turn("turn-2", "edit me")]) .mockResolvedValueOnce([retainedTurn]); const { archiveThread, control, controlFactory, forkThread } = forkControl(); const events: string[] = []; linkMocks.upsert.mockImplementation(() => { events.push("link"); return true; }); const mutate = vi.fn(async () => { events.push("bind"); return true; }); const runtime = createPluginRuntimeMock(); const createSessionEntry = vi.mocked(runtime.agent.session.createSessionEntry); const result = await forkCodexUpstreamSession(forkParams(), { bindingStore: { mutate } as unknown as CodexAppServerBindingStore, controlFactory, harnessRuntimeId: "codex-custom", resolveConfig: () => ({}), runtime, }); expect(forkThread).toHaveBeenCalledWith({ threadId: "thread-source", beforeTurnId: "turn-2", excludeTurns: true, }); expect(boundaryMocks.listTurns).toHaveBeenLastCalledWith(control, "thread-forked"); expect(transcriptMocks.importHistory).toHaveBeenCalledWith( expect.objectContaining({ sessionKey: "agent:main:dashboard:forked", thread: expect.objectContaining({ id: "thread-forked", turns: [retainedTurn] }), throughTurnId: "turn-1", }), ); expect(linkMocks.upsert).toHaveBeenCalledWith( expect.objectContaining({ marker: { turnId: "turn-1", userMessageCount: 1 }, sessionKey: "agent:main:dashboard:forked", threadId: "thread-forked", }), ); expect(runtime.agent.session.createSessionEntry).toHaveBeenCalledWith( expect.objectContaining({ initialEntry: expect.objectContaining({ agentHarnessId: "codex-custom" }), }), ); expect(createSessionEntry.mock.calls[0]?.[0]).not.toHaveProperty("recoverMatchingInitialEntry"); expect(events).toEqual(["link", "bind"]); expect(result).toEqual({ status: "created", key: "agent:main:dashboard:forked", editorText: "edit me", }); expect(archiveThread).not.toHaveBeenCalled(); }); it("forks through the secondary home selected by the upstream fingerprint", async () => { boundaryMocks.listTurns .mockResolvedValueOnce([turn("turn-2", "edit me")]) .mockResolvedValueOnce([turn("turn-1", "one")]); const primary = forkControl(undefined, "primary-fingerprint"); const secondary = forkControl(undefined, "secondary-fingerprint"); const forUpstream = vi.fn((_agentId: string, fingerprint: string) => fingerprint === secondary.control.connectionFingerprint ? secondary.control : undefined, ); const factory = { ...primary.controlFactory, forUpstream }; const params = forkParams(); params.upstream.ref = { connectionFingerprint: "secondary-fingerprint", threadId: params.upstream.threadId, }; await expect( forkCodexUpstreamSession(params, { bindingStore: { mutate: vi.fn(async () => true) } as unknown as CodexAppServerBindingStore, controlFactory: factory, harnessRuntimeId: "codex", resolveConfig: () => ({}), runtime: createPluginRuntimeMock(), }), ).resolves.toMatchObject({ status: "created", key: params.targetKey }); expect(forUpstream).toHaveBeenCalledWith("main", "secondary-fingerprint"); expect(secondary.forkThread).toHaveBeenCalledOnce(); expect(primary.forkThread).not.toHaveBeenCalled(); }); it("fails closed when the upstream fingerprint does not resolve to a known home", async () => { const { controlFactory, forkThread } = forkControl(); const params = forkParams(); params.upstream.ref = { connectionFingerprint: "unknown-fingerprint", threadId: params.upstream.threadId, }; await expect( forkCodexUpstreamSession(params, { bindingStore: {} as CodexAppServerBindingStore, controlFactory, harnessRuntimeId: "codex", runtime: createPluginRuntimeMock(), }), ).resolves.toEqual({ status: "failed", code: "upstream-unavailable", message: "This Codex thread is not available on the current connection. Reconnect to its host and try again.", }); expect(forkThread).not.toHaveBeenCalled(); expect(boundaryMocks.listTurns).not.toHaveBeenCalled(); }); it("forks incognito sessions ephemerally and imports history from the live response", async () => { const retainedTurn = turn("turn-1", "one"); boundaryMocks.listTurns.mockResolvedValueOnce([turn("turn-2", "edit me")]); const forkThread = vi.fn(async () => { const response = forkResponse(); return { ...response, thread: { ...response.thread, ephemeral: true, turns: [retainedTurn] }, }; }); const { controlFactory } = forkControl(forkThread); const params = forkParams(); params.targetKey = "agent:main:dashboard:incognito-forked"; const mutate = vi.fn(async () => true); await expect( forkCodexUpstreamSession(params, { bindingStore: { mutate } as unknown as CodexAppServerBindingStore, controlFactory, harnessRuntimeId: "codex", resolveConfig: () => ({}), runtime: createPluginRuntimeMock(), }), ).resolves.toMatchObject({ status: "created", key: params.targetKey }); expect(forkThread).toHaveBeenCalledWith({ threadId: "thread-source", beforeTurnId: "turn-2", ephemeral: true, excludeTurns: false, }); expect(boundaryMocks.listTurns).toHaveBeenCalledTimes(1); expect(transcriptMocks.importHistory).toHaveBeenCalledWith( expect.objectContaining({ sessionKey: params.targetKey, thread: expect.objectContaining({ turns: expect.arrayContaining([expect.objectContaining({ id: retainedTurn.id })]), }), }), ); expect(mutate).toHaveBeenCalledWith( expect.anything(), expect.objectContaining({ kind: "set", binding: expect.objectContaining({ clientId: "client-pinned" }), }), ); }); it("archives a fork whose read-back history proves beforeTurnId was ignored", async () => { boundaryMocks.listTurns .mockResolvedValueOnce([turn("turn-2", "edit me")]) .mockResolvedValueOnce([turn("turn-1", "one"), turn("turn-2", "edit me")]); const { archiveThread, controlFactory } = forkControl(); const runtime = createPluginRuntimeMock(); const result = await forkCodexUpstreamSession(forkParams(), { bindingStore: { mutate: vi.fn() } as unknown as CodexAppServerBindingStore, controlFactory, harnessRuntimeId: "codex", runtime, }); expect(result).toMatchObject({ status: "failed", code: "upstream-unavailable", message: expect.stringContaining("Codex version"), }); expect(archiveThread).toHaveBeenCalledWith("thread-forked"); expect(runtime.agent.session.createSessionEntry).not.toHaveBeenCalled(); expect(linkMocks.upsert).not.toHaveBeenCalled(); }); it("cleans the link and archives the fork when binding materialization fails", async () => { boundaryMocks.listTurns .mockResolvedValueOnce([turn("turn-2", "edit me")]) .mockResolvedValueOnce([turn("turn-1", "one")]); const { archiveThread, controlFactory } = forkControl(); const mutate = vi.fn(async () => false); const result = await forkCodexUpstreamSession(forkParams(), { bindingStore: { mutate } as unknown as CodexAppServerBindingStore, controlFactory, harnessRuntimeId: "codex", runtime: createPluginRuntimeMock(), }); expect(result).toMatchObject({ status: "failed", code: "upstream-unavailable" }); expect(linkMocks.delete).toHaveBeenCalledWith("agent:main:dashboard:forked", "main"); expect(mutate).toHaveBeenLastCalledWith(expect.anything(), { kind: "clear", threadId: "thread-forked", }); expect(archiveThread).toHaveBeenCalledWith("thread-forked"); }); it("archives a recoverable orphan id when the fork response is invalid", async () => { boundaryMocks.listTurns.mockResolvedValueOnce([turn("turn-2", "edit me")]); const { archiveThread, controlFactory } = forkControl( vi.fn(async () => ({ thread: { id: "thread-orphan" } })), ); const result = await forkCodexUpstreamSession(forkParams(), { bindingStore: {} as CodexAppServerBindingStore, controlFactory, harnessRuntimeId: "codex", runtime: createPluginRuntimeMock(), }); expect(result).toMatchObject({ status: "failed", code: "upstream-unavailable" }); expect(archiveThread).toHaveBeenCalledWith("thread-orphan"); }); it("rejects a fork response that reuses the source thread id", async () => { boundaryMocks.listTurns.mockResolvedValueOnce([turn("turn-2", "edit me")]); const { archiveThread, controlFactory } = forkControl( vi.fn(async () => forkResponse("thread-source")), ); const result = await forkCodexUpstreamSession(forkParams(), { bindingStore: { mutate: vi.fn() } as unknown as CodexAppServerBindingStore, controlFactory, harnessRuntimeId: "codex", runtime: createPluginRuntimeMock(), }); expect(result).toMatchObject({ status: "failed", code: "upstream-unavailable" }); expect(archiveThread).not.toHaveBeenCalled(); }); });