From 714cb42e29fbd6132ef82806498956d8ade1412f Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Wed, 12 Aug 2026 23:52:47 +0000 Subject: [PATCH 1/6] feat(dashboard): Cancel queued mailbox messages Add a participant-only pending-messages cancel API and a Cancel queue control in the dashboard mailbox stack so deferred follow-ups can be dropped before the worker runs them. Co-Authored-By: PDPM --- .../client/conversations/ConversationPage.tsx | 8 + .../conversations/PendingMailboxStack.tsx | 28 ++- .../src/client/conversations/queries.ts | 59 ++++- packages/junior-dashboard/src/client/http.ts | 17 ++ .../conversations/cancel-pending-messages.ts | 50 ++++ .../junior/src/api/conversations/routes.ts | 31 +++ packages/junior/src/api/schema.ts | 4 + .../junior/src/api/schema/conversation.ts | 22 ++ .../junior/src/chat/task-execution/state.ts | 62 +++++ .../junior/src/chat/task-execution/store.ts | 15 ++ .../cancel-pending-messages.test.ts | 216 ++++++++++++++++++ 11 files changed, 508 insertions(+), 4 deletions(-) create mode 100644 packages/junior/src/api/conversations/cancel-pending-messages.ts create mode 100644 packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts diff --git a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx index 5a763ab5f0..e1befdf750 100644 --- a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx +++ b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx @@ -7,6 +7,7 @@ import type { import { useAppendConversationMessage, useArchiveConversation, + useCancelConversationPendingMessages, useConversationData, type PendingArchiveConversationUpdate, } from "./queries"; @@ -53,6 +54,8 @@ export function ConversationPage(props: { const detail = useConversationData(conversationId); const archive = useArchiveConversation(conversationId); const appendMessage = useAppendConversationMessage(conversationId); + const cancelPendingMessages = + useCancelConversationPendingMessages(conversationId); const feedConversation = conversations.find( (item) => item.id === conversationId, ); @@ -174,8 +177,13 @@ export function ConversationPage(props: { ) : null} {detail.data ? ( { + cancelPendingMessages.mutate({}); + }} /> ) : null} void; }): ReactNode { const rows = unresolvedPendingTranscriptMessages( conversationTranscriptMessages(props.conversation), @@ -156,16 +160,34 @@ export function PendingMailboxStack(props: { ? rows.slice(0, COLLAPSED_PENDING_ROW_COUNT) : rows; const collapsedCount = rows.length - visibleRows.length; + const showCancel = Boolean(props.onCancelQueue); return (
-
- {countLabel} +
+
+ {countLabel} +
+ {showCancel ? ( + + ) : null}
-
+ {props.cancelError ? ( +
+ Could not cancel queued messages. Try again. +
+ ) : null} +
{visibleRows.map((message, index) => ( + del( + cancelConversationPendingMessagesResponseSchema, + `/api/conversations/${encodeURIComponent(conversationId)}/pending-messages`, + args ?? {}, + ), + onMutate: async () => { + await queryClient.cancelQueries({ + exact: true, + queryKey: conversationPendingMessagesQueryKey(conversationId), + }); + const previousPending = + queryClient.getQueryData( + conversationPendingMessagesQueryKey(conversationId), + ); + if (previousPending) { + queryClient.setQueryData( + conversationPendingMessagesQueryKey(conversationId), + { + ...previousPending, + messages: [], + }, + ); + } + return { previousPending }; + }, + onError: (_error, _args, context) => { + if (context?.previousPending) { + queryClient.setQueryData( + conversationPendingMessagesQueryKey(conversationId), + context.previousPending, + ); + } + }, + onSettled: async () => { + await Promise.all([ + queryClient.invalidateQueries({ + queryKey: ["dashboard", "conversations"], + }), + queryClient.invalidateQueries({ + exact: true, + queryKey: conversationDetailQueryKey(conversationId), + }), + queryClient.invalidateQueries({ + exact: true, + queryKey: conversationPendingMessagesQueryKey(conversationId), + }), + ]); + }, + }); +} + /** Archive or restore one conversation with an immediate reversible cache update. */ export function useArchiveConversation( conversationId: string, diff --git a/packages/junior-dashboard/src/client/http.ts b/packages/junior-dashboard/src/client/http.ts index 18ec435655..9eb2a8b43e 100644 --- a/packages/junior-dashboard/src/client/http.ts +++ b/packages/junior-dashboard/src/client/http.ts @@ -76,6 +76,23 @@ export async function deleteDashboardResource(path: string): Promise { if (!response.ok) throw new DashboardApiError(path, response.status); } +/** Send one authenticated DELETE request with JSON body and validate its response. */ +export async function del( + schema: ZodType, + path: string, + body: unknown = {}, +): Promise { + const response = await fetch(path, { + body: JSON.stringify(body), + credentials: "same-origin", + headers: { "content-type": "application/json" }, + method: "DELETE", + }); + if (response.status === 401) restartDashboardSignIn(); + if (!response.ok) throw new DashboardApiError(path, response.status); + return schema.parse(await response.json()); +} + /** Fetch one authenticated dashboard JSON resource and validate its response. */ export async function fetchDashboardJson( schema: ZodType, diff --git a/packages/junior/src/api/conversations/cancel-pending-messages.ts b/packages/junior/src/api/conversations/cancel-pending-messages.ts new file mode 100644 index 0000000000..dc56109acd --- /dev/null +++ b/packages/junior/src/api/conversations/cancel-pending-messages.ts @@ -0,0 +1,50 @@ +import type { User } from "@sentry/junior-plugin-api"; +import { getConversationStore, getDb } from "@/chat/db"; +import { cancelHumanFacingPendingMessages } from "@/chat/task-execution/store"; +import { throwApiError } from "../http"; +import type { + CancelConversationPendingMessagesBody, + CancelConversationPendingMessagesResponse, +} from "../schema/conversation"; +import { readConversationAccessFromSql } from "./access"; + +/** Cancel accepted human-facing mailbox rows for one conversation participant. */ +export async function cancelConversationPendingMessagesForViewer( + viewer: User, + conversationId: string, + body: CancelConversationPendingMessagesBody = {}, +): Promise { + const conversation = await getConversationStore().get({ conversationId }); + if (!conversation) { + throwApiError(404, "Conversation not found."); + } + + const access = await readConversationAccessFromSql( + getDb(), + [conversationId], + viewer, + ); + if (!access.get(conversationId)?.isParticipant) { + throwApiError( + 403, + "Only conversation participants can cancel queued messages.", + ); + } + + try { + const result = await cancelHumanFacingPendingMessages({ + conversationId, + ...(body.inboundMessageIds + ? { inboundMessageIds: body.inboundMessageIds } + : {}), + conversationStore: getConversationStore(), + }); + return { + cancelledCount: result.cancelledInboundMessageIds.length, + cancelledInboundMessageIds: result.cancelledInboundMessageIds, + conversationId, + }; + } catch (error) { + throwApiError(500, "Unable to cancel queued messages.", error); + } +} diff --git a/packages/junior/src/api/conversations/routes.ts b/packages/junior/src/api/conversations/routes.ts index 5947c1d713..d5bfbd2fd3 100644 --- a/packages/junior/src/api/conversations/routes.ts +++ b/packages/junior/src/api/conversations/routes.ts @@ -5,6 +5,8 @@ import { acceptedConversationMessageSchema, archiveConversationBodySchema, archiveConversationResponseSchema, + cancelConversationPendingMessagesBodySchema, + cancelConversationPendingMessagesResponseSchema, conversationDetailQuerySchema, conversationDetailReportSchema, conversationEventPageSchema, @@ -27,6 +29,7 @@ import { import { readConversationDetail } from "./detail"; import { readConversationEvents } from "./event-list"; import { readConversationFeed } from "./list"; +import { cancelConversationPendingMessagesForViewer } from "./cancel-pending-messages"; import { requireConversationPendingMessages } from "./pending-messages"; import { readConversationStats } from "./stats"; @@ -166,6 +169,34 @@ export function createConversationRoutes(): Hono { }, ); + app.delete( + "/:conversationId/pending-messages", + requireViewer, + validateRequest( + "param", + conversationParamsSchema, + "Invalid route parameters.", + ), + validateRequest( + "json", + cancelConversationPendingMessagesBodySchema, + "Invalid request body.", + ), + async (context) => { + const viewer = context.get("viewer"); + const { conversationId } = context.req.valid("param"); + const body = context.req.valid("json"); + return jsonResponse( + cancelConversationPendingMessagesResponseSchema, + await cancelConversationPendingMessagesForViewer( + viewer, + conversationId, + body, + ), + ); + }, + ); + app.get( "/:conversationId", validateRequest( diff --git a/packages/junior/src/api/schema.ts b/packages/junior/src/api/schema.ts index be24fce028..02326405ca 100644 --- a/packages/junior/src/api/schema.ts +++ b/packages/junior/src/api/schema.ts @@ -4,6 +4,8 @@ export { acceptedConversationMessageSchema, archiveConversationBodySchema, archiveConversationResponseSchema, + cancelConversationPendingMessagesBodySchema, + cancelConversationPendingMessagesResponseSchema, conversationAuxiliaryCostsSchema, conversationDetailQuerySchema, conversationDetailReportSchema, @@ -28,6 +30,8 @@ export type { AcceptedConversationMessage, ArchiveConversationBody, ArchiveConversationResponse, + CancelConversationPendingMessagesBody, + CancelConversationPendingMessagesResponse, ActorIdentity, ConversationAuxiliaryCosts, ConversationCost, diff --git a/packages/junior/src/api/schema/conversation.ts b/packages/junior/src/api/schema/conversation.ts index f545262f45..f4c4862b2b 100644 --- a/packages/junior/src/api/schema/conversation.ts +++ b/packages/junior/src/api/schema/conversation.ts @@ -146,6 +146,22 @@ export const conversationPendingMessagesReportSchema = z }) .strict(); +/** Optional filters for cancelling accepted mailbox rows. */ +export const cancelConversationPendingMessagesBodySchema = z + .object({ + inboundMessageIds: z.array(z.string().min(1)).min(1).optional(), + }) + .strict(); + +/** Result of cancelling accepted human-facing mailbox rows. */ +export const cancelConversationPendingMessagesResponseSchema = z + .object({ + cancelledCount: z.number().int().nonnegative(), + cancelledInboundMessageIds: z.array(z.string().min(1)), + conversationId: z.string().min(1), + }) + .strict(); + export const conversationAuxiliaryCostsSchema = z .object({ costUsd: z.number().finite().nonnegative(), @@ -792,3 +808,9 @@ export type ConversationPendingMessage = z.infer< export type ConversationPendingMessagesReport = z.infer< typeof conversationPendingMessagesReportSchema >; +export type CancelConversationPendingMessagesBody = z.infer< + typeof cancelConversationPendingMessagesBodySchema +>; +export type CancelConversationPendingMessagesResponse = z.infer< + typeof cancelConversationPendingMessagesResponseSchema +>; diff --git a/packages/junior/src/chat/task-execution/state.ts b/packages/junior/src/chat/task-execution/state.ts index 84f996e49d..d514960119 100644 --- a/packages/junior/src/chat/task-execution/state.ts +++ b/packages/junior/src/chat/task-execution/state.ts @@ -1612,6 +1612,68 @@ export async function ackMessages(args: { }); } +/** Cancel human-facing pending mailbox rows without requiring a worker lease. */ +export async function cancelHumanFacingPendingMessages(args: { + conversationId: string; + inboundMessageIds?: readonly string[]; + nowMs?: number; + state?: StateAdapter; +}): Promise<{ cancelledInboundMessageIds: string[] }> { + const nowMs = args.nowMs ?? now(); + const requestedIds = + args.inboundMessageIds === undefined + ? undefined + : new Set(args.inboundMessageIds); + return await withConversationMutation(args, async (state, lock) => { + const current = await readConversation(state, args.conversationId); + if (!current) { + return { cancelledInboundMessageIds: [] }; + } + + const cancelledInboundMessageIds: string[] = []; + const pendingMessages: InboundMessage[] = []; + for (const message of current.execution.pendingMessages) { + const isHumanFacing = message.source === "web" || message.source === "slack"; + const isRequested = + requestedIds === undefined || requestedIds.has(message.inboundMessageId); + if (isHumanFacing && isRequested) { + cancelledInboundMessageIds.push(message.inboundMessageId); + continue; + } + pendingMessages.push(message); + } + + if (cancelledInboundMessageIds.length === 0) { + return { cancelledInboundMessageIds }; + } + + const nextStatus = + current.execution.status === "pending" && pendingMessages.length === 0 + ? "idle" + : current.execution.status; + + await writeConversation( + state, + lock, + withExecutionUpdate( + current, + { + ...current.execution, + lastEnqueuedAtMs: + pendingMessages.length === 0 + ? undefined + : current.execution.lastEnqueuedAtMs, + pendingMessages, + status: nextStatus, + }, + nowMs, + ), + ); + + return { cancelledInboundMessageIds }; + }); +} + /** Mark the leased conversation as needing another queue-delivered slice. */ export async function requestAnotherSlice(args: { conversationId: string; diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index 17c481f557..9fa176e721 100644 --- a/packages/junior/src/chat/task-execution/store.ts +++ b/packages/junior/src/chat/task-execution/store.ts @@ -413,6 +413,21 @@ export async function ackMessages(args: { return result; } +/** Cancel human-facing pending mailbox rows without requiring a worker lease. */ +export async function cancelHumanFacingPendingMessages(args: { + conversationId: string; + inboundMessageIds?: readonly string[]; + conversationStore?: ConversationStore; + nowMs?: number; + state?: StateAdapter; +}) { + const result = await workState.cancelHumanFacingPendingMessages(args); + if (result.cancelledInboundMessageIds.length > 0) { + await recordExecutionMetadata(args); + } + return result; +} + /** Mark the leased conversation as needing another queue-delivered slice. */ export async function requestAnotherSlice(args: { conversationId: string; diff --git a/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts b/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts new file mode 100644 index 0000000000..ebdc13ccdf --- /dev/null +++ b/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts @@ -0,0 +1,216 @@ +import { afterEach, describe, expect, it } from "vitest"; +import { Hono } from "hono"; +import { createJuniorApi, type JuniorApiVariables } from "@/api"; +import { + cancelConversationPendingMessagesResponseSchema, + conversationPendingMessagesReportSchema, +} from "@/api/schema"; +import { + appendAndEnqueueApiConversationMessage, + createAndEnqueueApiConversation, +} from "@/chat/api-turns/work"; +import { closeDb } from "@/chat/db"; +import { + appendInboundMessage, + getConversation, +} from "@/chat/task-execution/store"; +import { + closeApiTurnWorkFixture, + createApiTurnWorkFixture, +} from "../../../fixtures/api-turn"; +import { testViewer } from "../../../fixtures/user"; + +describe("conversation cancel pending messages API", () => { + afterEach(async () => { + await closeApiTurnWorkFixture(); + await closeDb(); + }); + + it("cancels accepted web mailbox rows for participants", async () => { + const { actor, conversationStore, queue, state } = + await createApiTurnWorkFixture(); + const created = await createAndEnqueueApiConversation( + { + actor, + idempotencyKey: "cancel-root", + message: "first", + }, + { conversationStore, queue, state }, + ); + const continued = await appendAndEnqueueApiConversationMessage( + { + actor, + conversationId: created.conversationId, + idempotencyKey: "cancel-second", + message: "second", + }, + { conversationStore, queue, state }, + ); + + const app = new Hono<{ Variables: JuniorApiVariables }>(); + app.use("*", async (context, next) => { + context.set("viewer", testViewer(actor.email)); + await next(); + }); + app.route("/", createJuniorApi()); + + const response = await app.request( + `http://localhost/api/conversations/${encodeURIComponent(created.conversationId)}/pending-messages`, + { + body: JSON.stringify({}), + headers: { "content-type": "application/json" }, + method: "DELETE", + }, + ); + expect(response.status).toBe(200); + const cancelled = cancelConversationPendingMessagesResponseSchema.parse( + await response.json(), + ); + expect(cancelled.conversationId).toBe(created.conversationId); + expect(cancelled.cancelledCount).toBe(2); + expect(cancelled.cancelledInboundMessageIds.sort()).toEqual( + [created.messageId, continued.messageId].sort(), + ); + + const pending = await app.request( + `http://localhost/api/conversations/${encodeURIComponent(created.conversationId)}/pending-messages`, + ); + expect(pending.status).toBe(200); + const report = conversationPendingMessagesReportSchema.parse( + await pending.json(), + ); + expect(report.messages).toEqual([]); + + const work = await getConversation({ + conversationId: created.conversationId, + state, + }); + expect(work?.execution.pendingMessages).toEqual([]); + expect(work?.execution.status).toBe("idle"); + expect(work?.execution.inboundMessageIds).toEqual( + expect.arrayContaining([created.messageId, continued.messageId]), + ); + }); + + it("rejects cancel from non-participants", async () => { + const { actor, conversationStore, queue, state } = + await createApiTurnWorkFixture(); + const created = await createAndEnqueueApiConversation( + { + actor, + idempotencyKey: "cancel-forbidden", + message: "secret", + visibility: "private", + }, + { conversationStore, queue, state }, + ); + + const app = new Hono<{ Variables: JuniorApiVariables }>(); + app.use("*", async (context, next) => { + context.set("viewer", testViewer("stranger@example.com")); + await next(); + }); + app.route("/", createJuniorApi()); + + const response = await app.request( + `http://localhost/api/conversations/${encodeURIComponent(created.conversationId)}/pending-messages`, + { + body: JSON.stringify({}), + headers: { "content-type": "application/json" }, + method: "DELETE", + }, + ); + expect(response.status).toBe(403); + + const work = await getConversation({ + conversationId: created.conversationId, + state, + }); + expect(work?.execution.pendingMessages).toHaveLength(1); + }); + + it("keeps internal mailbox work when cancelling human-facing rows", async () => { + const { actor, conversationStore, queue, state } = + await createApiTurnWorkFixture(); + const created = await createAndEnqueueApiConversation( + { + actor, + idempotencyKey: "cancel-keep-internal-root", + message: "human", + }, + { conversationStore, queue, state }, + ); + + const existing = await getConversation({ + conversationId: created.conversationId, + state, + }); + await appendInboundMessage({ + message: { + conversationId: created.conversationId, + createdAtMs: Date.now(), + delivery: "defer", + ...(existing?.destination ? { destination: existing.destination } : {}), + inboundMessageId: "internal:keep-me", + input: { + authorId: "system", + text: "internal wake", + }, + receivedAtMs: Date.now(), + publishExternally: false, + source: "internal", + }, + state, + }); + + const app = new Hono<{ Variables: JuniorApiVariables }>(); + app.use("*", async (context, next) => { + context.set("viewer", testViewer(actor.email)); + await next(); + }); + app.route("/", createJuniorApi()); + + const response = await app.request( + `http://localhost/api/conversations/${encodeURIComponent(created.conversationId)}/pending-messages`, + { + body: JSON.stringify({}), + headers: { "content-type": "application/json" }, + method: "DELETE", + }, + ); + expect(response.status).toBe(200); + const cancelled = cancelConversationPendingMessagesResponseSchema.parse( + await response.json(), + ); + expect(cancelled.cancelledInboundMessageIds).toEqual([created.messageId]); + + const work = await getConversation({ + conversationId: created.conversationId, + state, + }); + expect(work?.execution.pendingMessages.map((m) => m.inboundMessageId)).toEqual([ + "internal:keep-me", + ]); + expect(work?.execution.status).toBe("pending"); + }); + + it("returns 404 for unknown conversations", async () => { + await createApiTurnWorkFixture(); + const app = new Hono<{ Variables: JuniorApiVariables }>(); + app.use("*", async (context, next) => { + context.set("viewer", testViewer("owner@example.com")); + await next(); + }); + app.route("/", createJuniorApi()); + + const response = await app.request( + "http://localhost/api/conversations/missing/pending-messages", + { + body: JSON.stringify({}), + headers: { "content-type": "application/json" }, + method: "DELETE", + }, + ); + expect(response.status).toBe(404); + }); +}); From 85d16fc2d661ca895a3659cee8f9eb52cb1e727f Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 13 Aug 2026 12:56:18 +0000 Subject: [PATCH 2/6] fix(chat): Clear cancelled run metadata Co-Authored-By: PDPM --- .../junior/src/chat/task-execution/state.ts | 16 +++++----- .../cancel-pending-messages.test.ts | 29 +++++++++++++++++-- 2 files changed, 35 insertions(+), 10 deletions(-) diff --git a/packages/junior/src/chat/task-execution/state.ts b/packages/junior/src/chat/task-execution/state.ts index d514960119..443761e95d 100644 --- a/packages/junior/src/chat/task-execution/state.ts +++ b/packages/junior/src/chat/task-execution/state.ts @@ -1633,9 +1633,11 @@ export async function cancelHumanFacingPendingMessages(args: { const cancelledInboundMessageIds: string[] = []; const pendingMessages: InboundMessage[] = []; for (const message of current.execution.pendingMessages) { - const isHumanFacing = message.source === "web" || message.source === "slack"; + const isHumanFacing = + message.source === "web" || message.source === "slack"; const isRequested = - requestedIds === undefined || requestedIds.has(message.inboundMessageId); + requestedIds === undefined || + requestedIds.has(message.inboundMessageId); if (isHumanFacing && isRequested) { cancelledInboundMessageIds.push(message.inboundMessageId); continue; @@ -1647,10 +1649,8 @@ export async function cancelHumanFacingPendingMessages(args: { return { cancelledInboundMessageIds }; } - const nextStatus = - current.execution.status === "pending" && pendingMessages.length === 0 - ? "idle" - : current.execution.status; + const becomesIdle = + current.execution.status === "pending" && pendingMessages.length === 0; await writeConversation( state, @@ -1664,7 +1664,9 @@ export async function cancelHumanFacingPendingMessages(args: { ? undefined : current.execution.lastEnqueuedAtMs, pendingMessages, - status: nextStatus, + retryCount: becomesIdle ? 0 : current.execution.retryCount, + runId: becomesIdle ? undefined : current.execution.runId, + status: becomesIdle ? "idle" : current.execution.status, }, nowMs, ), diff --git a/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts b/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts index ebdc13ccdf..13768a6e05 100644 --- a/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts +++ b/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts @@ -13,6 +13,8 @@ import { closeDb } from "@/chat/db"; import { appendInboundMessage, getConversation, + releaseConversationWork, + startConversationWork, } from "@/chat/task-execution/store"; import { closeApiTurnWorkFixture, @@ -47,6 +49,25 @@ describe("conversation cancel pending messages API", () => { { conversationStore, queue, state }, ); + const lease = await startConversationWork({ + conversationId: created.conversationId, + nowMs: 2_000, + state, + }); + expect(lease.status).toBe("acquired"); + if (lease.status !== "acquired") throw new Error("Expected work lease"); + await releaseConversationWork({ + conversationId: created.conversationId, + leaseToken: lease.leaseToken, + nowMs: 3_000, + state, + }); + await expect( + getConversation({ conversationId: created.conversationId, state }), + ).resolves.toMatchObject({ + execution: { runId: expect.any(String), status: "pending" }, + }); + const app = new Hono<{ Variables: JuniorApiVariables }>(); app.use("*", async (context, next) => { context.set("viewer", testViewer(actor.email)); @@ -87,6 +108,8 @@ describe("conversation cancel pending messages API", () => { }); expect(work?.execution.pendingMessages).toEqual([]); expect(work?.execution.status).toBe("idle"); + expect(work?.execution.retryCount).toBe(0); + expect(work?.execution.runId).toBeUndefined(); expect(work?.execution.inboundMessageIds).toEqual( expect.arrayContaining([created.messageId, continued.messageId]), ); @@ -188,9 +211,9 @@ describe("conversation cancel pending messages API", () => { conversationId: created.conversationId, state, }); - expect(work?.execution.pendingMessages.map((m) => m.inboundMessageId)).toEqual([ - "internal:keep-me", - ]); + expect( + work?.execution.pendingMessages.map((m) => m.inboundMessageId), + ).toEqual(["internal:keep-me"]); expect(work?.execution.status).toBe("pending"); }); From 1e334539c511cdaf475a32469d335da1fa3e228d Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 13 Aug 2026 13:30:24 +0000 Subject: [PATCH 3/6] fix(dashboard): Scope queue cancellation to accepted rows Co-Authored-By: PDPM --- .../client/conversations/ConversationPage.tsx | 23 ++++- .../conversations/PendingMailboxStack.tsx | 7 +- .../src/client/conversations/queries.ts | 17 +++- .../tests/pending-mailbox-stack.test.tsx | 89 +++++++++++++++++++ 4 files changed, 129 insertions(+), 7 deletions(-) create mode 100644 packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx diff --git a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx index 287161804b..517a6fbd32 100644 --- a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx +++ b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx @@ -262,6 +262,8 @@ function ConversationReplyFooter(props: { // composer does not re-render while the reader is still typing. const appendMessageRef = useRef(appendMessage); appendMessageRef.current = appendMessage; + const cancelPendingMessagesRef = useRef(cancelPendingMessages); + cancelPendingMessagesRef.current = cancelPendingMessages; const onPinRequestRef = useRef(props.onPinRequest); onPinRequestRef.current = props.onPinRequest; const onSubmit = useCallback( @@ -284,8 +286,25 @@ function ConversationReplyFooter(props: { onPinRequestRef.current(); }, []); const onSubmitStart = useCallback(() => { + cancelPendingMessagesRef.current.reset(); onPinRequestRef.current(); }, []); + const cancellableMessageIds = props.pendingMessages + .filter((message) => message.clientStatus === undefined) + .map((message) => message.inboundMessageId); + const cancellableMessageIdsRef = useRef(cancellableMessageIds); + cancellableMessageIdsRef.current = cancellableMessageIds; + const onCancelQueue = useCallback(() => { + const inboundMessageIds = cancellableMessageIdsRef.current; + if (inboundMessageIds.length === 0) return; + cancelPendingMessagesRef.current.mutate({ inboundMessageIds }); + }, []); + const cancelError = Boolean( + cancelPendingMessages.error && + cancelPendingMessages.variables?.inboundMessageIds?.some((id) => + cancellableMessageIds.includes(id), + ), + ); return (
@@ -303,11 +322,11 @@ function ConversationReplyFooter(props: { ) : null} cancelPendingMessages.mutate({})} + onCancelQueue={onCancelQueue} onRetry={onRetry} /> setExpanded((value) => !value); - const showCancel = Boolean(props.onCancelQueue); + const cancellableCount = rows.filter( + (message) => message.clientStatus === undefined, + ).length; + const showCancel = cancellableCount > 0 && Boolean(props.onCancelQueue); const countLabel = rows.length === 1 ? "1 queued message" : `${rows.length} queued messages`; @@ -236,7 +239,7 @@ export function PendingMailboxStack(props: {
) : null} - {props.cancelError ? ( + {showCancel && props.cancelError ? (
Could not cancel queued messages. Try again.
diff --git a/packages/junior-dashboard/src/client/conversations/queries.ts b/packages/junior-dashboard/src/client/conversations/queries.ts index e1a0bb6a55..8997c8c8dc 100644 --- a/packages/junior-dashboard/src/client/conversations/queries.ts +++ b/packages/junior-dashboard/src/client/conversations/queries.ts @@ -23,7 +23,13 @@ import { conversationPendingMessagesReportSchema, } from "@sentry/junior/api/schema"; -import { DashboardApiError, del, fetchDashboardJson, patch, post } from "../http"; +import { + DashboardApiError, + del, + fetchDashboardJson, + patch, + post, +} from "../http"; import { conversationOutboxMessageForSubmit, conversationOutboxQueryKey, @@ -224,7 +230,7 @@ export function useCancelConversationPendingMessages(conversationId: string) { `/api/conversations/${encodeURIComponent(conversationId)}/pending-messages`, args ?? {}, ), - onMutate: async () => { + onMutate: async (args) => { await queryClient.cancelQueries({ exact: true, queryKey: conversationPendingMessagesQueryKey(conversationId), @@ -238,7 +244,12 @@ export function useCancelConversationPendingMessages(conversationId: string) { conversationPendingMessagesQueryKey(conversationId), { ...previousPending, - messages: [], + messages: args?.inboundMessageIds + ? previousPending.messages.filter( + (message) => + !args.inboundMessageIds!.includes(message.inboundMessageId), + ) + : [], }, ); } diff --git a/packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx b/packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx new file mode 100644 index 0000000000..9d775abce5 --- /dev/null +++ b/packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx @@ -0,0 +1,89 @@ +import { renderToStaticMarkup } from "react-dom/server"; +import { describe, expect, it } from "vitest"; +import type { ConversationDetailReport } from "@sentry/junior/api/schema"; + +import { PendingMailboxStack } from "../src/client/conversations/PendingMailboxStack"; +import type { ConversationMailboxMessage } from "../src/client/conversations/conversationOutbox"; + +function conversation(): ConversationDetailReport { + return { + annotations: [], + conversationId: "local:web:pending-stack", + cumulativeDurationMs: 0, + displayTitle: "Pending stack", + eventHistory: { status: "available" }, + events: [], + generatedAt: new Date(0).toISOString(), + isParticipant: true, + lastProgressAt: new Date(0).toISOString(), + lastSeenAt: new Date(0).toISOString(), + startedAt: new Date(0).toISOString(), + status: "active", + surface: "api", + visibility: "public", + }; +} + +function message( + overrides: Partial = {}, +): ConversationMailboxMessage { + return { + createdAt: new Date(1_000).toISOString(), + delivery: "defer", + inboundMessageId: "accepted-1", + messageId: "accepted-1", + receivedAt: new Date(1_000).toISOString(), + role: "user", + source: "web", + text: "queued", + ...overrides, + }; +} + +describe("PendingMailboxStack cancel control", () => { + it("shows cancel only when an accepted mailbox row exists", () => { + const accepted = renderToStaticMarkup( + undefined} + />, + ); + const localOnly = renderToStaticMarkup( + undefined} + />, + ); + + expect(accepted).toContain("Cancel queue"); + expect(localOnly).not.toContain("Cancel queue"); + }); + + it("hides a stale cancel error when only local outbox rows remain", () => { + const html = renderToStaticMarkup( + undefined} + />, + ); + + expect(html).not.toContain("Could not cancel queued messages"); + }); +}); From aa8aa5b06abf14fb9ae7f04f78de6814d3163d52 Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 13 Aug 2026 13:43:08 +0000 Subject: [PATCH 4/6] fix(dashboard): Cancel the accepted queue atomically Co-Authored-By: PDPM --- .../client/conversations/ConversationPage.tsx | 13 ++++++++----- .../conversations/PendingMailboxStack.tsx | 6 +++++- .../src/client/conversations/queries.ts | 16 ++++++++-------- .../tests/pending-mailbox-stack.test.tsx | 19 +++++++++++++++++++ 4 files changed, 40 insertions(+), 14 deletions(-) diff --git a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx index 517a6fbd32..073b7109d7 100644 --- a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx +++ b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx @@ -294,14 +294,17 @@ function ConversationReplyFooter(props: { .map((message) => message.inboundMessageId); const cancellableMessageIdsRef = useRef(cancellableMessageIds); cancellableMessageIdsRef.current = cancellableMessageIds; + const hasSendingOutboxMessage = props.pendingMessages.some( + (message) => message.clientStatus === "sending", + ); const onCancelQueue = useCallback(() => { - const inboundMessageIds = cancellableMessageIdsRef.current; - if (inboundMessageIds.length === 0) return; - cancelPendingMessagesRef.current.mutate({ inboundMessageIds }); + const snapshotInboundMessageIds = cancellableMessageIdsRef.current; + if (snapshotInboundMessageIds.length === 0) return; + cancelPendingMessagesRef.current.mutate({ snapshotInboundMessageIds }); }, []); const cancelError = Boolean( cancelPendingMessages.error && - cancelPendingMessages.variables?.inboundMessageIds?.some((id) => + cancelPendingMessages.variables?.snapshotInboundMessageIds.some((id) => cancellableMessageIds.includes(id), ), ); @@ -326,7 +329,7 @@ function ConversationReplyFooter(props: { cancelPending={cancelPendingMessages.isPending} conversation={props.conversation} messages={props.pendingMessages} - onCancelQueue={onCancelQueue} + onCancelQueue={hasSendingOutboxMessage ? undefined : onCancelQueue} onRetry={onRetry} /> message.clientStatus === undefined, ).length; - const showCancel = cancellableCount > 0 && Boolean(props.onCancelQueue); + const hasSendingRow = rows.some( + (message) => message.clientStatus === "sending", + ); + const showCancel = + cancellableCount > 0 && !hasSendingRow && Boolean(props.onCancelQueue); const countLabel = rows.length === 1 ? "1 queued message" : `${rows.length} queued messages`; diff --git a/packages/junior-dashboard/src/client/conversations/queries.ts b/packages/junior-dashboard/src/client/conversations/queries.ts index 8997c8c8dc..ac2bd240e5 100644 --- a/packages/junior-dashboard/src/client/conversations/queries.ts +++ b/packages/junior-dashboard/src/client/conversations/queries.ts @@ -224,11 +224,11 @@ export function useAppendConversationMessage(conversationId: string) { export function useCancelConversationPendingMessages(conversationId: string) { const queryClient = useQueryClient(); return useMutation({ - mutationFn: (args?: { inboundMessageIds?: string[] }) => + mutationFn: (_args: { snapshotInboundMessageIds: string[] }) => del( cancelConversationPendingMessagesResponseSchema, `/api/conversations/${encodeURIComponent(conversationId)}/pending-messages`, - args ?? {}, + {}, ), onMutate: async (args) => { await queryClient.cancelQueries({ @@ -244,12 +244,12 @@ export function useCancelConversationPendingMessages(conversationId: string) { conversationPendingMessagesQueryKey(conversationId), { ...previousPending, - messages: args?.inboundMessageIds - ? previousPending.messages.filter( - (message) => - !args.inboundMessageIds!.includes(message.inboundMessageId), - ) - : [], + messages: previousPending.messages.filter( + (message) => + !args.snapshotInboundMessageIds.includes( + message.inboundMessageId, + ), + ), }, ); } diff --git a/packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx b/packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx index 9d775abce5..d051b679f4 100644 --- a/packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx +++ b/packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx @@ -67,6 +67,25 @@ describe("PendingMailboxStack cancel control", () => { expect(localOnly).not.toContain("Cancel queue"); }); + it("hides cancel while a local send can still become accepted", () => { + const html = renderToStaticMarkup( + undefined} + />, + ); + + expect(html).not.toContain("Cancel queue"); + }); + it("hides a stale cancel error when only local outbox rows remain", () => { const html = renderToStaticMarkup( Date: Thu, 13 Aug 2026 13:55:09 +0000 Subject: [PATCH 5/6] fix(dashboard): Cancel only the visible queue Co-Authored-By: PDPM --- .../src/client/conversations/ConversationComposer.tsx | 4 +++- .../src/client/conversations/ConversationPage.tsx | 9 +++++---- .../junior-dashboard/src/client/conversations/queries.ts | 8 +++----- 3 files changed, 11 insertions(+), 10 deletions(-) diff --git a/packages/junior-dashboard/src/client/conversations/ConversationComposer.tsx b/packages/junior-dashboard/src/client/conversations/ConversationComposer.tsx index ff30539db9..cf93084f9e 100644 --- a/packages/junior-dashboard/src/client/conversations/ConversationComposer.tsx +++ b/packages/junior-dashboard/src/client/conversations/ConversationComposer.tsx @@ -44,6 +44,7 @@ export function conversationAttemptForSubmit( } type ConversationComposerProps = { + disabled?: boolean; draftId: string; error?: string; label: string; @@ -85,7 +86,8 @@ export const ConversationComposer = memo(function ConversationComposer( const submittingRef = useRef(false); // Monotonic token so a late failed create never restores over a newer submit. const submitTokenRef = useRef(0); - const sendLocked = props.restoreDraftOnError && createPending; + const sendLocked = + Boolean(props.disabled) || (props.restoreDraftOnError && createPending); draftRef.current = draft; // Persist drafts after typing settles so storage never contends with keystrokes. diff --git a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx index 073b7109d7..b948fcec8d 100644 --- a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx +++ b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx @@ -298,13 +298,13 @@ function ConversationReplyFooter(props: { (message) => message.clientStatus === "sending", ); const onCancelQueue = useCallback(() => { - const snapshotInboundMessageIds = cancellableMessageIdsRef.current; - if (snapshotInboundMessageIds.length === 0) return; - cancelPendingMessagesRef.current.mutate({ snapshotInboundMessageIds }); + const inboundMessageIds = cancellableMessageIdsRef.current; + if (inboundMessageIds.length === 0) return; + cancelPendingMessagesRef.current.mutate({ inboundMessageIds }); }, []); const cancelError = Boolean( cancelPendingMessages.error && - cancelPendingMessages.variables?.snapshotInboundMessageIds.some((id) => + cancelPendingMessages.variables?.inboundMessageIds.some((id) => cancellableMessageIds.includes(id), ), ); @@ -333,6 +333,7 @@ function ConversationReplyFooter(props: { onRetry={onRetry} /> + mutationFn: (args: { inboundMessageIds: string[] }) => del( cancelConversationPendingMessagesResponseSchema, `/api/conversations/${encodeURIComponent(conversationId)}/pending-messages`, - {}, + args, ), onMutate: async (args) => { await queryClient.cancelQueries({ @@ -246,9 +246,7 @@ export function useCancelConversationPendingMessages(conversationId: string) { ...previousPending, messages: previousPending.messages.filter( (message) => - !args.snapshotInboundMessageIds.includes( - message.inboundMessageId, - ), + !args.inboundMessageIds.includes(message.inboundMessageId), ), }, ); From 8d47e5d22b669d4d48e408c88466d9bc840edaea Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 13 Aug 2026 14:09:30 +0000 Subject: [PATCH 6/6] fix(dashboard): Preserve later queued messages Co-Authored-By: PDPM --- .../client/conversations/ConversationPage.tsx | 12 +++- .../src/client/conversations/queries.ts | 6 +- .../conversations/cancel-pending-messages.ts | 3 + .../junior/src/api/schema/conversation.ts | 1 + .../junior/src/chat/task-execution/state.ts | 6 +- .../junior/src/chat/task-execution/store.ts | 1 + .../cancel-pending-messages.test.ts | 55 +++++++++++++++++++ 7 files changed, 80 insertions(+), 4 deletions(-) diff --git a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx index b948fcec8d..266a8b7612 100644 --- a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx +++ b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx @@ -231,6 +231,7 @@ export function ConversationPage(props: { live={live} onPinRequest={requestPin} pendingAuthorization={detail.pendingAuthorization} + pendingGeneratedAt={detail.pendingGeneratedAt} pendingMessages={detail.pendingMessages} /> ) : null} @@ -252,6 +253,7 @@ function ConversationReplyFooter(props: { live: boolean; onPinRequest: () => void; pendingAuthorization?: ConversationPendingMessagesReport["authorization"]; + pendingGeneratedAt?: string; pendingMessages: ConversationMailboxMessage[]; }) { const appendMessage = useAppendConversationMessage(props.conversationId); @@ -297,10 +299,16 @@ function ConversationReplyFooter(props: { const hasSendingOutboxMessage = props.pendingMessages.some( (message) => message.clientStatus === "sending", ); + const pendingGeneratedAtRef = useRef(props.pendingGeneratedAt); + pendingGeneratedAtRef.current = props.pendingGeneratedAt; const onCancelQueue = useCallback(() => { const inboundMessageIds = cancellableMessageIdsRef.current; - if (inboundMessageIds.length === 0) return; - cancelPendingMessagesRef.current.mutate({ inboundMessageIds }); + const receivedBefore = pendingGeneratedAtRef.current; + if (!receivedBefore || inboundMessageIds.length === 0) return; + cancelPendingMessagesRef.current.mutate({ + inboundMessageIds, + receivedBefore, + }); }, []); const cancelError = Boolean( cancelPendingMessages.error && diff --git a/packages/junior-dashboard/src/client/conversations/queries.ts b/packages/junior-dashboard/src/client/conversations/queries.ts index 0bc102d5b1..ec87badc21 100644 --- a/packages/junior-dashboard/src/client/conversations/queries.ts +++ b/packages/junior-dashboard/src/client/conversations/queries.ts @@ -224,7 +224,10 @@ export function useAppendConversationMessage(conversationId: string) { export function useCancelConversationPendingMessages(conversationId: string) { const queryClient = useQueryClient(); return useMutation({ - mutationFn: (args: { inboundMessageIds: string[] }) => + mutationFn: (args: { + inboundMessageIds: string[]; + receivedBefore: string; + }) => del( cancelConversationPendingMessagesResponseSchema, `/api/conversations/${encodeURIComponent(conversationId)}/pending-messages`, @@ -534,6 +537,7 @@ export function useConversationData(conversationId: string | undefined) { isPending: detail.isPending, isLoadingPreviousPage, pendingAuthorization: pending.data?.authorization, + pendingGeneratedAt: pending.data?.generatedAt, pendingMessages, loadCompleteTranscript: () => { if (!conversationId || !detail.data) { diff --git a/packages/junior/src/api/conversations/cancel-pending-messages.ts b/packages/junior/src/api/conversations/cancel-pending-messages.ts index dc56109acd..cf476453ec 100644 --- a/packages/junior/src/api/conversations/cancel-pending-messages.ts +++ b/packages/junior/src/api/conversations/cancel-pending-messages.ts @@ -37,6 +37,9 @@ export async function cancelConversationPendingMessagesForViewer( ...(body.inboundMessageIds ? { inboundMessageIds: body.inboundMessageIds } : {}), + ...(body.receivedBefore + ? { receivedBeforeMs: Date.parse(body.receivedBefore) } + : {}), conversationStore: getConversationStore(), }); return { diff --git a/packages/junior/src/api/schema/conversation.ts b/packages/junior/src/api/schema/conversation.ts index b437c244b0..02778252ad 100644 --- a/packages/junior/src/api/schema/conversation.ts +++ b/packages/junior/src/api/schema/conversation.ts @@ -157,6 +157,7 @@ export const conversationPendingMessagesReportSchema = z export const cancelConversationPendingMessagesBodySchema = z .object({ inboundMessageIds: z.array(z.string().min(1)).min(1).optional(), + receivedBefore: z.string().datetime().optional(), }) .strict(); diff --git a/packages/junior/src/chat/task-execution/state.ts b/packages/junior/src/chat/task-execution/state.ts index 443761e95d..f4264d830f 100644 --- a/packages/junior/src/chat/task-execution/state.ts +++ b/packages/junior/src/chat/task-execution/state.ts @@ -1616,6 +1616,7 @@ export async function ackMessages(args: { export async function cancelHumanFacingPendingMessages(args: { conversationId: string; inboundMessageIds?: readonly string[]; + receivedBeforeMs?: number; nowMs?: number; state?: StateAdapter; }): Promise<{ cancelledInboundMessageIds: string[] }> { @@ -1638,7 +1639,10 @@ export async function cancelHumanFacingPendingMessages(args: { const isRequested = requestedIds === undefined || requestedIds.has(message.inboundMessageId); - if (isHumanFacing && isRequested) { + const isInSnapshot = + args.receivedBeforeMs === undefined || + message.receivedAtMs <= args.receivedBeforeMs; + if (isHumanFacing && isRequested && isInSnapshot) { cancelledInboundMessageIds.push(message.inboundMessageId); continue; } diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index 9fa176e721..9a36e9cd78 100644 --- a/packages/junior/src/chat/task-execution/store.ts +++ b/packages/junior/src/chat/task-execution/store.ts @@ -417,6 +417,7 @@ export async function ackMessages(args: { export async function cancelHumanFacingPendingMessages(args: { conversationId: string; inboundMessageIds?: readonly string[]; + receivedBeforeMs?: number; conversationStore?: ConversationStore; nowMs?: number; state?: StateAdapter; diff --git a/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts b/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts index 13768a6e05..a3ab24b26d 100644 --- a/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts +++ b/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts @@ -115,6 +115,61 @@ describe("conversation cancel pending messages API", () => { ); }); + it("keeps messages received after the requested snapshot", async () => { + const { actor, conversationStore, queue, state } = + await createApiTurnWorkFixture(); + const created = await createAndEnqueueApiConversation( + { + actor, + idempotencyKey: "cancel-snapshot-root", + message: "first", + }, + { conversationStore, nowMs: 1_000, queue, state }, + ); + const later = await appendAndEnqueueApiConversationMessage( + { + actor, + conversationId: created.conversationId, + idempotencyKey: "cancel-snapshot-later", + message: "later", + }, + { conversationStore, nowMs: 3_000, queue, state }, + ); + + const app = new Hono<{ Variables: JuniorApiVariables }>(); + app.use("*", async (context, next) => { + context.set("viewer", testViewer(actor.email)); + await next(); + }); + app.route("/", createJuniorApi()); + + const response = await app.request( + `http://localhost/api/conversations/${encodeURIComponent(created.conversationId)}/pending-messages`, + { + body: JSON.stringify({ + inboundMessageIds: [created.messageId, later.messageId], + receivedBefore: new Date(2_000).toISOString(), + }), + headers: { "content-type": "application/json" }, + method: "DELETE", + }, + ); + expect(response.status).toBe(200); + const cancelled = cancelConversationPendingMessagesResponseSchema.parse( + await response.json(), + ); + expect(cancelled.cancelledInboundMessageIds).toEqual([created.messageId]); + + const work = await getConversation({ + conversationId: created.conversationId, + state, + }); + expect( + work?.execution.pendingMessages.map((message) => message.inboundMessageId), + ).toEqual([later.messageId]); + expect(work?.execution.status).toBe("pending"); + }); + it("rejects cancel from non-participants", async () => { const { actor, conversationStore, queue, state } = await createApiTurnWorkFixture();