refactor(gateway): extract non-agent reply finalization (#106740)

This commit is contained in:
Peter Steinberger
2026-07-13 12:00:27 -07:00
committed by GitHub
parent 72b42f0ed7
commit 8dd57279cd
3 changed files with 303 additions and 262 deletions
@@ -0,0 +1,33 @@
import { describe, expect, it } from "vitest";
import { buildChatSendBtwSideResult } from "./chat-send-nonagent-finalization.js";
describe("buildChatSendBtwSideResult", () => {
it("combines BTW replies and preserves their error state", () => {
expect(
buildChatSendBtwSideResult([
{ kind: "block", payload: { text: "ignored" } },
{
kind: "final",
payload: { text: "first", btw: { question: " why? " } },
},
{
kind: "final",
payload: { text: "second", btw: { question: "why?" }, isError: true },
},
]),
).toEqual({
question: "why?",
text: "first\n\nsecond",
isError: true,
});
});
it("ignores empty or absent BTW replies", () => {
expect(
buildChatSendBtwSideResult([
{ kind: "final", payload: { text: "answer" } },
{ kind: "final", payload: { text: " ", btw: { question: "why?" } } },
]),
).toBeUndefined();
});
});
@@ -0,0 +1,257 @@
import { expectDefined } from "@openclaw/normalization-core";
import type { ReplyPayload } from "../../auto-reply/reply-payload.js";
import {
appendLocalMediaParentRoots,
getAgentScopedMediaLocalRoots,
} from "../../media/local-roots.js";
import { stripInlineDirectiveTagsForDisplay } from "../../utils/directive-tags.js";
import { attachManagedOutgoingImagesToMessage } from "../managed-image-attachments.js";
import { loadSessionEntry } from "../session-utils.js";
import { formatForLog } from "../ws-log.js";
import {
buildAssistantDisplayContentFromReplyPayloads,
extractAssistantDisplayText,
extractAssistantDisplayTextFromContent,
hasAssistantDisplayMediaContent,
hasSensitiveMediaPayload,
hasVisibleAssistantFinalMessage,
replaceAssistantContentTextBlocks,
stripManagedOutgoingAssistantContentBlocks,
} from "./chat-assistant-content.js";
import { broadcastChatFinal, broadcastSideResult, isBtwReplyPayload } from "./chat-broadcast.js";
import { normalizeWebchatReplyMediaPathsForDisplay } from "./chat-reply-media.js";
import { selectChatSendFinalReplyPayloads } from "./chat-send-command-replies.js";
import { buildTranscriptReplyText } from "./chat-send-reply-dispatch.js";
import type { PreparedChatSendSession } from "./chat-send-session.js";
import type { GatewayInjectedTtsSupplementMarker } from "./chat-transcript-inject.js";
import { appendAssistantTranscriptMessage } from "./chat-transcript-persistence.js";
import { buildMediaOnlyTtsSupplementTranscriptMarker } from "./chat-tts-markers.js";
import { buildWebchatAssistantMessageFromReplyPayloads } from "./chat-webchat-media.js";
import type { GatewayRequestContext } from "./types.js";
type DeliveredReply = {
payload: ReplyPayload;
kind: "block" | "final";
};
export function buildChatSendBtwSideResult(deliveredReplies: readonly DeliveredReply[]) {
const replies = deliveredReplies.map((entry) => entry.payload).filter(isBtwReplyPayload);
const text = replies
.map((payload) => payload.text.trim())
.filter(Boolean)
.join("\n\n")
.trim();
if (replies.length === 0 || !text) {
return undefined;
}
return {
question: expectDefined(replies[0], "btw replies entry at 0").btw.question.trim(),
text,
isError: replies.some((payload) => payload.isError),
};
}
/** Persist and broadcast replies produced without a runtime-owned agent assistant turn. */
export async function finalizeChatSendNonAgentReplies(params: {
accountId: string | undefined;
context: GatewayRequestContext;
deliveredReplies: readonly DeliveredReply[];
emitFirstAssistantServerTiming: () => void;
foldCommandBlocks: boolean;
persistUserTurnTranscript: () => Promise<void>;
session: Pick<
PreparedChatSendSession,
"agentId" | "backingSessionId" | "cfg" | "clientRunId" | "sessionKey" | "sessionLoadOptions"
>;
suppressReplies: boolean;
}): Promise<void> {
const {
accountId,
context,
deliveredReplies,
emitFirstAssistantServerTiming,
foldCommandBlocks,
persistUserTurnTranscript,
session,
suppressReplies,
} = params;
const { agentId, backingSessionId, cfg, clientRunId, sessionKey, sessionLoadOptions } = session;
const btwResult = buildChatSendBtwSideResult(deliveredReplies);
if (btwResult) {
broadcastSideResult({
context,
payload: {
kind: "btw",
runId: clientRunId,
sessionKey,
...(sessionKey === "global" && agentId ? { agentId } : {}),
...btwResult,
ts: Date.now(),
},
});
broadcastChatFinal({
context,
runId: clientRunId,
sessionKey,
agentId,
});
return;
}
const rawFinalPayloads = selectChatSendFinalReplyPayloads({
deliveredReplies,
foldCommandBlocks,
suppressReplies,
});
const finalPayloads = await normalizeWebchatReplyMediaPathsForDisplay({
cfg,
sessionKey,
agentId,
accountId,
payloads: rawFinalPayloads,
});
const { storePath: latestStorePath, entry: latestEntry } = loadSessionEntry(
sessionKey,
sessionLoadOptions,
);
const sessionId = latestEntry?.sessionId ?? backingSessionId ?? clientRunId;
const mediaLocalRoots = appendLocalMediaParentRoots(
getAgentScopedMediaLocalRoots(cfg, agentId),
latestStorePath ? [latestStorePath] : undefined,
);
const assistantContent = await buildAssistantDisplayContentFromReplyPayloads({
sessionKey,
agentId,
payloads: finalPayloads,
managedImageLocalRoots: mediaLocalRoots,
includeSensitiveMedia: false,
includeSensitiveDisplay: true,
onLocalAudioAccessDenied: (message) => {
context.logGateway.warn(`webchat audio embedding denied local path: ${message}`);
},
onManagedImagePrepareError: (message) => {
context.logGateway.warn(`webchat image embedding skipped attachment: ${message}`);
},
onSensitiveDisplayPrepareError: (message) => {
context.logGateway.warn(`webchat sensitive display skipped attachment: ${message}`);
},
});
const mediaMessage = await buildWebchatAssistantMessageFromReplyPayloads(finalPayloads, {
localRoots: mediaLocalRoots,
onLocalAudioAccessDenied: (err) => {
context.logGateway.warn(`webchat audio embedding denied local path: ${formatForLog(err)}`);
},
});
const hasSensitiveMedia = hasSensitiveMediaPayload(finalPayloads);
const ttsSupplementMarker = finalPayloads
.map((payload) => buildMediaOnlyTtsSupplementTranscriptMarker(payload))
.find((marker): marker is GatewayInjectedTtsSupplementMarker => Boolean(marker));
const persistedAssistantContent = replaceAssistantContentTextBlocks(
hasSensitiveMedia
? await buildAssistantDisplayContentFromReplyPayloads({
sessionKey,
agentId,
payloads: finalPayloads,
managedImageLocalRoots: mediaLocalRoots,
includeSensitiveMedia: false,
onLocalAudioAccessDenied: (message) => {
context.logGateway.warn(`webchat audio embedding denied local path: ${message}`);
},
onManagedImagePrepareError: (message) => {
context.logGateway.warn(`webchat image embedding skipped attachment: ${message}`);
},
})
: assistantContent,
mediaMessage,
);
const persistedContentForAppend = hasAssistantDisplayMediaContent(persistedAssistantContent)
? persistedAssistantContent
: undefined;
const broadcastAssistantContent = hasAssistantDisplayMediaContent(assistantContent)
? assistantContent
: hasAssistantDisplayMediaContent(mediaMessage?.content)
? mediaMessage?.content
: assistantContent;
const displayReply =
extractAssistantDisplayTextFromContent(assistantContent) ??
buildTranscriptReplyText(finalPayloads);
const transcriptDisplayReply = displayReply
? stripInlineDirectiveTagsForDisplay(displayReply).text.trim()
: "";
const transcriptReply =
mediaMessage?.transcriptText ||
buildTranscriptReplyText(finalPayloads) ||
transcriptDisplayReply;
let message: Record<string, unknown> | undefined;
const shouldAppendAssistantTranscript = Boolean(
transcriptReply || persistedContentForAppend?.length,
);
await persistUserTurnTranscript();
if (shouldAppendAssistantTranscript) {
const appended = await appendAssistantTranscriptMessage({
sessionKey,
message: transcriptReply,
...(persistedContentForAppend?.length ? { content: persistedContentForAppend } : {}),
sessionId,
storePath: latestStorePath,
sessionFile: latestEntry?.sessionFile,
agentId,
createIfMissing: true,
idempotencyKey: clientRunId,
ttsSupplement: ttsSupplementMarker,
cfg,
});
if (appended.ok) {
if (appended.messageId && assistantContent?.length) {
await attachManagedOutgoingImagesToMessage({
messageId: appended.messageId,
blocks: assistantContent,
});
}
message = broadcastAssistantContent?.length
? { ...appended.message, content: broadcastAssistantContent }
: appended.message;
} else {
context.logGateway.warn(
`webchat transcript append failed: ${appended.error ?? "unknown error"}`,
);
const fallbackAssistantContent =
stripManagedOutgoingAssistantContentBlocks(persistedAssistantContent) ??
stripManagedOutgoingAssistantContentBlocks(assistantContent);
const fallbackText = extractAssistantDisplayText(fallbackAssistantContent) ?? displayReply;
message = {
role: "assistant",
...(fallbackAssistantContent?.length
? { content: fallbackAssistantContent }
: fallbackText
? { content: [{ type: "text", text: fallbackText }] }
: {}),
...(fallbackText ? { text: fallbackText } : {}),
timestamp: Date.now(),
...(ttsSupplementMarker ? { openclawTtsSupplement: ttsSupplementMarker } : {}),
// Keep compatible with runner stopReason enums when transcript persistence fails.
stopReason: "stop",
usage: { input: 0, output: 0, totalTokens: 0 },
};
}
} else if (broadcastAssistantContent?.length) {
message = {
role: "assistant",
content: broadcastAssistantContent,
text: extractAssistantDisplayText(broadcastAssistantContent) ?? "",
timestamp: Date.now(),
stopReason: "stop",
usage: { input: 0, output: 0, totalTokens: 0 },
};
}
if (hasVisibleAssistantFinalMessage(message)) {
emitFirstAssistantServerTiming();
}
broadcastChatFinal({
context,
runId: clientRunId,
sessionKey,
agentId,
message,
});
}
+13 -262
View File
@@ -2,7 +2,6 @@
// bridge UI RPCs to agent dispatch, transcripts, media, and streaming state.
import { createHash } from "node:crypto";
import { performance } from "node:perf_hooks";
import { expectDefined } from "@openclaw/normalization-core";
import {
GATEWAY_CLIENT_CAPS,
hasGatewayClientCap,
@@ -27,7 +26,6 @@ import type { ModelCatalogEntry, ModelCatalogSnapshot } from "../../agents/model
import { resolveProviderIdForAuth } from "../../agents/provider-auth-aliases.js";
import { createAgentRunRestartAbortError } from "../../agents/run-termination.js";
import { dispatchInboundMessage } from "../../auto-reply/dispatch.js";
import type { ReplyPayload } from "../../auto-reply/reply-payload.js";
import { resolveSessionWorkStartError } from "../../config/sessions.js";
import { resolveTranscriptSessionKeyBySessionId } from "../../config/sessions/session-accessor.js";
import type { OpenClawConfig } from "../../config/types.openclaw.js";
@@ -41,10 +39,6 @@ import {
} from "../../infra/diagnostics-timeline.js";
import { formatErrorMessage } from "../../infra/errors.js";
import { jsonUtf8Bytes } from "../../infra/json-utf8-bytes.js";
import {
appendLocalMediaParentRoots,
getAgentScopedMediaLocalRoots,
} from "../../media/local-roots.js";
import {
retainGatewayRootWorkAdmissionContinuation,
runWithGatewayIndependentRootWorkContinuation,
@@ -74,7 +68,6 @@ import {
isDashboardSessionTitleCandidate,
maybeGenerateDashboardSessionTitle,
} from "../dashboard-session-title.js";
import { attachManagedOutgoingImagesToMessage } from "../managed-image-attachments.js";
import type { ChatRunTiming } from "../server-chat-state.js";
import { getMaxChatHistoryMessagesBytes, MAX_PAYLOAD_BYTES } from "../server-constants.js";
import { persistGatewaySessionLifecycleEvent } from "../session-lifecycle-state.js";
@@ -95,22 +88,10 @@ import { formatForLog } from "../ws-log.js";
import { setGatewayDedupeEntry } from "./agent-job.js";
import { handleChatAbortRequest } from "./chat-abort-handler.js";
import { ensureChatQueuedTurns } from "./chat-abort-runtime.js";
import {
buildAssistantDisplayContentFromReplyPayloads,
extractAssistantDisplayText,
extractAssistantDisplayTextFromContent,
hasAssistantDisplayMediaContent,
hasSensitiveMediaPayload,
hasVisibleAssistantFinalMessage,
replaceAssistantContentTextBlocks,
scheduleChatHistoryManagedImageCleanup,
stripManagedOutgoingAssistantContentBlocks,
} from "./chat-assistant-content.js";
import { scheduleChatHistoryManagedImageCleanup } from "./chat-assistant-content.js";
import {
broadcastChatError,
broadcastChatFinal,
broadcastSideResult,
isBtwReplyPayload,
sendGlobalAwareNodeChatPayload,
} from "./chat-broadcast.js";
import {
@@ -130,19 +111,15 @@ import {
resolveRequestedChatAgentId,
validateChatSelectedAgent,
} from "./chat-origin-routing.js";
import { normalizeWebchatReplyMediaPathsForDisplay } from "./chat-reply-media.js";
import { terminalizeRestartSafeChatAdmission } from "./chat-restart-recovery.js";
import { admitChatSend } from "./chat-send-admission.js";
import { prepareChatSendAttachments } from "./chat-send-attachments.js";
import { selectChatSendFinalReplyPayloads } from "./chat-send-command-replies.js";
import { finalizeChatSendNonAgentReplies } from "./chat-send-nonagent-finalization.js";
import {
respondChatSessionRoutingChanged,
runChatSendPreAdmission,
} from "./chat-send-pre-admission.js";
import {
buildTranscriptReplyText,
createChatSendReplyDispatch,
} from "./chat-send-reply-dispatch.js";
import { createChatSendReplyDispatch } from "./chat-send-reply-dispatch.js";
import { normalizeChatSendRequest } from "./chat-send-request.js";
import { prepareChatSendSession } from "./chat-send-session.js";
import { finalizeChatSendSourceReplies } from "./chat-send-source-finalization.js";
@@ -155,11 +132,8 @@ import {
type ChatSendServerTimingPhase,
} from "./chat-server-timing.js";
import { normalizeOptionalChatText as normalizeOptionalText } from "./chat-text-normalization.js";
import type { GatewayInjectedTtsSupplementMarker } from "./chat-transcript-inject.js";
import { appendAssistantTranscriptMessage } from "./chat-transcript-persistence.js";
import { buildMediaOnlyTtsSupplementTranscriptMarker } from "./chat-tts-markers.js";
import { createGatewayChatUserTurnController } from "./chat-user-turn-recorder.js";
import { buildWebchatAssistantMessageFromReplyPayloads } from "./chat-webchat-media.js";
import {
loadOptionalServerMethodModelCatalog,
loadOptionalServerMethodModelCatalogSnapshot,
@@ -372,21 +346,6 @@ function resolveWebchatPromptCacheKey(params: {
return `openclaw-webchat-${digest}`;
}
async function buildWebchatAssistantMediaMessage(
payloads: ReplyPayload[],
options?: {
localRoots?: readonly string[];
onLocalAudioAccessDenied?: (message: string) => void;
},
): Promise<{ content: Array<Record<string, unknown>>; transcriptText: string } | null> {
return buildWebchatAssistantMessageFromReplyPayloads(payloads, {
localRoots: options?.localRoots,
onLocalAudioAccessDenied: (err) => {
options?.onLocalAudioAccessDenied?.(formatForLog(err));
},
});
}
export {
augmentChatHistoryWithCanvasBlocks,
DEFAULT_CHAT_HISTORY_TEXT_MAX_CHARS,
@@ -1479,224 +1438,16 @@ export const chatHandlers: GatewayRequestHandlers = {
// runtime-owned assistant turn, so it appends a gateway-injected assistant entry before
// broadcasting the final UI event.
if (!agentRunStarted && !queuedFollowupEnqueued) {
const btwReplies = deliveredReplies
.map((entryScoped) => entryScoped.payload)
.filter(isBtwReplyPayload);
const btwText = btwReplies
.map((payload) => payload.text.trim())
.filter(Boolean)
.join("\n\n")
.trim();
if (btwReplies.length > 0 && btwText) {
broadcastSideResult({
context,
payload: {
kind: "btw",
runId: clientRunId,
sessionKey,
...(sessionKey === "global" && agentId ? { agentId } : {}),
question: expectDefined(
btwReplies[0],
"btw replies entry at 0",
).btw.question.trim(),
text: btwText,
isError: btwReplies.some((payload) => payload.isError),
ts: Date.now(),
},
});
broadcastChatFinal({
context,
runId: clientRunId,
sessionKey,
agentId,
});
} else {
const rawFinalPayloads = selectChatSendFinalReplyPayloads({
deliveredReplies,
foldCommandBlocks: isInternalTextSlashCommandTurn,
suppressReplies: hasAppendedWebchatAgentMedia(),
});
const finalPayloads = await normalizeWebchatReplyMediaPathsForDisplay({
cfg,
sessionKey,
agentId,
accountId,
payloads: rawFinalPayloads,
});
const { storePath: latestStorePath, entry: latestEntry } = loadSessionEntry(
sessionKey,
sessionLoadOptions,
);
const sessionId = latestEntry?.sessionId ?? backingSessionId ?? clientRunId;
const mediaLocalRoots = appendLocalMediaParentRoots(
getAgentScopedMediaLocalRoots(cfg, agentId),
latestStorePath ? [latestStorePath] : undefined,
);
const assistantContent = await buildAssistantDisplayContentFromReplyPayloads({
sessionKey,
agentId,
payloads: finalPayloads,
managedImageLocalRoots: mediaLocalRoots,
includeSensitiveMedia: false,
includeSensitiveDisplay: true,
onLocalAudioAccessDenied: (message) => {
context.logGateway.warn(
`webchat audio embedding denied local path: ${message}`,
);
},
onManagedImagePrepareError: (message) => {
context.logGateway.warn(
`webchat image embedding skipped attachment: ${message}`,
);
},
onSensitiveDisplayPrepareError: (message) => {
context.logGateway.warn(
`webchat sensitive display skipped attachment: ${message}`,
);
},
});
const mediaMessage = await buildWebchatAssistantMediaMessage(finalPayloads, {
localRoots: mediaLocalRoots,
onLocalAudioAccessDenied: (message) => {
context.logGateway.warn(
`webchat audio embedding denied local path: ${message}`,
);
},
});
const hasSensitiveMedia = hasSensitiveMediaPayload(finalPayloads);
const ttsSupplementMarker = finalPayloads
.map((payload) => buildMediaOnlyTtsSupplementTranscriptMarker(payload))
.find((marker): marker is GatewayInjectedTtsSupplementMarker =>
Boolean(marker),
);
const persistedAssistantContent = replaceAssistantContentTextBlocks(
hasSensitiveMedia
? await buildAssistantDisplayContentFromReplyPayloads({
sessionKey,
agentId,
payloads: finalPayloads,
managedImageLocalRoots: mediaLocalRoots,
includeSensitiveMedia: false,
onLocalAudioAccessDenied: (message) => {
context.logGateway.warn(
`webchat audio embedding denied local path: ${message}`,
);
},
onManagedImagePrepareError: (message) => {
context.logGateway.warn(
`webchat image embedding skipped attachment: ${message}`,
);
},
})
: assistantContent,
mediaMessage,
);
const persistedContentForAppend = hasAssistantDisplayMediaContent(
persistedAssistantContent,
)
? persistedAssistantContent
: undefined;
const broadcastAssistantContent = hasAssistantDisplayMediaContent(
assistantContent,
)
? assistantContent
: hasAssistantDisplayMediaContent(mediaMessage?.content)
? mediaMessage?.content
: assistantContent;
const displayReply =
extractAssistantDisplayTextFromContent(assistantContent) ??
buildTranscriptReplyText(finalPayloads);
const transcriptDisplayReply = displayReply
? stripInlineDirectiveTagsForDisplay(displayReply).text.trim()
: "";
const transcriptReply =
mediaMessage?.transcriptText ||
buildTranscriptReplyText(finalPayloads) ||
transcriptDisplayReply;
let message: Record<string, unknown> | undefined;
const shouldAppendAssistantTranscript = Boolean(
transcriptReply || persistedContentForAppend?.length,
);
if (shouldAppendAssistantTranscript) {
await persistGatewayUserTurnTranscriptBestEffort();
} else {
await persistGatewayUserTurnTranscriptBestEffort();
}
if (shouldAppendAssistantTranscript) {
const appended = await appendAssistantTranscriptMessage({
sessionKey,
message: transcriptReply,
...(persistedContentForAppend?.length
? { content: persistedContentForAppend }
: {}),
sessionId,
storePath: latestStorePath,
sessionFile: latestEntry?.sessionFile,
agentId,
createIfMissing: true,
idempotencyKey: clientRunId,
ttsSupplement: ttsSupplementMarker,
cfg,
});
if (appended.ok) {
if (appended.messageId && assistantContent?.length) {
await attachManagedOutgoingImagesToMessage({
messageId: appended.messageId,
blocks: assistantContent,
});
}
message = broadcastAssistantContent?.length
? { ...appended.message, content: broadcastAssistantContent }
: appended.message;
} else {
context.logGateway.warn(
`webchat transcript append failed: ${appended.error ?? "unknown error"}`,
);
const fallbackAssistantContent =
stripManagedOutgoingAssistantContentBlocks(persistedAssistantContent) ??
stripManagedOutgoingAssistantContentBlocks(assistantContent);
const fallbackText =
extractAssistantDisplayText(fallbackAssistantContent) ?? displayReply;
const nowValue = Date.now();
message = {
role: "assistant",
...(fallbackAssistantContent?.length
? { content: fallbackAssistantContent }
: fallbackText
? { content: [{ type: "text", text: fallbackText }] }
: {}),
...(fallbackText ? { text: fallbackText } : {}),
timestamp: nowValue,
...(ttsSupplementMarker
? { openclawTtsSupplement: ttsSupplementMarker }
: {}),
// Keep this compatible with runner stopReason enums even though this message isn't
// persisted to the transcript due to the append failure.
stopReason: "stop",
usage: { input: 0, output: 0, totalTokens: 0 },
};
}
} else if (broadcastAssistantContent?.length) {
message = {
role: "assistant",
content: broadcastAssistantContent,
text: extractAssistantDisplayText(broadcastAssistantContent) ?? "",
timestamp: Date.now(),
stopReason: "stop",
usage: { input: 0, output: 0, totalTokens: 0 },
};
}
if (hasVisibleAssistantFinalMessage(message)) {
emitFirstAssistantServerTiming();
}
broadcastChatFinal({
context,
runId: clientRunId,
sessionKey,
agentId,
message,
});
}
await finalizeChatSendNonAgentReplies({
accountId,
context,
deliveredReplies,
emitFirstAssistantServerTiming,
foldCommandBlocks: isInternalTextSlashCommandTurn,
persistUserTurnTranscript: persistGatewayUserTurnTranscriptBestEffort,
session: preparedSession.value,
suppressReplies: hasAppendedWebchatAgentMedia(),
});
} else {
broadcastedSourceReplyFinal = await finalizeChatSendSourceReplies({
accountId,