diff --git a/config/max-lines-baseline.txt b/config/max-lines-baseline.txt index 955e2032c77d..13164b0a326b 100644 --- a/config/max-lines-baseline.txt +++ b/config/max-lines-baseline.txt @@ -92,7 +92,6 @@ extensions/discord/src/channel.ts extensions/discord/src/monitor.test.ts extensions/discord/src/monitor/message-handler.preflight.test.ts extensions/discord/src/monitor/message-handler.preflight.ts -extensions/discord/src/monitor/message-handler.process.ts extensions/discord/src/monitor/message-utils.test.ts extensions/discord/src/monitor/model-picker.test.ts extensions/discord/src/monitor/model-picker.view.ts diff --git a/extensions/discord/src/monitor/message-handler.process-progress.ts b/extensions/discord/src/monitor/message-handler.process-progress.ts new file mode 100644 index 000000000000..4585f8bc17e1 --- /dev/null +++ b/extensions/discord/src/monitor/message-handler.process-progress.ts @@ -0,0 +1,341 @@ +import type { StatusReactionController } from "openclaw/plugin-sdk/channel-feedback"; +import type { ChannelInboundTurnPlan } from "openclaw/plugin-sdk/channel-inbound"; +// Discord plugin module owns progress-window state and agent-event rendering. +import { + buildChannelProgressDraftLine, + buildChannelProgressDraftLineForEntry, + isChannelProgressDraftWorkToolName, +} from "openclaw/plugin-sdk/channel-outbound"; +import { getSessionEntry, resolveStorePath } from "openclaw/plugin-sdk/session-store-runtime"; +import type { createDiscordDraftPreviewController } from "./message-handler.draft-preview.js"; +import type { DiscordMessagePreflightContext } from "./message-handler.preflight.js"; + +type ReplyOptions = NonNullable; +type CallbackPayload = + NonNullable extends (...args: infer Args) => unknown ? Args[0] : never; +type DraftPreview = ReturnType; + +function isProcessAborted(abortSignal?: AbortSignal): boolean { + return Boolean(abortSignal?.aborted); +} + +function isFailedProgress(payload: { + phase?: string; + status?: string; + exitCode?: number | null; +}): boolean { + return ( + payload.phase === "error" || + payload.status === "failed" || + payload.status === "error" || + (typeof payload.exitCode === "number" && payload.exitCode !== 0) + ); +} + +export function createDiscordMessageProgressRuntime(params: { + ctx: DiscordMessagePreflightContext; + sessionKey?: string; + sourceRepliesAreToolOnly: boolean; + draftPreview: DraftPreview; + reactions: { + statusReactionsExplicitlyEnabled: boolean; + statusReactionsEnabled: boolean; + readonly controller: StatusReactionController; + maybeBindToToolReaction: (payload: CallbackPayload<"onToolStart">) => Promise; + }; + onTurnReset: () => void; +}) { + const { ctx, draftPreview } = params; + const { cfg, discordConfig, route, abortSignal } = ctx; + // Reasoning delivery follows the session /reasoning level, not streaming config. + const reasoningLevel = ((): "on" | "stream" | "off" => { + const normalizedAgentId = (route.agentId ?? "").trim().toLowerCase() || "main"; + const agentEntryDefault = cfg.agents?.list?.find( + (entry) => ((entry?.id ?? "").trim().toLowerCase() || "main") === normalizedAgentId, + )?.reasoningDefault; + const cfgDefault = agentEntryDefault ?? cfg.agents?.defaults?.reasoningDefault; + const configDefault: "on" | "stream" | "off" = + cfgDefault === "on" || cfgDefault === "stream" ? cfgDefault : "off"; + if (!params.sessionKey) { + return configDefault; + } + try { + const storePath = resolveStorePath(cfg.session?.store, { agentId: route.agentId }); + const level = getSessionEntry({ + agentId: route.agentId, + sessionKey: params.sessionKey, + storePath, + })?.reasoningLevel; + if (level === "on" || level === "stream" || level === "off") { + return level; + } + } catch { + return "off"; + } + return configDefault; + })(); + const reasoningDurableEnabled = reasoningLevel === "on"; + const reasoningWindowEnabled = reasoningLevel === "stream"; + let shouldYieldDraftProgress: () => boolean = () => false; + let progressTurnStartedAt = Date.now(); + let progressReasoningSteps = 0; + let progressToolCalls = 0; + let progressCommentaryNotes = 0; + // Preamble updates can re-fire; count each item id or id-less text once. + const seenCommentaryIds = new Set(); + let lastCommentaryNoteText = ""; + const noteWindowCommentary = (itemId?: string, noteText?: string) => { + const trimmed = noteText?.trim(); + if (!trimmed) { + return; + } + if (itemId) { + if (seenCommentaryIds.has(itemId)) { + return; + } + seenCommentaryIds.add(itemId); + progressCommentaryNotes += 1; + return; + } + if (trimmed !== lastCommentaryNoteText) { + lastCommentaryNoteText = trimmed; + progressCommentaryNotes += 1; + } + }; + // DeepSeek does not always emit a thinking_end, so tool/final boundaries also close bursts. + let windowReasoningOpen = false; + const closePendingWindowThought = () => { + if (windowReasoningOpen) { + windowReasoningOpen = false; + progressReasoningSteps += 1; + } + }; + const resetTurnState = () => { + progressTurnStartedAt = Date.now(); + progressReasoningSteps = 0; + progressToolCalls = 0; + progressCommentaryNotes = 0; + seenCommentaryIds.clear(); + lastCommentaryNoteText = ""; + windowReasoningOpen = false; + }; + const handleAssistantMessageBoundary = () => { + if (draftPreview.handleAssistantMessageBoundary()) { + resetTurnState(); + params.onTurnReset(); + } + }; + const buildProgressSummaryLine = () => { + closePendingWindowThought(); + const seconds = Math.max(1, Math.round((Date.now() - progressTurnStartedAt) / 1000)); + const parts = [ + ...(progressReasoningSteps > 0 + ? [`🧠 ${progressReasoningSteps} thought${progressReasoningSteps === 1 ? "" : "s"}`] + : []), + ...(progressCommentaryNotes > 0 + ? [`💬 ${progressCommentaryNotes} note${progressCommentaryNotes === 1 ? "" : "s"}`] + : []), + ...(progressToolCalls > 0 + ? [`🛠️ ${progressToolCalls} tool call${progressToolCalls === 1 ? "" : "s"}`] + : []), + `⏱️ ${seconds}s`, + ]; + return `-# ${parts.join(" · ")}`; + }; + + const replyOptions: Partial = { + onAssistantMessageStart: draftPreview.draftStream ? handleAssistantMessageBoundary : undefined, + onReasoningEnd: draftPreview.draftStream + ? () => { + closePendingWindowThought(); + handleAssistantMessageBoundary(); + } + : undefined, + suppressDefaultToolProgressMessages: + (params.sourceRepliesAreToolOnly && params.reactions.statusReactionsExplicitlyEnabled) || + draftPreview.suppressDefaultToolProgressMessages + ? true + : undefined, + allowToolLifecycleWhenProgressHidden: params.reactions.statusReactionsEnabled + ? true + : undefined, + commentaryProgressEnabled: draftPreview.isProgressMode + ? draftPreview.commentaryProgressEnabled + : undefined, + progressPreambleEnabled: + draftPreview.draftStream && draftPreview.isProgressMode ? true : undefined, + commentaryPayloadsEnabled: draftPreview.isProgressMode + ? draftPreview.commentaryProgressEnabled + : undefined, + reasoningPayloadsEnabled: reasoningDurableEnabled, + onVerboseProgressVisibility: (isActive) => { + shouldYieldDraftProgress = isActive; + }, + onNarrationUpdate: draftPreview.narrationProgressEnabled + ? async (payload) => { + if (isProcessAborted(abortSignal) || shouldYieldDraftProgress()) { + return; + } + await draftPreview.pushNarrationProgress(payload.text); + } + : undefined, + onProgressNarratorLifecycle: draftPreview.narrationProgressEnabled + ? (lifecycle) => draftPreview.setProgressNarratorLifecycle(lifecycle) + : undefined, + isProgressDraftVisible: draftPreview.narrationProgressEnabled + ? () => draftPreview.isProgressDraftVisible + : undefined, + narrationHideCommandText: draftPreview.narrationHideCommandText ? true : undefined, + onReasoningStream: async (payload) => { + if (payload?.requiresReasoningProgressOptIn === true && !reasoningWindowEnabled) { + return; + } + if (payload?.text) { + windowReasoningOpen = true; + } + await params.reactions.controller.setThinking(); + await draftPreview.pushReasoningProgress(payload?.text, { + snapshot: payload?.isReasoningSnapshot === true, + }); + }, + streamReasoningInNonStreamModes: reasoningWindowEnabled, + onToolStart: async (payload) => { + if (isProcessAborted(abortSignal)) { + return; + } + await params.reactions.maybeBindToToolReaction(payload); + await params.reactions.controller.setTool(payload.name); + if (payload.phase === "start") { + closePendingWindowThought(); + } + if (shouldYieldDraftProgress()) { + return; + } + // Match the compositor: message/react/typing are not work-tool lines. + if (payload.phase === "start" && isChannelProgressDraftWorkToolName(payload.name)) { + progressToolCalls += 1; + } + await draftPreview.pushToolProgress( + buildChannelProgressDraftLineForEntry( + discordConfig, + { + event: "tool", + itemId: payload.itemId, + toolCallId: payload.toolCallId, + name: payload.name, + phase: payload.phase, + args: payload.args, + }, + payload.detailMode ? { detailMode: payload.detailMode } : undefined, + ), + { toolName: payload.name }, + ); + }, + onItemEvent: async (payload) => { + if (isFailedProgress(payload)) { + return false; + } + if (payload.kind === "preamble") { + if (shouldYieldDraftProgress()) { + return undefined; + } + return await draftPreview.pushPreambleItemEvent(payload, noteWindowCommentary); + } + if (shouldYieldDraftProgress()) { + return undefined; + } + await draftPreview.pushToolProgress( + buildChannelProgressDraftLineForEntry(discordConfig, { + event: "item", + itemId: payload.itemId, + toolCallId: payload.toolCallId, + itemKind: payload.kind, + title: payload.title, + name: payload.name, + phase: payload.phase, + status: payload.status, + summary: payload.summary, + progressText: payload.progressText, + meta: payload.meta, + }), + ); + }, + onPlanUpdate: async (payload) => { + if (payload.phase === "update") { + await draftPreview.pushPlanProgress(payload.steps, { + explanation: payload.explanation, + }); + } + }, + onApprovalEvent: async (payload) => { + if (payload.phase === "requested") { + await draftPreview.pushToolProgress( + buildChannelProgressDraftLine({ + event: "approval", + phase: payload.phase, + title: payload.title, + command: payload.command, + reason: payload.reason, + message: payload.message, + }), + ); + } + }, + onCommandOutput: async (payload) => { + if (isFailedProgress(payload)) { + return false; + } + if (payload.phase !== "end" || shouldYieldDraftProgress()) { + return undefined; + } + await draftPreview.pushToolProgress( + buildChannelProgressDraftLine({ + event: "command-output", + itemId: payload.itemId, + toolCallId: payload.toolCallId, + phase: payload.phase, + title: payload.title, + name: payload.name, + status: payload.status, + exitCode: payload.exitCode, + }), + ); + return undefined; + }, + onPatchSummary: async (payload) => { + if (payload.phase !== "end" || shouldYieldDraftProgress()) { + return; + } + await draftPreview.pushToolProgress( + buildChannelProgressDraftLine({ + event: "patch", + itemId: payload.itemId, + toolCallId: payload.toolCallId, + phase: payload.phase, + title: payload.title, + name: payload.name, + added: payload.added, + modified: payload.modified, + deleted: payload.deleted, + summary: payload.summary, + }), + ); + }, + onCompactionStart: async () => { + if (!isProcessAborted(abortSignal)) { + await params.reactions.controller.setCompacting(); + } + }, + onCompactionEnd: async () => { + if (!isProcessAborted(abortSignal)) { + params.reactions.controller.cancelPending(); + await params.reactions.controller.setThinking(); + } + }, + }; + + return { + replyOptions, + buildProgressSummaryLine, + }; +} diff --git a/extensions/discord/src/monitor/message-handler.process-reactions.ts b/extensions/discord/src/monitor/message-handler.process-reactions.ts new file mode 100644 index 000000000000..31090bb3e611 --- /dev/null +++ b/extensions/discord/src/monitor/message-handler.process-reactions.ts @@ -0,0 +1,287 @@ +// Discord plugin module owns inbound ack and status-reaction lifecycle. +import { resolveAckReaction } from "openclaw/plugin-sdk/agent-runtime"; +import { + createStatusReactionController, + DEFAULT_TIMING, + logAckFailure, + shouldAckReaction as shouldAckReactionGate, + type StatusReactionController, +} from "openclaw/plugin-sdk/channel-feedback"; +import { logVerbose, sleep } from "openclaw/plugin-sdk/runtime-env"; +import { createDiscordRestClient } from "../client.js"; +import { removeReactionDiscord } from "../send.js"; +import { resolveDiscordTargetChannelId } from "../send.shared.js"; +import { resolveDiscordChannelId } from "../targets.js"; +import { + createDiscordAckReactionAdapter, + createDiscordAckReactionContext, + queueInitialDiscordAckReaction, +} from "./ack-reactions.js"; +import type { DiscordMessagePreflightContext } from "./message-handler.preflight.js"; + +type ToolStartPayload = { + name?: string; + phase?: string; + args?: Record; +}; + +function readToolStringArg(args: Record, key: string): string | undefined { + const value = args[key]; + return typeof value === "string" && value.trim() ? value.trim() : undefined; +} + +function readToolBooleanArg(args: Record, key: string): boolean { + return args[key] === true; +} + +export function createDiscordMessageReactionRuntime(params: { + ctx: DiscordMessagePreflightContext; + sourceRepliesAreToolOnly: boolean; + isRoomEvent: boolean; +}) { + const { ctx } = params; + const { + cfg, + accountId, + token, + ackReactionScope, + message, + messageChannelId, + isGuildMessage, + isDirectMessage, + isGroupDm, + shouldRequireMention, + canDetectMention, + effectiveWasMentioned, + shouldBypassMention, + route, + } = ctx; + const ackReaction = resolveAckReaction(cfg, route.agentId, { + channel: "discord", + accountId, + }); + const removeAckAfterReply = cfg.messages?.removeAckAfterReply ?? false; + const shouldSendAckReaction = Boolean( + ackReaction && + shouldAckReactionGate({ + scope: ackReactionScope, + inboundEventKind: ctx.inboundEventKind, + isDirect: isDirectMessage, + isGroup: isGuildMessage || isGroupDm, + isMentionableGroup: isGuildMessage, + requireMention: shouldRequireMention, + canDetectMention, + effectiveWasMentioned, + shouldBypassMention, + }), + ); + const statusReactionsExplicitlyEnabled = cfg.messages?.statusReactions?.enabled === true; + const statusReactionsEnabled = + !params.isRoomEvent && + shouldSendAckReaction && + cfg.messages?.statusReactions?.enabled !== false && + (!params.sourceRepliesAreToolOnly || statusReactionsExplicitlyEnabled); + const feedbackRest = createDiscordRestClient({ cfg, token, accountId }).rest; + const deliveryRest = createDiscordRestClient({ cfg, token, accountId }).rest; + // Discord outbound helpers expect the internal REST client shape explicitly. + const ackReactionContext = createDiscordAckReactionContext({ + rest: feedbackRest, + cfg, + accountId, + }); + const discordAdapter = createDiscordAckReactionAdapter({ + channelId: messageChannelId, + messageId: message.id, + reactionContext: ackReactionContext, + }); + const statusReactionTiming = { + ...DEFAULT_TIMING, + ...cfg.messages?.statusReactions?.timing, + }; + let statusReactionTarget = `${messageChannelId}/${message.id}`; + let statusReactionsActive = statusReactionsEnabled; + let statusReactions: StatusReactionController = createStatusReactionController({ + enabled: statusReactionsEnabled, + adapter: discordAdapter, + initialEmoji: ackReaction, + emojis: cfg.messages?.statusReactions?.emojis, + timing: statusReactionTiming, + onError: (err) => { + logAckFailure({ + log: logVerbose, + channel: "discord", + target: statusReactionTarget, + error: err, + }); + }, + }); + + const resolveTrackedReactionChannelId = async ( + args: Record, + ): Promise => { + const target = + readToolStringArg(args, "channelId") ?? + readToolStringArg(args, "channel_id") ?? + readToolStringArg(args, "to"); + if (!target) { + return messageChannelId; + } + try { + return resolveDiscordChannelId(target); + } catch { + return ( + await resolveDiscordTargetChannelId(target, { + cfg, + token, + accountId, + }) + ).channelId; + } + }; + + const maybeBindToToolReaction = async (payload: ToolStartPayload) => { + if ( + params.sourceRepliesAreToolOnly || + cfg.messages?.statusReactions?.enabled === false || + payload.phase !== "start" || + payload.name !== "message" || + !payload.args + ) { + return; + } + const args = payload.args; + if (readToolStringArg(args, "action")?.toLowerCase() !== "react") { + return; + } + const shouldTrack = + readToolBooleanArg(args, "trackToolCalls") || readToolBooleanArg(args, "track_tool_calls"); + if (!shouldTrack) { + return; + } + const emoji = readToolStringArg(args, "emoji"); + if (!emoji || readToolBooleanArg(args, "remove")) { + return; + } + const trackedMessageId = + readToolStringArg(args, "messageId") ?? readToolStringArg(args, "message_id") ?? message.id; + let trackedChannelId: string; + try { + trackedChannelId = await resolveTrackedReactionChannelId(args); + } catch (err) { + logAckFailure({ + log: logVerbose, + channel: "discord", + target: `${readToolStringArg(args, "to") ?? readToolStringArg(args, "channelId") ?? messageChannelId}/${trackedMessageId}`, + error: err, + }); + return; + } + statusReactionTarget = `${trackedChannelId}/${trackedMessageId}`; + if (statusReactionsActive) { + void statusReactions.clear(); + } + statusReactions = createStatusReactionController({ + enabled: true, + adapter: createDiscordAckReactionAdapter({ + channelId: trackedChannelId, + messageId: trackedMessageId, + reactionContext: ackReactionContext, + }), + initialEmoji: emoji, + emojis: cfg.messages?.statusReactions?.emojis, + timing: statusReactionTiming, + onError: (err) => { + logAckFailure({ + log: logVerbose, + channel: "discord", + target: statusReactionTarget, + error: err, + }); + }, + }); + statusReactionsActive = true; + void statusReactions.setQueued(); + }; + + let initialAckReactionQueued = false; + const queueInitialAckReactionAfterRecord = () => { + if (initialAckReactionQueued) { + return; + } + initialAckReactionQueued = true; + if (statusReactionsEnabled) { + statusReactionsActive = true; + } + queueInitialDiscordAckReaction({ + enabled: statusReactionsEnabled, + shouldSendAckReaction, + ackReaction, + statusReactions, + reactionAdapter: discordAdapter, + target: `${messageChannelId}/${message.id}`, + }); + }; + + const finish = async (result: { + dispatchAborted: boolean; + dispatchError: boolean; + finalDeliveryFailed: boolean; + }) => { + if (statusReactionsActive) { + if (result.dispatchAborted) { + if (removeAckAfterReply) { + void statusReactions.clear(); + } else { + void statusReactions.restoreInitial(); + } + return; + } + if (result.dispatchError || result.finalDeliveryFailed) { + await statusReactions.setError(); + } else { + await statusReactions.setDone(); + } + if (removeAckAfterReply) { + void (async () => { + await sleep( + result.dispatchError || result.finalDeliveryFailed + ? statusReactionTiming.errorHoldMs + : statusReactionTiming.doneHoldMs, + ); + await statusReactions.clear(); + })(); + } else { + void statusReactions.restoreInitial(); + } + return; + } + if (shouldSendAckReaction && ackReaction && removeAckAfterReply) { + void removeReactionDiscord( + messageChannelId, + message.id, + ackReaction, + ackReactionContext, + ).catch((err: unknown) => { + logAckFailure({ + log: logVerbose, + channel: "discord", + target: `${messageChannelId}/${message.id}`, + error: err, + }); + }); + } + }; + + return { + feedbackRest, + deliveryRest, + statusReactionsExplicitlyEnabled, + statusReactionsEnabled, + get controller() { + return statusReactions; + }, + maybeBindToToolReaction, + queueInitialAckReactionAfterRecord, + finish, + }; +} diff --git a/extensions/discord/src/monitor/message-handler.process-reply-runtime.ts b/extensions/discord/src/monitor/message-handler.process-reply-runtime.ts new file mode 100644 index 000000000000..59861c39b20f --- /dev/null +++ b/extensions/discord/src/monitor/message-handler.process-reply-runtime.ts @@ -0,0 +1,169 @@ +// Discord plugin module owns the reply pipeline, draft preview, and delivery correlation setup. +import { + createChannelMessageReplyPipeline, + resolveChannelStreamingBlockEnabled, +} from "openclaw/plugin-sdk/channel-outbound"; +import { resolveMarkdownTableMode } from "openclaw/plugin-sdk/markdown-table-runtime"; +import { resolveChunkMode } from "openclaw/plugin-sdk/reply-chunking"; +import { createChannelHistoryWindow } from "openclaw/plugin-sdk/reply-history"; +import { logVerbose } from "openclaw/plugin-sdk/runtime-env"; +import { getSessionEntry, resolveStorePath } from "openclaw/plugin-sdk/session-store-runtime"; +import { readLatestAssistantTextByIdentity } from "openclaw/plugin-sdk/session-transcript-runtime"; +import { resolveDiscordMaxLinesPerMessage } from "../accounts.js"; +import { beginDiscordInboundEventDeliveryCorrelation } from "../inbound-event-delivery.js"; +import type { RequestClient } from "../internal/discord.js"; +import { buildDiscordMessageProcessContext } from "./message-handler.context.js"; +import { createDiscordDraftPreviewController } from "./message-handler.draft-preview.js"; +import type { DiscordMessagePreflightContext } from "./message-handler.preflight.js"; +import { createDiscordReplyTypingFeedback } from "./reply-typing-feedback.js"; + +type DiscordMessageProcessContext = NonNullable< + Awaited> +>; + +export function createDiscordMessageReplyRuntime(params: { + ctx: DiscordMessagePreflightContext; + processContext: DiscordMessageProcessContext; + sourceRepliesAreToolOnly: boolean; + shouldDisableCoreTypingKeepalive: boolean; + isRoomEvent: boolean; + dispatchStartedAt: number; + feedbackRest: RequestClient; + deliveryRest: RequestClient; +}) { + const { ctx, processContext } = params; + const { + cfg, + discordConfig, + accountId, + token, + guildHistories, + historyLimit, + textLimit, + messageChannelId, + isDirectMessage, + route, + } = ctx; + const { ctxPayload, deliverTarget, replyReference } = processContext; + const typingChannelId = deliverTarget.startsWith("channel:") + ? deliverTarget.slice("channel:".length) + : messageChannelId; + let typingFeedback: ReturnType | undefined; + const getTypingFeedback = () => + (typingFeedback ??= createDiscordReplyTypingFeedback({ + cfg, + token, + accountId, + channelId: typingChannelId, + rest: params.feedbackRest, + log: logVerbose, + keepaliveIntervalMs: params.shouldDisableCoreTypingKeepalive ? undefined : 0, + })); + + const { onModelSelected, ...replyPipeline } = createChannelMessageReplyPipeline({ + cfg, + agentId: route.agentId, + channel: "discord", + accountId: route.accountId, + // The core lifecycle reaches this callback only after reply admission. + // Silent pre-dispatch outcomes therefore never allocate or emit feedback. + typingCallbacks: { + onReplyStart: () => getTypingFeedback().onReplyStart(), + onIdle: () => typingFeedback?.onIdle?.(), + onCleanup: () => typingFeedback?.onCleanup?.(), + }, + }); + const tableMode = resolveMarkdownTableMode({ cfg, channel: "discord", accountId }); + const maxLinesPerMessage = resolveDiscordMaxLinesPerMessage({ + cfg, + discordConfig, + accountId, + }); + const chunkMode = resolveChunkMode(cfg, "discord", accountId); + const clearGroupHistory = () => { + if (isDirectMessage) { + return; + } + createChannelHistoryWindow({ historyMap: guildHistories }).clear({ + historyKey: messageChannelId, + limit: historyLimit, + }); + }; + const beginDeliveryCorrelation = () => + params.isRoomEvent + ? beginDiscordInboundEventDeliveryCorrelation( + ctxPayload.SessionKey, + { + outboundTo: messageChannelId, + outboundAccountId: route.accountId, + markInboundEventDelivered: clearGroupHistory, + }, + { inboundEventKind: ctxPayload.InboundEventKind }, + ) + : () => {}; + const endDeliveryCorrelation = beginDeliveryCorrelation(); + + const resolveCurrentTurnTranscriptFinalText = async (): Promise => { + const sessionKey = ctxPayload.SessionKey; + if (!sessionKey) { + return undefined; + } + try { + const storePath = resolveStorePath(cfg.session?.store, { agentId: route.agentId }); + const sessionEntry = getSessionEntry({ + agentId: route.agentId, + sessionKey, + storePath, + }); + if (!sessionEntry?.sessionId) { + return undefined; + } + const latest = await readLatestAssistantTextByIdentity({ + agentId: route.agentId, + sessionId: sessionEntry.sessionId, + sessionKey, + storePath, + }); + if (!latest?.timestamp || latest.timestamp < params.dispatchStartedAt) { + return undefined; + } + return latest.text; + } catch (err) { + logVerbose(`discord transcript final candidate lookup failed: ${String(err)}`); + return undefined; + } + }; + + const deliverChannelId = deliverTarget.startsWith("channel:") + ? deliverTarget.slice("channel:".length) + : messageChannelId; + const draftPreview = createDiscordDraftPreviewController({ + cfg, + discordConfig, + accountId, + sourceRepliesAreToolOnly: params.sourceRepliesAreToolOnly, + textLimit, + deliveryRest: params.deliveryRest, + deliverChannelId, + replyReference, + tableMode, + maxLinesPerMessage, + chunkMode, + log: logVerbose, + }); + const resolvedBlockStreamingEnabled = resolveChannelStreamingBlockEnabled(discordConfig); + + return { + replyPipeline, + onModelSelected, + tableMode, + maxLinesPerMessage, + chunkMode, + beginQueuedDeliveryCorrelation: beginDeliveryCorrelation, + endDeliveryCorrelation, + resolveCurrentTurnTranscriptFinalText, + deliverChannelId, + draftPreview, + resolvedBlockStreamingEnabled, + }; +} diff --git a/extensions/discord/src/monitor/message-handler.process.ts b/extensions/discord/src/monitor/message-handler.process.ts index f3344b1bc5ad..f5c98386fc6d 100644 --- a/extensions/discord/src/monitor/message-handler.process.ts +++ b/extensions/discord/src/monitor/message-handler.process.ts @@ -1,34 +1,18 @@ // Discord plugin module implements message handler.process behavior. import type { APIAllowedMentions } from "discord-api-types/v10"; -import { resolveAckReaction, resolveHumanDelayConfig } from "openclaw/plugin-sdk/agent-runtime"; -import { - createStatusReactionController, - DEFAULT_TIMING, - logAckFailure, - shouldAckReaction as shouldAckReactionGate, -} from "openclaw/plugin-sdk/channel-feedback"; +import { resolveHumanDelayConfig } from "openclaw/plugin-sdk/agent-runtime"; import { dispatchChannelInboundTurn, hasFinalInboundReplyDispatch, } from "openclaw/plugin-sdk/channel-inbound"; import { bindIngressLifecycleToReplyOptions, - createChannelMessageReplyPipeline, defineFinalizableLivePreviewAdapter, deliverWithFinalizableLivePreviewAdapter, resolveChannelMessageSourceReplyDeliveryMode, } from "openclaw/plugin-sdk/channel-outbound"; -import { - buildChannelProgressDraftLine, - buildChannelProgressDraftLineForEntry, - isChannelProgressDraftWorkToolName, - resolveChannelStreamingBlockEnabled, - resolveTranscriptBackedChannelFinalText, -} from "openclaw/plugin-sdk/channel-outbound"; -import { resolveMarkdownTableMode } from "openclaw/plugin-sdk/markdown-table-runtime"; +import { resolveTranscriptBackedChannelFinalText } from "openclaw/plugin-sdk/channel-outbound"; import { getAgentScopedMediaLocalRoots } from "openclaw/plugin-sdk/media-runtime"; -import { resolveChunkMode } from "openclaw/plugin-sdk/reply-chunking"; -import { createChannelHistoryWindow } from "openclaw/plugin-sdk/reply-history"; import { getReplyPayloadTtsSupplement, isReplyPayloadNonTerminalToolErrorWarning, @@ -39,33 +23,24 @@ import { danger, logVerbose, shouldLogVerbose, - sleep, sleepWithAbort, } from "openclaw/plugin-sdk/runtime-env"; -import { getSessionEntry, resolveStorePath } from "openclaw/plugin-sdk/session-store-runtime"; -import { readLatestAssistantTextByIdentity } from "openclaw/plugin-sdk/session-transcript-runtime"; -import { resolveDiscordMaxLinesPerMessage } from "../accounts.js"; import { chunkDiscordTextWithMode } from "../chunk.js"; -import { createDiscordRestClient } from "../client.js"; -import { beginDiscordInboundEventDeliveryCorrelation } from "../inbound-event-delivery.js"; import { discordTextHasBroadcastMention } from "../mentions.js"; -import { removeReactionDiscord } from "../send.js"; import { editMessageDiscord } from "../send.messages.js"; -import { resolveDiscordTargetChannelId } from "../send.shared.js"; import type { DiscordMessageEdit } from "../send.types.js"; -import { resolveDiscordChannelId } from "../targets.js"; -import { - createDiscordAckReactionAdapter, - createDiscordAckReactionContext, - queueInitialDiscordAckReaction, -} from "./ack-reactions.js"; import { buildDiscordMessageProcessContext } from "./message-handler.context.js"; -import { createDiscordDraftPreviewController } from "./message-handler.draft-preview.js"; import type { DiscordMessagePreflightContext } from "./message-handler.preflight.js"; +import { createDiscordMessageProgressRuntime } from "./message-handler.process-progress.js"; +import { createDiscordMessageReactionRuntime } from "./message-handler.process-reactions.js"; +import { createDiscordMessageReplyRuntime } from "./message-handler.process-reply-runtime.js"; import { completeDiscordSessionConflict } from "./message-handler.retry.js"; -import { deliverDiscordReply, formatDiscordReplyDeliveryFailure } from "./reply-delivery.js"; +import { + deliverDiscordReply, + formatDiscordReplyDeliveryFailure, + formatDiscordReplySkip, +} from "./reply-delivery.js"; import { sanitizeDiscordFrontChannelReplyPayloads } from "./reply-safety.js"; -import { createDiscordReplyTypingFeedback } from "./reply-typing-feedback.js"; const TARGETED_ONLY_ALLOWED_MENTIONS = { parse: ["users", "roles"], @@ -82,35 +57,7 @@ function isFallbackOnlyToolWarningFinal(payload: ReplyPayload): boolean { return !resolveSendableOutboundReplyParts(payload).hasMedia; } -function isFailedProgress(payload: { - phase?: string; - status?: string; - exitCode?: number | null; -}): boolean { - return ( - payload.phase === "error" || - payload.status === "failed" || - payload.status === "error" || - (typeof payload.exitCode === "number" && payload.exitCode !== 0) - ); -} - -type DiscordReplySkipReason = "aborted before delivery" | "internal-only payload"; - -export function formatDiscordReplySkip(params: { - kind: "tool" | "block" | "final"; - reason: DiscordReplySkipReason; - target: string; - sessionKey?: string; -}) { - const context = [ - `target=${params.target}`, - params.sessionKey ? `session=${params.sessionKey}` : undefined, - ] - .filter(Boolean) - .join(" "); - return `discord ${params.kind} reply skipped (${params.reason}): ${context}`; -} +export { formatDiscordReplySkip } from "./reply-delivery.js"; type DiscordMessageProcessObserver = { onFinalReplyStart?: () => void; @@ -118,22 +65,6 @@ type DiscordMessageProcessObserver = { onReplyPlanResolved?: (params: { createdThreadId?: string; sessionKey?: string }) => void; }; -type ToolStartPayload = { - name?: string; - phase?: string; - args?: Record; - detailMode?: "explain" | "raw"; -}; - -function readToolStringArg(args: Record, key: string): string | undefined { - const value = args[key]; - return typeof value === "string" && value.trim() ? value.trim() : undefined; -} - -function readToolBooleanArg(args: Record, key: string): boolean { - return args[key] === true; -} - export async function processDiscordMessage( ctx: DiscordMessagePreflightContext, observer?: DiscordMessageProcessObserver, @@ -148,7 +79,6 @@ async function processDiscordMessageInner( const dispatchStartedAt = Date.now(); const { cfg, - discordConfig, accountId, token, runtime, @@ -156,17 +86,12 @@ async function processDiscordMessageInner( historyLimit, textLimit, replyToMode, - ackReactionScope, message, messageChannelId, isGuildMessage, isDirectMessage, isGroupDm, messageText, - shouldRequireMention, - canDetectMention, - effectiveWasMentioned, - shouldBypassMention, channelConfig, threadBindings, route, @@ -208,183 +133,13 @@ async function processDiscordMessageInner( sourceRepliesAreToolOnly && configuredTypingMode === undefined && configuredTypingInterval === undefined; - const ackReaction = resolveAckReaction(cfg, route.agentId, { - channel: "discord", - accountId, - }); - const removeAckAfterReply = cfg.messages?.removeAckAfterReply ?? false; const mediaLocalRoots = getAgentScopedMediaLocalRoots(cfg, route.agentId); const isRoomEvent = ctx.inboundEventKind === "room_event"; - const shouldAckReaction = () => - Boolean( - ackReaction && - shouldAckReactionGate({ - scope: ackReactionScope, - inboundEventKind: ctx.inboundEventKind, - isDirect: isDirectMessage, - isGroup: isGuildMessage || isGroupDm, - isMentionableGroup: isGuildMessage, - requireMention: shouldRequireMention, - canDetectMention, - effectiveWasMentioned, - shouldBypassMention, - }), - ); - const shouldSendAckReaction = shouldAckReaction(); - const statusReactionsExplicitlyEnabled = cfg.messages?.statusReactions?.enabled === true; - const statusReactionsEnabled = - !isRoomEvent && - shouldSendAckReaction && - cfg.messages?.statusReactions?.enabled !== false && - (!sourceRepliesAreToolOnly || statusReactionsExplicitlyEnabled); - const feedbackRest = createDiscordRestClient({ - cfg, - token, - accountId, - }).rest; - const deliveryRest = createDiscordRestClient({ - cfg, - token, - accountId, - }).rest; - // Discord outbound helpers expect the internal REST client shape explicitly. - const ackReactionContext = createDiscordAckReactionContext({ - rest: feedbackRest, - cfg, - accountId, + const reactions = createDiscordMessageReactionRuntime({ + ctx, + sourceRepliesAreToolOnly, + isRoomEvent, }); - const discordAdapter = createDiscordAckReactionAdapter({ - channelId: messageChannelId, - messageId: message.id, - reactionContext: ackReactionContext, - }); - const statusReactionTiming = { - ...DEFAULT_TIMING, - ...cfg.messages?.statusReactions?.timing, - }; - let statusReactionTarget = `${messageChannelId}/${message.id}`; - let statusReactionsActive = statusReactionsEnabled; - let statusReactions = createStatusReactionController({ - enabled: statusReactionsEnabled, - adapter: discordAdapter, - initialEmoji: ackReaction, - emojis: cfg.messages?.statusReactions?.emojis, - timing: statusReactionTiming, - onError: (err) => { - logAckFailure({ - log: logVerbose, - channel: "discord", - target: statusReactionTarget, - error: err, - }); - }, - }); - const resolveTrackedReactionChannelId = async ( - args: Record, - ): Promise => { - const target = - readToolStringArg(args, "channelId") ?? - readToolStringArg(args, "channel_id") ?? - readToolStringArg(args, "to"); - if (!target) { - return messageChannelId; - } - try { - return resolveDiscordChannelId(target); - } catch { - return ( - await resolveDiscordTargetChannelId(target, { - cfg, - token, - accountId, - }) - ).channelId; - } - }; - const maybeBindStatusReactionsToToolReaction = async (payload: ToolStartPayload) => { - if ( - sourceRepliesAreToolOnly || - cfg.messages?.statusReactions?.enabled === false || - payload.phase !== "start" || - payload.name !== "message" || - !payload.args - ) { - return; - } - const args = payload.args; - const action = readToolStringArg(args, "action")?.toLowerCase(); - if (action !== "react") { - return; - } - const shouldTrack = - readToolBooleanArg(args, "trackToolCalls") || readToolBooleanArg(args, "track_tool_calls"); - if (!shouldTrack) { - return; - } - const emoji = readToolStringArg(args, "emoji"); - const remove = readToolBooleanArg(args, "remove"); - if (!emoji || remove) { - return; - } - const trackedMessageId = - readToolStringArg(args, "messageId") ?? readToolStringArg(args, "message_id") ?? message.id; - let trackedChannelId: string; - try { - trackedChannelId = await resolveTrackedReactionChannelId(args); - } catch (err) { - logAckFailure({ - log: logVerbose, - channel: "discord", - target: `${readToolStringArg(args, "to") ?? readToolStringArg(args, "channelId") ?? messageChannelId}/${trackedMessageId}`, - error: err, - }); - return; - } - statusReactionTarget = `${trackedChannelId}/${trackedMessageId}`; - if (statusReactionsActive) { - void statusReactions.clear(); - } - const trackedAdapter = createDiscordAckReactionAdapter({ - channelId: trackedChannelId, - messageId: trackedMessageId, - reactionContext: ackReactionContext, - }); - statusReactions = createStatusReactionController({ - enabled: true, - adapter: trackedAdapter, - initialEmoji: emoji, - emojis: cfg.messages?.statusReactions?.emojis, - timing: statusReactionTiming, - onError: (err) => { - logAckFailure({ - log: logVerbose, - channel: "discord", - target: statusReactionTarget, - error: err, - }); - }, - }); - statusReactionsActive = true; - void statusReactions.setQueued(); - }; - let initialAckReactionQueued = false; - const queueInitialAckReactionAfterRecord = () => { - if (initialAckReactionQueued) { - return; - } - initialAckReactionQueued = true; - if (statusReactionsEnabled) { - statusReactionsActive = true; - } - queueInitialDiscordAckReaction({ - enabled: statusReactionsEnabled, - shouldSendAckReaction, - ackReaction, - statusReactions, - reactionAdapter: discordAdapter, - target: `${messageChannelId}/${message.id}`, - }); - }; const processContext = await buildDiscordMessageProcessContext({ ctx, text, @@ -407,116 +162,29 @@ async function processDiscordMessageInner( sessionKey: persistedSessionKey, }); - const typingChannelId = deliverTarget.startsWith("channel:") - ? deliverTarget.slice("channel:".length) - : messageChannelId; - let typingFeedback: ReturnType | undefined; - const getTypingFeedback = () => - (typingFeedback ??= createDiscordReplyTypingFeedback({ - cfg, - token, - accountId, - channelId: typingChannelId, - rest: feedbackRest, - log: logVerbose, - keepaliveIntervalMs: shouldDisableCoreTypingKeepalive ? undefined : 0, - })); - - const { onModelSelected, ...replyPipeline } = createChannelMessageReplyPipeline({ - cfg, - agentId: route.agentId, - channel: "discord", - accountId: route.accountId, - // The core lifecycle reaches this callback only after reply admission. - // Silent pre-dispatch outcomes therefore never allocate or emit feedback. - typingCallbacks: { - onReplyStart: () => getTypingFeedback().onReplyStart(), - onIdle: () => typingFeedback?.onIdle?.(), - onCleanup: () => typingFeedback?.onCleanup?.(), - }, - }); - const tableMode = resolveMarkdownTableMode({ - cfg, - channel: "discord", - accountId, - }); - const maxLinesPerMessage = resolveDiscordMaxLinesPerMessage({ - cfg, - discordConfig, - accountId, - }); - const chunkMode = resolveChunkMode(cfg, "discord", accountId); - const clearGroupHistory = () => { - if (isDirectMessage) { - return; - } - createChannelHistoryWindow({ historyMap: guildHistories }).clear({ - historyKey: messageChannelId, - limit: historyLimit, - }); - }; - const beginDeliveryCorrelation = () => - isRoomEvent - ? beginDiscordInboundEventDeliveryCorrelation( - ctxPayload.SessionKey, - { - outboundTo: messageChannelId, - outboundAccountId: route.accountId, - markInboundEventDelivered: clearGroupHistory, - }, - { inboundEventKind: ctxPayload.InboundEventKind }, - ) - : () => {}; - const endDiscordInboundEventDeliveryCorrelation = beginDeliveryCorrelation(); - const resolveCurrentTurnTranscriptFinalText = async (): Promise => { - const sessionKey = ctxPayload.SessionKey; - if (!sessionKey) { - return undefined; - } - try { - const storePath = resolveStorePath(cfg.session?.store, { agentId: route.agentId }); - const sessionEntry = getSessionEntry({ - agentId: route.agentId, - sessionKey, - storePath, - }); - if (!sessionEntry?.sessionId) { - return undefined; - } - const latest = await readLatestAssistantTextByIdentity({ - agentId: route.agentId, - sessionId: sessionEntry.sessionId, - sessionKey, - storePath, - }); - if (!latest?.timestamp || latest.timestamp < dispatchStartedAt) { - return undefined; - } - return latest.text; - } catch (err) { - logVerbose(`discord transcript final candidate lookup failed: ${String(err)}`); - return undefined; - } - }; - - const deliverChannelId = deliverTarget.startsWith("channel:") - ? deliverTarget.slice("channel:".length) - : messageChannelId; - const draftPreview = createDiscordDraftPreviewController({ - cfg, - discordConfig, - accountId, + const replyRuntime = createDiscordMessageReplyRuntime({ + ctx, + processContext, sourceRepliesAreToolOnly, - textLimit, - deliveryRest, - deliverChannelId, - replyReference, + shouldDisableCoreTypingKeepalive, + isRoomEvent, + dispatchStartedAt, + feedbackRest: reactions.feedbackRest, + deliveryRest: reactions.deliveryRest, + }); + const { + replyPipeline, + onModelSelected, tableMode, maxLinesPerMessage, chunkMode, - log: logVerbose, - }); - let shouldYieldDraftProgress: () => boolean = () => false; + beginQueuedDeliveryCorrelation, + endDeliveryCorrelation, + resolveCurrentTurnTranscriptFinalText, + deliverChannelId, + draftPreview, + resolvedBlockStreamingEnabled, + } = replyRuntime; let finalReplyStartNotified = false; const notifyFinalReplyStart = () => { if (finalReplyStartNotified) { @@ -550,110 +218,26 @@ async function processDiscordMessageInner( lines[0] = `🧠 ${lines[0]}`; return lines.map((line) => `> ${line}`).join("\n"); }; - // Reasoning delivery follows the session /reasoning level, not streaming config. - const reasoningLevel = ((): "on" | "stream" | "off" => { - const normalizedAgentId = (route.agentId ?? "").trim().toLowerCase() || "main"; - const agentEntryDefault = cfg.agents?.list?.find( - (entry) => ((entry?.id ?? "").trim().toLowerCase() || "main") === normalizedAgentId, - )?.reasoningDefault; - const cfgDefault = agentEntryDefault ?? cfg.agents?.defaults?.reasoningDefault; - const configDefault: "on" | "stream" | "off" = - cfgDefault === "on" || cfgDefault === "stream" ? cfgDefault : "off"; - const sessionKey = ctxPayload.SessionKey; - if (!sessionKey) { - return configDefault; - } - try { - const storePath = resolveStorePath(cfg.session?.store, { agentId: route.agentId }); - const level = getSessionEntry({ - agentId: route.agentId, - sessionKey, - storePath, - })?.reasoningLevel; - if (level === "on" || level === "stream" || level === "off") { - return level; - } - } catch { - return "off"; - } - return configDefault; - })(); - const reasoningDurableEnabled = reasoningLevel === "on"; - const reasoningWindowEnabled = reasoningLevel === "stream"; - let progressTurnStartedAt = Date.now(); - let progressReasoningSteps = 0; - let progressToolCalls = 0; - let progressCommentaryNotes = 0; // Set when a progress draft collapses: the receipt appends to the final // answer text and the draft message deletes once that answer delivered. let progressReceiptLine: string | undefined; let clearProgressDraftAfterFinalDelivery = false; - // Preamble updates can re-fire; count each item id or id-less text once. - const seenCommentaryIds = new Set(); - let lastCommentaryNoteText = ""; - const noteWindowCommentary = (itemId?: string, noteText?: string) => { - const trimmed = noteText?.trim(); - if (!trimmed) { - return; - } - if (itemId) { - if (seenCommentaryIds.has(itemId)) { - return; - } - seenCommentaryIds.add(itemId); - progressCommentaryNotes += 1; - return; - } - if (trimmed !== lastCommentaryNoteText) { - lastCommentaryNoteText = trimmed; - progressCommentaryNotes += 1; - } - }; - // DeepSeek does not always emit a thinking_end, so tool/final boundaries also close bursts. - let windowReasoningOpen = false; - const closePendingWindowThought = () => { - if (windowReasoningOpen) { - windowReasoningOpen = false; - progressReasoningSteps += 1; - } - }; - const resetProgressTurnState = () => { + const resetDeliveryState = () => { finalReplyStartNotified = false; userFacingFinalDelivered = false; userFacingFinalDeliveryFailed = false; pendingToolWarningFinal = undefined; - progressTurnStartedAt = Date.now(); - progressReasoningSteps = 0; - progressToolCalls = 0; - progressCommentaryNotes = 0; progressReceiptLine = undefined; clearProgressDraftAfterFinalDelivery = false; - seenCommentaryIds.clear(); - lastCommentaryNoteText = ""; - windowReasoningOpen = false; - }; - const handleAssistantMessageBoundary = () => { - if (draftPreview.handleAssistantMessageBoundary()) { - resetProgressTurnState(); - } - }; - const buildProgressSummaryLine = () => { - closePendingWindowThought(); - const seconds = Math.max(1, Math.round((Date.now() - progressTurnStartedAt) / 1000)); - const parts = [ - ...(progressReasoningSteps > 0 - ? [`🧠 ${progressReasoningSteps} thought${progressReasoningSteps === 1 ? "" : "s"}`] - : []), - ...(progressCommentaryNotes > 0 - ? [`💬 ${progressCommentaryNotes} note${progressCommentaryNotes === 1 ? "" : "s"}`] - : []), - ...(progressToolCalls > 0 - ? [`🛠️ ${progressToolCalls} tool call${progressToolCalls === 1 ? "" : "s"}`] - : []), - `⏱️ ${seconds}s`, - ]; - return `-# ${parts.join(" · ")}`; }; + const progress = createDiscordMessageProgressRuntime({ + ctx, + sessionKey: ctxPayload.SessionKey, + sourceRepliesAreToolOnly, + draftPreview, + reactions, + onTurnReset: resetDeliveryState, + }); let replyLifecycleStarted = false; const onDiscordReplyStart = async () => { if (isProcessAborted(abortSignal)) { @@ -661,7 +245,7 @@ async function processDiscordMessageInner( } replyLifecycleStarted = true; await replyPipeline.typingCallbacks?.onReplyStart(); - await statusReactions.setThinking(); + await reactions.controller.setThinking(); }; const beforeDiscordPayloadDelivery = ( payload: ReplyPayload, @@ -737,7 +321,7 @@ async function processDiscordMessageInner( target: deliverTarget, token, accountId, - rest: deliveryRest, + rest: reactions.deliveryRest, runtime, replyToId: replyReference.use(), replyToMode, @@ -815,7 +399,7 @@ async function processDiscordMessageInner( // deletes after that answer lands, so busy channels keep no orphaned // tool log above the reply. Error finals skip both and keep the draft // as the visible record of the failed turn. - progressReceiptLine = buildProgressSummaryLine(); + progressReceiptLine = progress.buildProgressSummaryLine(); clearProgressDraftAfterFinalDelivery = true; // Fall through to the generic fresh send below for the final itself. } @@ -848,7 +432,7 @@ async function processDiscordMessageInner( await editMessageDiscord(deliverChannelId, previewMessageId, edit, { cfg, accountId, - rest: deliveryRest, + rest: reactions.deliveryRest, }); }, onPreviewFinalized: () => { @@ -885,7 +469,7 @@ async function processDiscordMessageInner( target: deliverTarget, token, accountId, - rest: deliveryRest, + rest: reactions.deliveryRest, runtime, replyToId, replyToMode, @@ -944,7 +528,7 @@ async function processDiscordMessageInner( target: deliverTarget, token, accountId, - rest: deliveryRest, + rest: reactions.deliveryRest, runtime, replyToId, replyToMode, @@ -991,7 +575,6 @@ async function processDiscordMessageInner( ), ); }; - const resolvedBlockStreamingEnabled = resolveChannelStreamingBlockEnabled(discordConfig); let dispatchResult: { queuedFinal: boolean; counts: Record; @@ -1026,7 +609,7 @@ async function processDiscordMessageInner( accountId: route.accountId, route: { agentId: route.agentId, sessionKey: persistedSessionKey }, ctxPayload, - afterRecord: queueInitialAckReactionAfterRecord, + afterRecord: reactions.queueInitialAckReactionAfterRecord, sessionInitRetry: { delaysMs: [250, 1_000, 2_500], signal: abortSignal, @@ -1058,7 +641,11 @@ async function processDiscordMessageInner( skillFilter: channelConfig?.skills, sourceReplyDeliveryMode, typingKeepalive: shouldDisableCoreTypingKeepalive ? false : undefined, - queuedDeliveryCorrelations: isRoomEvent ? [{ begin: beginDeliveryCorrelation }] : undefined, + // The primary turn already owns one correlation; each queued followup + // needs a fresh owner so its eventual delivery clears room history. + queuedDeliveryCorrelations: isRoomEvent + ? [{ begin: beginQueuedDeliveryCorrelation }] + : undefined, suppressTyping: isRoomEvent ? true : undefined, allowProgressCallbacksWhenSourceDeliverySuppressed: sourceRepliesAreToolOnly && draftPreview.draftStream && draftPreview.isProgressMode @@ -1074,205 +661,8 @@ async function processDiscordMessageInner( draftPreview.draftStream && !draftPreview.isProgressMode ? (payload) => draftPreview.updateFromPartial(payload.text) : undefined, - onAssistantMessageStart: draftPreview.draftStream - ? handleAssistantMessageBoundary - : undefined, - onReasoningEnd: draftPreview.draftStream - ? () => { - closePendingWindowThought(); - handleAssistantMessageBoundary(); - } - : undefined, + ...progress.replyOptions, onModelSelected, - suppressDefaultToolProgressMessages: - (sourceRepliesAreToolOnly && statusReactionsExplicitlyEnabled) || - draftPreview.suppressDefaultToolProgressMessages - ? true - : undefined, - allowToolLifecycleWhenProgressHidden: statusReactionsEnabled ? true : undefined, - commentaryProgressEnabled: draftPreview.isProgressMode - ? draftPreview.commentaryProgressEnabled - : undefined, - progressPreambleEnabled: - draftPreview.draftStream && draftPreview.isProgressMode ? true : undefined, - commentaryPayloadsEnabled: draftPreview.isProgressMode - ? draftPreview.commentaryProgressEnabled - : undefined, - reasoningPayloadsEnabled: reasoningDurableEnabled, - onVerboseProgressVisibility: (isActive) => { - shouldYieldDraftProgress = isActive; - }, - onNarrationUpdate: draftPreview.narrationProgressEnabled - ? async (payload) => { - if (isProcessAborted(abortSignal) || shouldYieldDraftProgress()) { - return; - } - await draftPreview.pushNarrationProgress(payload.text); - } - : undefined, - onProgressNarratorLifecycle: draftPreview.narrationProgressEnabled - ? (lifecycle) => draftPreview.setProgressNarratorLifecycle(lifecycle) - : undefined, - isProgressDraftVisible: draftPreview.narrationProgressEnabled - ? () => draftPreview.isProgressDraftVisible - : undefined, - narrationHideCommandText: draftPreview.narrationHideCommandText ? true : undefined, - onReasoningStream: async (payload) => { - if (payload?.requiresReasoningProgressOptIn === true && !reasoningWindowEnabled) { - return; - } - if (payload?.text) { - windowReasoningOpen = true; - } - await statusReactions.setThinking(); - await draftPreview.pushReasoningProgress(payload?.text, { - snapshot: payload?.isReasoningSnapshot === true, - }); - }, - streamReasoningInNonStreamModes: reasoningWindowEnabled, - onToolStart: async (payload) => { - if (isProcessAborted(abortSignal)) { - return; - } - await maybeBindStatusReactionsToToolReaction(payload); - await statusReactions.setTool(payload.name); - if (payload.phase === "start") { - closePendingWindowThought(); - } - if (shouldYieldDraftProgress()) { - return; - } - // Match the compositor: message/react/typing are not work-tool lines. - if (payload.phase === "start" && isChannelProgressDraftWorkToolName(payload.name)) { - progressToolCalls += 1; - } - await draftPreview.pushToolProgress( - buildChannelProgressDraftLineForEntry( - discordConfig, - { - event: "tool", - itemId: payload.itemId, - toolCallId: payload.toolCallId, - name: payload.name, - phase: payload.phase, - args: payload.args, - }, - payload.detailMode ? { detailMode: payload.detailMode } : undefined, - ), - { toolName: payload.name }, - ); - }, - onItemEvent: async (payload) => { - if (isFailedProgress(payload)) { - return false; - } - if (payload.kind === "preamble") { - if (shouldYieldDraftProgress()) { - return undefined; - } - return await draftPreview.pushPreambleItemEvent(payload, noteWindowCommentary); - } - if (shouldYieldDraftProgress()) { - return undefined; - } - await draftPreview.pushToolProgress( - buildChannelProgressDraftLineForEntry(discordConfig, { - event: "item", - itemId: payload.itemId, - toolCallId: payload.toolCallId, - itemKind: payload.kind, - title: payload.title, - name: payload.name, - phase: payload.phase, - status: payload.status, - summary: payload.summary, - progressText: payload.progressText, - meta: payload.meta, - }), - ); - }, - onPlanUpdate: async (payload) => { - if (payload.phase !== "update") { - return; - } - await draftPreview.pushPlanProgress(payload.steps, { - explanation: payload.explanation, - }); - }, - onApprovalEvent: async (payload) => { - if (payload.phase !== "requested") { - return; - } - await draftPreview.pushToolProgress( - buildChannelProgressDraftLine({ - event: "approval", - phase: payload.phase, - title: payload.title, - command: payload.command, - reason: payload.reason, - message: payload.message, - }), - ); - }, - onCommandOutput: async (payload) => { - if (isFailedProgress(payload)) { - return false; - } - if (payload.phase !== "end") { - return undefined; - } - if (shouldYieldDraftProgress()) { - return undefined; - } - await draftPreview.pushToolProgress( - buildChannelProgressDraftLine({ - event: "command-output", - itemId: payload.itemId, - toolCallId: payload.toolCallId, - phase: payload.phase, - title: payload.title, - name: payload.name, - status: payload.status, - exitCode: payload.exitCode, - }), - ); - return undefined; - }, - onPatchSummary: async (payload) => { - if (payload.phase !== "end") { - return; - } - if (shouldYieldDraftProgress()) { - return; - } - await draftPreview.pushToolProgress( - buildChannelProgressDraftLine({ - event: "patch", - itemId: payload.itemId, - toolCallId: payload.toolCallId, - phase: payload.phase, - title: payload.title, - name: payload.name, - added: payload.added, - modified: payload.modified, - deleted: payload.deleted, - summary: payload.summary, - }), - ); - }, - onCompactionStart: async () => { - if (isProcessAborted(abortSignal)) { - return; - } - await statusReactions.setCompacting(); - }, - onCompactionEnd: async () => { - if (isProcessAborted(abortSignal)) { - return; - } - statusReactions.cancelPending(); - await statusReactions.setThinking(); - }, }, }); if (!preparedResult.dispatched) { @@ -1295,50 +685,10 @@ async function processDiscordMessageInner( } throw err; } finally { - endDiscordInboundEventDeliveryCorrelation(); + endDeliveryCorrelation(); await draftPreview.cleanup(); const finalDeliveryFailed = (dispatchResult?.failedCounts?.final ?? 0) > 0; - if (statusReactionsActive) { - if (dispatchAborted) { - if (removeAckAfterReply) { - void statusReactions.clear(); - } else { - void statusReactions.restoreInitial(); - } - } else { - if (dispatchError || finalDeliveryFailed) { - await statusReactions.setError(); - } else { - await statusReactions.setDone(); - } - if (removeAckAfterReply) { - void (async () => { - await sleep( - dispatchError || finalDeliveryFailed - ? statusReactionTiming.errorHoldMs - : statusReactionTiming.doneHoldMs, - ); - await statusReactions.clear(); - })(); - } else { - void statusReactions.restoreInitial(); - } - } - } else if (shouldSendAckReaction && ackReaction && removeAckAfterReply) { - void removeReactionDiscord( - messageChannelId, - message.id, - ackReaction, - ackReactionContext, - ).catch((err: unknown) => { - logAckFailure({ - log: logVerbose, - channel: "discord", - target: `${messageChannelId}/${message.id}`, - error: err, - }); - }); - } + await reactions.finish({ dispatchAborted, dispatchError, finalDeliveryFailed }); } if (dispatchAborted) { return; @@ -1355,4 +705,3 @@ async function processDiscordMessageInner( ); } } -/* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */ diff --git a/extensions/discord/src/monitor/reply-delivery.ts b/extensions/discord/src/monitor/reply-delivery.ts index f916e8b2ea9c..ded94de03c8e 100644 --- a/extensions/discord/src/monitor/reply-delivery.ts +++ b/extensions/discord/src/monitor/reply-delivery.ts @@ -53,6 +53,23 @@ export function formatDiscordReplyDeliveryFailure(params: { return `discord ${params.kind} reply failed (${context}): ${String(params.err)}`; } +type DiscordReplySkipReason = "aborted before delivery" | "internal-only payload"; + +export function formatDiscordReplySkip(params: { + kind: "tool" | "block" | "final"; + reason: DiscordReplySkipReason; + target: string; + sessionKey?: string; +}) { + const context = [ + `target=${params.target}`, + params.sessionKey ? `session=${params.sessionKey}` : undefined, + ] + .filter(Boolean) + .join(" "); + return `discord ${params.kind} reply skipped (${params.reason}): ${context}`; +} + function resolveTargetChannelId(target: string): string | undefined { if (!target.startsWith("channel:")) { return undefined;