Skip to content

Commit 95e37f8

Browse files
authored
refactor: guard reply session initialization (#96218)
* refactor: guard reply session initialization * refactor: tighten reply session initialization boundary * test: satisfy reply session accessor lint
1 parent 6f2869c commit 95e37f8

12 files changed

Lines changed: 579 additions & 118 deletions

src/auto-reply/reply/body.ts

Lines changed: 29 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22
import type { SessionEntry } from "../../config/sessions/types.js";
33
import { createLazyImportLoader } from "../../shared/lazy-promise.js";
44
import { setAbortMemory } from "./abort-primitives.js";
5+
import type { ReplySessionEntryHandle } from "./session-entry-handle.js";
56

67
const sessionAccessorRuntimeLoader = createLazyImportLoader(
78
() => import("../../config/sessions/session-accessor.js"),
@@ -16,6 +17,7 @@ export async function applySessionHints(params: {
1617
baseBody: string;
1718
abortedLastRun: boolean;
1819
sessionEntry?: SessionEntry;
20+
sessionEntryHandle?: ReplySessionEntryHandle;
1921
sessionStore?: Record<string, SessionEntry>;
2022
sessionKey?: string;
2123
storePath?: string;
@@ -28,11 +30,13 @@ export async function applySessionHints(params: {
2830
if (abortedHint) {
2931
prefixedBodyBase = `${abortedHint}\n\n${prefixedBodyBase}`;
3032
// The abort hint is one-shot; clear durable state once it is added.
31-
if (params.sessionEntry && params.sessionStore && params.sessionKey) {
33+
const sessionEntry = params.sessionEntryHandle?.getCurrent() ?? params.sessionEntry;
34+
if (sessionEntry && params.sessionEntryHandle && params.sessionKey) {
3235
const updatedAt = Date.now();
33-
params.sessionEntry.abortedLastRun = false;
34-
params.sessionEntry.updatedAt = updatedAt;
35-
params.sessionStore[params.sessionKey] = params.sessionEntry;
36+
params.sessionEntryHandle.patchCurrent({
37+
abortedLastRun: false,
38+
updatedAt,
39+
});
3640
if (params.storePath) {
3741
const sessionKey = params.sessionKey;
3842
const { patchSessionEntry } = await loadSessionAccessorRuntime();
@@ -45,7 +49,27 @@ export async function applySessionHints(params: {
4549
abortedLastRun: false,
4650
updatedAt,
4751
}),
48-
{ fallbackEntry: params.sessionEntry },
52+
{ fallbackEntry: params.sessionEntryHandle.getCurrent() ?? sessionEntry },
53+
);
54+
}
55+
} else if (sessionEntry && params.sessionStore && params.sessionKey) {
56+
const updatedAt = Date.now();
57+
sessionEntry.abortedLastRun = false;
58+
sessionEntry.updatedAt = updatedAt;
59+
params.sessionStore[params.sessionKey] = sessionEntry;
60+
if (params.storePath) {
61+
const sessionKey = params.sessionKey;
62+
const { patchSessionEntry } = await loadSessionAccessorRuntime();
63+
await patchSessionEntry(
64+
{
65+
storePath: params.storePath,
66+
sessionKey,
67+
},
68+
() => ({
69+
abortedLastRun: false,
70+
updatedAt,
71+
}),
72+
{ fallbackEntry: sessionEntry },
4973
);
5074
}
5175
} else if (params.abortKey) {

src/auto-reply/reply/get-reply-fast-path.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import { isFormattedGoalContinuationPrompt } from "./commands-goal.js";
1919
import { parseSoftResetCommand } from "./commands-reset-mode.js";
2020
import type { CommandContext } from "./commands-types.js";
2121
import { stripMentions, stripStructuralPrefixes } from "./mentions.js";
22+
import { createReplySessionEntryHandle } from "./session-entry-handle.js";
2223
import type { SessionInitResult } from "./session.js";
2324

2425
const COMPLETE_REPLY_CONFIG_SYMBOL = Symbol.for("openclaw.reply.complete-config");
@@ -284,6 +285,11 @@ export function initFastReplySessionState(params: {
284285
: {}),
285286
};
286287
sessionStore[sessionKey] = sessionEntry;
288+
const sessionEntryHandle = createReplySessionEntryHandle({
289+
sessionEntry,
290+
sessionKey,
291+
sessionStore,
292+
});
287293
const sessionCtx: TemplateContext = {
288294
...ctx,
289295
SessionKey: sessionKey,
@@ -294,6 +300,7 @@ export function initFastReplySessionState(params: {
294300
return {
295301
sessionCtx,
296302
sessionEntry,
303+
sessionEntryHandle,
297304
sessionStore,
298305
sessionKey,
299306
sessionId,

src/auto-reply/reply/get-reply-run.ts

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -100,6 +100,7 @@ import {
100100
import { resolveReplyToMode } from "./reply-threading.js";
101101
import { resolveRoutedDeliveryThreadId } from "./routed-delivery-thread.js";
102102
import { resolveRuntimePolicySessionKey } from "./runtime-policy-session-key.js";
103+
import type { ReplySessionEntryHandle } from "./session-entry-handle.js";
103104
import { resolveBareSessionResetPromptState } from "./session-reset-prompt.js";
104105
import { resolveBareResetBootstrapFileAccess } from "./session-reset-prompt.js";
105106
import { drainFormattedSystemEvents } from "./session-system-events.js";
@@ -444,6 +445,7 @@ type RunPreparedReplyParams = {
444445
resetTriggered: boolean;
445446
systemSent: boolean;
446447
sessionEntry?: SessionEntry;
448+
sessionEntryHandle?: ReplySessionEntryHandle;
447449
sessionStore?: Record<string, SessionEntry>;
448450
sessionKey: string;
449451
sessionId?: string;
@@ -491,6 +493,7 @@ export async function runPreparedReply(
491493
sessionId,
492494
storePath,
493495
workspaceDir,
496+
sessionEntryHandle,
494497
sessionStore,
495498
} = params;
496499
const runtimePolicySessionKey = resolveRuntimePolicySessionKey({
@@ -770,11 +773,13 @@ export async function runPreparedReply(
770773
baseBody: effectiveBaseBody,
771774
abortedLastRun,
772775
sessionEntry,
776+
sessionEntryHandle,
773777
sessionStore,
774778
sessionKey,
775779
storePath,
776780
abortKey: command.abortKey,
777781
});
782+
sessionEntry = sessionEntryHandle?.getCurrent() ?? sessionEntry;
778783
const isGroupSession = sessionEntry?.chatType === "group" || sessionEntry?.chatType === "channel";
779784
const isMainSession = !isGroupSession && sessionKey === normalizeMainKey(sessionCfg?.mainKey);
780785
// Extract first-token think hint from the user body BEFORE prepending system events.
@@ -854,6 +859,7 @@ export async function runPreparedReply(
854859
const { ensureSkillSnapshot } = await loadSessionUpdatesRuntime();
855860
return await ensureSkillSnapshot({
856861
sessionEntry,
862+
sessionEntryHandle,
857863
sessionStore,
858864
sessionKey,
859865
storePath,
@@ -865,6 +871,9 @@ export async function runPreparedReply(
865871
});
866872
});
867873
sessionEntry = skillResult.sessionEntry ?? sessionEntry;
874+
if (sessionEntry) {
875+
sessionEntryHandle?.replaceCurrent(sessionEntry);
876+
}
868877
const skillsSnapshot = skillResult.skillsSnapshot;
869878
let {
870879
prefixedCommandBody,

src/auto-reply/reply/get-reply.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -494,6 +494,7 @@ export async function getReplyFromConfig(
494494
const {
495495
sessionCtx,
496496
sessionEntry,
497+
sessionEntryHandle,
497498
previousSessionEntry,
498499
sessionStore,
499500
sessionKey,
@@ -534,6 +535,7 @@ export async function getReplyFromConfig(
534535
sessionEntry.pendingFinalDeliveryAttemptCount = undefined;
535536
sessionEntry.pendingFinalDeliveryLastError = undefined;
536537
sessionEntry.pendingFinalDeliveryContext = undefined;
538+
sessionEntryHandle.replaceCurrent(sessionEntry);
537539
if (sessionKey && sessionStore) {
538540
sessionStore[sessionKey] = sessionEntry;
539541
}
@@ -570,6 +572,7 @@ export async function getReplyFromConfig(
570572
sessionCtx,
571573
ctx: finalized,
572574
sessionEntry,
575+
sessionEntryHandle,
573576
sessionStore,
574577
sessionKey,
575578
storePath,
@@ -741,6 +744,7 @@ export async function getReplyFromConfig(
741744
resetTriggered,
742745
systemSent,
743746
sessionEntry,
747+
sessionEntryHandle,
744748
sessionStore,
745749
sessionKey,
746750
sessionId,

src/auto-reply/reply/session-delivery.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -184,7 +184,7 @@ export function resolveLastToRaw(params: {
184184
export function maybeRetireLegacyMainDeliveryRoute(params: {
185185
sessionCfg: { dmScope?: string } | undefined;
186186
sessionKey: string;
187-
sessionStore: Record<string, SessionEntry>;
187+
legacyMain?: SessionEntry;
188188
agentId: string;
189189
mainKey: string;
190190
isGroup: boolean;
@@ -201,7 +201,7 @@ export function maybeRetireLegacyMainDeliveryRoute(params: {
201201
if (params.sessionKey === canonicalMainSessionKey) {
202202
return undefined;
203203
}
204-
const legacyMain = params.sessionStore[canonicalMainSessionKey];
204+
const legacyMain = params.legacyMain;
205205
if (!legacyMain) {
206206
return undefined;
207207
}
Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
// Narrow mutable handle for the active reply session entry.
2+
import type { SessionEntry } from "../../config/sessions.js";
3+
4+
export type ReplySessionEntryHandle = {
5+
get(sessionKey: string): SessionEntry | undefined;
6+
getCurrent(): SessionEntry | undefined;
7+
patchCurrent(patch: Partial<SessionEntry>): SessionEntry | undefined;
8+
replaceCurrent(entry: SessionEntry): void;
9+
set(sessionKey: string, entry: SessionEntry): void;
10+
toCompatSessionStore(): Record<string, SessionEntry>;
11+
};
12+
13+
export function createReplySessionEntryHandle(params: {
14+
sessionEntry: SessionEntry;
15+
sessionKey: string;
16+
sessionStore?: Record<string, SessionEntry>;
17+
}): ReplySessionEntryHandle {
18+
const entries = params.sessionStore ?? { [params.sessionKey]: params.sessionEntry };
19+
let currentEntry = params.sessionEntry;
20+
entries[params.sessionKey] = currentEntry;
21+
22+
return {
23+
get: (sessionKey) => entries[sessionKey],
24+
getCurrent: () => currentEntry,
25+
patchCurrent: (patch) => {
26+
currentEntry = { ...currentEntry, ...patch };
27+
entries[params.sessionKey] = currentEntry;
28+
return currentEntry;
29+
},
30+
replaceCurrent: (entry) => {
31+
currentEntry = entry;
32+
entries[params.sessionKey] = entry;
33+
},
34+
set: (sessionKey, entry) => {
35+
entries[sessionKey] = entry;
36+
if (sessionKey === params.sessionKey) {
37+
currentEntry = entry;
38+
}
39+
},
40+
toCompatSessionStore: () => entries,
41+
};
42+
}
Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,62 @@
1+
// Prepares parent-context fork metadata for guarded reply session initialization.
2+
import path from "node:path";
3+
import type { SessionEntry } from "../../config/sessions.js";
4+
import { forkSessionFromParent, resolveParentForkDecision } from "./session-fork.js";
5+
6+
export async function prepareReplySessionParentFork(params: {
7+
agentId: string;
8+
alreadyForked: boolean;
9+
parentSessionKey?: string;
10+
readEntry: (sessionKey: string) => SessionEntry | undefined;
11+
sessionEntry: SessionEntry;
12+
sessionKey: string;
13+
storePath: string;
14+
warn: (message: string) => void;
15+
}): Promise<SessionEntry> {
16+
if (
17+
!params.parentSessionKey ||
18+
params.parentSessionKey === params.sessionKey ||
19+
params.alreadyForked
20+
) {
21+
return params.sessionEntry;
22+
}
23+
const parentEntry = params.readEntry(params.parentSessionKey);
24+
if (!parentEntry?.sessionId) {
25+
return params.sessionEntry;
26+
}
27+
const decision = await resolveParentForkDecision({
28+
parentEntry,
29+
agentId: params.agentId,
30+
storePath: params.storePath,
31+
});
32+
if (decision.status === "skip") {
33+
// The parent branch is too large to inherit usefully. Start fresh and
34+
// mark as handled so the thread does not retry this decision every turn.
35+
params.warn(
36+
`skipping parent fork (parent too large): parentKey=${params.parentSessionKey} → sessionKey=${params.sessionKey} ` +
37+
`parentTokens=${decision.parentTokens} maxTokens=${decision.maxTokens}`,
38+
);
39+
return { ...params.sessionEntry, forkedFromParent: true };
40+
}
41+
const fork = await forkSessionFromParent({
42+
parentEntry,
43+
agentId: params.agentId,
44+
sessionsDir: path.dirname(params.storePath),
45+
});
46+
if (!fork) {
47+
return params.sessionEntry;
48+
}
49+
params.warn(
50+
`forking from parent session: parentKey=${params.parentSessionKey} → sessionKey=${params.sessionKey} ` +
51+
`parentTokens=${decision.parentTokens ?? "unknown"}`,
52+
);
53+
params.warn(`forked session created: file=${fork.sessionFile}`);
54+
return {
55+
...params.sessionEntry,
56+
sessionId: fork.sessionId,
57+
sessionFile: fork.sessionFile,
58+
forkedFromParent: true,
59+
totalTokens: undefined,
60+
totalTokensFresh: false,
61+
};
62+
}

src/auto-reply/reply/session-reset-model.ts

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import {
1818
type ModelAliasIndex,
1919
type ModelDirectiveSelection,
2020
} from "./model-selection-directive.js";
21+
import type { ReplySessionEntryHandle } from "./session-entry-handle.js";
2122

2223
/** Result of applying a reset-message model override. */
2324
type ResetModelResult = {
@@ -109,12 +110,14 @@ function buildSelectionFromExplicit(params: {
109110
function applySelectionToSession(params: {
110111
selection: ModelDirectiveSelection;
111112
sessionEntry?: SessionEntry;
113+
sessionEntryHandle?: ReplySessionEntryHandle;
112114
sessionStore?: Record<string, SessionEntry>;
113115
sessionKey?: string;
114116
storePath?: string;
115117
}) {
116-
const { selection, sessionEntry, sessionStore, sessionKey, storePath } = params;
117-
if (!sessionEntry || !sessionStore || !sessionKey) {
118+
const { selection, sessionEntryHandle, sessionStore, sessionKey, storePath } = params;
119+
const sessionEntry = sessionEntryHandle?.getCurrent() ?? params.sessionEntry;
120+
if (!sessionEntry || !sessionKey) {
118121
return;
119122
}
120123
const { updated } = applyModelOverrideToSessionEntry({
@@ -124,7 +127,11 @@ function applySelectionToSession(params: {
124127
if (!updated) {
125128
return;
126129
}
127-
sessionStore[sessionKey] = sessionEntry;
130+
if (sessionEntryHandle) {
131+
sessionEntryHandle.replaceCurrent(sessionEntry);
132+
} else if (sessionStore) {
133+
sessionStore[sessionKey] = sessionEntry;
134+
}
128135
if (storePath) {
129136
void import("../../config/sessions/session-accessor.js")
130137
.then(({ replaceSessionEntry }) =>
@@ -146,6 +153,7 @@ export async function applyResetModelOverride(params: {
146153
sessionCtx: TemplateContext;
147154
ctx: MsgContext;
148155
sessionEntry?: SessionEntry;
156+
sessionEntryHandle?: ReplySessionEntryHandle;
149157
sessionStore?: Record<string, SessionEntry>;
150158
sessionKey?: string;
151159
storePath?: string;
@@ -248,6 +256,7 @@ export async function applyResetModelOverride(params: {
248256
applySelectionToSession({
249257
selection,
250258
sessionEntry: params.sessionEntry,
259+
sessionEntryHandle: params.sessionEntryHandle,
251260
sessionStore: params.sessionStore,
252261
sessionKey: params.sessionKey,
253262
storePath: params.storePath,

0 commit comments

Comments
 (0)