Skip to content

Commit 052bde7

Browse files
author
ai-hpc
committed
fix(agents): persist CLI user turns before attempts
1 parent 5b6d409 commit 052bde7

11 files changed

Lines changed: 205 additions & 3 deletions

src/agents/agent-command.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1611,6 +1611,7 @@ async function agentCommandInternal(
16111611
sessionCwd: workspaceDir,
16121612
config: cfg,
16131613
embeddedAssistantGapFill,
1614+
skipUserMessage: attemptLifecycleState.currentTurnUserMessagePersisted,
16141615
});
16151616
if (suppressVisibleSessionEffects) {
16161617
sessionEntry = prepared.sessionEntry;

src/agents/cli-runner.reliability.test.ts

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -222,6 +222,16 @@ function expectTextMessage(value: unknown, fields: { role: string; content: stri
222222
expect(message.timestamp).toBeTypeOf("number");
223223
}
224224

225+
function readTranscriptMessages(sessionFile: string): Array<Record<string, unknown>> {
226+
return fs
227+
.readFileSync(sessionFile, "utf-8")
228+
.trim()
229+
.split("\n")
230+
.map((line) => JSON.parse(line) as { type?: string; message?: unknown })
231+
.filter((entry) => entry.type === "message")
232+
.map((entry) => requireRecord(entry.message, "transcript message"));
233+
}
234+
225235
describe("runCliAgent reliability", () => {
226236
afterEach(() => {
227237
replyRunTesting.resetReplyRunRegistry();
@@ -418,6 +428,54 @@ describe("runCliAgent reliability", () => {
418428
}
419429
});
420430

431+
it("persists the CLI user turn before backend failures can drop it", async () => {
432+
supervisorSpawnMock.mockClear();
433+
supervisorSpawnMock.mockResolvedValueOnce(
434+
createManagedRun({
435+
reason: "exit",
436+
exitCode: 1,
437+
exitSignal: null,
438+
durationMs: 150,
439+
stdout: "",
440+
stderr: "rate limit exceeded",
441+
timedOut: false,
442+
noOutputTimedOut: false,
443+
}),
444+
);
445+
const { dir, sessionFile } = createSessionFile();
446+
const onUserMessagePersisted = vi.fn();
447+
448+
try {
449+
await expect(
450+
runPreparedCliAgent({
451+
...buildPreparedContext({
452+
sessionKey: "agent:main:main",
453+
runId: "run-cli-user-before-failure",
454+
}),
455+
params: {
456+
...buildPreparedContext({
457+
sessionKey: "agent:main:main",
458+
runId: "run-cli-user-before-failure",
459+
}).params,
460+
agentId: "main",
461+
sessionFile,
462+
workspaceDir: dir,
463+
transcriptPrompt: "transcript-safe prompt",
464+
onUserMessagePersisted,
465+
},
466+
}),
467+
).rejects.toThrow("rate limit exceeded");
468+
469+
expect(supervisorSpawnMock).toHaveBeenCalledTimes(1);
470+
expect(onUserMessagePersisted).toHaveBeenCalledTimes(1);
471+
const messages = readTranscriptMessages(sessionFile);
472+
expect(messages).toHaveLength(1);
473+
expectTextMessage(messages[0], { role: "user", content: "transcript-safe prompt" });
474+
} finally {
475+
fs.rmSync(dir, { recursive: true, force: true });
476+
}
477+
});
478+
421479
it("returns the assembled CLI prompt in meta for raw trace consumers", async () => {
422480
supervisorSpawnMock.mockResolvedValueOnce(
423481
createManagedRun({

src/agents/cli-runner.ts

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import type { AgentMessage } from "@earendil-works/pi-agent-core";
22
import { SessionManager } from "@earendil-works/pi-coding-agent";
33
import type { ReplyPayload } from "../auto-reply/reply-payload.js";
44
import { SILENT_REPLY_TOKEN } from "../auto-reply/tokens.js";
5+
import { appendSessionTranscriptMessage } from "../config/sessions/transcript-append.js";
56
import { formatErrorMessage } from "../infra/errors.js";
67
import { createSubsystemLogger } from "../logging/subsystem.js";
78
import { buildAgentHookContextChannelFields } from "../plugins/hook-agent-context.js";
@@ -109,6 +110,35 @@ function buildCliContextEngineAssistantMessage(params: {
109110
return buildCliHookAssistantMessage(params) as AgentMessage;
110111
}
111112

113+
async function persistCliCurrentUserMessage(params: RunCliAgentParams): Promise<boolean> {
114+
if (params.suppressNextUserMessagePersistence === true) {
115+
return false;
116+
}
117+
const prompt = params.transcriptPrompt ?? params.prompt;
118+
if (!prompt.trim()) {
119+
return false;
120+
}
121+
122+
try {
123+
const { message } = await appendSessionTranscriptMessage({
124+
transcriptPath: params.sessionFile,
125+
sessionId: params.sessionId,
126+
cwd: params.workspaceDir,
127+
config: params.config,
128+
message: {
129+
role: "user",
130+
content: prompt,
131+
timestamp: Date.now(),
132+
},
133+
});
134+
params.onUserMessagePersisted?.(message as Extract<AgentMessage, { role: "user" }>);
135+
return true;
136+
} catch (err) {
137+
log.warn(`CLI run: failed to persist user transcript message: ${formatErrorMessage(err)}`);
138+
return false;
139+
}
140+
}
141+
112142
type CliAgentEndHookParams = Parameters<typeof runAgentHarnessAgentEndHook>[0];
113143

114144
function shouldAwaitCliAgentEndHook(params: RunCliAgentParams): boolean {
@@ -368,6 +398,14 @@ export async function runPreparedCliAgent(
368398
},
369399
});
370400

401+
let currentUserMessagePersisted = false;
402+
const persistCurrentUserMessage = async (): Promise<void> => {
403+
if (currentUserMessagePersisted) {
404+
return;
405+
}
406+
currentUserMessagePersisted = await persistCliCurrentUserMessage(params);
407+
};
408+
371409
const persistBlockedBeforeAgentRun = async (block: {
372410
message: string;
373411
pluginId: string;
@@ -620,6 +658,7 @@ export async function runPreparedCliAgent(
620658
}
621659
}
622660

661+
await persistCurrentUserMessage();
623662
runAgentHarnessLlmInputHook({
624663
event: llmInputEvent,
625664
ctx: hookContext,
@@ -658,6 +697,7 @@ export async function runPreparedCliAgent(
658697

659698
// For now, retry without the session ID to create a new session
660699
try {
700+
await persistCurrentUserMessage();
661701
const { output, lastAssistant } = await executeCliAttempt(undefined);
662702
const assistantText = output.text.trim();
663703
const effectiveCliSessionId = output.sessionId;
@@ -743,6 +783,8 @@ export function buildRunClaudeCliAgentParams(params: RunClaudeCliAgentParams): R
743783
images: params.images,
744784
messageChannel: params.messageChannel,
745785
messageProvider: params.messageProvider,
786+
suppressNextUserMessagePersistence: params.suppressNextUserMessagePersistence,
787+
onUserMessagePersisted: params.onUserMessagePersisted,
746788
};
747789
}
748790

src/agents/cli-runner/types.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
import type { AgentMessage } from "@earendil-works/pi-agent-core";
12
import type { ImageContent } from "@earendil-works/pi-ai";
23
import type { SourceReplyDeliveryMode } from "../../auto-reply/get-reply-options.types.js";
34
import type { ReplyOperation } from "../../auto-reply/reply/reply-run-registry.js";
@@ -77,6 +78,8 @@ export type RunCliAgentParams = {
7778
source?: string;
7879
firstModelCallStarted?: boolean;
7980
}) => void;
81+
suppressNextUserMessagePersistence?: boolean;
82+
onUserMessagePersisted?: (message: Extract<AgentMessage, { role: "user" }>) => void;
8083
replyOperation?: ReplyOperation;
8184
/**
8285
* Close any long-lived CLI live session created for this run after the run

src/agents/command/attempt-execution.cli.test.ts

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -562,6 +562,39 @@ describe("CLI attempt execution", () => {
562562
});
563563
});
564564

565+
it("can append only the CLI assistant when the user turn is already persisted", async () => {
566+
const sessionKey = "agent:main:subagent:cli-transcript-assistant-only";
567+
const sessionEntry: SessionEntry = {
568+
sessionId: "session-cli-transcript-assistant-only",
569+
updatedAt: Date.now(),
570+
};
571+
const sessionStore: Record<string, SessionEntry> = { [sessionKey]: sessionEntry };
572+
await fs.writeFile(storePath, JSON.stringify(sessionStore, null, 2), "utf-8");
573+
574+
const updatedEntry = await persistCliTurnTranscript({
575+
body: "already durable",
576+
result: makeCliResult("assistant only"),
577+
sessionId: sessionEntry.sessionId,
578+
sessionKey,
579+
sessionEntry,
580+
sessionStore,
581+
storePath,
582+
sessionAgentId: "main",
583+
sessionCwd: tmpDir,
584+
config: {},
585+
skipUserMessage: true,
586+
});
587+
588+
const messages = await readSessionMessages(updatedEntry?.sessionFile ?? "");
589+
expect(messages).toHaveLength(1);
590+
expectRecordFields(requireRecord(messages[0], "assistant message"), {
591+
role: "assistant",
592+
provider: "claude-cli",
593+
model: "opus",
594+
content: [{ type: "text", text: "assistant only" }],
595+
});
596+
});
597+
565598
it("embedded assistant gap-fill skips user mirror and dedupes identical assistant tails", async () => {
566599
const sessionKey = "agent:main:subagent:embedded-gap-fill";
567600
const sessionEntry: SessionEntry = {

src/agents/command/attempt-execution.ts

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,7 @@ type PersistTextTurnTranscriptParams = {
111111
model: string;
112112
usage?: TranscriptUsage;
113113
};
114+
skipUserMessage?: boolean;
114115
};
115116

116117
type HarnessAuthProfileSelection = {
@@ -229,7 +230,7 @@ async function persistTextTurnTranscript(
229230
allowReentrant: true,
230231
});
231232
try {
232-
if (promptText) {
233+
if (promptText && !params.skipUserMessage) {
233234
await appendSessionTranscriptMessage({
234235
transcriptPath: sessionFile,
235236
sessionId: params.sessionId,
@@ -335,6 +336,7 @@ export async function persistCliTurnTranscript(params: {
335336
sessionCwd: string;
336337
config: OpenClawConfig;
337338
embeddedAssistantGapFill?: boolean;
339+
skipUserMessage?: boolean;
338340
}): Promise<SessionEntry | undefined> {
339341
const replyText = resolveCliTranscriptReplyText(params.result);
340342
const provider = params.result.meta.agentMeta?.provider?.trim() ?? "cli";
@@ -355,6 +357,7 @@ export async function persistCliTurnTranscript(params: {
355357
sessionCwd: params.sessionCwd,
356358
config: params.config,
357359
embeddedAssistantGapFill: gapFill,
360+
skipUserMessage: gapFill || params.skipUserMessage === true,
358361
assistant: {
359362
api: "cli",
360363
provider,
@@ -549,6 +552,8 @@ export function runAgentAttempt(params: {
549552
agentAccountId: params.runContext.accountId,
550553
senderIsOwner: params.opts.senderIsOwner,
551554
toolsAllow: params.opts.toolsAllow,
555+
suppressNextUserMessagePersistence: params.suppressPromptPersistenceOnRetry === true,
556+
onUserMessagePersisted: params.onUserMessagePersisted,
552557
cleanupBundleMcpOnRunEnd: params.opts.cleanupBundleMcpOnRunEnd,
553558
cleanupCliLiveSessionOnRunEnd: params.opts.cleanupCliLiveSessionOnRunEnd,
554559
});

src/auto-reply/reply/agent-runner-execution.test.ts

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5445,7 +5445,12 @@ describe("runAgentTurnWithFallback", () => {
54455445
attempts: [],
54465446
};
54475447
});
5448-
state.runCliAgentMock.mockRejectedValueOnce(new Error("cli failed"));
5448+
state.runCliAgentMock.mockImplementationOnce(
5449+
async (args: { onUserMessagePersisted?: (m: { role: "user"; content: string }) => void }) => {
5450+
args.onUserMessagePersisted?.({ role: "user", content: "queued" });
5451+
throw new Error("cli failed");
5452+
},
5453+
);
54495454
state.runEmbeddedPiAgentMock.mockResolvedValueOnce({
54505455
payloads: [{ text: "ok" }],
54515456
meta: {},
@@ -5457,6 +5462,7 @@ describe("runAgentTurnWithFallback", () => {
54575462
expect(state.runCliAgentMock).toHaveBeenCalledOnce();
54585463
expect(state.runEmbeddedPiAgentMock).toHaveBeenCalledOnce();
54595464
expectMockCallArgFields(state.runEmbeddedPiAgentMock, 0, "embedded fallback candidate", {
5465+
suppressNextUserMessagePersistence: true,
54605466
suppressAssistantErrorPersistence: false,
54615467
});
54625468
});

src/auto-reply/reply/agent-runner-execution.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1802,6 +1802,10 @@ export async function runAgentTurnWithFallback(params: {
18021802
disableTools: params.opts?.disableTools,
18031803
abortSignal: params.replyOperation?.abortSignal ?? params.opts?.abortSignal,
18041804
replyOperation: params.replyOperation,
1805+
suppressNextUserMessagePersistence: suppressQueuedUserPersistenceForCandidate,
1806+
onUserMessagePersisted: () => {
1807+
queuedUserMessagePersistedAcrossFallback = true;
1808+
},
18051809
},
18061810
transformResult: (rawResult) =>
18071811
isRoomEventCliRun && rawResult.meta.agentMeta

src/auto-reply/reply/followup-runner.test.ts

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -853,7 +853,14 @@ describe("createFollowupRunner runtime config", () => {
853853
};
854854
},
855855
);
856-
runCliAgentMock.mockRejectedValueOnce(new Error("cli failed"));
856+
runCliAgentMock.mockImplementationOnce(
857+
async (args: {
858+
onUserMessagePersisted?: (message: { role: "user"; content: string }) => void;
859+
}) => {
860+
args.onUserMessagePersisted?.({ role: "user", content: "queued" });
861+
throw new Error("cli failed");
862+
},
863+
);
857864
runEmbeddedPiAgentMock.mockImplementationOnce(async (params: { runId: string }) => {
858865
realAgentEvents.emitAgentEvent({
859866
runId: params.runId,
@@ -898,6 +905,7 @@ describe("createFollowupRunner runtime config", () => {
898905
expect(runCliAgentMock).toHaveBeenCalledTimes(1);
899906
expect(runEmbeddedPiAgentMock).toHaveBeenCalledTimes(1);
900907
const embeddedCall = requireLastMockCallArg(runEmbeddedPiAgentMock, "run embedded pi agent");
908+
expect(embeddedCall.suppressNextUserMessagePersistence).toBe(true);
901909
expect(embeddedCall.suppressAssistantErrorPersistence).toBe(false);
902910
expect(lifecyclePhases).toEqual(["start", "start", "end"]);
903911
});

src/auto-reply/reply/followup-runner.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -757,6 +757,10 @@ export function createFollowupRunner(params: {
757757
agentAccountId: run.agentAccountId,
758758
disableTools: opts?.disableTools,
759759
abortSignal: queued.abortSignal,
760+
suppressNextUserMessagePersistence: suppressQueuedUserPersistenceForCandidate,
761+
onUserMessagePersisted: () => {
762+
queuedUserMessagePersistedAcrossFallback = true;
763+
},
760764
},
761765
transformResult: (rawResult) =>
762766
isRoomEventCliRun && rawResult.meta.agentMeta

0 commit comments

Comments
 (0)