diff --git a/tests/e2e/fake_llm.py b/tests/e2e/fake_llm.py index 9ea10e99..3f491bdd 100644 --- a/tests/e2e/fake_llm.py +++ b/tests/e2e/fake_llm.py @@ -10,6 +10,7 @@ the preceding tool result, exactly as a real model would. from __future__ import annotations import re +import time from typing import Any from e2e_env import ( @@ -239,6 +240,11 @@ def _latest_attribution(messages: list[BaseMessage]) -> str | None: def _step_followup(messages: list[BaseMessage]) -> AIMessage: + if any( + isinstance(msg, HumanMessage) and "Please queue this follow-up" in _text(msg.content) + for msg in messages + ): + time.sleep(2) attribution = _latest_attribution(messages) suffix = f" I saw this follow-up was from {attribution}." if attribution else "" return AIMessage(content=f"{FOLLOW_UP_REPLY}{suffix}") diff --git a/tests/e2e/screenshots/queued-messages-dashboard.png b/tests/e2e/screenshots/queued-messages-dashboard.png new file mode 100644 index 00000000..956baabe Binary files /dev/null and b/tests/e2e/screenshots/queued-messages-dashboard.png differ diff --git a/tests/e2e/tests/dashboard.spec.ts b/tests/e2e/tests/dashboard.spec.ts index 7b1ebb51..cce7a994 100644 --- a/tests/e2e/tests/dashboard.spec.ts +++ b/tests/e2e/tests/dashboard.spec.ts @@ -10,15 +10,34 @@ async function loginAs(page: Page, user: { login: string; email: string }) { expect(res.ok()).toBeTruthy(); } +async function openRunningThreadViaSlackLink(page: Page) { + await page.goto("/mock/slack"); + await page.locator("#reset").click(); + await expect(page.locator("#thread")).toContainText("No messages yet"); + await page + .locator("#text") + .fill("<@U0BOT> please add a greet() helper and open a PR"); + await page.locator("#send").click(); + + const webLink = page.locator('.msg.bot a[href*="/agents/"]').first(); + await expect(webLink).toBeVisible(); + await webLink.click(); + await expect(page).toHaveURL(/\/agents\//); +} + // Run the Slack flow so a thread + PR exist, then click the bot's real // "Open in Web" link, landing on the actual dashboard app. async function openThreadViaSlackLink(page: Page) { await page.goto("/mock/slack"); await page.locator("#reset").click(); await expect(page.locator("#thread")).toContainText("No messages yet"); - await page.locator("#text").fill("<@U0BOT> please add a greet() helper and open a PR"); + await page + .locator("#text") + .fill("<@U0BOT> please add a greet() helper and open a PR"); await page.locator("#send").click(); - await expect(page.locator(".msg.bot").filter({ hasText: "Add greet() helper" })).toBeVisible(); + await expect( + page.locator(".msg.bot").filter({ hasText: "Add greet() helper" }), + ).toBeVisible(); const webLink = page.locator('.msg.bot a[href*="/agents/"]').first(); await expect(webLink).toBeVisible(); @@ -38,25 +57,61 @@ async function expectTranscriptVisible(page: Page) { } test.describe("Slack → web handoff (real dashboard UI)", () => { - test("the SAME user continues the conversation in the web app", async ({ page }) => { + test("the SAME user continues the conversation in the web app", async ({ + page, + }) => { await loginAs(page, SAME_USER); await openThreadViaSlackLink(page); // The owner sees the composer (either the follow-up bar once the transcript // hydrates, or the empty-state bar before it — both mean they can type). - const composer = page.getByPlaceholder(/Add a follow up|Send the first message/); + const composer = page.getByPlaceholder( + /Add a follow up|Send the first message/, + ); await expect(composer).toBeVisible(); // Continue from the web — a new agent reply streams into the same thread. await composer.fill("Looks good — can you also add a docstring?"); await composer.press("Enter"); - await expect(page.getByText(/anything else you'd like changed/)).toBeVisible(); + await expect( + page.getByText(/anything else you'd like changed/), + ).toBeVisible(); // The transcript that started in Slack is here too (incl. the PR link). - await expect(page.getByRole("link", { name: "Add greet() helper" }).first()).toBeVisible(); + await expect( + page.getByRole("link", { name: "Add greet() helper" }).first(), + ).toBeVisible(); }); - test("a DIFFERENT user can post, and their message is attributed", async ({ page }) => { + test("shows follow-ups queued while the agent is still running", async ({ + page, + }, testInfo) => { + await loginAs(page, SAME_USER); + await openRunningThreadViaSlackLink(page); + + const queuedText = "Please queue this follow-up while you finish the PR."; + const busyComposer = page.getByPlaceholder( + "Send a message to queue next...", + ); + await expect(busyComposer).toBeVisible(); + await busyComposer.fill(queuedText); + await busyComposer.press("Enter"); + + const queuedMessage = page + .getByTestId("queued-message") + .filter({ hasText: queuedText }); + await expect(queuedMessage).toBeVisible(); + const screenshotPath = testInfo.outputPath("queued-messages-dashboard.png"); + await page.screenshot({ path: screenshotPath, fullPage: true }); + await testInfo.attach("queued-messages-dashboard", { + path: screenshotPath, + contentType: "image/png", + }); + }); + + test("a DIFFERENT user can post, and their message is attributed", async ({ + page, + }) => { await loginAs(page, OTHER_USER); await openThreadViaSlackLink(page); @@ -64,16 +119,22 @@ test.describe("Slack → web handoff (real dashboard UI)", () => { await expectTranscriptVisible(page); // …and a non-owner now gets a composer too (owner-only restriction removed). - const composer = page.getByPlaceholder(/Add a follow up|Send the first message/); + const composer = page.getByPlaceholder( + /Add a follow up|Send the first message/, + ); await expect(composer).toBeVisible(); // Posting starts a new run — the agent's follow-up reply streams in. await composer.fill("Can you also add a docstring?"); await composer.press("Enter"); - await expect(page.getByText(/anything else you'd like changed/)).toBeVisible(); + await expect( + page.getByText(/anything else you'd like changed/), + ).toBeVisible(); // The non-owner's message is tagged server-side with their GitHub login, so // the owner can tell who sent it. - await expect(page.getByText(new RegExp(`@${OTHER_USER.login}`)).first()).toBeVisible(); + await expect( + page.getByText(new RegExp(`@${OTHER_USER.login}`)).first(), + ).toBeVisible(); }); }); diff --git a/ui/src/components/agents/AgentThreadView.tsx b/ui/src/components/agents/AgentThreadView.tsx index 57c7f792..a6e93c5d 100644 --- a/ui/src/components/agents/AgentThreadView.tsx +++ b/ui/src/components/agents/AgentThreadView.tsx @@ -3,7 +3,11 @@ import { Link } from "@tanstack/react-router" import { useStreamContext as useAgentThreadStream } from "@langchain/react" import { Map as MapIcon } from "lucide-react" -import type { AgentThread, Message } from "@/lib/agents/types" +import type { + AgentThread, + Message, + QueuedThreadMessage, +} from "@/lib/agents/types" import type { ModelSelection } from "@/lib/agents/provider/useModelOptions" import { AgentGitPanel, @@ -24,6 +28,44 @@ interface AgentThreadViewProps { thread: AgentThread } +function messageText(message: Message): string { + return message.chunks + .map((chunk) => (chunk.kind === "text" ? chunk.text : "")) + .join("\n") + .trim() +} + +function visibleQueuedMessages( + queuedMessages: Array | undefined, + messages: Array +): Array { + const queued = queuedMessages ?? [] + if (queued.length === 0) return queued + + const userMessages = messages + .filter((message) => message.author === "user") + .map((message) => ({ + text: messageText(message), + timestamp: Date.parse(message.timestamp), + consumed: false, + })) + + return queued.filter((queuedMessage) => { + const queuedText = queuedMessage.content.trim() + if (!queuedText) return true + + const match = userMessages.find((message) => { + if (message.consumed || !message.text.includes(queuedText)) return false + if (!Number.isFinite(message.timestamp)) return true + return message.timestamp >= queuedMessage.createdAt - 1000 + }) + if (!match) return true + + match.consumed = true + return false + }) +} + // The stream lives at the `/agents` layout (one persistent provider that // survives the home → thread navigation), so this view only consumes it. export function AgentThreadView({ thread }: AgentThreadViewProps) { @@ -71,8 +113,13 @@ export function AgentThreadView({ thread }: AgentThreadViewProps) { return live }, [stream.messages, stream.toolCalls, stream.subagents, thread.messages]) - const hasMessages = baseMessages.length > 0 const isStreaming = thread.status === "running" || stream.isLoading + const queuedMessages = useMemo( + () => visibleQueuedMessages(thread.queuedMessages, baseMessages), + [baseMessages, thread.queuedMessages] + ) + const hasMessages = baseMessages.length > 0 + const hasConversation = hasMessages || queuedMessages.length > 0 const isThinking = stream.isLoading const settingUpSandbox = isThinking && baseMessages.length === 0 // The transcript hydrates from the SDK (`GET …/state` → `stream.messages`). @@ -118,10 +165,11 @@ export function AgentThreadView({ thread }: AgentThreadViewProps) { )} - {hasMessages ? ( + {hasConversation ? (
-
+
; +}) { + if (queuedMessages.length === 0) return null; + + return ( +
+ {queuedMessages.map((message, index) => { + const imageCount = message.images?.length ?? 0; + return ( +
+
+ + {queuedMessages.length > 1 ? `Queued next #${index + 1}` : "Queued next"} + + +
+ {message.content &&
{message.content}
} + {imageCount > 0 && ( +
+ {imageCount} image{imageCount === 1 ? "" : "s"} attached +
+ )} +
+ ); + })} +
+ ); +} + +export const Messages = memo(function MessagesComponent({ messages, + queuedMessages = [], isStreaming, streamIsLoading, isThinking, @@ -212,6 +249,7 @@ export const Messages = memo(function Messages({ /> ); })} + ; + queuedMessages?: Array; isStreaming: boolean; /** Live run signal from `useStream().isLoading` — drives Streamdown token animation. */ streamIsLoading?: boolean; diff --git a/ui/src/lib/agents/provider/useSubmitAgentMessage.ts b/ui/src/lib/agents/provider/useSubmitAgentMessage.ts index 5a33442a..11ff403d 100644 --- a/ui/src/lib/agents/provider/useSubmitAgentMessage.ts +++ b/ui/src/lib/agents/provider/useSubmitAgentMessage.ts @@ -1,9 +1,13 @@ -import { useMutation, useQueryClient } from "@tanstack/react-query"; -import { useStreamContext as useAgentThreadStream } from "@langchain/react"; +import { useMutation, useQueryClient } from "@tanstack/react-query" +import { useStreamContext as useAgentThreadStream } from "@langchain/react" -import type { SendAgentMessageVariables } from "@/lib/agents/queries"; -import { AgentsApiError, agentsApi } from "@/lib/agents/api"; -import { agentThreadKeys, invalidateAgentThreadLists } from "@/lib/agents/queries"; +import type { SendAgentMessageVariables } from "@/lib/agents/queries" +import type { AgentThread } from "@/lib/agents/types" +import { AgentsApiError, agentsApi } from "@/lib/agents/api" +import { + agentThreadKeys, + invalidateAgentThreadLists, +} from "@/lib/agents/queries" /** * Construct the message content for the LangGraph run. @@ -12,14 +16,44 @@ import { agentThreadKeys, invalidateAgentThreadLists } from "@/lib/agents/querie * @returns The message content. */ function messageContent(vars: SendAgentMessageVariables) { - const text = vars.content.trim(); - const imageBlocks = vars.images?.map((image) => ({ - type: "image", - base64: image.base64, - mime_type: image.mimeType, - ...(image.fileName ? { file_name: image.fileName } : {}), - })) ?? []; - return [...imageBlocks, ...(text ? [{ type: "text", text }] : [])]; + const text = vars.content.trim() + const imageBlocks = + vars.images?.map((image) => ({ + type: "image", + base64: image.base64, + mime_type: image.mimeType, + ...(image.fileName ? { file_name: image.fileName } : {}), + })) ?? [] + return [...imageBlocks, ...(text ? [{ type: "text", text }] : [])] +} + +function appendQueuedMessage( + thread: AgentThread, + vars: SendAgentMessageVariables, + id: string, + createdAt: number +): AgentThread { + return { + ...thread, + queuedMessages: [ + ...(thread.queuedMessages ?? []), + { + id, + content: vars.content.trim(), + images: vars.images, + createdAt, + }, + ], + } +} + +function removeQueuedMessage(thread: AgentThread, id: string): AgentThread { + return { + ...thread, + queuedMessages: thread.queuedMessages?.filter( + (message) => message.id !== id + ), + } } /** @@ -32,49 +66,65 @@ function messageContent(vars: SendAgentMessageVariables) { * That endpoint writes to the thread store; `check_message_queue_before_model` * injects the message into the *current* run before the next model call — the * same mid-run follow-up path used by Slack, Linear, and GitHub webhooks. - * + * * @param threadId - The ID of the thread to submit the message to. * @returns The mutation object. */ export function useSubmitAgentMessage(threadId: string) { - const queryClient = useQueryClient(); - const stream = useAgentThreadStream(); + const queryClient = useQueryClient() + const stream = useAgentThreadStream() return useMutation({ mutationFn: async (vars: SendAgentMessageVariables) => { - const queue = () => - agentsApi.queueMessage(threadId, { - content: vars.content, - images: vars.images, - model_id: vars.model_id, - effort: vars.effort, - plan_mode: vars.plan_mode, - }); - - if (stream.isLoading) { - await queue(); - return; - } - - try { - await queue(); - return; - } catch (error) { - if (!(error instanceof AgentsApiError) || error.status !== 409) { - throw error; + const queue = async () => { + const queuedAt = Date.now() + const queuedId = `queued-${queuedAt}-${Math.random().toString(36).slice(2)}` + queryClient.setQueryData( + agentThreadKeys.detail(threadId), + (prev) => + prev ? appendQueuedMessage(prev, vars, queuedId, queuedAt) : prev + ) + try { + await agentsApi.queueMessage(threadId, { + content: vars.content, + images: vars.images, + model_id: vars.model_id, + effort: vars.effort, + plan_mode: vars.plan_mode, + }) + } catch (error) { + queryClient.setQueryData( + agentThreadKeys.detail(threadId), + (prev) => (prev ? removeQueuedMessage(prev, queuedId) : prev) + ) + throw error } } - const configurable: Record = {}; + if (stream.isLoading) { + await queue() + return + } + + try { + await queue() + return + } catch (error) { + if (!(error instanceof AgentsApiError) || error.status !== 409) { + throw error + } + } + + const configurable: Record = {} if (vars.model_id && vars.effort) { - configurable.agent_model_id = vars.model_id; - configurable.agent_effort = vars.effort; + configurable.agent_model_id = vars.model_id + configurable.agent_effort = vars.effort } if (vars.plan_mode) { - configurable.plan_mode = true; + configurable.plan_mode = true } const config = - Object.keys(configurable).length > 0 ? { configurable } : undefined; + Object.keys(configurable).length > 0 ? { configurable } : undefined // Don't await: `stream.submit` resolves only when the run *finishes*, so // awaiting would keep the mutation `isPending` (and the prompt bar @@ -83,7 +133,7 @@ export function useSubmitAgentMessage(threadId: string) { void stream .submit( { messages: [{ type: "human", content: messageContent(vars) }] }, - { config }, + { config } ) .catch(() => { // The run failed to start (e.g. expired OAuth token → 401, or a @@ -91,16 +141,16 @@ export function useSubmitAgentMessage(threadId: string) { // `status: "running"`. Surface the failure and clear the busy state // instead of leaving the thread falsely running. queryClient.setQueryData(agentThreadKeys.detail(threadId), (prev) => - prev ? { ...prev, status: "error" as const } : prev, - ); - invalidateAgentThreadLists(queryClient); - }); + prev ? { ...prev, status: "error" as const } : prev + ) + invalidateAgentThreadLists(queryClient) + }) }, onSuccess: () => { queryClient.setQueryData(agentThreadKeys.detail(threadId), (prev) => - prev ? { ...prev, status: "running" as const } : prev, - ); - invalidateAgentThreadLists(queryClient); + prev ? { ...prev, status: "running" as const } : prev + ) + invalidateAgentThreadLists(queryClient) }, - }); + }) } diff --git a/ui/src/lib/agents/types.ts b/ui/src/lib/agents/types.ts index 2e88365f..8c664b47 100644 --- a/ui/src/lib/agents/types.ts +++ b/ui/src/lib/agents/types.ts @@ -175,6 +175,13 @@ export interface AgentSchedule { updatedAt?: string | null } +export interface QueuedThreadMessage { + id: string + content: string + images?: Array + createdAt: number +} + export interface AgentThread { id: string title: string @@ -196,6 +203,7 @@ export interface AgentThread { updatedAt: number traceUrl?: string | null messages: Array + queuedMessages?: Array pr?: { number: number title: string