diff --git a/docs/plugins/sdk-channel-plugins/message-adapter.md b/docs/plugins/sdk-channel-plugins/message-adapter.md index 9356e7cb3a10..0f03378f3976 100644 --- a/docs/plugins/sdk-channel-plugins/message-adapter.md +++ b/docs/plugins/sdk-channel-plugins/message-adapter.md @@ -127,6 +127,17 @@ When transport cleanup policy must change, pass a synchronous `prepareCleanup` callback to `cleanupPending`; it runs in order with clears, before deletion. Rejected deletions remain owned for a later cleanup attempt. +For synchronous turn rotation, `reset()` advances `generation`, reopens delivery, +and resets the current message and pending updates. Published messages remain the +adapter's responsibility. `reset("discard")` also retires +creates from earlier generations when they settle. Call `createMessage(send, publish)` +inside the serialized send loop; its synchronous `publish` callback installs only +current-generation receipts. Capture `generation` before edits and recheck it before +publishing their results. `retireCurrent(stopForClear)` similarly fences awaited +cleanup. Adapters that already settle their sends before rotation can use +`resetMessage()` to reset identity and pending/throttle state without reopening delivery +or advancing the generation. + ### Commentary delivery ownership Set `commentaryPayloadsEnabled: true` when the channel supports durable commentary messages. diff --git a/extensions/discord/src/draft-stream.test.ts b/extensions/discord/src/draft-stream.test.ts index 2626c9082e91..102e32e55cc0 100644 --- a/extensions/discord/src/draft-stream.test.ts +++ b/extensions/discord/src/draft-stream.test.ts @@ -343,57 +343,38 @@ describe("createDiscordDraftStream", () => { expect(stream.messageId()).toBe("1002"); }); - it("preserves an in-flight block while starting the next block", async () => { - let finishFirstCreate: ((value: { id: string }) => void) | undefined; - const firstCreate = new Promise<{ id: string }>((resolve) => { - finishFirstCreate = resolve; - }); + it.each([ + { modes: ["preserve"] as const, discarded: false }, + { modes: ["discard"] as const, discarded: true }, + { modes: ["discard", "preserve"] as const, discarded: true }, + { modes: ["preserve", "discard"] as const, discarded: true }, + ])("settles an in-flight create after $modes rotations", async ({ modes, discarded }) => { + const firstCreate = createDeferred<{ id: string }>(); const rest = { - post: vi.fn().mockReturnValueOnce(firstCreate).mockResolvedValueOnce({ id: "1002" }), + post: vi.fn().mockReturnValueOnce(firstCreate.promise).mockResolvedValueOnce({ id: "1002" }), patch: vi.fn(async () => undefined), delete: vi.fn(async () => undefined), }; const stream = createDraftStream(rest); stream.update("old turn draft"); - await vi.waitFor(() => expect(rest.post).toHaveBeenCalledTimes(1)); - stream.forceNewMessage(); + expect(rest.post).toHaveBeenCalledTimes(1); + for (const mode of modes) { + stream.forceNewMessage(mode); + } stream.update("queued turn draft"); - finishFirstCreate?.({ id: "1001" }); + firstCreate.resolve({ id: "1001" }); await stream.flush(); expect(rest.post).toHaveBeenCalledTimes(2); expect(rest.post.mock.calls[1]?.[1]).toMatchObject({ body: { content: "queued turn draft" }, }); - expect(rest.delete).not.toHaveBeenCalled(); - expect(stream.messageId()).toBe("1002"); - }); - - it("discards an in-flight progress draft while starting the queued turn", async () => { - let finishFirstCreate: ((value: { id: string }) => void) | undefined; - const firstCreate = new Promise<{ id: string }>((resolve) => { - finishFirstCreate = resolve; - }); - const rest = { - post: vi.fn().mockReturnValueOnce(firstCreate).mockResolvedValueOnce({ id: "1002" }), - patch: vi.fn(async () => undefined), - delete: vi.fn(async () => undefined), - }; - const stream = createDraftStream(rest); - - stream.update("old progress draft"); - await vi.waitFor(() => expect(rest.post).toHaveBeenCalledTimes(1)); - stream.forceNewMessage("discard"); - stream.update("queued turn draft"); - finishFirstCreate?.({ id: "1001" }); - await stream.flush(); - - expect(rest.post).toHaveBeenCalledTimes(2); - expect(rest.post.mock.calls[1]?.[1]).toMatchObject({ - body: { content: "queued turn draft" }, - }); - expect(rest.delete).toHaveBeenCalledWith(Routes.channelMessage("c1", "1001")); + if (discarded) { + expect(rest.delete).toHaveBeenCalledExactlyOnceWith(Routes.channelMessage("c1", "1001")); + } else { + expect(rest.delete).not.toHaveBeenCalled(); + } expect(stream.messageId()).toBe("1002"); }); diff --git a/extensions/discord/src/draft-stream.ts b/extensions/discord/src/draft-stream.ts index b37d13b3de99..cb39f9c23f79 100644 --- a/extensions/discord/src/draft-stream.ts +++ b/extensions/discord/src/draft-stream.ts @@ -42,15 +42,12 @@ export function createDiscordDraftStream(params: { const streamState = { stopped: false, final: false }; let streamMessage: DiscordDraftMessage | undefined; let lastSentText = ""; - let streamGeneration = 0; - let activeCreateGeneration: number | undefined; - let discardActiveCreate = false; const sendOrEditStreamMessage = async ({ text, complete, }: DiscordDraftUpdate): Promise => { - const generation = streamGeneration; + const generation = lifecycle.generation; const targetChannelId = channelId; // Allow final flush even if stopped (e.g., after clear()). if (streamState.stopped && !streamState.final) { @@ -88,7 +85,7 @@ export function createDiscordDraftStream(params: { await editChannelMessage(rest, streamMessage.channelId, streamMessage.messageId, { body, }); - if (generation === streamGeneration) { + if (generation === lifecycle.generation) { lastSentText = trimmed; } return true; @@ -97,37 +94,31 @@ export function createDiscordDraftStream(params: { const messageReference = replyToMessageId ? { message_id: replyToMessageId, fail_if_not_exists: false } : undefined; - activeCreateGeneration = generation; - const sent = (await rest.post(Routes.channelMessages(targetChannelId), { - body: { - ...body, - ...(messageReference ? { message_reference: messageReference } : {}), + return await lifecycle.createMessage( + async () => { + const sent = (await rest.post(Routes.channelMessages(targetChannelId), { + body: { + ...body, + ...(messageReference ? { message_reference: messageReference } : {}), + }, + })) as { id?: string }; // SAFETY: The create response's ID is checked before use. + return typeof sent?.id === "string" && sent.id + ? { channelId: targetChannelId, messageId: sent.id } + : undefined; }, - })) as { id?: string }; // SAFETY: The create response's ID is checked before use. - const sentMessageId = sent?.id; - const shouldDiscardStaleCreate = activeCreateGeneration === generation && discardActiveCreate; - activeCreateGeneration = undefined; - discardActiveCreate = false; - if (generation !== streamGeneration) { - if (shouldDiscardStaleCreate && typeof sentMessageId === "string" && sentMessageId) { - await lifecycle.retire({ channelId: targetChannelId, messageId: sentMessageId }); - } - return true; - } - if (typeof sentMessageId !== "string" || !sentMessageId) { - streamState.stopped = true; - params.warn?.("discord stream preview stopped (missing message id from send)"); - return false; - } - streamMessage = { channelId: targetChannelId, messageId: sentMessageId }; - lastSentText = trimmed; - return true; + (message) => { + if (!message) { + streamState.stopped = true; + params.warn?.("discord stream preview stopped (missing message id from send)"); + return false; + } + streamMessage = message; + lastSentText = trimmed; + return true; + }, + ); } catch (err) { - if (activeCreateGeneration === generation) { - activeCreateGeneration = undefined; - discardActiveCreate = false; - } - if (generation !== streamGeneration) { + if (generation !== lifecycle.generation) { return true; } streamState.stopped = true; @@ -159,22 +150,6 @@ export function createDiscordDraftStream(params: { const update = (text: string, options?: { complete?: boolean }) => updateDraft({ text, complete: options?.complete === true }); - /** Reset the draft so the next update creates a new message. */ - const forceNewMessage = (mode: "preserve" | "discard" = "preserve") => { - // In-flight REST calls may finish after a turn boundary. Advance identity - // synchronously so their result cannot overwrite the next turn's state. - // Block mode preserves the prior block; progress mode discards its draft. - if (mode === "discard" && activeCreateGeneration !== undefined) { - discardActiveCreate = true; - } - streamGeneration += 1; - streamState.stopped = false; - streamState.final = false; - streamMessage = undefined; - lastSentText = ""; - loop.resetPending(); - loop.resetThrottleWindow(); - }; /** Move the draft to another channel, preserving its current text. */ const retarget = async (nextChannelId: string) => { const normalized = nextChannelId.trim(); @@ -185,13 +160,8 @@ export function createDiscordDraftStream(params: { const pending = loop.takePending(); const previousMessage = streamMessage; const previousText = pending.text || lastSentText; - streamGeneration += 1; channelId = normalized; - streamMessage = undefined; - lastSentText = ""; - streamState.stopped = false; - streamState.final = false; - loop.resetThrottleWindow(); + lifecycle.reset(); if (previousText) { update(previousText, { complete: pending.text ? pending.complete : true }); await loop.flush(); @@ -204,18 +174,6 @@ export function createDiscordDraftStream(params: { await lifecycle.retire(previousMessage); } }; - const retireCurrentMessage = async (stopForClear: () => Promise) => { - const generation = streamGeneration; - await stopForClear(); - if (generation !== streamGeneration) { - return; - } - const message = streamMessage; - clearMessageId(); - if (message) { - await lifecycle.retire(message); - } - }; params.log?.(`discord stream preview ready (maxChars=${maxChars}, throttleMs=${throttleMs})`); @@ -224,9 +182,9 @@ export function createDiscordDraftStream(params: { flush: loop.flush, messageId: () => streamMessage?.messageId, lastDeliveredText: () => lastSentText, - clear: () => retireCurrentMessage(discardPending), + clear: () => lifecycle.retireCurrent(discardPending), deleteCurrentMessage: () => - retireCurrentMessage(async () => { + lifecycle.retireCurrent(async () => { loop.resetPending(); await loop.waitForInFlight(); }), @@ -234,9 +192,7 @@ export function createDiscordDraftStream(params: { seal, stop, retarget, - cleanupPendingMessages: async () => { - await lifecycle.cleanupPending(); - }, - forceNewMessage, + cleanupPendingMessages: lifecycle.cleanupPending, + forceNewMessage: lifecycle.reset, }; } diff --git a/extensions/matrix/src/matrix/draft-stream.ts b/extensions/matrix/src/matrix/draft-stream.ts index 19bb67ea17ca..7ec4f7de910e 100644 --- a/extensions/matrix/src/matrix/draft-stream.ts +++ b/extensions/matrix/src/matrix/draft-stream.ts @@ -114,6 +114,7 @@ export function createMatrixDraftStream(params: { clear, retire, cleanupPending, + resetMessage, } = createFinalizableDraftLifecycle({ throttleMs: DEFAULT_THROTTLE_MS, state: streamState, @@ -166,14 +167,10 @@ export function createMatrixDraftStream(params: { }; const resetCurrentMessage = (): void => { - currentEventId = undefined; - lastSentText = ""; - lastSentContent = ""; + resetMessage(); sendFailed = false; finalizeInPlaceBlocked = false; liveFinalized = false; - loop.resetPending(); - loop.resetThrottleWindow(); }; const reset = (): void => { // A new block consumes the first-only reply reference; retraction does not. diff --git a/src/channels/draft-stream-controls.ts b/src/channels/draft-stream-controls.ts index de25a9e5e9c6..0e61b1f64ad0 100644 --- a/src/channels/draft-stream-controls.ts +++ b/src/channels/draft-stream-controls.ts @@ -1,9 +1,6 @@ import { formatErrorMessage } from "../infra/errors.js"; import { createDraftStreamLoop } from "./draft-stream-loop.js"; -/** - * Mutable finalization flags shared by draft stream controls and channel adapters. - */ export type FinalizableDraftStreamState = { stopped: boolean; final: boolean; @@ -101,17 +98,13 @@ export function createFinalizableDraftStreamControls( }; } -/** - * Creates finalizable draft controls backed by a shared mutable state object. - */ export function createFinalizableDraftStreamControlsForState( params: DraftStreamOptions & { state: FinalizableDraftStreamState; }, ) { return createFinalizableDraftStreamControls({ - throttleMs: params.throttleMs, - coalesceInFlight: params.coalesceInFlight, + ...params, isStopped: () => params.state.stopped, isFinal: () => params.state.final, markStopped: () => { @@ -120,15 +113,9 @@ export function createFinalizableDraftStreamControlsForState( markFinal: () => { params.state.final = true; }, - sendOrEditStreamMessage: params.sendOrEditStreamMessage, - ...(params.emptyValue !== undefined ? { emptyValue: params.emptyValue } : {}), - ...(params.isEmpty ? { isEmpty: params.isEmpty } : {}), }); } -/** - * Stops a draft stream, reads the current preview message id, then clears the stored id. - */ export async function takeMessageIdAfterStop( params: StopAndClearMessageIdParams, ): Promise { @@ -180,9 +167,6 @@ export async function clearFinalizableDraftMessage( } } -/** - * Builds the standard draft lifecycle used by channel streaming preview implementations. - */ export function createFinalizableDraftLifecycle( params: FinalizableDraftLifecycleParams, ) { @@ -193,6 +177,8 @@ export function createFinalizableDraftLifecycle( }; const pending = new Map(); let clearTail = Promise.resolve(); + let generation = 0; + let discardThroughGeneration = -1; const claim = ( messageId: TMessageId, @@ -275,8 +261,57 @@ export function createFinalizableDraftLifecycle( clearTail = stopRun; return stopRun; }; + const resetMessage = () => { + params.clearMessageId(); + controls.loop.resetPending(); + controls.loop.resetThrottleWindow(); + }; + const reset = (mode: "preserve" | "discard" = "preserve") => { + // A later rotation cannot revoke an earlier request to discard an in-flight create. + if (mode === "discard") { + discardThroughGeneration = generation; + } + generation += 1; + params.state.stopped = false; + params.state.final = false; + resetMessage(); + }; + const createMessage = async ( + send: () => Promise, + publish: (messageId: TMessageId | undefined) => boolean, + ): Promise => { + const startedGeneration = generation; + const messageId = await send(); + if (startedGeneration === generation) { + // Publish in the same turn as the generation check so rotation cannot lose the receipt. + return publish(messageId); + } + if (startedGeneration <= discardThroughGeneration && params.isValidMessageId(messageId)) { + await retire(messageId); + } + return true; + }; + const retireCurrent = async (stopForClear: () => Promise): Promise => { + const startedGeneration = generation; + await stopForClear(); + if (startedGeneration !== generation) { + return; + } + const messageId = params.readMessageId(); + params.clearMessageId(); + if (params.isValidMessageId(messageId)) { + await retire(messageId); + } + }; return { ...controls, + get generation() { + return generation; + }, + reset, + resetMessage, + createMessage, + retireCurrent, stop, clear: () => clearWithStop(controls.stopForClear), clearWithStop,