mirror of
https://github.com/getpaseo/paseo.git
synced 2026-07-29 12:01:31 +00:00
Show Fork chat for every agent provider (#2022)
* fix(chat): show fork controls for every provider Real assistant messages require stable identity so provider-agnostic chat actions can address the selected turn. Preserve provider IDs where available and generate stable adapter-owned IDs otherwise. * test(opencode): expect assistant message identities * test(pi): cover assistant updates without message start * fix(providers): preserve late assistant identities * fix(opencode): track fallback compaction identities
This commit is contained in:
@@ -2090,11 +2090,16 @@ describe("ACPAgentSession", () => {
|
||||
|
||||
test("emits assistant and reasoning chunks as deltas while user chunks stay accumulated", async () => {
|
||||
const session = createSession();
|
||||
const events: Array<{ type: string; item?: { type: string; text?: string } }> = [];
|
||||
const events: Array<{
|
||||
type: string;
|
||||
item?: { type: string; text?: string; messageId?: string };
|
||||
}> = [];
|
||||
asInternals<ACPSessionInternals>(session).sessionId = "session-1";
|
||||
|
||||
session.subscribe((event) => {
|
||||
events.push(event as { type: string; item?: { type: string; text?: string } });
|
||||
events.push(
|
||||
event as { type: string; item?: { type: string; text?: string; messageId?: string } },
|
||||
);
|
||||
});
|
||||
|
||||
await session.sessionUpdate({
|
||||
@@ -2152,8 +2157,8 @@ describe("ACPAgentSession", () => {
|
||||
.filter(Boolean);
|
||||
|
||||
expect(timeline).toEqual([
|
||||
{ type: "assistant_message", text: "Hey!" },
|
||||
{ type: "assistant_message", text: " How are you?" },
|
||||
{ type: "assistant_message", text: "Hey!", messageId: "assistant-1" },
|
||||
{ type: "assistant_message", text: " How are you?", messageId: "assistant-1" },
|
||||
{ type: "reasoning", text: "Thinking" },
|
||||
{ type: "reasoning", text: " more" },
|
||||
{ type: "user_message", text: "hel", messageId: "user-1" },
|
||||
@@ -2161,6 +2166,48 @@ describe("ACPAgentSession", () => {
|
||||
]);
|
||||
});
|
||||
|
||||
test("assigns one fallback ID per contiguous assistant message", async () => {
|
||||
const session = createSession();
|
||||
const assistantMessages: Array<{ text: string; messageId?: string }> = [];
|
||||
asInternals<ACPSessionInternals>(session).sessionId = "session-1";
|
||||
|
||||
session.subscribe((event) => {
|
||||
if (event.type === "timeline" && event.item.type === "assistant_message") {
|
||||
assistantMessages.push(event.item);
|
||||
}
|
||||
});
|
||||
|
||||
for (const text of ["First", " message"]) {
|
||||
await session.sessionUpdate({
|
||||
sessionId: "session-1",
|
||||
update: {
|
||||
sessionUpdate: "agent_message_chunk",
|
||||
content: { type: "text", text },
|
||||
} as SessionUpdate,
|
||||
});
|
||||
}
|
||||
await session.sessionUpdate({
|
||||
sessionId: "session-1",
|
||||
update: {
|
||||
sessionUpdate: "agent_thought_chunk",
|
||||
content: { type: "text", text: "Next response" },
|
||||
} as SessionUpdate,
|
||||
});
|
||||
await session.sessionUpdate({
|
||||
sessionId: "session-1",
|
||||
update: {
|
||||
sessionUpdate: "agent_message_chunk",
|
||||
content: { type: "text", text: "Second message" },
|
||||
} as SessionUpdate,
|
||||
});
|
||||
|
||||
expect(assistantMessages).toHaveLength(3);
|
||||
expect(assistantMessages[0].messageId).toEqual(expect.any(String));
|
||||
expect(assistantMessages[1].messageId).toBe(assistantMessages[0].messageId);
|
||||
expect(assistantMessages[2].messageId).toEqual(expect.any(String));
|
||||
expect(assistantMessages[2].messageId).not.toBe(assistantMessages[0].messageId);
|
||||
});
|
||||
|
||||
test("startTurn returns before the ACP prompt settles and completes later via subscribers", async () => {
|
||||
const session = createSession();
|
||||
const events: Array<{ type: string; turnId?: string }> = [];
|
||||
@@ -2727,6 +2774,103 @@ describe("ACP session/load invariant — cwd and mcpServers always passed", () =
|
||||
});
|
||||
});
|
||||
|
||||
test("preserves assistant message IDs from loadSession replay", async () => {
|
||||
let session!: ACPAgentSession;
|
||||
const loadSession = async () => {
|
||||
await session.sessionUpdate({
|
||||
sessionId: "session-1",
|
||||
update: {
|
||||
sessionUpdate: "agent_message_chunk",
|
||||
messageId: "assistant-replay-1",
|
||||
content: { type: "text", text: "Welcome back" },
|
||||
} as SessionUpdate,
|
||||
});
|
||||
return {
|
||||
sessionId: "session-1",
|
||||
modes: null,
|
||||
models: null,
|
||||
configOptions: [],
|
||||
};
|
||||
};
|
||||
({ session } = makeTestSession({
|
||||
capabilities: { loadSession: true },
|
||||
handle: { sessionId: "session-1", provider: "claude-acp" },
|
||||
loadSession,
|
||||
}));
|
||||
|
||||
await session.initializeResumedSession();
|
||||
|
||||
const history: AgentStreamEvent[] = [];
|
||||
for await (const event of session.streamHistory()) {
|
||||
history.push(event);
|
||||
}
|
||||
expect(history).toEqual([
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "claude-acp",
|
||||
item: {
|
||||
type: "assistant_message",
|
||||
text: "Welcome back",
|
||||
messageId: "assistant-replay-1",
|
||||
},
|
||||
},
|
||||
]);
|
||||
});
|
||||
|
||||
test("assigns stable fallback IDs to ID-less assistant messages during loadSession replay", async () => {
|
||||
let session!: ACPAgentSession;
|
||||
const loadSession = async () => {
|
||||
for (const text of ["Loaded", " response"]) {
|
||||
await session.sessionUpdate({
|
||||
sessionId: "session-1",
|
||||
update: {
|
||||
sessionUpdate: "agent_message_chunk",
|
||||
content: { type: "text", text },
|
||||
} as SessionUpdate,
|
||||
});
|
||||
}
|
||||
await session.sessionUpdate({
|
||||
sessionId: "session-1",
|
||||
update: {
|
||||
sessionUpdate: "user_message_chunk",
|
||||
content: { type: "text", text: "Follow up" },
|
||||
} as SessionUpdate,
|
||||
});
|
||||
await session.sessionUpdate({
|
||||
sessionId: "session-1",
|
||||
update: {
|
||||
sessionUpdate: "agent_message_chunk",
|
||||
content: { type: "text", text: "Loaded second response" },
|
||||
} as SessionUpdate,
|
||||
});
|
||||
return {
|
||||
sessionId: "session-1",
|
||||
modes: null,
|
||||
models: null,
|
||||
configOptions: [],
|
||||
};
|
||||
};
|
||||
({ session } = makeTestSession({
|
||||
capabilities: { loadSession: true },
|
||||
handle: { sessionId: "session-1", provider: "claude-acp" },
|
||||
loadSession,
|
||||
}));
|
||||
|
||||
await session.initializeResumedSession();
|
||||
|
||||
const assistantMessages: Array<{ text: string; messageId?: string }> = [];
|
||||
for await (const event of session.streamHistory()) {
|
||||
if (event.type === "timeline" && event.item.type === "assistant_message") {
|
||||
assistantMessages.push(event.item);
|
||||
}
|
||||
}
|
||||
expect(assistantMessages).toHaveLength(3);
|
||||
expect(assistantMessages[0].messageId).toEqual(expect.any(String));
|
||||
expect(assistantMessages[1].messageId).toBe(assistantMessages[0].messageId);
|
||||
expect(assistantMessages[2].messageId).toEqual(expect.any(String));
|
||||
expect(assistantMessages[2].messageId).not.toBe(assistantMessages[0].messageId);
|
||||
});
|
||||
|
||||
test("loadSession is always called with mcpServers even when supportsMcpServers is false", async () => {
|
||||
const { session, loadSession } = makeTestSession({
|
||||
capabilities: { loadSession: true, supportsMcpServers: false },
|
||||
|
||||
@@ -1331,6 +1331,7 @@ export class ACPAgentSession implements AgentSession, ACPClient {
|
||||
private readonly extensionCommandsParser?: ACPExtensionCommandsParser;
|
||||
private currentTurnUsage: AgentUsage | undefined;
|
||||
private activeForegroundTurnId: string | null = null;
|
||||
private fallbackAssistantMessageId: string | null = null;
|
||||
private closed = false;
|
||||
private historyPending = false;
|
||||
private replayingHistory = false;
|
||||
@@ -1475,6 +1476,7 @@ export class ACPAgentSession implements AgentSession, ACPClient {
|
||||
const turnId = randomUUID();
|
||||
const messageId = options?.messageId ?? randomUUID();
|
||||
this.activeForegroundTurnId = turnId;
|
||||
this.fallbackAssistantMessageId = null;
|
||||
this.activeSubmittedUserMessage = null;
|
||||
this.emitBootstrapThreadEvent();
|
||||
this.pushEvent({ type: "turn_started", provider: this.provider, turnId });
|
||||
@@ -2443,6 +2445,7 @@ export class ACPAgentSession implements AgentSession, ACPClient {
|
||||
private translateSessionUpdate(update: SessionUpdate): AgentStreamEvent[] {
|
||||
switch (update.sessionUpdate) {
|
||||
case "user_message_chunk": {
|
||||
this.fallbackAssistantMessageId = null;
|
||||
const item = this.createMessageTimelineItem("user_message", update);
|
||||
if (!item) {
|
||||
return [];
|
||||
@@ -2460,10 +2463,12 @@ export class ACPAgentSession implements AgentSession, ACPClient {
|
||||
return item ? [this.wrapTimeline(item)] : [];
|
||||
}
|
||||
case "agent_thought_chunk": {
|
||||
this.fallbackAssistantMessageId = null;
|
||||
const item = this.createMessageTimelineItem("reasoning", update);
|
||||
return item ? [this.wrapTimeline(item)] : [];
|
||||
}
|
||||
case "tool_call":
|
||||
this.fallbackAssistantMessageId = null;
|
||||
return this.handleToolCallUpdate(update.toolCallId, update, undefined);
|
||||
case "tool_call_update":
|
||||
return this.handleToolCallUpdate(
|
||||
@@ -2472,6 +2477,7 @@ export class ACPAgentSession implements AgentSession, ACPClient {
|
||||
this.toolCalls.get(update.toolCallId),
|
||||
);
|
||||
case "plan":
|
||||
this.fallbackAssistantMessageId = null;
|
||||
return [this.wrapTimeline(mapPlanToTimeline(update))];
|
||||
case "current_mode_update":
|
||||
this.handleCurrentModeUpdate(update);
|
||||
@@ -2526,7 +2532,7 @@ export class ACPAgentSession implements AgentSession, ACPClient {
|
||||
>,
|
||||
):
|
||||
| { type: "user_message"; text: string; messageId?: string }
|
||||
| { type: "assistant_message"; text: string }
|
||||
| { type: "assistant_message"; text: string; messageId: string }
|
||||
| { type: "reasoning"; text: string }
|
||||
| null {
|
||||
const chunkText = contentBlockToText(update.content);
|
||||
@@ -2542,11 +2548,24 @@ export class ACPAgentSession implements AgentSession, ACPClient {
|
||||
return { type: "user_message", text: state.text, messageId: update.messageId ?? undefined };
|
||||
}
|
||||
if (type === "assistant_message") {
|
||||
return { type: "assistant_message", text: chunkText };
|
||||
return {
|
||||
type: "assistant_message",
|
||||
text: chunkText,
|
||||
messageId: this.resolveAssistantMessageId(update.messageId),
|
||||
};
|
||||
}
|
||||
return { type: "reasoning", text: chunkText };
|
||||
}
|
||||
|
||||
private resolveAssistantMessageId(messageId: string | null | undefined): string {
|
||||
if (messageId) {
|
||||
this.fallbackAssistantMessageId = null;
|
||||
return messageId;
|
||||
}
|
||||
this.fallbackAssistantMessageId ??= randomUUID();
|
||||
return this.fallbackAssistantMessageId;
|
||||
}
|
||||
|
||||
private messageAssemblyKey(
|
||||
type: "user_message" | "assistant_message" | "reasoning",
|
||||
messageId: string | null | undefined,
|
||||
@@ -2701,6 +2720,7 @@ export class ACPAgentSession implements AgentSession, ACPClient {
|
||||
event: Extract<AgentStreamEvent, { type: "turn_completed" | "turn_failed" | "turn_canceled" }>,
|
||||
): void {
|
||||
this.activeForegroundTurnId = null;
|
||||
this.fallbackAssistantMessageId = null;
|
||||
if (this.activeSubmittedUserMessage?.turnId === event.turnId) {
|
||||
this.activeSubmittedUserMessage = null;
|
||||
}
|
||||
|
||||
@@ -215,12 +215,13 @@ describe("OpenCodeAgentClient adapter smoke tests", () => {
|
||||
|
||||
expect(turn.turnCompleted).toBe(true);
|
||||
expect(turn.turnFailed).toBe(false);
|
||||
expect(turn.assistantMessages.length).toBeGreaterThan(0);
|
||||
for (const msg of turn.assistantMessages) {
|
||||
expect(msg.text.length).toBeGreaterThan(0);
|
||||
}
|
||||
const fullResponse = turn.assistantMessages.map((m) => m.text).join("");
|
||||
expect(fullResponse).toBe("Hello from OpenCode");
|
||||
expect(turn.assistantMessages).toEqual([
|
||||
{
|
||||
type: "assistant_message",
|
||||
text: "Hello from OpenCode",
|
||||
messageId: "msg_assistant",
|
||||
},
|
||||
]);
|
||||
expect(openCodeClient.calls.sessionPromptAsync).toEqual([
|
||||
expect.objectContaining({
|
||||
sessionID: "session-1",
|
||||
@@ -234,6 +235,74 @@ describe("OpenCodeAgentClient adapter smoke tests", () => {
|
||||
rmSync(cwd, { recursive: true, force: true });
|
||||
}, 120_000);
|
||||
|
||||
test("completed and structured assistant messages preserve OpenCode message IDs", async () => {
|
||||
const cwd = tmpCwd();
|
||||
const runtime = new TestOpenCodeHarness();
|
||||
const openCodeClient = new TestOpenCodeClient();
|
||||
openCodeClient.sessionPromptAsyncEvents = [
|
||||
{
|
||||
type: "message.updated",
|
||||
properties: {
|
||||
info: {
|
||||
id: "msg_structured",
|
||||
sessionID: "session-1",
|
||||
role: "assistant",
|
||||
structured: "structured reply",
|
||||
time: { completed: 1 },
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
type: "message.updated",
|
||||
properties: {
|
||||
info: {
|
||||
id: "msg_completed",
|
||||
sessionID: "session-1",
|
||||
role: "assistant",
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
type: "message.part.updated",
|
||||
properties: {
|
||||
part: {
|
||||
id: "prt_completed",
|
||||
sessionID: "session-1",
|
||||
messageID: "msg_completed",
|
||||
type: "text",
|
||||
text: "completed reply",
|
||||
time: { start: 1, end: 2 },
|
||||
},
|
||||
},
|
||||
},
|
||||
{ type: "session.idle", properties: { sessionID: "session-1" } },
|
||||
];
|
||||
runtime.enqueueClient(openCodeClient);
|
||||
const client = new OpenCodeAgentClient(logger, undefined, {
|
||||
serverManager: runtime,
|
||||
createClient: runtime.createClient,
|
||||
});
|
||||
const session = await client.createSession(buildConfig(cwd));
|
||||
|
||||
const turn = await collectTurnEvents(streamSession(session, "Reply twice"));
|
||||
|
||||
expect(turn.assistantMessages).toEqual([
|
||||
{
|
||||
type: "assistant_message",
|
||||
text: "structured reply",
|
||||
messageId: "msg_structured",
|
||||
},
|
||||
{
|
||||
type: "assistant_message",
|
||||
text: "completed reply",
|
||||
messageId: "msg_completed",
|
||||
},
|
||||
]);
|
||||
|
||||
await session.close();
|
||||
rmSync(cwd, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
test("manual compact hides the generated summary text", async () => {
|
||||
const cwd = tmpCwd();
|
||||
const runtime = new TestOpenCodeHarness();
|
||||
@@ -1470,7 +1539,11 @@ describe("OpenCode adapter startTurn error handling", () => {
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
timestamp: "2026-05-14T12:41:23.612Z",
|
||||
item: { type: "assistant_message", text: "probe ok" },
|
||||
item: {
|
||||
type: "assistant_message",
|
||||
text: "probe ok",
|
||||
messageId: "msg_assistant",
|
||||
},
|
||||
},
|
||||
]);
|
||||
});
|
||||
@@ -1522,7 +1595,11 @@ describe("OpenCode adapter startTurn error handling", () => {
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text: "no clocks here" },
|
||||
item: {
|
||||
type: "assistant_message",
|
||||
text: "no clocks here",
|
||||
messageId: "msg_assistant",
|
||||
},
|
||||
},
|
||||
]);
|
||||
});
|
||||
|
||||
@@ -989,12 +989,16 @@ function buildOpenCodeReplayTimelineEvent(params: {
|
||||
|
||||
function buildOpenCodeReplayPartTimelineEvent(params: {
|
||||
part: OpenCodePart;
|
||||
message: { structured?: unknown; time?: { created?: number; completed?: number } | undefined };
|
||||
message: {
|
||||
id: string;
|
||||
structured?: unknown;
|
||||
time?: { created?: number; completed?: number } | undefined;
|
||||
};
|
||||
}): Extract<AgentStreamEvent, { type: "timeline" }> | null {
|
||||
const { part, message } = params;
|
||||
if (part.type === "text" && part.text) {
|
||||
return buildOpenCodeReplayTimelineEvent({
|
||||
item: { type: "assistant_message", text: part.text },
|
||||
item: { type: "assistant_message", text: part.text, messageId: message.id },
|
||||
message,
|
||||
part,
|
||||
});
|
||||
@@ -1183,7 +1187,7 @@ function buildOpenCodeReplayTimelineEvents(
|
||||
if (text) {
|
||||
events.push(
|
||||
buildOpenCodeReplayTimelineEvent({
|
||||
item: { type: "assistant_message", text },
|
||||
item: { type: "assistant_message", text, messageId: info.id },
|
||||
message: info,
|
||||
}),
|
||||
);
|
||||
@@ -2329,7 +2333,7 @@ function appendOpenCodeMessageUpdated(
|
||||
events.push({
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text },
|
||||
item: { type: "assistant_message", text, messageId: info.id },
|
||||
});
|
||||
}
|
||||
|
||||
@@ -2457,7 +2461,7 @@ function appendOpenCodeTextPart(
|
||||
events.push({
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text: part.text },
|
||||
item: { type: "assistant_message", text: part.text, messageId: part.messageID },
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -2523,8 +2527,12 @@ function appendOpenCodeMessagePartDelta(
|
||||
if (messageRole === "user") {
|
||||
return;
|
||||
}
|
||||
if (messageID && state.suppressAssistantMessagesUntilIdle?.active === true) {
|
||||
state.compactionSummaryMessageIds.add(messageID);
|
||||
const assistantMessageId = messageID || partID;
|
||||
if (!assistantMessageId) {
|
||||
return;
|
||||
}
|
||||
if (state.suppressAssistantMessagesUntilIdle?.active === true) {
|
||||
state.compactionSummaryMessageIds.add(assistantMessageId);
|
||||
return;
|
||||
}
|
||||
if (partID) {
|
||||
@@ -2533,7 +2541,11 @@ function appendOpenCodeMessagePartDelta(
|
||||
events.push({
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text: delta },
|
||||
item: {
|
||||
type: "assistant_message",
|
||||
text: delta,
|
||||
messageId: assistantMessageId,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -124,7 +124,11 @@ describe("translateOpenCodeEvent", () => {
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text: "hey! what can I help with?" },
|
||||
item: {
|
||||
type: "assistant_message",
|
||||
text: "hey! what can I help with?",
|
||||
messageId: "message-1",
|
||||
},
|
||||
},
|
||||
]);
|
||||
});
|
||||
@@ -167,7 +171,7 @@ describe("translateOpenCodeEvent", () => {
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text: "final text" },
|
||||
item: { type: "assistant_message", text: "final text", messageId: "message-2" },
|
||||
},
|
||||
]);
|
||||
});
|
||||
@@ -269,16 +273,66 @@ describe("translateOpenCodeEvent", () => {
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text: "hey! " },
|
||||
item: { type: "assistant_message", text: "hey! ", messageId: "msg-d1" },
|
||||
},
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text: "what's up?" },
|
||||
item: { type: "assistant_message", text: "what's up?", messageId: "msg-d1" },
|
||||
},
|
||||
]);
|
||||
});
|
||||
|
||||
it("uses the part id when an assistant delta omits its message id", () => {
|
||||
const state = createState();
|
||||
|
||||
const events = translateOpenCodeEvent(
|
||||
{
|
||||
type: "message.part.delta",
|
||||
properties: {
|
||||
sessionID: "session-1",
|
||||
partID: "part-without-message",
|
||||
field: "text",
|
||||
delta: "still visible",
|
||||
},
|
||||
},
|
||||
state,
|
||||
);
|
||||
|
||||
expect(events).toEqual([
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: {
|
||||
type: "assistant_message",
|
||||
text: "still visible",
|
||||
messageId: "part-without-message",
|
||||
},
|
||||
},
|
||||
]);
|
||||
});
|
||||
|
||||
it("suppresses part-id-only assistant deltas by their resolved identity", () => {
|
||||
const state = createState();
|
||||
state.suppressAssistantMessagesUntilIdle = { active: true };
|
||||
|
||||
const events = translateOpenCodeEvent(
|
||||
{
|
||||
type: "message.part.delta",
|
||||
properties: {
|
||||
sessionID: "session-1",
|
||||
partID: "compaction-part",
|
||||
field: "text",
|
||||
delta: "hidden summary",
|
||||
},
|
||||
},
|
||||
state,
|
||||
);
|
||||
|
||||
expect(events).toEqual([]);
|
||||
expect(state.compactionSummaryMessageIds).toEqual(new Set(["compaction-part"]));
|
||||
});
|
||||
|
||||
it("humanizes permission requests and includes shell detail when command metadata exists", () => {
|
||||
const state = createState();
|
||||
|
||||
@@ -1073,7 +1127,7 @@ describe("translateOpenCodeEvent", () => {
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text: "hello there" },
|
||||
item: { type: "assistant_message", text: "hello there", messageId: "msg-dd1" },
|
||||
},
|
||||
]);
|
||||
});
|
||||
@@ -1274,7 +1328,11 @@ describe("translateOpenCodeEvent", () => {
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "opencode",
|
||||
item: { type: "assistant_message", text: '{"summary":"hello"}' },
|
||||
item: {
|
||||
type: "assistant_message",
|
||||
text: '{"summary":"hello"}',
|
||||
messageId: "message-structured-1",
|
||||
},
|
||||
},
|
||||
]);
|
||||
expect(second).toEqual([]);
|
||||
|
||||
@@ -426,10 +426,19 @@ describe("PiRpcAgentSession", () => {
|
||||
const fakeSession = pi.latestSession();
|
||||
|
||||
await session.startTurn("hello");
|
||||
fakeSession.emit({
|
||||
type: "message_start",
|
||||
message: { role: "assistant", content: [], responseId: "response-1" },
|
||||
});
|
||||
fakeSession.emit({
|
||||
type: "message_update",
|
||||
message: { role: "assistant", content: [] },
|
||||
assistantMessageEvent: { type: "text_delta", delta: "hello" },
|
||||
message: { role: "assistant", content: [], responseId: "response-1" },
|
||||
assistantMessageEvent: { type: "text_delta", delta: "hel" },
|
||||
});
|
||||
fakeSession.emit({
|
||||
type: "message_update",
|
||||
message: { role: "assistant", content: [], responseId: "response-1" },
|
||||
assistantMessageEvent: { type: "text_delta", delta: "lo" },
|
||||
});
|
||||
fakeSession.emit({
|
||||
type: "message_update",
|
||||
@@ -454,7 +463,8 @@ describe("PiRpcAgentSession", () => {
|
||||
await events.nextTurnCompletion();
|
||||
|
||||
expect(events.timelineItems()).toEqual([
|
||||
{ type: "assistant_message", text: "hello" },
|
||||
{ type: "assistant_message", text: "hel", messageId: "response-1" },
|
||||
{ type: "assistant_message", text: "lo", messageId: "response-1" },
|
||||
{ type: "reasoning", text: "thinking" },
|
||||
{
|
||||
type: "tool_call",
|
||||
@@ -475,6 +485,60 @@ describe("PiRpcAgentSession", () => {
|
||||
]);
|
||||
});
|
||||
|
||||
test("keeps one generated message id when Pi omits message start and response id", async () => {
|
||||
const { pi, session, events } = await createSession();
|
||||
const fakeSession = pi.latestSession();
|
||||
|
||||
await session.startTurn("hello");
|
||||
fakeSession.emit({
|
||||
type: "message_update",
|
||||
message: { role: "assistant", content: [] },
|
||||
assistantMessageEvent: { type: "text_delta", delta: "hel" },
|
||||
});
|
||||
fakeSession.emit({
|
||||
type: "message_update",
|
||||
message: { role: "assistant", content: [] },
|
||||
assistantMessageEvent: { type: "text_delta", delta: "lo" },
|
||||
});
|
||||
|
||||
const [firstChunk, secondChunk] = events.timelineItems();
|
||||
expect(firstChunk).toMatchObject({
|
||||
type: "assistant_message",
|
||||
text: "hel",
|
||||
messageId: expect.any(String),
|
||||
});
|
||||
const firstMessageId = (firstChunk as { messageId: string }).messageId;
|
||||
expect(secondChunk).toEqual({
|
||||
type: "assistant_message",
|
||||
text: "lo",
|
||||
messageId: firstMessageId,
|
||||
});
|
||||
});
|
||||
|
||||
test("uses a response id that first appears on the assistant update", async () => {
|
||||
const { pi, session, events } = await createSession();
|
||||
const fakeSession = pi.latestSession();
|
||||
|
||||
await session.startTurn("hello");
|
||||
fakeSession.emit({
|
||||
type: "message_start",
|
||||
message: { role: "assistant", content: [] },
|
||||
});
|
||||
fakeSession.emit({
|
||||
type: "message_update",
|
||||
message: { role: "assistant", content: [], responseId: "late-response-id" },
|
||||
assistantMessageEvent: { type: "text_delta", delta: "hello" },
|
||||
});
|
||||
|
||||
expect(events.timelineItems()).toEqual([
|
||||
{
|
||||
type: "assistant_message",
|
||||
text: "hello",
|
||||
messageId: "late-response-id",
|
||||
},
|
||||
]);
|
||||
});
|
||||
|
||||
test("emits live user messages with captured Pi tree entry ids", async () => {
|
||||
const { pi, session, events } = await createSession();
|
||||
const fakeSession = pi.latestSession();
|
||||
|
||||
@@ -1064,6 +1064,7 @@ export class PiRpcAgentSession implements AgentSession {
|
||||
private activeAskUserDialog: ActiveAskUserDialog | null = null;
|
||||
private pendingCombinedAskUserResponse: PendingCombinedAskUserResponse | null = null;
|
||||
private activeTurnId: string | null = null;
|
||||
private activeAssistantMessageId: string | null = null;
|
||||
private lastKnownThinkingOptionId: string | null;
|
||||
currentLeafOverrideId: string | null | undefined;
|
||||
private readonly capturedUserEntries: PiCapturedEntry[] = [];
|
||||
@@ -1123,6 +1124,7 @@ export class PiRpcAgentSession implements AgentSession {
|
||||
const payload = convertPromptInput(prompt, { model: this.state.model });
|
||||
const turnId = randomUUID();
|
||||
this.activeTurnId = turnId;
|
||||
this.activeAssistantMessageId = null;
|
||||
|
||||
void this.runtimeSession.prompt(payload.text, payload.images).catch((error) => {
|
||||
if (this.activeTurnId !== turnId) {
|
||||
@@ -1711,6 +1713,7 @@ export class PiRpcAgentSession implements AgentSession {
|
||||
});
|
||||
return;
|
||||
case "message_start":
|
||||
this.handleMessageStart(event);
|
||||
return;
|
||||
case "message_end":
|
||||
this.handleMessageEnd(event, turnId);
|
||||
@@ -1818,6 +1821,8 @@ export class PiRpcAgentSession implements AgentSession {
|
||||
return;
|
||||
}
|
||||
if (event.assistantMessageEvent.type === "text_delta") {
|
||||
// Pi-compatible runtimes may emit updates without a preceding message_start.
|
||||
this.activeAssistantMessageId ??= event.message.responseId || randomUUID();
|
||||
this.emit({
|
||||
type: "timeline",
|
||||
provider: PI_PROVIDER,
|
||||
@@ -1825,6 +1830,7 @@ export class PiRpcAgentSession implements AgentSession {
|
||||
item: {
|
||||
type: "assistant_message",
|
||||
text: event.assistantMessageEvent.delta ?? "",
|
||||
messageId: this.activeAssistantMessageId,
|
||||
},
|
||||
});
|
||||
return;
|
||||
@@ -1842,10 +1848,20 @@ export class PiRpcAgentSession implements AgentSession {
|
||||
}
|
||||
}
|
||||
|
||||
private handleMessageStart(event: Extract<PiAgentSessionEvent, { type: "message_start" }>): void {
|
||||
if (event.message.role === "assistant") {
|
||||
this.activeAssistantMessageId = event.message.responseId || null;
|
||||
}
|
||||
}
|
||||
|
||||
private handleMessageEnd(
|
||||
event: Extract<PiAgentSessionEvent, { type: "message_end" }>,
|
||||
turnId: string | undefined,
|
||||
): void {
|
||||
if (event.message.role === "assistant") {
|
||||
this.activeAssistantMessageId = null;
|
||||
return;
|
||||
}
|
||||
if (event.message.role === "custom") {
|
||||
const text = getUserMessageText(event.message.content);
|
||||
if (text) {
|
||||
@@ -1906,6 +1922,7 @@ export class PiRpcAgentSession implements AgentSession {
|
||||
|
||||
private completeTurn(turnId: string | undefined, messages: PiAgentMessage[]): void {
|
||||
this.activeTurnId = null;
|
||||
this.activeAssistantMessageId = null;
|
||||
const errorMessage = latestPiErrorMessage(messages);
|
||||
if (typeof errorMessage === "string" && errorMessage.length > 0) {
|
||||
this.emit({
|
||||
|
||||
@@ -29,6 +29,7 @@ describe("Pi history mapper", () => {
|
||||
},
|
||||
{
|
||||
role: "assistant",
|
||||
responseId: "response-1",
|
||||
content: [
|
||||
{ type: "thinking", thinking: "checking file" },
|
||||
{ type: "toolCall", id: "tool-1", name: "read", arguments: { path: "note.txt" } },
|
||||
@@ -77,7 +78,7 @@ describe("Pi history mapper", () => {
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "pi",
|
||||
item: { type: "assistant_message", text: "done" },
|
||||
item: { type: "assistant_message", text: "done", messageId: "response-1" },
|
||||
},
|
||||
{
|
||||
type: "timeline",
|
||||
@@ -153,7 +154,11 @@ describe("Pi history mapper", () => {
|
||||
{
|
||||
type: "timeline",
|
||||
provider: "pi",
|
||||
item: { type: "assistant_message", text: "first answer" },
|
||||
item: {
|
||||
type: "assistant_message",
|
||||
text: "first answer",
|
||||
messageId: "pi-history-assistant-1",
|
||||
},
|
||||
},
|
||||
{
|
||||
type: "timeline",
|
||||
|
||||
@@ -44,6 +44,7 @@ export async function* streamPiHistory(
|
||||
): AsyncGenerator<AgentStreamEvent> {
|
||||
const pendingToolCalls = new Map<string, PiTrackedToolCall>();
|
||||
let userIndex = 0;
|
||||
let assistantIndex = 0;
|
||||
|
||||
for (const message of messages) {
|
||||
if (message.role === "user") {
|
||||
@@ -65,12 +66,14 @@ export async function* streamPiHistory(
|
||||
}
|
||||
|
||||
if (message.role === "assistant") {
|
||||
assistantIndex += 1;
|
||||
const messageId = message.responseId || `${provider}-history-assistant-${assistantIndex}`;
|
||||
for (const content of message.content) {
|
||||
if (content.type === "text" && content.text) {
|
||||
yield {
|
||||
type: "timeline",
|
||||
provider,
|
||||
item: { type: "assistant_message", text: content.text },
|
||||
item: { type: "assistant_message", text: content.text, messageId },
|
||||
};
|
||||
continue;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user