feat: show queued dashboard follow-ups (#1631)

* feat: show queued dashboard follow-ups

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>

* fix: de-dupe queued follow-ups while streaming

---------

Co-authored-by: Ramon Nogueira <270434257+ramon-langchain@users.noreply.github.com>
Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
Co-authored-by: Johannes du Plessis <johannes@langchain.dev>
This commit is contained in:
Ramon Nogueira 2026-06-29 14:23:45 -04:00 • committed by GitHub
parent db2ae58edb
commit 8e0788dc2a
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 278 additions and 66 deletions

View file

@ -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}")

Binary file not shown.

After

Width:  |  Height:  |  Size: 44 KiB

View file

@ -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();
});
});

View file

@ -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<QueuedThreadMessage> | undefined,
messages: Array<Message>
): Array<QueuedThreadMessage> {
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) {
</span>
</Link>
)}
{hasMessages ? (
{hasConversation ? (
<div className="relative flex min-h-0 flex-1 flex-col overflow-hidden">
<Messages
messages={baseMessages}
queuedMessages={queuedMessages}
isStreaming={isStreaming}
streamIsLoading={stream.isLoading}
isThinking={isThinking}
@ -129,7 +177,7 @@ export function AgentThreadView({ thread }: AgentThreadViewProps) {
contentWidthClass="max-w-3xl"
/>
<div className="shrink-0 px-4 pb-4">
<div className="mx-auto w-full min-w-0 max-w-3xl">
<div className="mx-auto w-full max-w-3xl min-w-0">
<AgentPromptBar
placeholder="Add a follow up"
compact

View file

@ -9,8 +9,45 @@ import { useLiveMarkdownMessageId } from "@/lib/agents/provider/useLiveMarkdownM
const BOTTOM_LOCK_THRESHOLD_PX = 24;
export const Messages = memo(function Messages({
function QueuedMessages({
queuedMessages,
}: {
queuedMessages: NonNullable<MessagesProps["queuedMessages"]>;
}) {
if (queuedMessages.length === 0) return null;
return (
<div className="mb-3 space-y-2" data-testid="queued-messages">
{queuedMessages.map((message, index) => {
const imageCount = message.images?.length ?? 0;
return (
<div
key={message.id}
className="ml-auto max-w-[85%] rounded-2xl border border-dashed border-[var(--ui-border)] bg-[var(--ui-panel)] px-3 py-2 text-[13px] text-[color:var(--ui-text)] shadow-sm"
data-testid="queued-message"
>
<div className="mb-1 flex items-center gap-2 text-[11px] font-medium uppercase tracking-wide text-[color:var(--ui-text-dim)]">
<span>
{queuedMessages.length > 1 ? `Queued next #${index + 1}` : "Queued next"}
</span>
<span className="h-1.5 w-1.5 rounded-full bg-[var(--ui-accent)]" />
</div>
{message.content && <div className="whitespace-pre-wrap break-words">{message.content}</div>}
{imageCount > 0 && (
<div className="mt-1 text-xs text-[color:var(--ui-text-muted)]">
{imageCount} image{imageCount === 1 ? "" : "s"} attached
</div>
)}
</div>
);
})}
</div>
);
}
export const Messages = memo(function MessagesComponent({
messages,
queuedMessages = [],
isStreaming,
streamIsLoading,
isThinking,
@ -212,6 +249,7 @@ export const Messages = memo(function Messages({
/>
);
})}
<QueuedMessages queuedMessages={queuedMessages} />
<ThinkingSpinner
isActive={isThinking ?? streamIsLoading ?? isStreaming}
settingUpSandbox={settingUpSandbox}

View file

@ -1,4 +1,4 @@
import type { Message, Project } from "@/lib/agents/types";
import type { Message, Project, QueuedThreadMessage } from "@/lib/agents/types";
export interface ChangedFileSummaryItem {
filePath: string;
@ -21,6 +21,7 @@ export type MessagesScrollControl = {
export interface MessagesProps extends ApprovalCallbacks {
messages: Array<Message>;
queuedMessages?: Array<QueuedThreadMessage>;
isStreaming: boolean;
/** Live run signal from `useStream().isLoading` — drives Streamdown token animation. */
streamIsLoading?: boolean;

View file

@ -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<AgentThread>(
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<AgentThread>(
agentThreadKeys.detail(threadId),
(prev) => (prev ? removeQueuedMessage(prev, queuedId) : prev)
)
throw error
}
}
const configurable: Record<string, unknown> = {};
if (stream.isLoading) {
await queue()
return
}
try {
await queue()
return
} catch (error) {
if (!(error instanceof AgentsApiError) || error.status !== 409) {
throw error
}
}
const configurable: Record<string, unknown> = {}
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)
},
});
})
}

View file

@ -175,6 +175,13 @@ export interface AgentSchedule {
updatedAt?: string | null
}
export interface QueuedThreadMessage {
id: string
content: string
images?: Array<ImageChunk>
createdAt: number
}
export interface AgentThread {
id: string
title: string
@ -196,6 +203,7 @@ export interface AgentThread {
updatedAt: number
traceUrl?: string | null
messages: Array<Message>
queuedMessages?: Array<QueuedThreadMessage>
pr?: {
number: number
title: string