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); + }); +});