Skip to content

Commit 0dfa22c

Browse files
authored
refactor: add embedded run session target seam (#90439)
1 parent 6f80552 commit 0dfa22c

23 files changed

Lines changed: 1409 additions & 338 deletions

src/agents/embedded-agent-runner/compact.queued.ts

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -14,9 +14,7 @@ import type {
1414
ContextEngineRuntimeSettings,
1515
} from "../../context-engine/types.js";
1616
import {
17-
captureCompactionCheckpointSnapshotAsync,
18-
cleanupCompactionCheckpointSnapshot,
19-
persistSessionCompactionCheckpoint,
17+
createFileBackedCompactionCheckpointStore,
2018
readSessionLeafStateFromTranscriptAsync,
2119
resolveCompactionCheckpointTranscriptPosition,
2220
resolveSessionCompactionCheckpointReason,
@@ -63,6 +61,8 @@ import { resolveModelAsync } from "./model.js";
6361
import type { EmbeddedAgentCompactResult } from "./types.js";
6462
import { normalizeContextTokenBudget } from "./utils.js";
6563

64+
const compactionCheckpointStore = createFileBackedCompactionCheckpointStore();
65+
6666
function shouldFallbackAfterHarnessCompaction(
6767
result: EmbeddedAgentCompactResult | undefined,
6868
): boolean {
@@ -352,7 +352,7 @@ export async function compactEmbeddedAgentSession(
352352
// are notified regardless of which engine is active.
353353
const engineOwnsCompaction = contextEngine.info.ownsCompaction === true;
354354
checkpointSnapshot = engineOwnsCompaction
355-
? await captureCompactionCheckpointSnapshotAsync({
355+
? await compactionCheckpointStore.captureSnapshot({
356356
sessionFile: params.sessionFile,
357357
})
358358
: null;
@@ -478,7 +478,7 @@ export async function compactEmbeddedAgentSession(
478478
preferredLeafId: postCompactionLeafId,
479479
transcriptState,
480480
});
481-
const storedCheckpoint = await persistSessionCompactionCheckpoint({
481+
const storedCheckpoint = await compactionCheckpointStore.persistCheckpoint({
482482
cfg: params.config,
483483
sessionKey: params.sessionKey,
484484
sessionId: postCompactionSessionId,
@@ -620,7 +620,7 @@ export async function compactEmbeddedAgentSession(
620620
};
621621
} finally {
622622
if (!checkpointSnapshotRetained) {
623-
await cleanupCompactionCheckpointSnapshot(checkpointSnapshot);
623+
await compactionCheckpointStore.cleanupSnapshot(checkpointSnapshot);
624624
}
625625
await contextEngine.dispose?.();
626626
}
Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,12 @@
11
/**
22
* Types for the lazy embedded-agent compaction runtime boundary.
33
*/
4-
import type { CompactEmbeddedAgentSessionParams } from "./compact.types.js";
4+
import type { CompactEmbeddedAgentSessionRuntimeParams } from "./compact.types.js";
55
import type { EmbeddedAgentCompactResult } from "./types.js";
66

77
/**
88
* Lazy-runtime signature for direct embedded session compaction.
99
*/
1010
export type CompactEmbeddedAgentSessionDirect = (
11-
params: CompactEmbeddedAgentSessionParams,
11+
params: CompactEmbeddedAgentSessionRuntimeParams,
1212
) => Promise<EmbeddedAgentCompactResult>;

src/agents/embedded-agent-runner/compact.ts

Lines changed: 25 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -8,9 +8,7 @@ import type { ThinkLevel } from "../../auto-reply/thinking.js";
88
import { resolveAgentModelFallbackValues } from "../../config/model-input.js";
99
import type { OpenClawConfig } from "../../config/types.openclaw.js";
1010
import {
11-
captureCompactionCheckpointSnapshotAsync,
12-
cleanupCompactionCheckpointSnapshot,
13-
persistSessionCompactionCheckpoint,
11+
createFileBackedCompactionCheckpointStore,
1412
readSessionLeafStateFromTranscriptAsync,
1513
resolveCompactionCheckpointTranscriptPosition,
1614
resolveSessionCompactionCheckpointReason,
@@ -107,6 +105,10 @@ import { wrapStreamFnTextTransforms } from "../plugin-text-transforms.js";
107105
import { resolveAgentPromptSurfaceForSessionKey } from "../prompt-surface.js";
108106
import { applyPreparedRuntimeAuthToModel } from "../provider-request-config.js";
109107
import { registerProviderStreamForModel } from "../provider-stream.js";
108+
import {
109+
applyAgentRunSessionTargetIdentity,
110+
resolveAgentRunSessionTarget,
111+
} from "../run-session-target.js";
110112
import { collectRuntimeChannelCapabilities } from "../runtime-capabilities.js";
111113
import { buildAgentRuntimePlan } from "../runtime-plan/build.js";
112114
import type { AgentRuntimePlan } from "../runtime-plan/types.js";
@@ -135,6 +137,7 @@ import {
135137
} from "./compact-reasons.js";
136138
import type {
137139
CompactEmbeddedAgentSessionParams,
140+
CompactEmbeddedAgentSessionRuntimeParams,
138141
CompactionMessageMetrics,
139142
} from "./compact.types.js";
140143
import { dedupeDuplicateUserMessagesForCompaction } from "./compaction-duplicate-user-messages.js";
@@ -192,6 +195,11 @@ import { mapThinkingLevel, normalizeContextTokenBudget } from "./utils.js";
192195
import { flushPendingToolResultsAfterIdle } from "./wait-for-idle-before-flush.js";
193196
export type { CompactEmbeddedAgentSessionParams } from "./compact.types.js";
194197

198+
const compactionCheckpointStore = createFileBackedCompactionCheckpointStore();
199+
type CompactEmbeddedAgentSessionParamsWithSessionFile = CompactEmbeddedAgentSessionRuntimeParams & {
200+
sessionFile: string;
201+
};
202+
195203
function hasRealConversationContent(
196204
msg: AgentMessage,
197205
messages: AgentMessage[],
@@ -464,8 +472,17 @@ function fallbackFailureToCompactionResult(err: unknown): EmbeddedAgentCompactRe
464472
* Use this when already inside a session/global lane to avoid deadlocks.
465473
*/
466474
export async function compactEmbeddedAgentSessionDirect(
467-
params: CompactEmbeddedAgentSessionParams,
475+
paramsInput: CompactEmbeddedAgentSessionRuntimeParams,
468476
): Promise<EmbeddedAgentCompactResult> {
477+
const paramsBase = applyAgentRunSessionTargetIdentity(paramsInput);
478+
const runSessionTarget = await resolveAgentRunSessionTarget(paramsBase);
479+
const params: CompactEmbeddedAgentSessionParamsWithSessionFile = {
480+
...paramsBase,
481+
agentId: paramsBase.agentId ?? runSessionTarget.agentId,
482+
sessionId: runSessionTarget.sessionId,
483+
sessionKey: paramsBase.sessionKey ?? runSessionTarget.sessionKey,
484+
sessionFile: runSessionTarget.sessionFile,
485+
};
469486
if (hasExplicitCompactionModel(params) || !hasCompactionModelFallbackCandidates(params)) {
470487
return await compactEmbeddedAgentSessionDirectOnce(params);
471488
}
@@ -530,7 +547,7 @@ export async function compactEmbeddedAgentSessionDirect(
530547
}
531548

532549
async function compactEmbeddedAgentSessionDirectOnce(
533-
params: CompactEmbeddedAgentSessionParams,
550+
params: CompactEmbeddedAgentSessionParamsWithSessionFile,
534551
): Promise<EmbeddedAgentCompactResult> {
535552
const startedAt = Date.now();
536553
const diagId = params.diagId?.trim() || createCompactionDiagId();
@@ -1190,7 +1207,7 @@ async function compactEmbeddedAgentSessionDirectOnce(
11901207
: undefined,
11911208
allowedToolNames,
11921209
});
1193-
checkpointSnapshot = await captureCompactionCheckpointSnapshotAsync({
1210+
checkpointSnapshot = await compactionCheckpointStore.captureSnapshot({
11941211
sessionManager,
11951212
sessionFile: params.sessionFile,
11961213
});
@@ -1546,7 +1563,7 @@ async function compactEmbeddedAgentSessionDirectOnce(
15461563
preferredLeafId: activePostLeafId,
15471564
transcriptState,
15481565
});
1549-
const storedCheckpoint = await persistSessionCompactionCheckpoint({
1566+
const storedCheckpoint = await compactionCheckpointStore.persistCheckpoint({
15501567
cfg: params.config,
15511568
sessionKey: params.sessionKey,
15521569
sessionId: activeSessionId,
@@ -1670,7 +1687,7 @@ async function compactEmbeddedAgentSessionDirectOnce(
16701687
return fail(reason, err);
16711688
} finally {
16721689
if (!checkpointSnapshotRetained) {
1673-
await cleanupCompactionCheckpointSnapshot(checkpointSnapshot);
1690+
await compactionCheckpointStore.cleanupSnapshot(checkpointSnapshot);
16741691
}
16751692
restoreSkillEnv?.();
16761693
}

src/agents/embedded-agent-runner/compact.types.ts

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,12 +9,15 @@ import type { ContextEngine, ContextEngineRuntimeContext } from "../../context-e
99
import type { CommandQueueEnqueueFn } from "../../process/command-queue.types.js";
1010
import type { SkillSnapshot } from "../../skills/types.js";
1111
import type { ExecElevatedDefaults, ExecToolDefaults } from "../bash-tools.exec-types.js";
12+
import type { AgentRunSessionTarget } from "../run-session-target.js";
1213
import type { AgentRuntimePlan } from "../runtime-plan/types.js";
1314

1415
export type CompactEmbeddedAgentSessionParams = {
1516
sessionId: string;
1617
runId?: string;
1718
sessionKey?: string;
19+
/** Storage-neutral transcript/session target. Defaults to sessionId/sessionKey/agentId. */
20+
sessionTarget?: AgentRunSessionTarget;
1821
/** Caller-resolved owner agent for global session aliases. */
1922
agentId?: string;
2023
/** Session key used only for runtime policy/sandbox resolution. Defaults to sessionKey. */
@@ -106,6 +109,14 @@ export type CompactEmbeddedAgentSessionParams = {
106109
oneShotCliRun?: boolean;
107110
};
108111

112+
export type CompactEmbeddedAgentSessionRuntimeParams = Omit<
113+
CompactEmbeddedAgentSessionParams,
114+
"sessionFile"
115+
> & {
116+
/** Deprecated file-backed artifact target. Prefer sessionTarget for new callers. */
117+
sessionFile?: string;
118+
};
119+
109120
export type CompactionMessageMetrics = {
110121
messages: number;
111122
historyTextChars: number;

src/agents/embedded-agent-runner/run.ts

Lines changed: 22 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,10 @@ import {
125125
import { resolveProviderIdForAuth } from "../provider-auth-aliases.js";
126126
import { hasOnlyAssistantReasoningContent } from "../replay-turn-classification.js";
127127
import { runAgentCleanupStep } from "../run-cleanup-timeout.js";
128+
import {
129+
applyAgentRunSessionTargetIdentity,
130+
resolveAgentRunSessionTarget,
131+
} from "../run-session-target.js";
128132
import { buildAgentRuntimeAuthPlan } from "../runtime-plan/auth.js";
129133
import { buildAgentRuntimePlan } from "../runtime-plan/build.js";
130134
import { ensureRuntimePluginsLoaded } from "../runtime-plugins.js";
@@ -255,6 +259,7 @@ const BEFORE_AGENT_FINALIZE_RETRY_PROMPT_PREFIX =
255259
"Before accepting the previous final answer, apply this revision request and produce the revised final answer. Do not repeat completed work or rerun tools unless the request explicitly requires it.";
256260
const MAX_BEFORE_AGENT_FINALIZE_REVISIONS = 3;
257261
type EmbeddedRunAttemptForRunner = Awaited<ReturnType<typeof runEmbeddedAttemptWithBackend>>;
262+
type RunEmbeddedAgentParamsWithSessionFile = RunEmbeddedAgentParams & { sessionFile: string };
258263

259264
function isNoRealConversationCompactionNoop(params: {
260265
ok?: boolean;
@@ -617,20 +622,28 @@ export function runEmbeddedAgent(
617622
async function runEmbeddedAgentInternal(
618623
paramsInput: RunEmbeddedAgentParams,
619624
): Promise<EmbeddedAgentRunResult> {
620-
let params = paramsInput;
621-
let lifecycleGeneration = params.lifecycleGeneration!;
625+
const paramsBase = applyAgentRunSessionTargetIdentity(paramsInput);
626+
let lifecycleGeneration = paramsBase.lifecycleGeneration!;
622627
const queuedLifecycleGeneration = getAgentEventLifecycleGeneration();
623628
// Resolve sessionKey early so all downstream consumers (hooks, LCM, compaction)
624629
// receive a non-null key even when callers omit it. See #60552.
625630
const effectiveSessionKey = backfillSessionKey({
626-
config: params.config,
627-
sessionId: params.sessionId,
628-
sessionKey: params.sessionKey,
629-
agentId: params.agentId,
631+
config: paramsBase.config,
632+
sessionId: paramsBase.sessionId,
633+
sessionKey: paramsBase.sessionKey,
634+
agentId: paramsBase.agentId,
630635
});
631-
if (effectiveSessionKey !== params.sessionKey) {
632-
params = { ...params, sessionKey: effectiveSessionKey };
633-
}
636+
const runSessionTarget = await resolveAgentRunSessionTarget({
637+
...paramsBase,
638+
sessionKey: effectiveSessionKey,
639+
});
640+
let params: RunEmbeddedAgentParamsWithSessionFile = {
641+
...paramsBase,
642+
agentId: paramsBase.agentId ?? runSessionTarget.agentId,
643+
sessionId: runSessionTarget.sessionId,
644+
sessionKey: effectiveSessionKey ?? runSessionTarget.sessionKey,
645+
sessionFile: runSessionTarget.sessionFile,
646+
};
634647
const sessionLane = resolveSessionLane(params.sessionKey?.trim() || params.sessionId);
635648
const globalLane = resolveGlobalLane(params.lane);
636649
// Outer fallback attempts defer session suspension only while another

src/agents/embedded-agent-runner/run/params.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ import type {
2929
} from "../../embedded-agent-subscribe.shared-types.js";
3030
import type { FastModeAutoProgressState } from "../../fast-mode.js";
3131
import type { AgentInternalEvent } from "../../internal-events.js";
32+
import type { AgentRunSessionTarget } from "../../run-session-target.js";
3233
import type { AgentMessage } from "../../runtime/index.js";
3334
import type { SilentReplyPromptMode } from "../../system-prompt.types.js";
3435
import type { PromptMode } from "../../system-prompt.types.js";
@@ -47,6 +48,8 @@ export type CurrentInboundPromptContext = {
4748
export type RunEmbeddedAgentParams = {
4849
sessionId: string;
4950
sessionKey?: string;
51+
/** Storage-neutral transcript/session target. Defaults to sessionId/sessionKey/agentId. */
52+
sessionTarget?: AgentRunSessionTarget;
5053
/** Immutable gateway lifecycle ownership captured when this execution was admitted. */
5154
lifecycleGeneration?: string;
5255
/** Provider prompt-cache affinity key; distinct from transcript/session identity. */
@@ -122,7 +125,8 @@ export type RunEmbeddedAgentParams = {
122125
forceHeartbeatTool?: boolean;
123126
/** Allow runtime plugins for this run to late-bind the gateway subagent. */
124127
allowGatewaySubagentBinding?: boolean;
125-
sessionFile: string;
128+
/** @deprecated Use sessionTarget plus sessionId/sessionKey/agentId for runtime identity. */
129+
sessionFile?: string;
126130
workspaceDir: string;
127131
/** Task working directory for tool/runtime execution. Defaults to workspaceDir. */
128132
cwd?: string;

src/agents/embedded-agent-runner/run/types.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ type EmbeddedRunAttemptBase = Omit<
4040
| "fastMode"
4141
| "lane"
4242
| "enqueue"
43+
| "sessionFile"
4344
>;
4445

4546
export type EmbeddedRunContextWindowInfo = {
@@ -51,6 +52,8 @@ export type EmbeddedRunContextWindowInfo = {
5152
export type EmbeddedRunFastModeParam = boolean | (() => boolean | undefined);
5253

5354
export type EmbeddedRunAttemptParams = EmbeddedRunAttemptBase & {
55+
/** Active file-backed artifact target resolved by the run/session target seam. */
56+
sessionFile: string;
5457
initialReplayState?: EmbeddedRunReplayState;
5558
/** Pluggable context engine for ingest/assemble/compact lifecycle. */
5659
contextEngine?: ContextEngine;

src/agents/openclaw-tools.subagents.sessions-spawn.test-harness.ts

Lines changed: 20 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -206,13 +206,26 @@ export async function getSessionsSpawnTool(opts: CreateOpenClawToolsOpts) {
206206
compact: async () => ({ ok: true, compacted: false }),
207207
ingest: async () => ({ ingested: false }),
208208
}),
209-
resolveParentForkDecision: async () => ({
210-
status: "fork",
211-
maxTokens: 100_000,
212-
}),
213-
forkSessionFromParent: async () => ({
214-
sessionId: "forked-session-id",
215-
sessionFile: "/tmp/forked-session.jsonl",
209+
forkSessionEntryFromParent: async () => ({
210+
status: "forked",
211+
fork: {
212+
sessionId: "forked-session-id",
213+
sessionFile: "/tmp/forked-session.jsonl",
214+
},
215+
parentEntry: {
216+
sessionId: "parent-session-id",
217+
updatedAt: Date.now(),
218+
},
219+
sessionEntry: {
220+
sessionId: "forked-session-id",
221+
sessionFile: "/tmp/forked-session.jsonl",
222+
forkedFromParent: true,
223+
updatedAt: Date.now(),
224+
},
225+
decision: {
226+
status: "fork",
227+
maxTokens: 100_000,
228+
},
216229
}),
217230
updateSessionStore: async (_storePath, mutator) => mutator({}),
218231
});
Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,59 @@
1+
import fs from "node:fs";
2+
import os from "node:os";
3+
import path from "node:path";
4+
import { afterEach, beforeEach, describe, expect, it } from "vitest";
5+
import { loadSessionStore } from "../config/sessions/store.js";
6+
import type { OpenClawConfig } from "../config/types.openclaw.js";
7+
import { resolveAgentRunSessionTarget } from "./run-session-target.js";
8+
9+
describe("agent run session target", () => {
10+
let tempDir: string;
11+
12+
beforeEach(() => {
13+
tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-run-session-target-"));
14+
});
15+
16+
afterEach(() => {
17+
fs.rmSync(tempDir, { recursive: true, force: true });
18+
});
19+
20+
it("resolves runtime identity through the run config store", async () => {
21+
const storePath = path.join(tempDir, "custom-sessions", "sessions.json");
22+
const sessionKey = "agent:helper:commitments:test-run";
23+
24+
const target = await resolveAgentRunSessionTarget({
25+
agentId: "helper",
26+
config: { session: { store: storePath } } as OpenClawConfig,
27+
sessionId: "test-run",
28+
sessionKey,
29+
});
30+
31+
expect(target).toMatchObject({
32+
agentId: "helper",
33+
sessionId: "test-run",
34+
sessionKey,
35+
});
36+
expect(path.dirname(target.sessionFile)).toBe(path.dirname(storePath));
37+
expect(loadSessionStore(storePath, { skipCache: true })[sessionKey]?.sessionFile).toBe(
38+
target.sessionFile,
39+
);
40+
});
41+
42+
it("uses the agent from an agent-scoped session key when agentId is omitted", async () => {
43+
const storeRoot = path.join(tempDir, "agents", "{agentId}", "sessions.json");
44+
const sessionKey = "agent:helper:main";
45+
46+
const target = await resolveAgentRunSessionTarget({
47+
config: { session: { store: storeRoot } } as OpenClawConfig,
48+
sessionId: "helper-session",
49+
sessionKey,
50+
});
51+
52+
const helperStorePath = path.join(tempDir, "agents", "helper", "sessions.json");
53+
expect(target.agentId).toBe("helper");
54+
expect(path.dirname(target.sessionFile)).toBe(path.dirname(helperStorePath));
55+
expect(loadSessionStore(helperStorePath, { skipCache: true })[sessionKey]?.sessionFile).toBe(
56+
target.sessionFile,
57+
);
58+
});
59+
});

0 commit comments

Comments
 (0)