import { join } from "node:path"; import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; import { createDeferred } from "../../../test/helpers/promise.js"; import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import type { ContextEngine, ContextEngineRuntimeContext } from "../../context-engine/types.js"; import { AsyncWorkScope } from "../../shared/async-work-scope.js"; import { createAssistant, createAssistantResultStream, testModel, } from "../sessions/agent-session-loop-correctness.test-support.js"; import type { ProviderConfigInput } from "../sessions/model-registry.js"; import { applyExtraParamsToAgentMock, hookRunner, limitHistoryTurnsMock, loadCompactHooksHarness, resetCompactHooksHarnessMocks, resolveContextEngineMock, resolveModelMock, } from "./compact.hooks.harness.js"; const { requestPreparedCompaction } = vi.hoisted(() => ({ requestPreparedCompaction: vi.fn(), })); vi.mock("@openclaw/ai/transports", async (importOriginal) => ({ ...(await importOriginal()), requestPreparedOpenAIResponsesCompaction: requestPreparedCompaction, })); let delegate: typeof import("../../context-engine/delegate.js").delegateCompactionToRuntime; let compactQueued: typeof import("./compact.queued.js").compactEmbeddedAgentSession; let sessions: typeof import("../sessions/index.js"); let accessor: typeof import("../../config/sessions/session-accessor.js"); let databases: typeof import("../../state/openclaw-agent-db.js"); let streamResolution: typeof import("./stream-resolution.js"); let replay: typeof import("../openai-transport-stream.test-support.js").testing; let accounting: typeof import("./run/compaction-accounting-bridge.js"); const tempDirs = useAutoCleanupTempDirTracker((cleanup) => afterEach(() => { databases.closeOpenClawAgentDatabasesForTest(); cleanup(); }), ); let workspaceDir: string; const model = { ...testModel, id: "compaction-fixture", provider: "openai", baseUrl: "https://api.openai.com/v1", contextWindow: 128_000, maxTokens: 1_024, }; const summary = "The deployment checklist was reviewed. Compare the remaining options."; beforeAll(async () => { // Reuse the hook harness's provider setup with the real constructor, SQLite // manager, transcript guard, and both compaction owners left intact. ({ compactEmbeddedAgentSession: compactQueued } = await loadCompactHooksHarness({ durableSession: true, })); [ { delegateCompactionToRuntime: delegate }, sessions, accessor, databases, streamResolution, { testing: replay }, accounting, ] = await Promise.all([ import("../../context-engine/delegate.js"), import("../sessions/index.js"), import("../../config/sessions/session-accessor.js"), import("../../state/openclaw-agent-db.js"), import("./stream-resolution.js"), import("../openai-transport-stream.test-support.js"), import("./run/compaction-accounting-bridge.js"), ]); }); beforeEach(async () => { workspaceDir = tempDirs.make("openclaw-compaction-delegate-"); resetCompactHooksHarnessMocks(workspaceDir); const actualScope = await vi.importActual("../agent-scope.js"); const scope = await import("../agent-scope.js"); vi.mocked(scope.listAgentEntries).mockImplementation(actualScope.listAgentEntries); vi.mocked(scope.resolveSessionAgentId).mockImplementation(actualScope.resolveSessionAgentId); vi.mocked(scope.resolveSessionAgentIds).mockImplementation(actualScope.resolveSessionAgentIds); requestPreparedCompaction.mockReset(); limitHistoryTurnsMock.mockImplementation((messages) => messages); hookRunner.hasHooks.mockImplementation( (name) => name === "before_compaction" || name === "after_compaction", ); }); async function createFixture(operation: "summary" | "endpoint", globalAlias = false) { const target = { agentId: globalAlias ? "marketing" : "main", sessionId: "compaction-session", sessionKey: globalAlias ? "global" : "agent:main:compaction-session", storePath: join(workspaceDir, "alternate", "openclaw-agent.sqlite"), }; const configuredStore = join(workspaceDir, "configured", "openclaw-agent.sqlite"); await accessor.upsertSessionEntryCore(target, { sessionId: target.sessionId, updatedAt: 1 }); // The same physical id in another store must not redirect partial-target resolution. const decoy = { ...target, storePath: configuredStore }; await accessor.upsertSessionEntryCore(decoy, { sessionId: target.sessionId, updatedAt: 1 }); const decoyManager = sessions.SessionManager.open(decoy, workspaceDir); decoyManager.appendMessage({ role: "user", content: "Unrelated store history", timestamp: 1 }); decoyManager.flushPendingPersistence(); const sessionManager = sessions.SessionManager.open(target, workspaceDir); sessionManager.appendModelChange(model.provider, model.id); sessionManager.appendThinkingLevelChange("off"); for (const content of [ "Review the deployment checklist.", "Compare the remaining options.", "Keep the rollout notes.", ]) { sessionManager.appendMessage({ role: "user", content, timestamp: 1 }); sessionManager.appendMessage( createAssistant(model, [{ type: "text", text: `Recorded: ${content}` }]), ); } sessionManager.flushPendingPersistence(); const originalMessages = sessionManager.buildSessionContext().messages; const originalEntries = sessionManager.getEntries(); const stream = vi.fn>((activeModel) => createAssistantResultStream(createAssistant(activeModel, [{ type: "text", text: summary }])), ); const authStorage = sessions.AuthStorage.inMemory(); authStorage.setRuntimeApiKey(model.provider, "test-api-key"); const modelRegistry = sessions.ModelRegistry.inMemory(authStorage); modelRegistry.registerProvider(model.provider, { api: model.api, streamSimple: stream }); resolveModelMock.mockImplementation((provider = model.provider, modelId = model.id) => ({ logicalRef: { provider, model: modelId }, model, error: null, authStorage, modelRegistry, })); vi.mocked(streamResolution.resolveEmbeddedAgentStream).mockReturnValue({ streamFn: stream, strategy: "session-custom", }); applyExtraParamsToAgentMock.mockReturnValue({ effectiveExtraParams: { responsesCompactEndpoint: operation === "endpoint" }, }); requestPreparedCompaction.mockResolvedValue({ item: { type: "compaction", id: "cmp_fixture", encrypted_content: "opaque-fixture" }, output: [{ type: "compaction", id: "cmp_fixture", encrypted_content: "opaque-fixture" }], historyMode: "compacted-prefix", usage: { input_tokens: 1_000, output_tokens: 200 }, model, replayMetadata: replay.buildOpenAIResponsesReasoningReplayMetadata(model, { sessionId: target.sessionId, }), }); const runtimeContext = { workspaceDir, provider: model.provider, model: model.id, trigger: "manual", thinkLevel: "off", config: { session: { store: configuredStore }, agents: { ownership: "explicit", list: [{ id: "main" }, { id: "marketing" }], defaults: { compaction: { mode: "default", keepRecentTokens: 1, postIndexSync: "off" } }, }, }, } satisfies ContextEngineRuntimeContext & { config: OpenClawConfig }; const recordUsage = vi.fn(); const recordCompaction = vi.fn(); accounting.attachCompactionAccountingRecorder(runtimeContext, { recordUsage, recordCompaction }); return { target, decoy, originalEntries, originalMessages, runtimeContext, stream, recordUsage, recordCompaction, }; } describe("direct compactor through the context-engine delegate", () => { it.each([ { operation: "summary", partial: false, threadId: "thread-route" }, { operation: "summary", partial: true, threadId: 0 }, { operation: "endpoint", partial: false, threadId: 0 }, { operation: "endpoint", partial: true, threadId: "thread-route" }, { operation: "summary", partial: false, threadId: undefined }, ] as const)( "returns durable $operation identity (partial=$partial, thread=$threadId)", async ({ operation, partial, threadId }) => { const fixture = await createFixture(operation, partial); const { target } = fixture; const sessionTarget = { ...(partial ? { agentId: target.agentId, storePath: target.storePath } : target), ...(threadId !== undefined ? { threadId } : {}), expectedWriterRunId: "caller-private-writer", expectedLifecycleRevision: "caller-private-revision", unrelatedCapability: "caller-private-capability", }; const result = await delegate({ sessionId: target.sessionId, sessionKey: target.sessionKey, sessionTarget, runtimeContext: fixture.runtimeContext, }); expect(result, JSON.stringify(result)).toMatchObject({ ok: true, compacted: true }); const returned = result.result?.sessionTarget; expect(returned).toEqual({ ...target, ...(threadId !== undefined ? { threadId } : {}) }); expect(result.result).not.toHaveProperty("sessionFile"); if ( !returned?.agentId || !returned.sessionId || !returned.sessionKey || !returned.storePath ) { throw new Error("Compactor must return its complete resolved identity"); } // Close the actual DB handles before observing the returned identity and history. databases.closeOpenClawAgentDatabasesForTest(); const reopened = sessions.SessionManager.open({ agentId: returned.agentId, sessionId: returned.sessionId, sessionKey: returned.sessionKey, storePath: returned.storePath, }); expect(reopened.getSessionId()).toBe(target.sessionId); expect(accessor.loadSessionEntry(target)?.sessionId).toBe(target.sessionId); const messages = reopened.buildSessionContext().messages; if (operation === "summary") { expect(result.result?.summary).toContain(summary); expect(reopened.getBranch().filter((entry) => entry.type === "compaction")).toHaveLength(1); expect(accessor.loadSessionEntry(target)?.compactionCheckpoints).toEqual([ expect.objectContaining({ sessionId: target.sessionId, sessionKey: target.sessionKey, summary: result.result?.summary, preCompaction: expect.objectContaining({ sessionId: target.sessionId }), postCompaction: expect.objectContaining({ sessionId: target.sessionId }), }), ]); const firstKeptIndex = fixture.originalEntries.findIndex( (entry) => entry.id === result.result?.firstKeptEntryId, ); expect(firstKeptIndex).toBeGreaterThan(0); const retained = fixture.originalEntries .slice(firstKeptIndex) .filter((entry) => entry.type === "message") .map((entry) => entry.message); expect(retained.length).toBeGreaterThan(0); expect(messages.slice(-retained.length)).toEqual(retained); for (const entry of fixture.originalEntries) { expect(reopened.getEntries()).toContainEqual(entry); } expect(fixture.stream).toHaveBeenCalled(); expect(requestPreparedCompaction).not.toHaveBeenCalled(); } else { expect(result.result?.summary).toBeUndefined(); expect(result.result?.details).toMatchObject({ compactionKind: "server-endpoint" }); expect(messages).toMatchObject(fixture.originalMessages); expect(messages.at(-1)).toMatchObject({ providerReplay: { data: "opaque-fixture", compactedWindow: { state: "ready" } }, }); expect(requestPreparedCompaction).toHaveBeenCalledOnce(); expect(fixture.stream).not.toHaveBeenCalled(); } expect(fixture.recordUsage).toHaveBeenCalled(); expect(fixture.recordCompaction).toHaveBeenCalledOnce(); expect(hookRunner.runBeforeCompaction).toHaveBeenCalledOnce(); expect(hookRunner.runAfterCompaction).toHaveBeenCalledOnce(); expect(sessions.SessionManager.open(fixture.decoy).buildSessionContext().messages).toEqual([ { role: "user", content: "Unrelated store history", timestamp: 1 }, ]); expect(accessor.loadSessionEntry(fixture.decoy)?.compactionCheckpoints).toBeUndefined(); }, ); it.each(["summary", "endpoint"] as const)( "cancels %s without publishing or rewriting history", async (operation) => { const fixture = await createFixture(operation); const started = createDeferred(); const stopped = createDeferred(); const controller = new AbortController(); const abortReason = new Error("caller cancelled compaction"); if (operation === "endpoint") { requestPreparedCompaction.mockImplementationOnce( async (_stream, _model, _context, options) => { started.resolve(); return await new Promise((_, reject) => { options.signal?.addEventListener( "abort", () => { stopped.resolve(); reject(abortReason); }, { once: true }, ); }); }, ); } else { const { createAssistantMessageEventStream } = await import("openclaw/plugin-sdk/llm"); fixture.stream.mockImplementationOnce((activeModel, _context, options) => { const stream = createAssistantMessageEventStream(); started.resolve(); options?.signal?.addEventListener( "abort", () => { stream.push({ type: "error", reason: "aborted", error: createAssistant(activeModel, [], "aborted"), }); stream.end(); stopped.resolve(); }, { once: true }, ); return stream; }); } const pending = delegate({ sessionId: fixture.target.sessionId, sessionKey: fixture.target.sessionKey, sessionTarget: { ...fixture.target, threadId: 0 }, runtimeContext: fixture.runtimeContext, abortSignal: controller.signal, }); await started.promise; controller.abort(abortReason); const result = await pending; await stopped.promise; expect(result).toMatchObject({ ok: false, compacted: false }); expect(result.result).toBeUndefined(); databases.closeOpenClawAgentDatabasesForTest(); const reopened = sessions.SessionManager.open(fixture.target); expect(reopened.getSessionId()).toBe(fixture.target.sessionId); expect(reopened.getEntries()).toEqual(fixture.originalEntries); expect(reopened.buildSessionContext().messages).toEqual(fixture.originalMessages); expect(fixture.recordCompaction).not.toHaveBeenCalled(); expect(hookRunner.runAfterCompaction).not.toHaveBeenCalled(); }, ); it("rejects contradictory physical identity without touching either durable transcript", async () => { const fixture = await createFixture("summary"); await expect( delegate({ sessionId: "contradictory-session", sessionKey: fixture.target.sessionKey, sessionTarget: fixture.target, runtimeContext: fixture.runtimeContext, }), ).rejects.toThrow("conflicts with the caller session identity"); expect(fixture.stream).not.toHaveBeenCalled(); expect(resolveModelMock).not.toHaveBeenCalled(); expect(hookRunner.runBeforeCompaction).not.toHaveBeenCalled(); databases.closeOpenClawAgentDatabasesForTest(); expect(sessions.SessionManager.open(fixture.target).getEntries()).toEqual( fixture.originalEntries, ); expect(sessions.SessionManager.open(fixture.decoy).buildSessionContext().messages).toEqual([ { role: "user", content: "Unrelated store history", timestamp: 1 }, ]); }); it.each(["summary", "endpoint"] as const)( "keeps queued manual %s compaction countable when cancellation follows its commit during a post-compaction hook", async (operation) => { const fixture = await createFixture(operation); const { markRuntimeCompactionDelegate } = await import("../../context-engine/compaction-watchdog.js"); const { incrementCompactionCount } = await import("../../auto-reply/reply/session-updates.js"); const backend = vi.fn(delegate); markRuntimeCompactionDelegate(backend); resolveContextEngineMock.mockResolvedValueOnce({ info: { ownsCompaction: false }, compact: backend, }); const expectedSession = accessor.loadSessionEntry(fixture.target); if (!expectedSession) { throw new Error("expected the queued compaction's persisted session"); } const hookEntered = createDeferred(); const releaseHook = createDeferred(); const expectCommittedTranscript = () => { const manager = sessions.SessionManager.open(fixture.target); if (operation === "summary") { const compactions = manager.getBranch().filter((entry) => entry.type === "compaction"); expect(compactions).toHaveLength(1); const committed = compactions[0]; if (!committed) { throw new Error("expected the persisted summary compaction"); } expect(committed.summary).toContain(summary); expect(committed.tokensBefore).toBeGreaterThan(0); return committed.tokensBefore; } expect(manager.buildSessionContext().messages.at(-1)).toMatchObject({ providerReplay: { data: "opaque-fixture", compactedWindow: { state: "ready" } }, }); return 1_000; }; hookRunner.runAfterCompaction.mockImplementationOnce(async () => { expectCommittedTranscript(); hookEntered.resolve(); await releaseHook.promise; }); const controller = new AbortController(); const work = new AsyncWorkScope(); const pending = work.run(() => compactQueued({ ...fixture.target, sessionTarget: fixture.target, sessionFile: fixture.target.sessionKey, workspaceDir, config: fixture.runtimeContext.config, provider: model.provider, model: model.id, agentHarnessId: "openclaw", trigger: "manual", abortSignal: controller.signal, enqueue: async (task) => await task(), }), ); try { await Promise.race([ hookEntered.promise, pending.then((result) => { throw new Error( `Queued compaction returned before its post-commit hook: ${JSON.stringify(result)}`, ); }), ]); controller.abort(new Error("caller stopped after compaction committed")); const result = await pending; expect.soft(result).toMatchObject({ ok: true, compacted: true, compactionKind: operation === "summary" ? "context-engine" : "server-endpoint", }); const count = result.ok && result.compacted ? await incrementCompactionCount({ agentId: fixture.target.agentId, sessionEntry: expectedSession, sessionStore: { [fixture.target.sessionKey]: expectedSession }, sessionKey: fixture.target.sessionKey, storePath: fixture.target.storePath, expectedSession, tokensAfter: result.result?.tokensAfter, compactionKind: result.compactionKind, }) : undefined; expect.soft(count).toBe(1); expect.soft(accessor.loadSessionEntry(fixture.target)?.compactionCount).toBe(1); expect(backend).toHaveBeenCalledOnce(); expect(result.result?.tokensBefore).toBe(expectCommittedTranscript()); expect(sessions.SessionManager.open(fixture.decoy).buildSessionContext().messages).toEqual([ { role: "user", content: "Unrelated store history", timestamp: 1 }, ]); } finally { releaseHook.resolve(); await pending.catch(() => undefined); await work.drain(); } }, ); });