diff --git a/packages/junior-dashboard/src/client/conversations/ConversationComposer.tsx b/packages/junior-dashboard/src/client/conversations/ConversationComposer.tsx
index ff30539db..cf93084f9 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 a196a469b..266a8b761 100644
--- a/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx
+++ b/packages/junior-dashboard/src/client/conversations/ConversationPage.tsx
@@ -8,6 +8,7 @@ import type {
import {
useAppendConversationMessage,
useArchiveConversation,
+ useCancelConversationPendingMessages,
useConversationData,
type PendingArchiveConversationUpdate,
} from "./queries";
@@ -230,6 +231,7 @@ export function ConversationPage(props: {
live={live}
onPinRequest={requestPin}
pendingAuthorization={detail.pendingAuthorization}
+ pendingGeneratedAt={detail.pendingGeneratedAt}
pendingMessages={detail.pendingMessages}
/>
) : null}
@@ -251,13 +253,19 @@ function ConversationReplyFooter(props: {
live: boolean;
onPinRequest: () => void;
pendingAuthorization?: ConversationPendingMessagesReport["authorization"];
+ pendingGeneratedAt?: string;
pendingMessages: ConversationMailboxMessage[];
}) {
const appendMessage = useAppendConversationMessage(props.conversationId);
+ const cancelPendingMessages = useCancelConversationPendingMessages(
+ props.conversationId,
+ );
// Keep submit identity stable across mutation status flips so the memoized
// 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(
@@ -280,8 +288,34 @@ 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 hasSendingOutboxMessage = props.pendingMessages.some(
+ (message) => message.clientStatus === "sending",
+ );
+ const pendingGeneratedAtRef = useRef(props.pendingGeneratedAt);
+ pendingGeneratedAtRef.current = props.pendingGeneratedAt;
+ const onCancelQueue = useCallback(() => {
+ const inboundMessageIds = cancellableMessageIdsRef.current;
+ const receivedBefore = pendingGeneratedAtRef.current;
+ if (!receivedBefore || inboundMessageIds.length === 0) return;
+ cancelPendingMessagesRef.current.mutate({
+ inboundMessageIds,
+ receivedBefore,
+ });
+ }, []);
+ const cancelError = Boolean(
+ cancelPendingMessages.error &&
+ cancelPendingMessages.variables?.inboundMessageIds.some((id) =>
+ cancellableMessageIds.includes(id),
+ ),
+ );
return (
@@ -299,11 +333,15 @@ function ConversationReplyFooter(props: {
) : null}
void;
onRetry?(message: ConversationMailboxMessage): void;
}): ReactNode {
const [expanded, setExpanded] = useState(false);
@@ -208,12 +212,42 @@ export function PendingMailboxStack(props: {
const visibleRows = showCollapsed ? previewRows : rows;
const hiddenCount = Math.max(0, rows.length - COLLAPSED_PENDING_ROW_COUNT);
const toggleExpanded = () => setExpanded((value) => !value);
+ const cancellableCount = rows.filter(
+ (message) => message.clientStatus === undefined,
+ ).length;
+ 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`;
return (
+ {showCancel ? (
+
+
+ {countLabel}
+
+
+
+ ) : null}
+ {showCancel && props.cancelError ? (
+
+ Could not cancel queued messages. Try again.
+
+ ) : null}
{showCollapsed ? (
// Desktop keeps a two-row preview; mobile collapses to the control only.
diff --git a/packages/junior-dashboard/src/client/conversations/queries.ts b/packages/junior-dashboard/src/client/conversations/queries.ts
index 9d26e1fd9..ec87badc2 100644
--- a/packages/junior-dashboard/src/client/conversations/queries.ts
+++ b/packages/junior-dashboard/src/client/conversations/queries.ts
@@ -17,12 +17,19 @@ import type {
import {
acceptedConversationMessageSchema,
archiveConversationResponseSchema,
+ cancelConversationPendingMessagesResponseSchema,
conversationDetailReportSchema,
conversationEventPageSchema,
conversationPendingMessagesReportSchema,
} from "@sentry/junior/api/schema";
-import { DashboardApiError, fetchDashboardJson, patch, post } from "../http";
+import {
+ DashboardApiError,
+ del,
+ fetchDashboardJson,
+ patch,
+ post,
+} from "../http";
import {
conversationOutboxMessageForSubmit,
conversationOutboxQueryKey,
@@ -213,6 +220,68 @@ export function useAppendConversationMessage(conversationId: string) {
});
}
+/** Cancel accepted human-facing mailbox rows for the open conversation. */
+export function useCancelConversationPendingMessages(conversationId: string) {
+ const queryClient = useQueryClient();
+ return useMutation({
+ mutationFn: (args: {
+ inboundMessageIds: string[];
+ receivedBefore: string;
+ }) =>
+ del(
+ cancelConversationPendingMessagesResponseSchema,
+ `/api/conversations/${encodeURIComponent(conversationId)}/pending-messages`,
+ args,
+ ),
+ onMutate: async (args) => {
+ await queryClient.cancelQueries({
+ exact: true,
+ queryKey: conversationPendingMessagesQueryKey(conversationId),
+ });
+ const previousPending =
+ queryClient.getQueryData
(
+ conversationPendingMessagesQueryKey(conversationId),
+ );
+ if (previousPending) {
+ queryClient.setQueryData(
+ conversationPendingMessagesQueryKey(conversationId),
+ {
+ ...previousPending,
+ messages: previousPending.messages.filter(
+ (message) =>
+ !args.inboundMessageIds.includes(message.inboundMessageId),
+ ),
+ },
+ );
+ }
+ 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,
@@ -468,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-dashboard/src/client/http.ts b/packages/junior-dashboard/src/client/http.ts
index 18ec43565..9eb2a8b43 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-dashboard/tests/pending-mailbox-stack.test.tsx b/packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx
new file mode 100644
index 000000000..d051b679f
--- /dev/null
+++ b/packages/junior-dashboard/tests/pending-mailbox-stack.test.tsx
@@ -0,0 +1,108 @@
+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 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(
+ undefined}
+ />,
+ );
+
+ expect(html).not.toContain("Could not cancel queued messages");
+ });
+});
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 000000000..cf476453e
--- /dev/null
+++ b/packages/junior/src/api/conversations/cancel-pending-messages.ts
@@ -0,0 +1,53 @@
+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 }
+ : {}),
+ ...(body.receivedBefore
+ ? { receivedBeforeMs: Date.parse(body.receivedBefore) }
+ : {}),
+ 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 6e952c554..a5ba27c92 100644
--- a/packages/junior/src/api/conversations/routes.ts
+++ b/packages/junior/src/api/conversations/routes.ts
@@ -6,6 +6,8 @@ import {
acceptedConversationMessageSchema,
archiveConversationBodySchema,
archiveConversationResponseSchema,
+ cancelConversationPendingMessagesBodySchema,
+ cancelConversationPendingMessagesResponseSchema,
conversationAttachmentParamsSchema,
conversationDetailQuerySchema,
conversationDetailReportSchema,
@@ -33,6 +35,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";
@@ -174,6 +177,34 @@ export function createConversationRoutes(options: {
},
);
+ 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/attachments/:attachmentId",
validateRequest(
diff --git a/packages/junior/src/api/schema.ts b/packages/junior/src/api/schema.ts
index be24fce02..02326405c 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 8614e049f..02778252a 100644
--- a/packages/junior/src/api/schema/conversation.ts
+++ b/packages/junior/src/api/schema/conversation.ts
@@ -153,6 +153,23 @@ 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(),
+ receivedBefore: z.string().datetime().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(),
@@ -810,3 +827,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 84f996e49..f4264d830 100644
--- a/packages/junior/src/chat/task-execution/state.ts
+++ b/packages/junior/src/chat/task-execution/state.ts
@@ -1612,6 +1612,74 @@ 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[];
+ receivedBeforeMs?: number;
+ 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);
+ const isInSnapshot =
+ args.receivedBeforeMs === undefined ||
+ message.receivedAtMs <= args.receivedBeforeMs;
+ if (isHumanFacing && isRequested && isInSnapshot) {
+ cancelledInboundMessageIds.push(message.inboundMessageId);
+ continue;
+ }
+ pendingMessages.push(message);
+ }
+
+ if (cancelledInboundMessageIds.length === 0) {
+ return { cancelledInboundMessageIds };
+ }
+
+ const becomesIdle =
+ current.execution.status === "pending" && pendingMessages.length === 0;
+
+ await writeConversation(
+ state,
+ lock,
+ withExecutionUpdate(
+ current,
+ {
+ ...current.execution,
+ lastEnqueuedAtMs:
+ pendingMessages.length === 0
+ ? undefined
+ : current.execution.lastEnqueuedAtMs,
+ pendingMessages,
+ retryCount: becomesIdle ? 0 : current.execution.retryCount,
+ runId: becomesIdle ? undefined : current.execution.runId,
+ status: becomesIdle ? "idle" : current.execution.status,
+ },
+ 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 17c481f55..9a36e9cd7 100644
--- a/packages/junior/src/chat/task-execution/store.ts
+++ b/packages/junior/src/chat/task-execution/store.ts
@@ -413,6 +413,22 @@ 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[];
+ receivedBeforeMs?: number;
+ 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 000000000..a3ab24b26
--- /dev/null
+++ b/packages/junior/tests/integration/api/conversations/cancel-pending-messages.test.ts
@@ -0,0 +1,294 @@
+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,
+ releaseConversationWork,
+ startConversationWork,
+} 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 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));
+ 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.retryCount).toBe(0);
+ expect(work?.execution.runId).toBeUndefined();
+ expect(work?.execution.inboundMessageIds).toEqual(
+ expect.arrayContaining([created.messageId, continued.messageId]),
+ );
+ });
+
+ 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();
+ 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);
+ });
+});