refactor(discord): split message process runtime (#111119)

This commit is contained in:
Peter Steinberger
2026-07-18 19:04:30 -07:00
committed by GitHub
parent 047232bd16
commit ced95b3fce
6 changed files with 872 additions and 710 deletions
-1
View File
@@ -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
@@ -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<ChannelInboundTurnPlan["replyOptions"]>;
type CallbackPayload<K extends keyof ReplyOptions> =
NonNullable<ReplyOptions[K]> extends (...args: infer Args) => unknown ? Args[0] : never;
type DraftPreview = ReturnType<typeof createDiscordDraftPreviewController>;
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<void>;
};
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<string>();
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<ReplyOptions> = {
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,
};
}
@@ -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<string, unknown>;
};
function readToolStringArg(args: Record<string, unknown>, key: string): string | undefined {
const value = args[key];
return typeof value === "string" && value.trim() ? value.trim() : undefined;
}
function readToolBooleanArg(args: Record<string, unknown>, 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<string, unknown>,
): Promise<string> => {
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,
};
}
@@ -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<ReturnType<typeof buildDiscordMessageProcessContext>>
>;
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<typeof createDiscordReplyTypingFeedback> | 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<string | undefined> => {
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,
};
}
@@ -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<string, unknown>;
detailMode?: "explain" | "raw";
};
function readToolStringArg(args: Record<string, unknown>, key: string): string | undefined {
const value = args[key];
return typeof value === "string" && value.trim() ? value.trim() : undefined;
}
function readToolBooleanArg(args: Record<string, unknown>, 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<string, unknown>,
): Promise<string> => {
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<typeof createDiscordReplyTypingFeedback> | 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<string | undefined> => {
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<string>();
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<ReplyDispatchKind, number>;
@@ -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. */
@@ -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;