Skip to content

Commit 14acb97

Browse files
authored
fix(server): preserve Pi message IDs after resume (#2313)
1 parent 4a4556f commit 14acb97

4 files changed

Lines changed: 252 additions & 74 deletions

File tree

packages/server/src/server/agent/providers/pi/agent.test.ts

Lines changed: 97 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -57,21 +57,30 @@ function readUtf8File(pathname: string): string {
5757
}
5858
}
5959

60-
async function applyPaseoExtensionSystemPrompt(
60+
type PaseoExtensionListener = (event: unknown, context?: unknown) => unknown;
61+
62+
async function loadPaseoExtensionListeners(
6163
extensionPath: string,
62-
systemPrompt: string,
63-
): Promise<string | undefined> {
64-
const listeners = new Map<string, (event: { systemPrompt: string }) => unknown>();
64+
): Promise<Map<string, PaseoExtensionListener>> {
65+
const listeners = new Map<string, PaseoExtensionListener>();
6566
const extension = (await import(pathToFileURL(extensionPath).href)) as {
6667
default: (piApi: {
67-
on: (event: string, listener: (event: { systemPrompt: string }) => unknown) => void;
68+
on: (event: string, listener: PaseoExtensionListener) => void;
6869
registerCommand: () => void;
6970
}) => void;
7071
};
7172
extension.default({
7273
on: (event, listener) => listeners.set(event, listener),
7374
registerCommand: () => undefined,
7475
});
76+
return listeners;
77+
}
78+
79+
async function applyPaseoExtensionSystemPrompt(
80+
extensionPath: string,
81+
systemPrompt: string,
82+
): Promise<string | undefined> {
83+
const listeners = await loadPaseoExtensionListeners(extensionPath);
7584
const result = await listeners.get("before_agent_start")?.({ systemPrompt });
7685
return (result as { systemPrompt?: string } | undefined)?.systemPrompt;
7786
}
@@ -594,15 +603,15 @@ describe("PiRpcAgentSession", () => {
594603
]);
595604
});
596605

597-
test("emits live user messages with captured Pi tree entry ids", async () => {
606+
test("emits live user messages with submitted Pi tree entry ids", async () => {
598607
const { pi, session, events } = await createSession();
599608
const fakeSession = pi.latestSession();
600609

601-
fakeSession.capturedUserEntries = [{ id: "entry-user-1", parentId: null, text: "hello" }];
602610
await session.startTurn("hello");
603-
fakeSession.emit({
604-
type: "message_end",
605-
message: { role: "user", content: "hello" },
611+
fakeSession.finishSubmittedUserMessage({
612+
id: "entry-user-1",
613+
parentId: null,
614+
text: "hello",
606615
});
607616

608617
await events.nextTimelineEvent();
@@ -612,6 +621,38 @@ describe("PiRpcAgentSession", () => {
612621
]);
613622
});
614623

624+
test("uses the Pi entry attached to a submitted prompt after resuming old history", async () => {
625+
const pi = new FakePi();
626+
const client = createClient(pi);
627+
const session = (await client.resumeSession({
628+
provider: "pi",
629+
sessionId: "pi-session-1",
630+
nativeHandle: "/tmp/native-pi-session",
631+
metadata: { cwd: "/workspace/project" },
632+
})) as PiRpcAgentSession;
633+
const events = new SessionEvents(session);
634+
const fakeSession = pi.latestSession();
635+
fakeSession.capturedUserEntries = [{ id: "entry-old", parentId: null, text: "old prompt" }];
636+
637+
await session.startTurn("new prompt", { clientMessageId: "client-new" });
638+
fakeSession.finishSubmittedUserMessage({
639+
id: "entry-new",
640+
parentId: "entry-old-assistant",
641+
text: "new prompt",
642+
});
643+
644+
await events.nextTimelineEvent();
645+
646+
expect(events.timelineItems()).toEqual([
647+
{
648+
type: "user_message",
649+
text: "new prompt",
650+
messageId: "entry-new",
651+
clientMessageId: "client-new",
652+
},
653+
]);
654+
});
655+
615656
test("surfaces Pi extension command messages and completes when no agent turn starts", async () => {
616657
const { pi, session, events } = await createSession();
617658
const fakeSession = pi.latestSession();
@@ -734,6 +775,52 @@ describe("PiRpcAgentSession", () => {
734775
]);
735776
});
736777

778+
test("reports the persisted Pi entry attached to the submitted message", async () => {
779+
const pi = new FakePi();
780+
const client = createClient(pi);
781+
const session = await client.createSession(createConfig());
782+
const extensionPath = pi.recordedLaunches[0]?.extensionPaths[0];
783+
expect(extensionPath).toBeDefined();
784+
const listeners = await loadPaseoExtensionListeners(extensionPath!);
785+
const submittedMessage = { role: "user", content: "new prompt" };
786+
const entries: Array<{
787+
type: string;
788+
id: string;
789+
parentId: string | null;
790+
message: { role: string; content: string };
791+
}> = [
792+
{
793+
type: "message",
794+
id: "entry-old",
795+
parentId: null,
796+
message: { role: "user", content: "old prompt" },
797+
},
798+
];
799+
const notifications: string[] = [];
800+
const context = {
801+
sessionManager: { getEntries: () => entries },
802+
ui: { notify: (message: string) => notifications.push(message) },
803+
};
804+
805+
await listeners.get("message_end")?.({ message: submittedMessage }, context);
806+
entries.push({
807+
type: "message",
808+
id: "entry-new",
809+
parentId: "entry-old-assistant",
810+
message: submittedMessage,
811+
});
812+
await listeners.get("message_start")?.(
813+
{ message: { role: "assistant", content: [] } },
814+
context,
815+
);
816+
817+
expect(notifications).toEqual([
818+
'PASEO_SUBMITTED_USER_ENTRY {"entry":{"id":"entry-new","parentId":"entry-old-assistant","text":"new prompt"}}',
819+
]);
820+
821+
await session.close();
822+
});
823+
737824
test("appends agent and daemon prompts after Pi's discovered system prompt", async () => {
738825
const pi = new FakePi();
739826
const client = createClient(pi);

packages/server/src/server/agent/providers/pi/agent.ts

Lines changed: 69 additions & 61 deletions
Original file line numberDiff line numberDiff line change
@@ -92,6 +92,7 @@ const PI_CATALOG_REQUEST_TIMEOUT_MS = 120_000;
9292
const PASEO_PI_TREE_EXTENSION_COMMAND = "paseo_tree";
9393
const PASEO_PI_CAPTURE_EXTENSION_COMMAND = "paseo_capture_entries";
9494
const PASEO_PI_ENTRY_CAPTURE_MARKER = "PASEO_ENTRY_CAPTURE";
95+
const PASEO_PI_SUBMITTED_USER_ENTRY_MARKER = "PASEO_SUBMITTED_USER_ENTRY";
9596
const PASEO_PI_COMMAND_RESULT_MARKER = "PASEO_COMMAND_RESULT";
9697
const DEFAULT_PI_EXTENSION_RESULT_TIMEOUT_MS = 30_000;
9798
const QUESTION_RESPONSE_HEADER = "Response";
@@ -256,12 +257,6 @@ interface PiCapturedEntry extends PiCapturedUserMessageEntry {
256257
parentId: string | null;
257258
}
258259

259-
interface PendingPiUserMessage {
260-
text: string;
261-
turnId: string | undefined;
262-
clientMessageId?: string;
263-
}
264-
265260
interface PendingExtensionResult {
266261
resolve: (value: unknown) => void;
267262
reject: (error: Error) => void;
@@ -653,11 +648,15 @@ function createPiPaseoExtensionFile(systemPrompt?: string): PiTempFile {
653648
return ctx.sessionManager
654649
.getEntries()
655650
.filter((entry) => entry.type === "message" && entry.message?.role === "user")
656-
.map((entry) => ({
651+
.map(toCapturedUserEntry);
652+
}
653+
654+
function toCapturedUserEntry(entry) {
655+
return {
657656
id: entry.id,
658657
parentId: entry.parentId ?? null,
659658
text: readTextContent(entry.message.content),
660-
}));
659+
};
661660
}
662661
663662
function emitEntryCapture(ctx, reason, requestId) {
@@ -676,6 +675,30 @@ function createPiPaseoExtensionFile(systemPrompt?: string): PiTempFile {
676675
}
677676
678677
export default function paseoIntegration(pi) {
678+
const submittedUserMessages = [];
679+
680+
function emitSubmittedUserEntries(ctx) {
681+
const entries = ctx.sessionManager.getEntries();
682+
for (let index = 0; index < submittedUserMessages.length; index += 1) {
683+
const message = submittedUserMessages[index];
684+
// Pi assigns the entry ID after message_end, then persists this same message object.
685+
// Reference equality preserves the exact association even when another extension edits it.
686+
const entry = entries.find(
687+
(candidate) => candidate.type === "message" && candidate.message === message,
688+
);
689+
if (!entry) {
690+
continue;
691+
}
692+
submittedUserMessages.splice(index, 1);
693+
index -= 1;
694+
ctx.ui.notify(
695+
"${PASEO_PI_SUBMITTED_USER_ENTRY_MARKER} " +
696+
JSON.stringify({ entry: toCapturedUserEntry(entry) }),
697+
"info",
698+
);
699+
}
700+
}
701+
679702
${
680703
systemPrompt
681704
? `pi.on("before_agent_start", async (event) => ({
@@ -688,7 +711,20 @@ function createPiPaseoExtensionFile(systemPrompt?: string): PiTempFile {
688711
emitEntryCapture(ctx, "session_start");
689712
});
690713
714+
pi.on("message_end", async (event) => {
715+
if (event.message?.role === "user") {
716+
submittedUserMessages.push(event.message);
717+
}
718+
});
719+
720+
pi.on("message_start", async (event, ctx) => {
721+
if (event.message?.role === "assistant") {
722+
emitSubmittedUserEntries(ctx);
723+
}
724+
});
725+
691726
pi.on("turn_end", async (_event, ctx) => {
727+
emitSubmittedUserEntries(ctx);
692728
emitEntryCapture(ctx, "turn_end");
693729
});
694730
@@ -1197,8 +1233,6 @@ export class PiRpcAgentSession implements AgentSession {
11971233
currentLeafOverrideId: string | null | undefined;
11981234
private readonly capturedUserEntries: PiCapturedEntry[] = [];
11991235
private readonly capturedUserEntriesById = new Map<string, PiCapturedEntry>();
1200-
private readonly seenUserEntryIds = new Set<string>();
1201-
private readonly pendingUserMessages: PendingPiUserMessage[] = [];
12021236
private readonly pendingExtensionResults = new Map<string, PendingExtensionResult>();
12031237
private outOfBandCompactionEmit: ((event: AgentStreamEvent) => void) | null = null;
12041238
private outOfBandCompactionStarted = false;
@@ -1776,42 +1810,34 @@ export class PiRpcAgentSession implements AgentSession {
17761810
}
17771811

17781812
private recordCapturedUserEntries(entries: PiCapturedEntry[]): void {
1779-
const previouslySeenEntryIds = new Set(this.seenUserEntryIds);
17801813
this.capturedUserEntries.splice(0, this.capturedUserEntries.length, ...entries);
17811814
this.capturedUserEntriesById.clear();
17821815
for (const entry of entries) {
17831816
this.capturedUserEntriesById.set(entry.id, entry);
17841817
}
1785-
this.flushPendingUserMessages(previouslySeenEntryIds);
1786-
for (const entry of entries) {
1787-
this.seenUserEntryIds.add(entry.id);
1788-
}
17891818
}
17901819

1791-
private flushPendingUserMessages(previouslySeenEntryIds: Set<string>): void {
1792-
for (let index = 0; index < this.pendingUserMessages.length; index += 1) {
1793-
const pending = this.pendingUserMessages[index]!;
1794-
const entry = this.capturedUserEntries.find(
1795-
(candidate) => !previouslySeenEntryIds.has(candidate.id),
1796-
);
1797-
if (!entry) {
1798-
continue;
1799-
}
1800-
previouslySeenEntryIds.add(entry.id);
1801-
this.pendingUserMessages.splice(index, 1);
1802-
index -= 1;
1803-
this.emit({
1804-
type: "timeline",
1805-
provider: this.provider,
1806-
turnId: pending.turnId,
1807-
item: {
1808-
type: "user_message",
1809-
text: pending.text,
1810-
messageId: entry.id,
1811-
...(pending.clientMessageId ? { clientMessageId: pending.clientMessageId } : {}),
1812-
},
1813-
});
1820+
private handleSubmittedUserEntryMarker(message: string): boolean {
1821+
const payload = parseExtensionMarkerPayload(message, PASEO_PI_SUBMITTED_USER_ENTRY_MARKER);
1822+
if (!payload) {
1823+
return false;
1824+
}
1825+
const [entry] = parseCapturedEntries([payload.entry]);
1826+
if (!entry) {
1827+
return true;
18141828
}
1829+
this.emit({
1830+
type: "timeline",
1831+
provider: this.provider,
1832+
turnId: this.currentTurnIdForEvent(),
1833+
item: {
1834+
type: "user_message",
1835+
text: entry.text,
1836+
messageId: entry.id,
1837+
...(this.activeClientMessageId ? { clientMessageId: this.activeClientMessageId } : {}),
1838+
},
1839+
});
1840+
return true;
18151841
}
18161842

18171843
private handleEntryCaptureMarker(message: string): boolean {
@@ -1849,7 +1875,11 @@ export class PiRpcAgentSession implements AgentSession {
18491875
): void {
18501876
const message = optionalString(event.message);
18511877
if (event.method === "notify" && message) {
1852-
if (this.handleEntryCaptureMarker(message) || this.handleCommandResultMarker(message)) {
1878+
if (
1879+
this.handleSubmittedUserEntryMarker(message) ||
1880+
this.handleEntryCaptureMarker(message) ||
1881+
this.handleCommandResultMarker(message)
1882+
) {
18531883
return;
18541884
}
18551885
this.bufferNoTurnOutput(message);
@@ -2176,28 +2206,6 @@ export class PiRpcAgentSession implements AgentSession {
21762206
this.completeTurn(turnId, []);
21772207
return;
21782208
}
2179-
2180-
if (event.message.role !== "user") {
2181-
return;
2182-
}
2183-
const text = getUserMessageText(event.message.content);
2184-
if (!text) {
2185-
return;
2186-
}
2187-
this.pendingUserMessages.push({
2188-
text,
2189-
turnId,
2190-
...(this.activeClientMessageId ? { clientMessageId: this.activeClientMessageId } : {}),
2191-
});
2192-
void this.requestEntryCapture("message_end").catch((error: unknown) => {
2193-
const message = error instanceof Error ? error.message : String(error);
2194-
this.emit({
2195-
type: "turn_failed",
2196-
provider: this.provider,
2197-
turnId,
2198-
error: message,
2199-
});
2200-
});
22012209
}
22022210

22032211
private emitToolCallEvent(

packages/server/src/server/agent/providers/pi/test-utils/fake-pi.ts

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,12 @@ export interface FakePiSubagentMessagesResult {
4545
messages: PiAgentMessage[];
4646
}
4747

48+
interface FakePiUserEntry {
49+
id: string;
50+
parentId: string | null;
51+
text: string;
52+
}
53+
4854
export class FakePi implements PiRuntime {
4955
readonly recordedLaunches: PiRuntimeLaunch[] = [];
5056
private readonly sessions: FakePiSession[] = [];
@@ -349,6 +355,19 @@ export class FakePiSession implements PiRuntimeSession {
349355
this.emit({ type: "agent_end", messages: this.messages });
350356
}
351357

358+
finishSubmittedUserMessage(entry: FakePiUserEntry): void {
359+
this.emit({
360+
type: "message_end",
361+
message: { role: "user", content: entry.text },
362+
});
363+
this.emit({
364+
type: "extension_ui_request",
365+
id: `submitted-user-${entry.id}`,
366+
method: "notify",
367+
message: `PASEO_SUBMITTED_USER_ENTRY ${JSON.stringify({ entry })}`,
368+
});
369+
}
370+
352371
private handleTreeNavigationCommand(message: string): void {
353372
const prefix = "/paseo_tree ";
354373
if (!message.startsWith(prefix)) {

0 commit comments

Comments
 (0)