Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions src/clients/http/RuntaCloudAgentsClient.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -470,4 +470,28 @@ describe("RuntaCloudAgentsClient", () => {
expect(JSON.stringify(result.messages)).not.toContain("Checking.");
expect(JSON.stringify(result.messages)).not.toContain("giant aggregate");
});

it("replays historical tool activity onto the completed agent message", async () => {
const request = vi.fn(async ({ path }: CloudRequest) => ({ status: 200, body: path.includes("/artifacts?") ? [] : [
{ id: "run-plain", agent_id: "agent-1", status: "finished", prompt: "Just reply", result: "No tools used", error: null, created_at: "2026-08-26T02:00:00Z", updated_at: "2026-08-26T02:00:01Z" },
{ id: "run-tools", agent_id: "agent-1", status: "finished", prompt: "Trace it", result: "All done", error: null, created_at: "2026-08-26T01:00:00Z", updated_at: "2026-08-26T01:00:01Z" },
] }));
const subscribe = vi.fn((path: string, listener: (event: CloudStreamEvent) => void) => {
queueMicrotask(() => {
if (path.includes("run-tools")) {
listener({ event: "pi.event", id: "1", data: { type: "tool_execution_start", toolCallId: "tool-1", toolName: "read", args: { path: "README.md" } } });
listener({ event: "pi.event", id: "2", data: { type: "tool_execution_end", toolCallId: "tool-1", toolName: "read", args: { path: "README.md" }, result: { content: [{ type: "text", text: "README contents" }] }, isError: false } });
listener({ event: "pi.event", id: "3", data: { type: "message_update", message: { id: "assistant-1" }, assistantMessageEvent: { type: "text_delta", delta: "All done" } } });
}
listener({ event: "stream.closed" });
});
return () => undefined;
});
window.runtaCrew = { cloud: { request, subscribe } } as unknown as DesktopBridge;

const { messages } = await new RuntaCloudAgentsClient().getConversation("conversation-agent-1");

expect(messages.find((message) => message.id.startsWith("run-tools:agent"))).toMatchObject({ role: "agent", parts: [{ type: "text", text: "All done" }], activities: [expect.objectContaining({ id: "tool:tool-1", kind: "file", title: "read README.md", output: "README contents", status: "completed" })] });
expect(messages.find((message) => message.id.startsWith("run-plain:agent"))?.activities).toBeUndefined();
});
});
19 changes: 12 additions & 7 deletions src/clients/http/RuntaCloudAgentsClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -246,21 +246,23 @@ export class RuntaCloudAgentsClient implements CloudAgentsClient {
let sawAssistant = false;
let settled = false;
let unsubscribe: () => void = () => undefined;
const toolActivities = new Map<string, ActivityEvent>();
const withActivities = (messages: Message[]): Message[] => toolActivities.size === 0 ? messages : messages.map((message) => message.role === "agent" ? { ...message, activities: [...toolActivities.values()] } : message);
const finish = (incomplete = false) => {
if (settled) return;
settled = true;
window.clearTimeout(timeout);
signal?.removeEventListener("abort", aborted);
unsubscribe();
if (run.status === "finished" && run.result?.trim() && (incomplete || replyError)) { resolve(fallback); return; }
if (replyError) { resolve(runMessages(run, id, replyError)); return; }
if (!sawAssistant) { resolve(fallback); return; }
if (run.status === "finished" && run.result?.trim() && (incomplete || replyError)) { resolve(withActivities(fallback)); return; }
if (replyError) { resolve(withActivities(runMessages(run, id, replyError))); return; }
if (!sawAssistant) { resolve(withActivities(fallback)); return; }
const recoveredReply = current?.parts.some((part) => part.type === "text" && part.text.trim());
resolve([
resolve(withActivities([
...fallback.filter((message) => message.role === "user"),
...(current ? [current] : []),
...fallback.filter((message) => message.role === "system" && !(recoveredReply && run.status === "finished" && !run.error?.trim())),
]);
]));
};
const aborted = () => finish(true);
const timeout = window.setTimeout(() => finish(true), 10_000);
Expand All @@ -281,8 +283,11 @@ export class RuntaCloudAgentsClient implements CloudAgentsClient {
if (failure) { replyError = failure; current = undefined; }
else if (assistantSucceeded(event)) replyError = undefined;
if (event.event === "pi.event" && event.data && typeof event.data === "object") {
const payload = event.data as { type?: string };
if (payload.type === "tool_execution_start") { current = undefined; return; }
const update = event.data as Record<string, unknown>;
const activityId = stringValue(update.toolCallId);
const activity = activityFromPiEvent(update, id, activityId ? toolActivities.get(`tool:${activityId}`) : undefined);
if (activity) toolActivities.set(activity.id, activity);
if (update.type === "tool_execution_start") { current = undefined; return; }
}
const chunk = assistantChunk(event);
if (!chunk) return;
Expand Down
2 changes: 1 addition & 1 deletion src/domain/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ export interface ActivityPart { type: "activity"; activityId: string }
export interface Attachment { id: string; name: string; size: number; mediaType: string; source: "local-selection" | "cloud"; agentId?: string }
export interface AttachmentPart { type: "attachment"; attachment: Attachment }
export type MessagePart = TextPart | ActivityPart | AttachmentPart;
export interface Message { id: string; conversationId: string; role: MessageRole; parts: MessagePart[]; createdAt: string; streaming?: boolean; interrupted?: boolean }
export interface Message { id: string; conversationId: string; role: MessageRole; parts: MessagePart[]; createdAt: string; streaming?: boolean; interrupted?: boolean; activities?: ActivityEvent[] }
export interface Conversation { id: string; agentId: string; title: string; updatedAt: string }
export type ActivityStatus = "running" | "completed" | "failed";
export interface ActivityEvent { id: string; conversationId: string; kind: "browser" | "terminal" | "file" | "handoff" | "status"; title: string; detail: string; output?: string; status: ActivityStatus; createdAt: string; updatedAt?: string }
Expand Down
42 changes: 42 additions & 0 deletions src/state/useCrewController.test.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -251,4 +251,46 @@ describe("useCrewController", () => {
await waitFor(() => expect(result.current.messages[0]?.parts[0]).toEqual({ type: "text", text: "Server reply" }));
expect(getConversation).toHaveBeenCalledTimes(2);
});

it("keeps the tool trace on the completed message and clears the live activity feed", async () => {
const listeners = new Map<string, (event: ConversationEvent) => void>();
const client: CloudAgentsClient = {
listModelProviders: async () => ({ organizationId: "org-test", providers: [] }), listAgents: async () => [agent("atlas", "Atlas")],
getAgent: async (id) => agent(id, id), createAgent: async () => agent("new", "New"), updateAgent: async (id) => agent(id, id), deleteAgent: async () => undefined, duplicateAgent: async (id) => agent(`${id}-copy`, id), setAgentUnread: async (id) => agent(id, id),
listConversations: async (id) => [{ id: `conversation-${id}`, agentId: id, title: id, updatedAt: new Date(0).toISOString() }], getConversation: async (id) => ({ conversation: { id, agentId: "atlas", title: "Atlas", updatedAt: new Date(0).toISOString() }, messages: [] }),
sendMessage: async () => { throw new Error("unused"); }, subscribeToConversationEvents: (id, listener) => { listeners.set(id, listener); return { unsubscribe: () => undefined }; }, listApprovalRequests: async () => [], respondToApproval: async () => { throw new Error("unused"); }, getComputer: async (id) => ({ id, agentId: id, runtimeName: id, status: "online", capabilities: ["open"] }), openComputer: async () => ({ mode: "remote", url: "https://example.test", protocols: ["binary", "vnc-ticket.test"] }), takeOverComputer: async () => ({ mode: "remote", url: "https://example.test", protocols: ["binary", "vnc-ticket.test"] }), reconnect: async () => undefined,
};
const { result } = renderHook(() => useCrewController(client));
await waitFor(() => expect(listeners.has("conversation-atlas")).toBe(true));
act(() => listeners.get("conversation-atlas")?.({ type: "message.created", message: { id: "run-1:agent", conversationId: "conversation-atlas", role: "agent", parts: [{ type: "text", text: "Working on it" }], createdAt: new Date(0).toISOString(), streaming: true } }));
act(() => listeners.get("conversation-atlas")?.({ type: "activity.updated", activity: { id: "tool:read", conversationId: "conversation-atlas", kind: "file", title: "read README.md", detail: "read README.md", output: "README contents", status: "completed", createdAt: new Date(0).toISOString() } }));
expect(result.current.activities).toHaveLength(1);
act(() => listeners.get("conversation-atlas")?.({ type: "message.completed", messageId: "run-1:agent", notify: true }));
expect(result.current.messages.at(-1)).toMatchObject({ id: "run-1:agent", streaming: false, activities: [expect.objectContaining({ id: "tool:read", output: "README contents" })] });
expect(result.current.activities).toEqual([]);
});

it("does not leak the live trace into another agent's conversation", async () => {
const listeners = new Map<string, (event: ConversationEvent) => void>();
const client: CloudAgentsClient = {
listModelProviders: async () => ({ organizationId: "org-test", providers: [] }), listAgents: async () => [agent("atlas", "Atlas"), agent("scout", "Scout")],
getAgent: async (id) => agent(id, id), createAgent: async () => agent("new", "New"), updateAgent: async (id) => agent(id, id), deleteAgent: async () => undefined, duplicateAgent: async (id) => agent(`${id}-copy`, id), setAgentUnread: async (id) => agent(id, id),
listConversations: async (id) => [{ id: `conversation-${id}`, agentId: id, title: id, updatedAt: new Date(0).toISOString() }], getConversation: async (id) => ({ conversation: { id, agentId: id.endsWith("atlas") ? "atlas" : "scout", title: id, updatedAt: new Date(0).toISOString() }, messages: [] }),
sendMessage: async () => { throw new Error("unused"); }, subscribeToConversationEvents: (id, listener) => { listeners.set(id, listener); return { unsubscribe: () => undefined }; }, listApprovalRequests: async () => [], respondToApproval: async () => { throw new Error("unused"); }, getComputer: async (id) => ({ id, agentId: id, runtimeName: id, status: "online", capabilities: ["open"] }), openComputer: async () => ({ mode: "remote", url: "https://example.test", protocols: ["binary", "vnc-ticket.test"] }), takeOverComputer: async () => ({ mode: "remote", url: "https://example.test", protocols: ["binary", "vnc-ticket.test"] }), reconnect: async () => undefined,
};
const { result } = renderHook(() => useCrewController(client));
await waitFor(() => expect(listeners.has("conversation-atlas")).toBe(true));
act(() => listeners.get("conversation-atlas")?.({ type: "message.created", message: { id: "run-1:agent", conversationId: "conversation-atlas", role: "agent", parts: [{ type: "text", text: "" }], createdAt: new Date(0).toISOString(), streaming: true } }));
act(() => listeners.get("conversation-atlas")?.({ type: "activity.updated", activity: { id: "tool:live", conversationId: "conversation-atlas", kind: "terminal", title: "npm test", detail: "npm test", status: "running", createdAt: new Date(0).toISOString() } }));
act(() => listeners.get("conversation-atlas")?.({ type: "message.completed", messageId: "run-1:agent", notify: true }));
act(() => result.current.setSelectedAgentId("scout"));
await waitFor(() => expect(result.current.selectedAgentId).toBe("scout"));
expect(result.current.activities).toEqual([]);
await waitFor(() => expect(result.current.conversationLoading).toBe(false));
expect(result.current.messages.every((message) => !message.activities)).toBe(true);
act(() => result.current.setSelectedAgentId("atlas"));
await waitFor(() => expect(result.current.messages.at(-1)?.id).toBe("run-1:agent"));
expect(result.current.messages.at(-1)?.activities).toEqual([expect.objectContaining({ id: "tool:live" })]);
expect(result.current.activities).toEqual([]);
});
});
Loading